"""Versioned physical layout for :class:`httk.store.backend.sql.store.SqlStore`."""
import dataclasses
from collections.abc import Mapping, Sequence
from types import MappingProxyType
from typing import Final, Literal
import sqlalchemy
from httk.store.backend.schema import TableSchema, resolve_schema
from httk.store.backend.sql.mapping import (
dispatch_table_for,
entry_dispatch_table_name,
identity_owner_tables,
table_for,
)
from httk.store.storage_layout import (
EntryFamilyDeclaration,
EntryFamilyLayout,
StorageLayout,
StorageLayoutUpgradeRequiredError,
_merge_storage_layouts,
_walk_closure,
declaration_json,
validate_entry_id_fields,
)
from httk.store.storage_layout import (
normalize_entry_families as _normalize_entry_families,
)
from httk.store.storage_layout import (
normalize_entry_records as _normalize_entry_records,
)
__all__ = [
"METADATA_TABLE_NAME",
"STORAGE_PROTOCOL_VERSION",
"WRITE_PROFILE_VOCABULARY",
"BackendFacts",
"EntryFamilyLayout",
"StorageLayout",
"StorageLayoutUpgradeRequiredError",
"StoreUnderConstructionError",
"actual_columns",
"actual_schema_objects",
"actual_table_names",
"backend_facts_for_dialect",
"declaration_json",
"expected_metadata",
"metadata_table_for",
"normalize_entry_declaration",
"normalize_entry_families",
"normalize_entry_records",
"read_store_metadata",
]
[docs]
STORAGE_PROTOCOL_VERSION: Final = "2"
# The value is the major store generation only, compared for strict equality on
# reopen; it is never parsed. Bump it to "3" solely on a breaking change to the
# physical SqlStore layout that an existing store could not be reopened against.
"""The persisted SqlStore layout protocol implemented by this package."""
"""Reserved key/value table holding the store protocol and entry declaration."""
_METADATA_PROTOCOL_KEY: Final = "protocol"
_METADATA_DECLARATION_KEY: Final = "entry_declaration"
_RESERVED_PREFIX: Final = "_httk_"
[docs]
WRITE_PROFILE_VOCABULARY: Final = frozenset({"transactional", "degraded", "bulk-fenced"})
[docs]
class StoreUnderConstructionError(RuntimeError):
"""A new open found an interrupted empty-store bulk ingest.
Crash window for new SQLite/DuckDB opens: before the marker commits the
old clean state remains accepted; after the marker and through ingest,
finalize, or before marker clear the store is rejected; after clear it is
accepted again. The marker is intentionally not a resume protocol.
ClickHouse marker residue is fail-closed: the default recovery is to drop
the database and re-ingest. Clearing the marker is valid only after an
operator has restored and verified the declared empty-store invariant.
"""
@dataclasses.dataclass(frozen=True)
[docs]
class BackendFacts:
"""Dialect capabilities used by the SQL storage protocol."""
[docs]
transactional_ddl: bool
[docs]
transactional_dml: bool
[docs]
supports_sequences: bool
[docs]
supports_deferred_finalize: bool
[docs]
supports_degraded: bool
[docs]
write_profiles: tuple[str, ...]
[docs]
supports_incremental_save: bool
[docs]
system_catalog: Literal["sqlite", "duckdb", "clickhouse", "postgresql"]
[docs]
stage_load: Literal["attach", "duckdb-views", "client-stream"]
[docs]
finalize_map_maintenance: Literal["update", "swap"]
[docs]
supports_adhoc_indexes: bool
_BACKEND_FACTS: Final[dict[str, BackendFacts]] = {
"sqlite": BackendFacts(
transactional_ddl=False,
transactional_dml=True,
supports_sequences=False,
atomic_upsert=True,
serial_stage_format="sqlite",
parallel_shard_format="sqlite",
supports_deferred_finalize=True,
supports_degraded=True,
write_profiles=("transactional", "degraded"),
metadata_backend="table",
supports_incremental_save=True,
system_catalog="sqlite",
stage_load="attach",
finalize_map_maintenance="update",
supports_adhoc_indexes=True,
),
"duckdb": BackendFacts(
transactional_ddl=True,
transactional_dml=True,
supports_sequences=True,
atomic_upsert=True,
serial_stage_format="duckdb-attach",
parallel_shard_format="parquet",
supports_deferred_finalize=True,
supports_degraded=False,
write_profiles=("transactional",),
metadata_backend="table",
supports_incremental_save=True,
system_catalog="duckdb",
stage_load="duckdb-views",
finalize_map_maintenance="update",
supports_adhoc_indexes=True,
),
"postgresql": BackendFacts(
transactional_ddl=True,
transactional_dml=True,
supports_sequences=True,
atomic_upsert=True,
serial_stage_format="sqlite",
parallel_shard_format="sqlite",
supports_deferred_finalize=True,
supports_degraded=False,
write_profiles=("transactional",),
metadata_backend="table",
supports_incremental_save=True,
system_catalog="postgresql",
stage_load="attach",
finalize_map_maintenance="update",
supports_adhoc_indexes=True,
),
"clickhousedb": BackendFacts(
transactional_ddl=False,
transactional_dml=False,
supports_sequences=False,
atomic_upsert=False,
serial_stage_format="parquet",
parallel_shard_format="parquet",
supports_deferred_finalize=True,
supports_degraded=False,
write_profiles=("bulk-fenced",),
metadata_backend="keepermap",
supports_incremental_save=False,
system_catalog="clickhouse",
stage_load="client-stream",
finalize_map_maintenance="swap",
supports_adhoc_indexes=False,
),
}
[docs]
def backend_facts_for_dialect(dialect_name: str) -> BackendFacts:
"""Resolve the hardcoded protocol facts for one supported dialect."""
try:
return _BACKEND_FACTS[dialect_name]
except KeyError as error:
raise ValueError(f"SqlStore layout validation does not support dialect {dialect_name!r}") from error
[docs]
def normalize_entry_records(entry_records: Mapping[type, type | tuple[type, ...]]) -> StorageLayout:
"""Normalize a declaration and apply SQL physical-name validation."""
layout = _normalize_entry_records(entry_records)
_validate_physical_names(layout)
return layout
[docs]
def normalize_entry_families(entry_families: Sequence[EntryFamilyDeclaration]) -> StorageLayout:
"""Normalize application-owned declarations and apply SQL physical-name validation."""
layout = _normalize_entry_families(entry_families)
_validate_physical_names(layout)
return layout
[docs]
def normalize_entry_declaration(
entry_records: Mapping[type, type | tuple[type, ...]] | None,
entry_families: Sequence[EntryFamilyDeclaration] | None,
) -> StorageLayout | None:
"""Merge registered and application-owned declarations and validate SQL names."""
layouts = []
if entry_records is not None:
layouts.append(_normalize_entry_records(entry_records))
if entry_families is not None:
layouts.append(_normalize_entry_families(entry_families))
if not layouts:
return None
layout = _merge_storage_layouts(*layouts)
validate_entry_id_fields(layout)
_validate_physical_names(layout)
return layout
def _layout_from_declaration(value: str) -> StorageLayout:
"""Parse a declaration and apply SQL physical-name validation."""
from httk.store.storage_layout import _layout_from_declaration as parse_declaration
layout = parse_declaration(value)
_validate_physical_names(layout)
return layout
[docs]
def actual_schema_objects(connection: sqlalchemy.Connection) -> Mapping[str, frozenset[str]]:
"""Return application schema-object names mapped to their stable object kinds.
The DuckDB SQLAlchemy inspector presently routes column inspection through
a PostgreSQL catalogue relation DuckDB does not expose, so the whole
layout path intentionally uses the dialect catalogues directly.
"""
facts = backend_facts_for_dialect(connection.dialect.name)
if facts.system_catalog == "sqlite":
rows = connection.execute(
sqlalchemy.text(
"SELECT name, type FROM sqlite_master WHERE type IN ('table', 'view') AND name NOT LIKE 'sqlite_%'"
)
)
elif facts.system_catalog == "duckdb":
rows = connection.execute(
sqlalchemy.text(
"SELECT table_name, lower(table_type) FROM information_schema.tables "
"WHERE table_catalog = current_database() AND table_schema = current_schema() "
"UNION ALL "
"SELECT sequence_name, 'sequence' FROM duckdb_sequences() "
"WHERE database_name = current_database() AND schema_name = current_schema()"
)
)
elif facts.system_catalog == "postgresql":
rows = connection.execute(
sqlalchemy.text(
"SELECT table_name, lower(table_type) FROM information_schema.tables "
"WHERE table_catalog = current_database() AND table_schema = current_schema() "
"UNION ALL "
"SELECT sequence_name, 'sequence' FROM information_schema.sequences "
"WHERE sequence_catalog = current_database() AND sequence_schema = current_schema() "
"UNION ALL "
# Materialized views are absent from information_schema.tables; a
# stray one must still count as a non-empty schema object.
"SELECT matviewname, 'view' FROM pg_matviews WHERE schemaname = current_schema()"
)
)
else:
from httk.store.backend.clickhouse.support import actual_schema_objects as clickhouse_schema_objects
return clickhouse_schema_objects(connection)
result: dict[str, set[str]] = {}
for name, kind in rows:
result.setdefault(str(name), set()).add(str(kind).lower().replace("base ", ""))
return MappingProxyType({name: frozenset(kinds) for name, kinds in result.items()})
[docs]
def actual_table_names(connection: sqlalchemy.Connection) -> frozenset[str]:
"""Return application base-table names without SQLAlchemy reflection."""
return frozenset(name for name, kinds in actual_schema_objects(connection).items() if "table" in kinds)
[docs]
def actual_columns(connection: sqlalchemy.Connection, table_name: str) -> frozenset[str]:
"""Return the physical column names of ``table_name`` without SQLAlchemy reflection.
The DuckDB SQLAlchemy inspector routes column inspection through a
PostgreSQL catalogue relation DuckDB does not expose (see
:func:`actual_schema_objects`), so this reads the dialect catalogues directly.
:param connection: The open connection whose dialect selects the catalogue.
:param table_name: The table whose columns are read.
:return: The set of physical column names, empty when the table is absent.
"""
facts = backend_facts_for_dialect(connection.dialect.name)
if facts.system_catalog == "sqlite":
# PRAGMA cannot bind parameters, so the name is quoted, not bound.
quoted = connection.dialect.identifier_preparer.quote(table_name)
pragma_rows = connection.execute(sqlalchemy.text(f"PRAGMA table_info({quoted})")).mappings().all()
return frozenset(str(row["name"]) for row in pragma_rows)
if facts.system_catalog in ("duckdb", "postgresql"):
rows = connection.execute(
sqlalchemy.text(
"SELECT column_name FROM information_schema.columns "
"WHERE table_catalog = current_database() AND table_schema = current_schema() "
"AND table_name = :name"
),
{"name": table_name},
).all()
return frozenset(str(row[0]) for row in rows)
from httk.store.backend.clickhouse.support import actual_columns as clickhouse_actual_columns
return frozenset(clickhouse_actual_columns(connection, table_name))
def _validate_physical_names(layout: StorageLayout) -> None:
owners: dict[str, type] = {}
def check(schema: TableSchema) -> None:
record = schema.cls
names = [schema.table_name]
names.extend(spec.child.table_name for spec in schema.fields if spec.child is not None)
for name in names:
if name.startswith(_RESERVED_PREFIX):
raise ValueError(f"record {record.__name__} claims reserved SqlStore table name {name!r}")
previous = owners.get(name)
if previous is not None and previous is not record:
raise ValueError(
f"records {previous.__name__} and {record.__name__} collide on physical table name {name!r}"
)
owners[name] = record
_walk_closure(layout, check)
for family in layout.families:
dispatch_name = entry_dispatch_table_name(family.name) if len(family.records) > 1 else None
if dispatch_name is not None:
if dispatch_name in owners:
raise ValueError(
f"entry family {family.name!r} dispatch table collides with record table {dispatch_name!r}"
)
owners[dispatch_name] = family.family