From ead7e91dabf1e963796ec834b196984a2fa44ff4 Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Wed, 9 Sep 2026 17:59:21 +0530 Subject: [PATCH] refactor(recovery): stream the salvaged population; table the shape rules MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Pass 1 no longer materialises every classified record (full `messages.content` included) until pass 2; `LayoutEvidence` keeps only the capped per-position value sets (+ sessions rows for the one cross-column invariant) and pass 2 re-streams the lost_and_found tables. A 276 MB corrupted store no longer has to fit in memory. - `_sentinel_holds` / `_text_shape_holds` if-ladders become rule tables. - Tests trimmed to the three that bind behaviour (upgraded store maps by name; verifier refuses when rows matched no layout; replayed history ends at the current schema — the drift guard). No behaviour change; reverting inference to "no layout" still fails the name-mapping test. --- hermes_cli/session_lost_and_found.py | 351 ++++++++++-------- hermes_cli/session_recovery.py | 9 +- hermes_cli/session_schema_history.py | 2 +- .../test_session_recovery_lost_and_found.py | 175 --------- .../hermes_cli/test_session_schema_history.py | 74 +--- 5 files changed, 205 insertions(+), 406 deletions(-) diff --git a/hermes_cli/session_lost_and_found.py b/hermes_cli/session_lost_and_found.py index 368e0c885f..ee2f36cd66 100644 --- a/hermes_cli/session_lost_and_found.py +++ b/hermes_cli/session_lost_and_found.py @@ -11,7 +11,7 @@ import sqlite3 import subprocess import tempfile from pathlib import Path -from typing import Any, Optional, Sequence +from typing import Any, Callable, Optional, Sequence from hermes_cli.session_schema_history import SCHEMA_HISTORY, reachable_physical_layouts @@ -312,10 +312,7 @@ def _insert_prefix_row( def _declared_types(conn: sqlite3.Connection, table: str) -> dict[str, str]: - return { - str(row[1]): str(row[2] or "") - for row in conn.execute(f'PRAGMA table_info("{table}")') - } + return {str(row[1]): str(row[2] or "") for row in conn.execute(f'PRAGMA table_info("{table}")')} def _type_conflicts(value: Any, declared: str) -> bool: @@ -356,67 +353,82 @@ _HANDOFF_STATES = frozenset({"pending", "running", "completed", "failed"}) _TITLE_SOURCES = frozenset({"derived", "llm", "user"}) +def _is_epoch(value: Any) -> bool: + return isinstance(value, (int, float)) and _EPOCH_LOW <= float(value) <= _EPOCH_HIGH + + +def _is_nonempty_str(value: Any) -> bool: + return isinstance(value, str) and bool(value) + + +# Sentinel columns: the cells that differ hardest between candidate layouts, so a wrong layout is +# rejected instead of silently shifting every field. messages.id is a rowid alias, never a sentinel. +_SENTINEL_RULES: dict[str, Callable[[Any], bool]] = { + "id": _is_session_id, + "source": _looks_like_source, + "started_at": _is_epoch, + "timestamp": _is_epoch, + "session_id": _is_nonempty_str, + "role": lambda value: value in MESSAGE_ROLES, + "model": _is_nonempty_str, +} + + def _sentinel_holds(name: str, value: Any) -> bool: - if name == "id": # sessions.id; messages.id is a rowid alias, never a sentinel - return _is_session_id(value) - if name == "source": - return _looks_like_source(value) - if name in ("started_at", "timestamp"): - return ( - isinstance(value, (int, float)) - and _EPOCH_LOW <= float(value) <= _EPOCH_HIGH - ) - if name == "session_id": - return isinstance(value, str) and bool(value) - if name == "role": - return value in MESSAGE_ROLES - if name == "model": - return isinstance(value, str) and bool(value) - return True + rule = _SENTINEL_RULES.get(name) + return rule(value) if rule else True + + +def _is_token(value: str) -> bool: + return bool(_TOKEN_PATTERN.match(value)) + + +def _is_json_start(value: str) -> bool: + return value[:1] in "{[" + + +def _is_path(value: str) -> bool: + return value[:1] in "/~" or (len(value) > 1 and value[1] == ":") + + +def _is_url(value: str) -> bool: + return "://" in value + + +def _blank_or(rule: Callable[[str], bool]) -> Callable[[str], bool]: + return lambda value: value == "" or rule(value) + + +# Cheap per-column shape rules for text cells, by table (see module comment above). +_TEXT_SHAPE_RULES: dict[str, dict[str, Callable[[str], bool]]] = { + "sessions": { + "session_key": lambda value: bool(_SESSION_KEY_PATTERN.match(value)), + **dict.fromkeys(("chat_type", "end_reason", "cost_status", "cost_source", "billing_mode", + "last_activity_provenance", "handoff_platform"), _is_token), + "pricing_version": lambda value: bool(re.fullmatch(r"[a-z0-9][a-z0-9._-]*", value)), + "title_source": lambda value: value in _TITLE_SOURCES, + "handoff_state": lambda value: value in _HANDOFF_STATES, + "parent_session_id": _is_session_id, + "system_prompt_hash": lambda value: bool(re.fullmatch(r"[0-9a-f]{64}", value)), + "model_config": _is_json_start, "origin_json": _is_json_start, + "cwd": _is_path, "git_repo_root": _is_path, + "billing_base_url": _is_url, + }, + "messages": { + **dict.fromkeys(("effect_disposition", "finish_reason", "display_kind"), _is_token), + **dict.fromkeys(("tool_calls", "reasoning_details", "codex_reasoning_items", "codex_message_items", + "api_content", "display_metadata"), _blank_or(_is_json_start)), + }, + "session_model_usage": { + "billing_base_url": _blank_or(_is_url), + **dict.fromkeys(("billing_mode", "cost_status", "cost_source"), _blank_or(_is_token)), + }, +} def _text_shape_holds(kind: str, name: str, value: str) -> bool: - """Cheap per-column shape rules for text cells (see module comment above).""" - - if kind == "sessions": - if name == "session_key": - return bool(_SESSION_KEY_PATTERN.match(value)) - if name in ("chat_type", "end_reason", "cost_status", "cost_source", - "billing_mode", "last_activity_provenance", "handoff_platform"): - return bool(_TOKEN_PATTERN.match(value)) - if name == "pricing_version": - return bool(re.fullmatch(r"[a-z0-9][a-z0-9._-]*", value)) - if name == "title_source": - return value in _TITLE_SOURCES - if name == "handoff_state": - return value in _HANDOFF_STATES - if name == "parent_session_id": - return _is_session_id(value) - if name == "system_prompt_hash": - return bool(re.fullmatch(r"[0-9a-f]{64}", value)) - if name in ("model_config", "origin_json"): - return value[:1] in "{[" - if name in ("cwd", "git_repo_root"): - return value[:1] in "/~" or (len(value) > 1 and value[1] == ":") - if name == "billing_base_url": - return "://" in value - return True - if kind == "messages": - if name == "effect_disposition": - return bool(_TOKEN_PATTERN.match(value)) - if name in ("tool_calls", "reasoning_details", "codex_reasoning_items", - "codex_message_items", "api_content", "display_metadata"): - return value[:1] in "{[" or value == "" - if name in ("finish_reason", "display_kind"): - return bool(_TOKEN_PATTERN.match(value)) - return True - if kind == "session_model_usage": - if name == "billing_base_url": - return "://" in value or value == "" - if name in ("billing_mode", "cost_status", "cost_source"): - return value == "" or bool(_TOKEN_PATTERN.match(value)) - return True - return True + rule = _TEXT_SHAPE_RULES.get(kind, {}).get(name) + return rule(value) if rule else True def _cell_fits(kind: str, name: str, value: Any, dest_types: dict[str, str]) -> bool: @@ -434,9 +446,7 @@ def _cell_fits(kind: str, name: str, value: Any, dest_types: dict[str, str]) -> return not isinstance(value, str) or _text_shape_holds(kind, name, value) -def _row_invariants_hold( - kind: str, layout: tuple[str, ...], rows: Sequence[tuple[Any, ...]] -) -> bool: +def _row_invariants_hold(kind: str, layout: tuple[str, ...], rows: Sequence[tuple[Any, ...]]) -> bool: """Cross-column invariants every writer honours, checked per record. ``handoff_error`` is only ever written together with ``handoff_state`` @@ -454,21 +464,41 @@ def _row_invariants_hold( state_at = positions.get("handoff_state") if error_at is None or state_at is None: return True - for row in rows: - if ( - error_at < len(row) - and row[error_at] is not None - and (state_at >= len(row) or row[state_at] is None) - ): - return False - return True + # rows are exactly len(layout) wide (bucketed by width), so both positions are in range. + return all(row[error_at] is None or row[state_at] is not None for row in rows) -def infer_physical_layouts( - kind: str, - rows: Sequence[tuple[Any, ...]], - dest_types: dict[str, str], -) -> dict[int, list[Optional[str]]]: +_SAMPLE_CAP = 512 + + +class LayoutEvidence: + """What layout inference needs from a population, gathered in one streaming pass. + + Distinct non-NULL values per position (a column's admissibility is a property of the value, not the + row, and salvaged populations repeat values heavily — a few hundred distinct values per position + discriminate as well as 150k) plus, for ``sessions`` only, the rows themselves for the cross-column + invariant. Keeping the whole population — every ``messages.content`` included — until pass 2 would + hold the entire corrupted store in memory. + """ + + def __init__(self, kind: str) -> None: + self.kind = kind + self.widths: set[int] = set() + self.by_position: list[set[Any]] = [] + self.rows_by_width: dict[int, list[tuple[Any, ...]]] = {} + + def add(self, cells: tuple[Any, ...]) -> None: + self.widths.add(len(cells)) + while len(self.by_position) < len(cells): + self.by_position.append(set()) + for index, value in enumerate(cells): + if value is not None and len(self.by_position[index]) < _SAMPLE_CAP: + self.by_position[index].add(value) + if self.kind == "sessions": + self.rows_by_width.setdefault(len(cells), []).append(cells) + + +def infer_physical_layouts(evidence: LayoutEvidence, dest_types: dict[str, str]) -> dict[int, list[Optional[str]]]: """Infer which source column each record position holds, per field count. Salvaged records carry no schema. Every layout a real store can have is, @@ -488,18 +518,10 @@ def infer_physical_layouts( Returns an empty dict when no known layout fits the records. """ - if kind not in SCHEMA_HISTORY or not rows: + kind = evidence.kind + if kind not in SCHEMA_HISTORY or not evidence.widths: return {} - widths = sorted({len(row) for row in rows}) - # Distinct non-NULL values per position: a column's admissibility is a - # property of the value, not the row, and salvaged populations repeat - # values heavily (flags, counters, roles). Cap the sample per position — - # a few hundred distinct values discriminate as well as 150k. - by_position: list[set[Any]] = [set() for _ in range(widths[-1])] - for row in rows: - for index, value in enumerate(row): - if value is not None and len(by_position[index]) < 512: - by_position[index].add(value) + by_position = evidence.by_position verdicts: dict[tuple[str, Any], bool] = {} def fits(name: str, value: Any) -> bool: @@ -509,11 +531,20 @@ def infer_physical_layouts( verdict = verdicts[key] = _cell_fits(kind, name, value, dest_types) return verdict - # Cross-column invariants need whole rows, not per-position value sets; - # a candidate of width ``w`` only ever names the rows of that width. - rows_by_width: dict[int, list[tuple[Any, ...]]] = {w: [] for w in widths} - for row in rows: - rows_by_width[len(row)].append(row) + # The invariant verdict depends only on where a candidate puts the two handoff columns, and most + # candidates of one width agree on that — memoise so the rows are not rescanned per candidate. + invariant_verdicts: dict[tuple[int, Optional[int], Optional[int]], bool] = {} + + def invariants_hold(layout: tuple[str, ...]) -> bool: + same_width = evidence.rows_by_width.get(len(layout)) + if not same_width: + return True + positions = {name: index for index, name in enumerate(layout)} + key = (len(layout), positions.get("handoff_error"), positions.get("handoff_state")) + verdict = invariant_verdicts.get(key) + if verdict is None: + verdict = invariant_verdicts[key] = _row_invariants_hold(kind, layout, same_width) + return verdict def accept(layout: tuple[str, ...], first_new: int) -> bool: for index in range(first_new, min(len(layout), len(by_position))): @@ -521,14 +552,13 @@ def infer_physical_layouts( for value in by_position[index]: if not fits(name, value): return False - same_width = rows_by_width.get(len(layout)) - return not same_width or _row_invariants_hold(kind, layout, same_width) + return invariants_hold(layout) # Enumerate once and bucket by width. A record of width ``k`` was written # while the table had exactly ``k`` columns (ADD COLUMN runs at startup, # before any row is written), so its layout is a chain state of exactly # that length. - survivors_by_width: dict[int, list[tuple[str, ...]]] = {w: [] for w in widths} + survivors_by_width: dict[int, list[tuple[str, ...]]] = {w: [] for w in sorted(evidence.widths)} for layout in reachable_physical_layouts(kind, accept): bucket = survivors_by_width.get(len(layout)) if bucket is not None and layout not in bucket: @@ -623,7 +653,8 @@ def map_lost_and_found_rows(lf_conn: sqlite3.Connection, dest: sqlite3.Connectio """Best-effort mapping of a .recover output DB into a fresh SessionDB.""" report: dict[str, Any] = { "direct_table_rows": {}, "mapped": {"sessions": 0, "messages": 0, "session_model_usage": 0}, - "legacy_minimal_sessions": 0, "mapped_by_layout": 0, "unrecognized_layout_rows": 0, "inferred_layouts": {}, + "legacy_minimal_sessions": 0, "mapped_by_layout": 0, "unrecognized_layout_rows": 0, + "unrecognized_layout_widths": {}, "inferred_layouts": {}, "unmapped_rows": 0, "insert_conflicts": 0, "lost_and_found_tables": [], } with _immediate_transaction(dest): @@ -644,30 +675,30 @@ def map_lost_and_found_rows(lf_conn: sqlite3.Connection, dest: sqlite3.Connectio ] report["lost_and_found_tables"] = lf_tables - # Pass 1: classify every record and group it by kind. The physical layout is a property of the - # whole population (one store wrote all of them), so it is inferred once per kind, not per row. - classified: dict[str, list[tuple[Any, int, tuple[Any, ...]]]] = {kind: [] for kind in targets} - for lf_table in lf_tables: - if _table_columns(lf_conn, lf_table)[:3] != ["rootpgno", "pgno", "nfield"]: - continue - for row in lf_conn.execute(f'SELECT * FROM "{lf_table}"'): - try: - nfield = int(row[2]) if row[2] is not None else 0 - except (TypeError, ValueError): - report["unmapped_rows"] += 1 + def records(): + """Yield (kind, lf_rowid, nfield, cells) for every classifiable lost_and_found row.""" + for lf_table in lf_tables: + if _table_columns(lf_conn, lf_table)[:3] != ["rootpgno", "pgno", "nfield"]: continue - lf_rowid = row[3] - cells = tuple(row[4 : 4 + max(nfield, 0)]) - kind = classify_lost_and_found_row(nfield, cells) - if kind is None: - report["unmapped_rows"] += 1 - continue - classified[kind].append((lf_rowid, nfield, cells)) + for row in lf_conn.execute(f'SELECT * FROM "{lf_table}"'): + try: + nfield = int(row[2]) if row[2] is not None else 0 + except (TypeError, ValueError): + yield None, None, 0, () + continue + cells = tuple(row[4 : 4 + max(nfield, 0)]) + yield classify_lost_and_found_row(nfield, cells), row[3], nfield, cells - layouts = { - kind: infer_physical_layouts(kind, [cells for _, _, cells in records], dest_types[kind]) - for kind, records in classified.items() - } + # Pass 1: stream the population once, keeping only what layout inference needs. The physical + # layout is a property of the whole population (one store wrote all of them), so it is inferred + # once per kind, not per row. + evidence = {kind: LayoutEvidence(kind) for kind in targets} + for kind, _, _, cells in records(): + if kind is None: + report["unmapped_rows"] += 1 + else: + evidence[kind].add(cells) + layouts = {kind: infer_physical_layouts(evidence[kind], dest_types[kind]) for kind in targets} # Per kind and record width, the column each position resolved to (None where the surviving # layouts disagreed and the cell was left to the destination default). report["inferred_layouts"] = { @@ -677,46 +708,50 @@ def map_lost_and_found_rows(lf_conn: sqlite3.Connection, dest: sqlite3.Connectio # Pass 2: insert. Records whose width resolved to a layout are mapped by column name (#101409); # the rest take the historical positional prefix, audited by the recovery verifier's plausibility gate. - for kind, records in classified.items(): + for kind, lf_rowid, nfield, cells in records(): + if kind is None: + continue # counted in pass 1 columns, defaults = targets[kind] - for lf_rowid, nfield, cells in records: - layout = layouts[kind].get(len(cells)) - legacy_minimal = kind == "sessions" and nfield == SESSIONS_LEGACY_MINIMAL_NFIELD - if layout is None and not legacy_minimal: - report["unrecognized_layout_rows"] += 1 - try: - if layout is not None: - # messages.id is a rowid alias: NULL in the record, carried by the lost_and_found row id. - inserted = _insert_named_row( - dest, kind, layout, cells, columns, defaults, - {"id": lf_rowid} if kind == "messages" else None, - ) - report["mapped_by_layout"] += int(inserted) - elif legacy_minimal: - # A 14-field record matching no known layout (torn cells, or a pre-history store): - # salvage identity + timing rather than guessing 14 positional meanings. - row_values = ( - cells[0], cells[1] if _looks_like_source(cells[1]) else "recovered", - _heuristic_started_at(cells), - f"{STUB_TITLE_PREFIX}] legacy session row (layout unknown)", - ) - inserted = dest.execute( - "INSERT OR IGNORE INTO sessions (id, source, started_at, title) VALUES (?, ?, ?, ?)", - row_values, - ).rowcount == 1 - report["legacy_minimal_sessions"] += int(inserted) - else: - values = list(cells[:len(columns)]) - if kind == "messages": - values[0] = lf_rowid - inserted = _insert_prefix_row(dest, kind, columns, values, defaults) - except sqlite3.DatabaseError: - report["unmapped_rows"] += 1 - continue - if inserted: - report["mapped"][kind] += 1 + layout = layouts[kind].get(len(cells)) + legacy_minimal = kind == "sessions" and nfield == SESSIONS_LEGACY_MINIMAL_NFIELD + if layout is None and not legacy_minimal: + report["unrecognized_layout_rows"] += 1 + widths = report["unrecognized_layout_widths"].setdefault(kind, []) + if len(cells) not in widths: + widths.append(len(cells)) + try: + if layout is not None: + # messages.id is a rowid alias: NULL in the record, carried by the lost_and_found row id. + inserted = _insert_named_row( + dest, kind, layout, cells, columns, defaults, + {"id": lf_rowid} if kind == "messages" else None, + ) + report["mapped_by_layout"] += int(inserted) + elif legacy_minimal: + # A 14-field record matching no known layout (torn cells, or a pre-history store): + # salvage identity + timing rather than guessing 14 positional meanings. + row_values = ( + cells[0], cells[1] if _looks_like_source(cells[1]) else "recovered", + _heuristic_started_at(cells), + f"{STUB_TITLE_PREFIX}] legacy session row (layout unknown)", + ) + inserted = dest.execute( + "INSERT OR IGNORE INTO sessions (id, source, started_at, title) VALUES (?, ?, ?, ?)", + row_values, + ).rowcount == 1 + report["legacy_minimal_sessions"] += int(inserted) else: - report["insert_conflicts"] += 1 + values = list(cells[:len(columns)]) + if kind == "messages": + values[0] = lf_rowid + inserted = _insert_prefix_row(dest, kind, columns, values, defaults) + except sqlite3.DatabaseError: + report["unmapped_rows"] += 1 + continue + if inserted: + report["mapped"][kind] += 1 + else: + report["insert_conflicts"] += 1 return report diff --git a/hermes_cli/session_recovery.py b/hermes_cli/session_recovery.py index 3d1f6ec208..73cf4645b1 100644 --- a/hermes_cli/session_recovery.py +++ b/hermes_cli/session_recovery.py @@ -1032,10 +1032,13 @@ def _recover_via_lost_and_found( # looks clean. guessed = int(mapping.get("unrecognized_layout_rows") or 0) if guessed: + widths = ", ".join( + f"{kind} rows with {'/'.join(map(str, sorted(ws)))} fields" + for kind, ws in sorted((mapping.get("unrecognized_layout_widths") or {}).items()) + ) plausibility_errors.append( - f"{guessed} salvaged row(s) matched no known physical column " - "layout and were mapped positionally; their fields may be " - "shifted. Inspect the affected sessions before trusting them." + f"{guessed} salvaged row(s) matched no known physical column layout and were mapped " + f"positionally ({widths}); their fields may be shifted. Inspect those sessions before trusting them." ) if plausibility_errors: verification["errors"].extend(plausibility_errors) diff --git a/hermes_cli/session_schema_history.py b/hermes_cli/session_schema_history.py index a792eb0da4..477b645abf 100644 --- a/hermes_cli/session_schema_history.py +++ b/hermes_cli/session_schema_history.py @@ -58,7 +58,7 @@ def declared_snapshots(table: str) -> list[tuple[str, ...]]: history = SCHEMA_HISTORY[table] columns = list(history.base) snapshots = [tuple(columns)] - for _date, edits in history.events: + for _label, edits in history.events: _apply(columns, edits) snapshots.append(tuple(columns)) return snapshots diff --git a/tests/hermes_cli/test_session_recovery_lost_and_found.py b/tests/hermes_cli/test_session_recovery_lost_and_found.py index fc34e55b87..0927ceaa15 100644 --- a/tests/hermes_cli/test_session_recovery_lost_and_found.py +++ b/tests/hermes_cli/test_session_recovery_lost_and_found.py @@ -1081,181 +1081,6 @@ def test_plausibility_gate_ignores_stub_only_sessions(tmp_path: Path) -> None: conn.close() -def _upgraded_store_layout(table: str, created_at: int, upgrades: tuple[int, ...]) -> tuple[str, ...]: - """Physical order of a store created at snapshot ``created_at`` and later - opened by the releases at ``upgrades`` (each appends the columns it had - not seen), exactly as ``_reconcile_columns`` does.""" - snapshots = session_schema_history.declared_snapshots(table) - columns = list(snapshots[created_at]) - for index in upgrades: - columns.extend(c for c in snapshots[index] if c not in columns) - return tuple(columns) - - -def _write_lost_and_found(lf_path: Path, records: list[tuple[int, tuple[str, ...], dict]]) -> None: - width = max(len(layout) for _, layout, _ in records) - conn = sqlite3.connect(str(lf_path), isolation_level=None) - try: - columns = ", ".join(f"c{index}" for index in range(width)) - conn.execute( - "CREATE TABLE lost_and_found (rootpgno INTEGER, pgno INTEGER, " - "nfield INTEGER, id INTEGER, " + columns + ")" - ) - for rowid, layout, values in records: - cells = [values.get(name) for name in layout] - cells += [None] * (width - len(cells)) - placeholders = ", ".join("?" for _ in range(4 + width)) - conn.execute( - "INSERT INTO lost_and_found VALUES (" + placeholders + ")", - [2, 5, len(layout), rowid, *cells], - ) - finally: - conn.close() - - -# (creation snapshot, later upgrade snapshots) — a spread of real histories: -# a Feb-2026 v1 store upgraded monthly, a May store upgraded once to current, -# a July store, and one created at the current schema. -_UPGRADE_PATHS = [ - (0, (2, 5, 9, 12, 16, 20, 26)), - (5, (26,)), - (12, (18, 26)), - (20, (23, 26)), - (26, ()), -] - - -@pytest.mark.parametrize("created_at,upgrades", _UPGRADE_PATHS) -def test_any_upgrade_history_maps_cells_by_name( - tmp_path: Path, created_at: int, upgrades: tuple[int, ...] -) -> None: - """#101409 generalised: whatever release a store was created at and - however it was upgraded since, salvaged cells land on the columns they - were written from. The layout is inferred from the schema history plus - the salvaged values, not looked up in a hardcoded table.""" - - sessions_layout = _upgraded_store_layout("sessions", created_at, upgrades) - # messages/usage histories are shorter; pick the snapshot by date parity. - n_messages = len(session_schema_history.declared_snapshots("messages")) - n_usage = len(session_schema_history.declared_snapshots("session_model_usage")) - messages_layout = _upgraded_store_layout( - "messages", min(created_at // 2, n_messages - 1), - tuple(min(u // 2, n_messages - 1) for u in upgrades), - ) - usage_layout = _upgraded_store_layout( - "session_model_usage", min(created_at // 10, n_usage - 1), - tuple(min(u // 10, n_usage - 1) for u in upgrades), - ) - - output = tmp_path / "mapped.db" - SessionDB(db_path=output).close() - declared = sqlite3.connect(str(output)) - try: - sessions_declared = [r[1] for r in declared.execute("PRAGMA table_info(sessions)")] - finally: - declared.close() - shifted = sessions_declared[: len(sessions_layout)] != list(sessions_layout) - - started = 1_754_000_000.0 - records = [] - for index in range(3): - sid = f"2026070{index + 1}_101010_abc00{index}" - records.append((index + 1, sessions_layout, { - "id": sid, "source": "telegram", "model": "gpt-4.1", - "started_at": started + index, "ended_at": started + 600 + index, - "end_reason": "completed", "message_count": 4, "tool_call_count": 1, - "input_tokens": 100 + index, "output_tokens": 50, - "title": f"planning notes {index}", "cwd": "/home/user/project", - "archived": 0, "pinned": 0, "rewind_count": 0, "hidden": 0, - "api_call_count": 2, "expiry_finalized": 0, - })) - for m in range(2): - records.append((100 + index * 2 + m, messages_layout, { - "id": None, "session_id": sid, "role": "user" if m == 0 else "assistant", - "content": f"payload {index} {m}", "timestamp": started + 10 + m, - "token_count": 12, "finish_reason": "stop" if m else None, - "observed": 1, "active": 1, "compacted": 0, "_compressed_summary": 0, - })) - records.append((200 + index, usage_layout, { - "session_id": sid, "model": "gpt-4.1", "billing_provider": "openai", - "billing_base_url": "https://api.openai.com/v1", "billing_mode": "api", - "task": "", "api_call_count": 2, "input_tokens": 100 + index, - "output_tokens": 50, "cache_read_tokens": 0, "cache_write_tokens": 0, - "reasoning_tokens": 0, "estimated_cost_usd": 0.01, "actual_cost_usd": 0.0, - "first_seen": started, "last_seen": started + 600, - })) - lf_path = tmp_path / "lost_and_found.db" - _write_lost_and_found(lf_path, records) - - lf_conn = sqlite3.connect(str(lf_path), isolation_level=None) - dest = sqlite3.connect(str(output), isolation_level=None) - try: - dest.execute("PRAGMA foreign_keys=OFF") - report = map_lost_and_found_rows(lf_conn, dest) - assert report["mapped"] == {"sessions": 3, "messages": 6, "session_model_usage": 3} - assert report["unrecognized_layout_rows"] == 0 - rows = dest.execute( - "SELECT id, source, model, started_at, ended_at, message_count, title, cwd " - "FROM sessions ORDER BY id" - ).fetchall() - for index, row in enumerate(rows): - assert row == ( - f"2026070{index + 1}_101010_abc00{index}", "telegram", "gpt-4.1", - started + index, started + 600 + index, 4, - f"planning notes {index}", "/home/user/project", - ), (shifted, sessions_layout) - messages = dest.execute( - "SELECT role, content, timestamp, token_count, effect_disposition " - "FROM messages ORDER BY id" - ).fetchall() - assert [m[0] for m in messages] == ["user", "assistant"] * 3 - assert all(m[2] >= started for m in messages) - assert all(m[3] == 12 and m[4] is None for m in messages) - usage = dest.execute( - "SELECT model, billing_mode, task, api_call_count, input_tokens " - "FROM session_model_usage ORDER BY session_id" - ).fetchall() - assert usage == [("gpt-4.1", "api", "", 2, 100 + i) for i in range(3)] - assert session_recovery._lost_and_found_plausibility_errors(dest) == [] - finally: - lf_conn.close() - dest.close() - - -def test_unknown_layout_rows_fall_back_positionally_and_are_counted(tmp_path: Path) -> None: - """Records no known history produced (here: a sessions row whose sentinel - cells contradict every chain) take the positional guess and are counted - so the verifier can refuse to call the salvage verified.""" - - output = tmp_path / "mapped.db" - SessionDB(db_path=output).close() - declared = sqlite3.connect(str(output)) - try: - sessions_declared = tuple(r[1] for r in declared.execute("PRAGMA table_info(sessions)")) - finally: - declared.close() - # A 40-wide record with a text value where EVERY chain state of width 40 - # has a numeric column: no layout fits. - bogus = list(sessions_declared[:40]) - values = {"id": "20260701_101010_abc000", "source": "cli", "started_at": 1_754_000_000.0} - for name in bogus[2:]: - values[name] = "not-a-number" - lf_path = tmp_path / "lost_and_found.db" - _write_lost_and_found(lf_path, [(1, tuple(bogus), values)]) - - lf_conn = sqlite3.connect(str(lf_path), isolation_level=None) - dest = sqlite3.connect(str(output), isolation_level=None) - try: - dest.execute("PRAGMA foreign_keys=OFF") - report = map_lost_and_found_rows(lf_conn, dest) - finally: - lf_conn.close() - dest.close() - assert report["unrecognized_layout_rows"] == 1 - assert report["mapped_by_layout"] == 0 - assert report["inferred_layouts"] == {} - - @pytest.mark.skipif( not HAVE_SQLITE3_CLI, reason="sqlite3 CLI not on PATH; .recover is a shell-only feature", diff --git a/tests/hermes_cli/test_session_schema_history.py b/tests/hermes_cli/test_session_schema_history.py index 7c4562ba20..98483af4d1 100644 --- a/tests/hermes_cli/test_session_schema_history.py +++ b/tests/hermes_cli/test_session_schema_history.py @@ -34,6 +34,11 @@ def test_replayed_history_ends_at_current_schema(table: str) -> None: older events — real stores were shaped by them). """ + labels = [int(label.split()[0]) for label, _ in history.SCHEMA_HISTORY[table].events] + assert labels == list(range(1, len(labels) + 1)), ( + f"SCHEMA_HISTORY[{table!r}].events is out of order: append new events at the END with the next " + f"sequence number, never insert mid-list (got {labels})" + ) replayed = history.current_declared_columns(table) declared = _declared_now(table) missing = [c for c in declared if c not in replayed] @@ -52,72 +57,3 @@ def test_replayed_history_ends_at_current_schema(table: str) -> None: f"{next(i for i, (a, b) in enumerate(zip(replayed, declared)) if a != b)}: " f"replayed {replayed} vs SCHEMA_SQL {declared} — record the move as ('-', c) + ('+', c, after) in {hint}" ) - - -@pytest.mark.parametrize("table", sorted(history.SCHEMA_HISTORY)) -def test_snapshots_are_distinct_and_replay_is_well_formed(table: str) -> None: - snapshots = history.declared_snapshots(table) - assert len(snapshots) == len(history.SCHEMA_HISTORY[table].events) + 1 - for snapshot in snapshots: - assert len(set(snapshot)) == len(snapshot), "duplicate column in a snapshot" - # Snapshots may repeat (a column was declared, reverted, then declared - # again) but consecutive ones must differ or the event was a no-op. - for before, after in zip(snapshots, snapshots[1:]): - assert before != after, "an event changed nothing" - sequence = [int(label.split()[0]) for label, _ in history.SCHEMA_HISTORY[table].events] - assert sequence == list(range(1, len(sequence) + 1)), ( - "event labels must be numbered 01, 02, ... in replay order: append new " - "events at the END with the next number, never insert mid-list" - ) - - -def test_reachable_layouts_include_the_upgraded_store_from_101409() -> None: - """The reporter's physical order (created ~May 2026, upgraded through - v27) is a chain over the recorded snapshots.""" - - reported = ( - "id", "source", "user_id", "model", "model_config", "system_prompt", - "parent_session_id", "started_at", "ended_at", "end_reason", - "message_count", "tool_call_count", "input_tokens", "output_tokens", - "cache_read_tokens", "cache_write_tokens", "reasoning_tokens", - "billing_provider", "billing_base_url", "billing_mode", - "estimated_cost_usd", "actual_cost_usd", "cost_status", "cost_source", - "pricing_version", "title", "api_call_count", "handoff_state", - "handoff_platform", "handoff_error", "cwd", "rewind_count", "archived", - "session_key", "chat_id", "chat_type", "thread_id", "git_branch", - "git_repo_root", "compression_failure_cooldown_until", - "compression_failure_error", "display_name", "origin_json", - "expiry_finalized", "compression_fallback_streak", "profile_name", - "compression_ineffective_count", "pinned", "system_prompt_hash", - "last_activity_at", "last_activity_description", - "last_activity_provenance", "git_metadata_generation", "title_source", - "hidden", "last_read_at", - ) - - def accept(layout: tuple[str, ...], first_new: int) -> bool: - return layout[first_new:len(reported)] == reported[first_new:len(layout)] - - assert any( - layout[: len(reported)] == reported - for layout in history.reachable_physical_layouts("sessions", accept) - ) - - -def test_pruning_callback_stops_branches() -> None: - """``accept`` returning False must prune, not just filter the output.""" - - offered: list[tuple[tuple[str, ...], int]] = [] - - def accept(layout: tuple[str, ...], first_new: int) -> bool: - offered.append((layout, first_new)) - return len(layout) <= 10 - - layouts = list(history.reachable_physical_layouts("messages", accept)) - assert layouts and all(len(layout) <= 10 for layout in layouts) - # Every extension offered must grow from an ACCEPTED parent: its prefix - # up to first_new is a layout that passed. Seeds (first_new == 0) are - # always offered. - accepted = set(layouts) - for layout, first_new in offered: - if first_new: - assert layout[:first_new] in accepted, "a rejected layout was extended"