fix(recovery): register lazy state.db tables in one schema map; cover .recover lane and count-mismatch loss
Follow-up to the salvaged #100350 commits: replace the per-table 'if table == "delivery_obligations"' branches in session_recovery.py and session_lost_and_found.py with a single _AUXILIARY_TABLE_SCHEMAS registry (table -> destination DDL initializer) that both the SQL-level and the lost_and_found lanes consume, so the next lazily-created state.db table is one entry, not three code paths. The .recover lane now iterates _CANONICAL_TABLES + _AUXILIARY_TABLES instead of a duplicated literal list. Tests: the .recover direct-copy lane creates the missing ledger on the destination; a source-vs-destination obligation count mismatch fails verification (complete=False) instead of reporting a clean salvage. Docs: state.db table inventory lists delivery_obligations. Addresses #100313
This commit is contained in:
@@ -325,25 +325,24 @@ def _copy_direct_tables(
|
||||
) -> dict[str, int]:
|
||||
"""Copy rows .recover managed to attribute to real canonical tables."""
|
||||
|
||||
# Lazy import: session_recovery imports this module inside a function, so
|
||||
# a module-level import here would be circular.
|
||||
from hermes_cli.session_recovery import (
|
||||
_AUXILIARY_TABLE_SCHEMAS,
|
||||
_AUXILIARY_TABLES,
|
||||
_CANONICAL_TABLES,
|
||||
)
|
||||
|
||||
copied: dict[str, int] = {}
|
||||
for table in (
|
||||
"system_prompts",
|
||||
"sessions",
|
||||
"messages",
|
||||
"session_model_usage",
|
||||
"compression_locks",
|
||||
"gateway_routing",
|
||||
"async_delegations",
|
||||
"delivery_obligations",
|
||||
):
|
||||
for table in (*_CANONICAL_TABLES, *_AUXILIARY_TABLES):
|
||||
source_columns = _table_columns(lf_conn, table)
|
||||
if not source_columns:
|
||||
continue
|
||||
dest_columns = _table_columns(dest, table)
|
||||
if table == "delivery_obligations" and not dest_columns:
|
||||
from gateway.delivery_ledger import _initialize_schema
|
||||
|
||||
_initialize_schema(dest)
|
||||
if not dest_columns and table in _AUXILIARY_TABLE_SCHEMAS:
|
||||
# Lazily-created gateway table: base SessionDB never made it on
|
||||
# the fresh destination, so create it before copying.
|
||||
_AUXILIARY_TABLE_SCHEMAS[table](dest)
|
||||
dest_columns = _table_columns(dest, table)
|
||||
columns = [c for c in dest_columns if c in source_columns]
|
||||
if not columns:
|
||||
|
||||
@@ -45,12 +45,26 @@ _TOPIC_TABLES = (
|
||||
"telegram_dm_topic_bindings",
|
||||
)
|
||||
|
||||
# Durable gateway outbox. Created lazily by gateway.delivery_ledger, so a
|
||||
# fresh SessionDB destination does not have the table until we initialize it.
|
||||
# Omitting it from the copy inventory drops owed replies (#100313, #86236).
|
||||
_AUXILIARY_TABLES = (
|
||||
"delivery_obligations",
|
||||
)
|
||||
|
||||
|
||||
def _init_delivery_ledger_schema(conn: sqlite3.Connection) -> None:
|
||||
from gateway.delivery_ledger import _initialize_schema
|
||||
|
||||
_initialize_schema(conn)
|
||||
|
||||
|
||||
# Tables that live in state.db but are created lazily by a gateway module on
|
||||
# first use, so base ``SessionDB`` never creates them on a fresh destination.
|
||||
# Every entry maps the table to the initializer that owns its DDL; recovery
|
||||
# creates the table on the destination before copying, so owed rows survive
|
||||
# instead of silently vanishing from a "complete" salvage (#100313, #86236).
|
||||
# Add new lazily-created state.db tables HERE, never as one-off ``if table ==``
|
||||
# branches.
|
||||
_AUXILIARY_TABLE_SCHEMAS: dict[str, Callable[[sqlite3.Connection], None]] = {
|
||||
"delivery_obligations": _init_delivery_ledger_schema,
|
||||
}
|
||||
|
||||
_AUXILIARY_TABLES = tuple(_AUXILIARY_TABLE_SCHEMAS)
|
||||
|
||||
_INVENTORY_TABLES = (
|
||||
*_CANONICAL_TABLES,
|
||||
@@ -410,10 +424,12 @@ def _ensure_auxiliary_destination_schema(
|
||||
would report ``missing`` / ``no compatible columns`` and drop the rows.
|
||||
"""
|
||||
|
||||
if table == "delivery_obligations":
|
||||
from gateway.delivery_ledger import _initialize_schema
|
||||
|
||||
_initialize_schema(destination)
|
||||
initialize = _AUXILIARY_TABLE_SCHEMAS.get(table)
|
||||
if initialize is None:
|
||||
raise SessionRecoverySafetyError(
|
||||
f"no destination schema initializer registered for table {table!r}"
|
||||
)
|
||||
initialize(destination)
|
||||
|
||||
|
||||
def _copy_table(
|
||||
|
||||
@@ -753,3 +753,83 @@ def test_recovery_without_delivery_ledger_is_not_lossy(tmp_path: Path) -> None:
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
def test_recovery_flags_delivery_obligation_count_mismatch_as_loss(
|
||||
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
|
||||
) -> None:
|
||||
"""A source-vs-destination ledger count mismatch must not verify as complete.
|
||||
|
||||
The destination table is created through the registered initializer; a
|
||||
real SQL trigger that silently drops one row stands in for the "rows went
|
||||
missing on the way over" failure the verifier has to catch.
|
||||
"""
|
||||
|
||||
from hermes_cli import session_recovery
|
||||
|
||||
source = tmp_path / "state.db"
|
||||
output = tmp_path / "recovered.db"
|
||||
_make_source(source)
|
||||
now = 1_720_000_000.0
|
||||
_insert_delivery_obligations(
|
||||
source,
|
||||
[
|
||||
("ob-a", "k", "telegram", "chat-1", None, "a", "pending", 0, now, now, None, None, None, "default"),
|
||||
("ob-b", "k", "telegram", "chat-1", None, "b", "pending", 0, now, now, None, None, None, "default"),
|
||||
],
|
||||
)
|
||||
|
||||
real_init = session_recovery._AUXILIARY_TABLE_SCHEMAS["delivery_obligations"]
|
||||
|
||||
def lossy_init(conn: sqlite3.Connection) -> None:
|
||||
real_init(conn)
|
||||
conn.execute(
|
||||
"""CREATE TRIGGER drop_ob_b BEFORE INSERT ON delivery_obligations
|
||||
WHEN NEW.obligation_id = 'ob-b' BEGIN SELECT RAISE(IGNORE); END"""
|
||||
)
|
||||
|
||||
monkeypatch.setitem(
|
||||
session_recovery._AUXILIARY_TABLE_SCHEMAS, "delivery_obligations", lossy_init
|
||||
)
|
||||
|
||||
report = recover_session_database(source, output, work_dir=tmp_path)
|
||||
assert report["verification"]["table_counts"]["delivery_obligations"] == 1
|
||||
assert report["complete"] is False
|
||||
assert any(
|
||||
"delivery_obligations count is 1, expected 2" in error
|
||||
for error in report["verification"]["errors"]
|
||||
)
|
||||
|
||||
|
||||
def test_lost_and_found_direct_copy_creates_lazy_delivery_ledger(tmp_path: Path) -> None:
|
||||
"""The .recover lane copies the ledger even though SessionDB never made it."""
|
||||
|
||||
from hermes_cli.session_lost_and_found import _copy_direct_tables
|
||||
|
||||
recovered_source = tmp_path / "lost_and_found.db"
|
||||
now = 1_720_000_000.0
|
||||
_insert_delivery_obligations(
|
||||
recovered_source,
|
||||
[
|
||||
("ob-1", "k", "telegram", "chat-1", None, "one", "pending", 0, now, now, None, None, None, "default"),
|
||||
("ob-2", "k", "telegram", "chat-1", None, "two", "failed", 3, now, now, None, None, "boom", "default"),
|
||||
],
|
||||
)
|
||||
output = tmp_path / "rebuilt.db"
|
||||
SessionDB(db_path=output).close()
|
||||
|
||||
lf_conn = sqlite3.connect(str(recovered_source), isolation_level=None)
|
||||
dest = sqlite3.connect(str(output), isolation_level=None)
|
||||
try:
|
||||
assert not dest.execute(
|
||||
"SELECT 1 FROM sqlite_master WHERE type='table' AND name='delivery_obligations'"
|
||||
).fetchall()
|
||||
copied = _copy_direct_tables(lf_conn, dest)
|
||||
assert copied["delivery_obligations"] == 2
|
||||
rows = dest.execute(
|
||||
"SELECT obligation_id, state, last_error FROM delivery_obligations ORDER BY obligation_id"
|
||||
).fetchall()
|
||||
finally:
|
||||
lf_conn.close()
|
||||
dest.close()
|
||||
assert rows == [("ob-1", "pending", None), ("ob-2", "failed", "boom")]
|
||||
|
||||
@@ -21,9 +21,15 @@ Source file: `hermes_state.py`
|
||||
├── gateway_routing — Gateway routing metadata
|
||||
├── compression_locks — Cross-process compression locking
|
||||
├── async_delegations — Async delegation bookkeeping
|
||||
├── delivery_obligations — Gateway outbox (owed replies); created lazily by gateway/delivery_ledger.py
|
||||
└── schema_version — Single-row table tracking migration state
|
||||
```
|
||||
|
||||
`hermes sessions recover` copies the row-bearing tables above into the
|
||||
recovered database (FTS indexes and `schema_version` are regenerated), including
|
||||
the lazily-created `delivery_obligations` ledger when the source has one — its
|
||||
row count is verified like `sessions`/`messages`.
|
||||
|
||||
Key design decisions:
|
||||
- **WAL mode** for concurrent readers + one writer (gateway multi-platform)
|
||||
- **FTS5 virtual table** for fast text search across all session messages
|
||||
|
||||
Reference in New Issue
Block a user