refactor(recovery): stream the salvaged population; table the shape rules
- 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.
This commit is contained in:
@@ -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
|
||||
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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"
|
||||
|
||||
Reference in New Issue
Block a user