Files

55 lines
2.3 KiB
Python

"""Direct PostgreSQL implementation of the DWH port."""
from tht.config import DatabaseConfig
from sqlalchemy.exc import OperationalError, SQLAlchemyError
from tht.db import execute, sampling
from tht.db.connection import can_create_in_schema, make_engine, ping, writable_tables
from tht.db.introspect import introspect
from tht.execute import ExecResult, PlanSummary
from tht.mschema.models import PhysicalSchema
from tht.ports.dwh import DistinctValues, DwhCapabilities, DwhHealth
class PostgresDwhAdapter:
capabilities = DwhCapabilities()
def __init__(self, config: DatabaseConfig, *, statement_timeout_ms: int = 30_000):
self._config = config
self._engine = make_engine(config)
self._statement_timeout_ms = statement_timeout_ms
def health(self) -> DwhHealth:
try:
ping(self._engine)
except OperationalError as exc:
return DwhHealth(ok=False, detail=str(exc.orig), error_kind="connection")
except SQLAlchemyError as exc:
return DwhHealth(ok=False, detail=str(exc), error_kind="connection")
writable = tuple(writable_tables(self._engine, self._config.db_schema))
can_create = can_create_in_schema(self._engine, self._config.db_schema)
return DwhHealth(ok=True, database=self._config.database, schema=self._config.db_schema,
read_only=not writable and not can_create,
writable_tables=writable, can_create=can_create)
def introspect(self) -> PhysicalSchema:
return introspect(self._engine, self._config.database, self._config.db_schema)
def run_query(self, sql: str, *, limit: int) -> ExecResult:
return execute.run_query(
self._engine, sql, limit=limit, timeout_ms=self._statement_timeout_ms
)
def explain(self, sql: str) -> PlanSummary:
return execute.explain(self._engine, sql, timeout_ms=self._statement_timeout_ms)
def sample_column(self, table: str, column: str, *, limit: int) -> list[object]:
return sampling.sample_column(
self._engine, self._config.db_schema, table, column, limit=limit
)
def distinct_values(self, table: str, column: str, *, limit: int) -> DistinctValues:
return sampling.distinct_values(
self._engine, self._config.db_schema, table, column, max_values=limit
)