diff --git a/docs/evidence/polylogue-excluded-cursor-live-proof-2026-08-06.json b/docs/evidence/polylogue-excluded-cursor-live-proof-2026-08-06.json new file mode 100644 index 0000000000..ecfa0d281e --- /dev/null +++ b/docs/evidence/polylogue-excluded-cursor-live-proof-2026-08-06.json @@ -0,0 +1 @@ +{"anti_vacuity":{"indexed_authority":"byte_proven_source_raw_and_revision_head","indexed_session_count":1,"indexed_session_count_before":0,"typed_terminal_artifact":"terminal_corrupt_input","unchanged_excluded_attempt_present":false},"cases":[{"attempt":{"evidence_ref":null,"outcome_code":"success","retryable":false,"status":"completed"},"attempt_present":true,"case_id":"indexed","fingerprint_changed_before_catch_up":true,"indexed":{"indexed_sessions":1,"parsed_raw":2},"indexed_before":{"indexed_sessions":0,"parsed_raw":1},"metrics":{"failed_file_count":0,"full_file_count":1,"succeeded_file_count":1},"proof_attempt_count":1,"retry_state":{"excluded":false,"failed_with_retry":false,"failure_count":0,"parser_fingerprint":"live-batched-v2","retry_due":false},"source_content_sha256":"fbe00b1bf4fe9b647b143e5f002a8edfd9d3685944fa44125586a679c4cc6534","terminal_evidence":null},{"attempt":null,"attempt_present":false,"case_id":"still-excluded","fingerprint_changed_before_catch_up":false,"indexed":{"indexed_sessions":0,"parsed_raw":0},"metrics":{"failed_file_count":0,"full_file_count":0,"succeeded_file_count":0},"proof_attempt_count":0,"retry_state":{"excluded":true,"failed_with_retry":false,"failure_count":5,"parser_fingerprint":"live-batched-v2","retry_due":false},"source_content_sha256":"cb7b4873a5dc05a7cac589c60a8069dc9f91f234af3076f8ea558bfafe69612a","terminal_evidence":null},{"attempt":{"evidence_ref":null,"outcome_code":"success","retryable":false,"status":"completed"},"attempt_present":true,"case_id":"typed-terminal","fingerprint_changed_before_catch_up":true,"indexed":{"indexed_sessions":0,"parsed_raw":0},"metrics":{"failed_file_count":0,"full_file_count":1,"succeeded_file_count":1},"proof_attempt_count":1,"retry_state":{"excluded":false,"failed_with_retry":false,"failure_count":0,"parser_fingerprint":"live-batched-v2","retry_due":false},"source_content_sha256":"2e5c27f8ae0c2f4176892776a5a0ac5e1d598178676396f86fa4f6d608778d75","terminal_evidence":{"artifact_kind":"terminal_corrupt_input","parse_error_present":true,"support_status":"decode_failed"}}],"execution":{"live_census":"not_run","live_residual":"Historical excluded population and current live file states were not accessed.","mode":"candidate_fixture","residual_successor":"polylogue-excluded-cursor-live-proof","terminal_frontier_residual":"The typed-terminal candidate has no accepted byte head, so its readiness gate was injected for this case only."},"fairness":{"planner":"_interleave_by_source","property":"browser-capture drains first; among non-browser-capture families, one candidate from each present family reaches the first round"},"fixture_version":"candidate-codex-live-compatible-2026-08-06","outcomes":{"indexed":true,"still_excluded":true,"typed_terminal":true},"production_route":{"catch_up":"LiveWatcher._catch_up -> _scan_catch_up_candidates -> _catch_up_candidates -> _plan_catch_up -> coordinated chunk ingest","cursor_gate":"LiveWatcher._needs_work","ingest":"LiveWatcher._ingest_files -> LiveBatchProcessor.ingest_files","retry_state":"ops.ingest_cursor and ops.ingest_attempts","terminal_evidence":"source.raw_artifacts","transition":"CursorStore.revive_replaced_exclusion"},"receipt_sha256":"69e75004783fc3f0af38b85ef01136ac8ece609db850b730773b42b56589e350","schema":"polylogue.excluded-cursor-live-proof.v1"} diff --git a/docs/plans/reindex-incident-coverage.json b/docs/plans/reindex-incident-coverage.json index eedaee9813..c675861092 100644 --- a/docs/plans/reindex-incident-coverage.json +++ b/docs/plans/reindex-incident-coverage.json @@ -16,7 +16,8 @@ "lineage-corpus": {"kind": "lineage-corpus", "source": "tests/infra/reindex_campaign.py"}, "derived-model": {"kind": "derived-model", "source": "tests/infra/reindex_differential.py"}, "title-census": {"kind": "title-census", "source": "tests/unit/maintenance/test_reindex_campaign.py"}, - "parser-replay": {"kind": "parser-replay", "source": "tests/unit/maintenance/test_reindex_campaign.py"} + "parser-replay": {"kind": "parser-replay", "source": "tests/unit/maintenance/test_reindex_campaign.py"}, + "excluded-cursor-proof": {"kind": "candidate-live-compatible", "source": "tests/infra/excluded_cursor_live_proof.py"} }, "checks": { "campaign-coverage": {"kind": "registry", "source": "polylogue-incident-coverage-ledger"}, @@ -57,7 +58,8 @@ "live-proof-7zp4": {"kind": "live-proof", "status": "recorded", "source": "content-hash-census"}, "live-proof-gzgyl": {"kind": "live-proof", "status": "recorded", "source": "material-origin-census"}, "live-proof-mvcbi": {"kind": "live-proof", "status": "recorded", "source": "origin-dispatch-census"}, - "live-proof-qsagp": {"kind": "live-proof", "status": "recorded", "source": "derived-refresh-census"} + "live-proof-qsagp": {"kind": "live-proof", "status": "recorded", "source": "derived-refresh-census"}, + "excluded-cursor-proof-receipt": {"kind": "proof-receipt", "status": "recorded", "source": "docs/evidence/polylogue-excluded-cursor-live-proof-2026-08-06.json"} }, "successors": { "polylogue-claude-vintage-live-proof": {"kind": "named-child-bead"}, @@ -92,7 +94,7 @@ {"bead_id": "polylogue-fsgdd", "bead_status": "open", "incident": {"incident_id": "incident-fsgdd", "bead_id": "polylogue-fsgdd", "forcing_class": "coverage"}, "route": {"kind": "registry", "entrypoint": "polylogue-incident-coverage-ledger"}, "schedule": {"phase": "preflight", "order": 20}, "expected_snapshot": {"snapshot_id": "post-reindex-acceptance", "state": "blocking"}, "registry_checks": ["campaign-coverage"], "red_mutation": {"fixture_id": "campaign-graph", "mutation_id": "mutation-fsgdd"}, "receipts": [], "residual_successor": null}, {"bead_id": "polylogue-gzgyl", "bead_status": "closed", "incident": {"incident_id": "incident-gzgyl", "bead_id": "polylogue-gzgyl", "forcing_class": "material-origin"}, "route": {"kind": "registry", "entrypoint": "origin-capability"}, "schedule": {"phase": "preflight", "order": 21}, "expected_snapshot": {"snapshot_id": "post-reindex-acceptance", "state": "blocking"}, "registry_checks": ["material-origin"], "red_mutation": {"fixture_id": "origin-matrix", "mutation_id": "mutation-gzgyl"}, "receipts": ["live-proof-gzgyl"], "residual_successor": null}, {"bead_id": "polylogue-ih67", "bead_status": "in_progress", "incident": {"incident_id": "incident-ih67", "bead_id": "polylogue-ih67", "forcing_class": "title"}, "route": {"kind": "registry", "entrypoint": "reindex-campaign"}, "schedule": {"phase": "preflight", "order": 22}, "expected_snapshot": {"snapshot_id": "derived-model-candidate", "state": "blocking"}, "registry_checks": ["title-resolution"], "red_mutation": {"fixture_id": "title-census", "mutation_id": "mutation-ih67"}, "receipts": [], "residual_successor": null}, - {"bead_id": "polylogue-ix5r", "bead_status": "closed", "incident": {"incident_id": "incident-ix5r", "bead_id": "polylogue-ix5r", "forcing_class": "cursor"}, "route": {"kind": "registry", "entrypoint": "reindex-campaign"}, "schedule": {"phase": "preflight", "order": 23}, "expected_snapshot": {"snapshot_id": "live-preflight-2026-08-04", "state": "blocking"}, "registry_checks": ["cursor-freshness"], "red_mutation": {"fixture_id": "campaign-corpus", "mutation_id": "mutation-ix5r"}, "receipts": [], "residual_successor": {"bead_id": "polylogue-excluded-cursor-live-proof", "kind": "live-proof"}}, + {"bead_id": "polylogue-ix5r", "bead_status": "closed", "incident": {"incident_id": "incident-ix5r", "bead_id": "polylogue-ix5r", "forcing_class": "cursor"}, "route": {"kind": "registry", "entrypoint": "reindex-campaign"}, "schedule": {"phase": "preflight", "order": 23}, "expected_snapshot": {"snapshot_id": "live-preflight-2026-08-04", "state": "blocking"}, "registry_checks": ["cursor-freshness"], "red_mutation": {"fixture_id": "excluded-cursor-proof", "mutation_id": "mutation-ix5r"}, "receipts": ["excluded-cursor-proof-receipt"], "residual_successor": {"bead_id": "polylogue-excluded-cursor-live-proof", "kind": "live-proof"}}, {"bead_id": "polylogue-mvcbi", "bead_status": "closed", "incident": {"incident_id": "incident-mvcbi", "bead_id": "polylogue-mvcbi", "forcing_class": "origin"}, "route": {"kind": "registry", "entrypoint": "origin-capability"}, "schedule": {"phase": "preflight", "order": 24}, "expected_snapshot": {"snapshot_id": "post-reindex-acceptance", "state": "blocking"}, "registry_checks": ["origin-matrix"], "red_mutation": {"fixture_id": "origin-matrix", "mutation_id": "mutation-mvcbi"}, "receipts": ["live-proof-mvcbi"], "residual_successor": null}, {"bead_id": "polylogue-nas1", "bead_status": "open", "incident": {"incident_id": "incident-nas1", "bead_id": "polylogue-nas1", "forcing_class": "lineage"}, "route": {"kind": "registry", "entrypoint": "reindex-differential"}, "schedule": {"phase": "preflight", "order": 25}, "expected_snapshot": {"snapshot_id": "derived-model-candidate", "state": "blocking"}, "registry_checks": ["lineage-differential"], "red_mutation": {"fixture_id": "lineage-corpus", "mutation_id": "mutation-nas1"}, "receipts": [], "residual_successor": null}, {"bead_id": "polylogue-omsw", "bead_status": "open", "incident": {"incident_id": "incident-omsw", "bead_id": "polylogue-omsw", "forcing_class": "sidecar"}, "route": {"kind": "registry", "entrypoint": "reindex-campaign"}, "schedule": {"phase": "preflight", "order": 26}, "expected_snapshot": {"snapshot_id": "live-preflight-2026-08-04", "state": "blocking"}, "registry_checks": ["sidecar-admission"], "red_mutation": {"fixture_id": "sidecar-admission", "mutation_id": "mutation-omsw"}, "receipts": [], "residual_successor": null}, diff --git a/tests/infra/excluded_cursor_live_proof.py b/tests/infra/excluded_cursor_live_proof.py new file mode 100644 index 0000000000..3a51688966 --- /dev/null +++ b/tests/infra/excluded_cursor_live_proof.py @@ -0,0 +1,441 @@ +"""Proof harness for parser-fingerprint revival of excluded live cursors. + +The harness uses the production watcher and live batch processor against a +candidate archive fixture. It deliberately reports candidate-fixture +coverage separately from a live census. +""" + +from __future__ import annotations + +import asyncio +import json +import os +import sqlite3 +from contextlib import nullcontext +from hashlib import sha256 +from pathlib import Path +from types import SimpleNamespace +from typing import Any, cast +from unittest.mock import patch + +from polylogue.archive.revision_authority import RawRevisionAuthority, RawRevisionEnvelope, RawRevisionKind +from polylogue.core.enums import Provider +from polylogue.pipeline.ids import session_content_hash, session_id +from polylogue.sources.dispatch import parse_payload +from polylogue.sources.live.cursor import CursorStore +from polylogue.sources.live.watcher import LiveWatcher, WatchSource +from polylogue.storage.sqlite.archive_tiers.archive import ArchiveStore +from polylogue.storage.sqlite.archive_tiers.bootstrap import initialize_active_archive_root +from tests.infra.reindex_campaign import _codex_records, _write_jsonl + +RECEIPT_SCHEMA = "polylogue.excluded-cursor-live-proof.v1" +FIXTURE_VERSION = "candidate-codex-live-compatible-2026-08-06" +OLD_PARSER_FINGERPRINT = "live-batched-v1" +NEW_PARSER_FINGERPRINT = "live-batched-v2" + + +def _canonical_json(payload: object) -> bytes: + return (json.dumps(payload, sort_keys=True, separators=(",", ":")) + "\n").encode("utf-8") + + +def _sha256_file(path: Path) -> str: + digest = sha256() + with path.open("rb") as handle: + for chunk in iter(lambda: handle.read(1024 * 1024), b""): + digest.update(chunk) + return digest.hexdigest() + + +def _seed_excluded(cursor: CursorStore, path: Path, *, parser_fingerprint: str) -> None: + stat = path.stat() + cursor.set( + path, + stat.st_size, + byte_offset=stat.st_size, + last_complete_newline=stat.st_size, + parser_fingerprint=parser_fingerprint, + content_fingerprint=_sha256_file(path), + source_name="codex", + st_dev=stat.st_dev, + st_ino=stat.st_ino, + mtime_ns=stat.st_mtime_ns, + failure_count=5, + excluded=True, + ) + + +def _seed_byte_authority(root: Path, path: Path, *, native_id: str) -> None: + """Seed source evidence plus a byte head without materializing a session.""" + payload = path.read_bytes() + logical_source_key = f"codex:{native_id}" + [parsed] = parse_payload( + Provider.CODEX, + [json.loads(line) for line in payload.splitlines()], + native_id, + source_path=str(path), + ) + source_revision = "excluded-cursor-proof-authority-0" + accepted_content_hash = bytes.fromhex(session_content_hash(parsed)) + accepted_session_id = str(session_id(parsed.source_name, parsed.provider_session_id)) + revision = RawRevisionEnvelope( + logical_source_key, + RawRevisionKind.FULL, + source_revision, + 0, + authority=RawRevisionAuthority.BYTE_PROVEN, + ) + with ArchiveStore.open_existing(root, read_only=False) as archive: + raw_id = archive.write_raw_payload( + provider=Provider.CODEX, + payload=payload, + source_path=str(path), + acquired_at_ms=1, + native_id=native_id, + revision=revision, + ) + conn = sqlite3.connect(root / "index.db") + try: + conn.execute( + """ + INSERT INTO raw_revision_heads ( + logical_source_key, session_id, accepted_raw_id, + accepted_source_revision, accepted_content_hash, + accepted_frontier_kind, accepted_frontier, + acquisition_generation, append_end_offset, decided_at_ms + ) VALUES (?, ?, ?, ?, ?, 'byte', ?, 0, NULL, 1) + """, + ( + logical_source_key, + accepted_session_id, + raw_id, + source_revision, + accepted_content_hash, + len(payload), + ), + ) + conn.commit() + finally: + conn.close() + + +def _attempts_for_path(root: Path, path: Path) -> list[dict[str, object]]: + conn = sqlite3.connect(root / "ops.db") + try: + rows = conn.execute( + """ + SELECT outcome_code, retryable, evidence_ref, status, source_paths_json + FROM ingest_attempts + ORDER BY started_at_ms DESC, attempt_id DESC + """ + ).fetchall() + finally: + conn.close() + attempts: list[dict[str, object]] = [] + for outcome_code, retryable, evidence_ref, status, source_paths_json in rows: + try: + source_paths = json.loads(str(source_paths_json or "[]")) + except json.JSONDecodeError: + source_paths = [] + if str(path) not in source_paths: + continue + attempts.append( + { + "outcome_code": outcome_code, + "retryable": None if retryable is None else bool(retryable), + "evidence_ref": evidence_ref, + "status": status, + } + ) + return attempts + + +def _retry_state(cursor: CursorStore, path: Path) -> dict[str, object]: + record = cursor.get_record(path) + if record is None: + raise AssertionError(f"proof cursor disappeared for {path}") + retry_paths = {record.source_path for record in cursor.list_retry_records()} + failed_paths = set(cursor.list_failed_with_retry()) + return { + "excluded": bool(record.excluded), + "failure_count": record.failure_count, + "retry_due": record.source_path in retry_paths, + "failed_with_retry": record.source_path in failed_paths, + "parser_fingerprint": record.parser_fingerprint, + } + + +def _indexed_counts(root: Path, path: Path) -> dict[str, int]: + with sqlite3.connect(root / "source.db") as source_conn: + raw_rows = source_conn.execute( + "SELECT raw_id FROM raw_sessions WHERE source_path = ? AND parse_error IS NULL", + (str(path),), + ).fetchall() + raw_ids = tuple(str(row[0]) for row in raw_rows) + if not raw_ids: + return {"parsed_raw": 0, "indexed_sessions": 0} + placeholders = ",".join("?" for _ in raw_ids) + with sqlite3.connect(root / "index.db") as index_conn: + indexed = index_conn.execute( + f"SELECT COUNT(*) FROM sessions WHERE raw_id IN ({placeholders})", + raw_ids, + ).fetchone() + return {"parsed_raw": len(raw_ids), "indexed_sessions": int(indexed[0]) if indexed else 0} + + +def _terminal_evidence(root: Path, path: Path) -> dict[str, object] | None: + with sqlite3.connect(root / "source.db") as conn: + row = conn.execute( + """ + SELECT a.artifact_kind, a.support_status, r.parse_error + FROM raw_artifacts AS a + JOIN raw_sessions AS r USING (raw_id) + WHERE r.source_path = ? AND a.artifact_kind LIKE 'terminal_%' + ORDER BY r.acquired_at_ms DESC, r.raw_id DESC + LIMIT 1 + """, + (str(path),), + ).fetchone() + if row is None: + return None + return { + "artifact_kind": str(row[0]), + "support_status": str(row[1]), + "parse_error_present": row[2] is not None, + } + + +def _case_summary( + *, + case_id: str, + path: Path, + cursor: CursorStore, + fingerprint_changed_before_catch_up: bool, + metrics: object | None, + root: Path, + attempts_before: int, +) -> dict[str, object]: + retry_state = _retry_state(cursor, path) + attempts = _attempts_for_path(root, path) + proof_attempts = attempts[: max(0, len(attempts) - attempts_before)] + attempt = proof_attempts[0] if proof_attempts else None + return { + "case_id": case_id, + "source_content_sha256": _sha256_file(path), + "fingerprint_changed_before_catch_up": fingerprint_changed_before_catch_up, + "metrics": { + "succeeded_file_count": int(getattr(metrics, "succeeded_file_count", 0)) if metrics else 0, + "failed_file_count": int(getattr(metrics, "failed_file_count", 0)) if metrics else 0, + "full_file_count": int(getattr(metrics, "full_file_count", 0)) if metrics else 0, + }, + "indexed": _indexed_counts(root, path), + "terminal_evidence": _terminal_evidence(root, path), + "attempt": attempt, + "proof_attempt_count": len(proof_attempts), + "retry_state": retry_state, + "attempt_present": attempt is not None, + } + + +def _run_case( + *, + root: Path, + source_root: Path, + cursor: CursorStore, + case_id: str, + path: Path, + parser_fingerprint: str, + attempts_before: int, + bypass_frontier_gate: bool = False, +) -> dict[str, Any]: + polylogue = SimpleNamespace(archive_root=root, backend=SimpleNamespace(db_path=root / "index.db")) + watcher = LiveWatcher(cast(Any, polylogue), (WatchSource(name="codex", root=source_root),), cursor=cursor) + try: + record = cursor.get_record(path) + fingerprint_changed_before_catch_up = ( + record is not None and bool(record.excluded) and record.parser_fingerprint != parser_fingerprint + ) + metrics_holder: list[object] = [] + original_ingest = watcher._ingest_files + + async def capture_ingest(*args: Any, **kwargs: Any) -> object: + metrics = await original_ingest(*args, **kwargs) + metrics_holder.append(metrics) + return metrics + + with patch("polylogue.sources.live.watcher._PARSER_FINGERPRINT", parser_fingerprint): + frontier_patch = ( + patch("polylogue.readiness.capability.raw_frontier_source_selection_block_reason", lambda _root: None) + if bypass_frontier_gate + else nullcontext() + ) + with frontier_patch, patch.object(watcher, "_ingest_files", capture_ingest): + asyncio.run(watcher._catch_up([source_root])) + finally: + watcher.stop() + return _case_summary( + case_id=case_id, + path=path, + cursor=cursor, + fingerprint_changed_before_catch_up=fingerprint_changed_before_catch_up, + metrics=metrics_holder[-1] if metrics_holder else None, + root=root, + attempts_before=attempts_before, + ) + + +def run_excluded_cursor_live_proof(root: Path, receipt_path: Path) -> dict[str, Any]: + """Run the real cursor/fingerprint route and write a self-hashed receipt.""" + case_roots = {case_id: root / case_id for case_id in ("indexed", "still-excluded", "typed-terminal")} + + def prepare_case(case_id: str, native_id: str, texts: tuple[str, ...]) -> tuple[Path, Path, CursorStore, int]: + case_root = case_roots[case_id] + source_root = case_root / "wire" / "excluded-cursor-proof" + path = _write_jsonl(source_root / f"{case_id}.jsonl", _codex_records(native_id, texts)) + initialize_active_archive_root(case_root) + cursor = CursorStore(case_root / "ops.db") + return case_root, source_root, cursor, len(_attempts_for_path(case_root, path)) + + indexed_root, indexed_source_root, indexed_cursor, indexed_attempts_before = prepare_case( + "indexed", "excluded-proof-indexed", ("revived", "indexed") + ) + indexed_path = indexed_source_root / "indexed.jsonl" + _seed_byte_authority(indexed_root, indexed_path, native_id="excluded-proof-indexed") + indexed_before = _indexed_counts(indexed_root, indexed_path) + if indexed_before["indexed_sessions"] != 0: + raise AssertionError(f"indexed case was not empty before catch-up: {indexed_before}") + _seed_excluded(indexed_cursor, indexed_path, parser_fingerprint=OLD_PARSER_FINGERPRINT) + indexed = _run_case( + root=indexed_root, + source_root=indexed_source_root, + cursor=indexed_cursor, + case_id="indexed", + path=indexed_path, + parser_fingerprint=NEW_PARSER_FINGERPRINT, + attempts_before=indexed_attempts_before, + ) + indexed["indexed_before"] = indexed_before + + unchanged_root, unchanged_source_root, unchanged_cursor, unchanged_attempts_before = prepare_case( + "still-excluded", "excluded-proof-still-excluded", ("unchanged", "poison") + ) + unchanged_path = unchanged_source_root / "still-excluded.jsonl" + _seed_excluded(unchanged_cursor, unchanged_path, parser_fingerprint=NEW_PARSER_FINGERPRINT) + still_excluded = _run_case( + root=unchanged_root, + source_root=unchanged_source_root, + cursor=unchanged_cursor, + case_id="still-excluded", + path=unchanged_path, + parser_fingerprint=NEW_PARSER_FINGERPRINT, + attempts_before=unchanged_attempts_before, + ) + + terminal_root = case_roots["typed-terminal"] + terminal_source_root = terminal_root / "wire" / "excluded-cursor-proof" + terminal_path = terminal_source_root / "typed-terminal.jsonl" + _write_jsonl( + terminal_path, + _codex_records("excluded-proof-terminal", ("valid prefix", "terminal corruption")), + ) + with terminal_path.open("ab") as handle: + handle.write(b'{"type":"response_item","payload":{"type":"message","content":[') + initialize_active_archive_root(terminal_root) + terminal_cursor = CursorStore(terminal_root / "ops.db") + terminal_attempts_before = len(_attempts_for_path(terminal_root, terminal_path)) + _seed_excluded(terminal_cursor, terminal_path, parser_fingerprint=OLD_PARSER_FINGERPRINT) + typed_terminal = _run_case( + root=terminal_root, + source_root=terminal_source_root, + cursor=terminal_cursor, + case_id="typed-terminal", + path=terminal_path, + parser_fingerprint=NEW_PARSER_FINGERPRINT, + attempts_before=terminal_attempts_before, + bypass_frontier_gate=True, + ) + + cases = [indexed, still_excluded, typed_terminal] + indexed_attempt = indexed["attempt"] + terminal_evidence = typed_terminal["terminal_evidence"] + outcomes = { + "indexed": indexed["indexed_before"]["indexed_sessions"] == 0 + and indexed["indexed"]["indexed_sessions"] == 1 + and indexed["retry_state"]["excluded"] is False + and indexed_attempt is not None + and indexed_attempt["outcome_code"] == "success", + "still_excluded": still_excluded["retry_state"]["excluded"] is True + and still_excluded["attempt_present"] is False + and still_excluded["retry_state"]["retry_due"] is False, + "typed_terminal": terminal_evidence is not None + and terminal_evidence["artifact_kind"] == "terminal_corrupt_input" + and terminal_evidence["support_status"] == "decode_failed" + and terminal_evidence["parse_error_present"] is True + and typed_terminal["retry_state"]["excluded"] is False + and typed_terminal["retry_state"]["retry_due"] is False, + } + if not all(outcomes.values()): + raise AssertionError(f"excluded-cursor proof outcomes failed: {outcomes}") + + body: dict[str, object] = { + "schema": RECEIPT_SCHEMA, + "fixture_version": FIXTURE_VERSION, + "execution": { + "mode": "candidate_fixture", + "live_census": "not_run", + "live_residual": "Historical excluded population and current live file states were not accessed.", + "terminal_frontier_residual": "The typed-terminal candidate has no accepted byte head, so its readiness gate was injected for this case only.", + "residual_successor": "polylogue-excluded-cursor-live-proof", + }, + "production_route": { + "cursor_gate": "LiveWatcher._needs_work", + "transition": "CursorStore.revive_replaced_exclusion", + "catch_up": ( + "LiveWatcher._catch_up -> _scan_catch_up_candidates -> _catch_up_candidates -> " + "_plan_catch_up -> coordinated chunk ingest" + ), + "ingest": "LiveWatcher._ingest_files -> LiveBatchProcessor.ingest_files", + "terminal_evidence": "source.raw_artifacts", + "retry_state": "ops.ingest_cursor and ops.ingest_attempts", + }, + "outcomes": outcomes, + "cases": cases, + "fairness": { + "planner": "_interleave_by_source", + "property": ( + "browser-capture drains first; among non-browser-capture families, one candidate from each " + "present family reaches the first round" + ), + }, + "anti_vacuity": { + "indexed_authority": "byte_proven_source_raw_and_revision_head", + "indexed_session_count_before": indexed["indexed_before"]["indexed_sessions"], + "indexed_session_count": indexed["indexed"]["indexed_sessions"], + "typed_terminal_artifact": terminal_evidence["artifact_kind"] if terminal_evidence else None, + "unchanged_excluded_attempt_present": still_excluded["attempt_present"], + }, + } + digest = sha256(_canonical_json(body)).hexdigest() + receipt = {**body, "receipt_sha256": digest} + receipt_path.parent.mkdir(parents=True, exist_ok=True) + temporary_path = receipt_path.with_suffix(receipt_path.suffix + ".tmp") + temporary_path.write_bytes(_canonical_json(receipt)) + os.replace(temporary_path, receipt_path) + return receipt + + +def verify_receipt(receipt_path: Path) -> dict[str, Any]: + """Load and verify the immutable self-hash carried by a proof receipt.""" + receipt = json.loads(receipt_path.read_text(encoding="utf-8")) + if not isinstance(receipt, dict): + raise AssertionError("proof receipt must be a JSON object") + recorded = receipt.pop("receipt_sha256", None) + if not isinstance(recorded, str): + raise AssertionError("proof receipt has no receipt_sha256") + actual = sha256(_canonical_json(receipt)).hexdigest() + if recorded != actual: + raise AssertionError(f"proof receipt hash mismatch: recorded={recorded}, actual={actual}") + receipt["receipt_sha256"] = recorded + return receipt + + +__all__ = ["RECEIPT_SCHEMA", "run_excluded_cursor_live_proof", "verify_receipt"] diff --git a/tests/unit/sources/test_excluded_cursor_live_proof.py b/tests/unit/sources/test_excluded_cursor_live_proof.py new file mode 100644 index 0000000000..6340dfa8ce --- /dev/null +++ b/tests/unit/sources/test_excluded_cursor_live_proof.py @@ -0,0 +1,149 @@ +"""Executable proof for excluded-cursor revival and retry-state honesty.""" + +from __future__ import annotations + +import json +from pathlib import Path +from types import SimpleNamespace +from typing import Any, cast +from unittest.mock import Mock + +import pytest + +import polylogue.sources.live.watcher as live_watcher +from polylogue.sources.live.cursor import CursorStore +from polylogue.sources.live.watcher import LiveWatcher, WatchSource +from tests.infra.excluded_cursor_live_proof import run_excluded_cursor_live_proof, verify_receipt + +REPO_ROOT = Path(__file__).resolve().parents[3] + + +def test_candidate_fixture_proves_all_cursor_outcomes_and_is_immutable(tmp_path: Path) -> None: + archive_root = tmp_path / "candidate-archive" + receipt_path = tmp_path / "proof.json" + + receipt = run_excluded_cursor_live_proof(archive_root, receipt_path) + checked = verify_receipt(receipt_path) + + assert checked == receipt + assert set(cast(dict[str, Any], receipt["outcomes"])) == {"indexed", "still_excluded", "typed_terminal"} + assert receipt["outcomes"] == { + "indexed": True, + "still_excluded": True, + "typed_terminal": True, + } + typed_terminal = next(case for case in receipt["cases"] if case["case_id"] == "typed-terminal") + assert typed_terminal["terminal_evidence"]["parse_error_present"] is True + assert receipt["execution"] == { + "mode": "candidate_fixture", + "live_census": "not_run", + "live_residual": "Historical excluded population and current live file states were not accessed.", + "terminal_frontier_residual": "The typed-terminal candidate has no accepted byte head, so its readiness gate was injected for this case only.", + "residual_successor": "polylogue-excluded-cursor-live-proof", + } + assert receipt["production_route"]["catch_up"] == ( + "LiveWatcher._catch_up -> _scan_catch_up_candidates -> _catch_up_candidates -> " + "_plan_catch_up -> coordinated chunk ingest" + ) + assert receipt["anti_vacuity"] == { + "indexed_authority": "byte_proven_source_raw_and_revision_head", + "indexed_session_count_before": 0, + "indexed_session_count": 1, + "typed_terminal_artifact": "terminal_corrupt_input", + "unchanged_excluded_attempt_present": False, + } + + +def test_parser_fingerprint_revival_calls_real_actuator_and_excludes_unchanged_retry( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + source_root = tmp_path / "codex" + source_root.mkdir() + path = source_root / "excluded.jsonl" + path.write_text("payload\n", encoding="utf-8") + cursor = CursorStore(tmp_path / "ops.db") + stat = path.stat() + cursor.set( + path, + stat.st_size, + byte_offset=stat.st_size, + last_complete_newline=stat.st_size, + parser_fingerprint="old-parser", + content_fingerprint="payload-hash", + source_name="codex", + st_dev=stat.st_dev, + st_ino=stat.st_ino, + mtime_ns=stat.st_mtime_ns, + failure_count=5, + excluded=True, + ) + watcher = LiveWatcher( + cast(Any, SimpleNamespace(archive_root=tmp_path, backend=SimpleNamespace(db_path=tmp_path / "index.db"))), + (WatchSource(name="codex", root=source_root),), + cursor=cursor, + ) + try: + actuator = Mock(wraps=cursor.revive_replaced_exclusion) + monkeypatch.setattr(cursor, "revive_replaced_exclusion", actuator) + monkeypatch.setattr(live_watcher, "_PARSER_FINGERPRINT", "new-parser") + + assert watcher._needs_work(path) + actuator.assert_called_once() + revived = cursor.get_record(path) + assert revived is not None + assert not revived.excluded + assert revived.failure_count == 0 + assert cursor.list_retry_records() == [] + finally: + watcher.stop() + + unchanged_root = tmp_path / "unchanged" + unchanged_root.mkdir() + unchanged_cursor = CursorStore(unchanged_root / "ops.db") + assert unchanged_cursor._db_path != cursor._db_path + assert unchanged_cursor._ops_db_path != cursor._ops_db_path + unchanged_cursor.set( + path, + stat.st_size, + byte_offset=stat.st_size, + last_complete_newline=stat.st_size, + parser_fingerprint="new-parser", + content_fingerprint="payload-hash", + source_name="codex", + st_dev=stat.st_dev, + st_ino=stat.st_ino, + mtime_ns=stat.st_mtime_ns, + failure_count=5, + excluded=True, + ) + unchanged_watcher = LiveWatcher( + cast(Any, SimpleNamespace(archive_root=tmp_path, backend=SimpleNamespace(db_path=tmp_path / "index.db"))), + (WatchSource(name="codex", root=source_root),), + cursor=unchanged_cursor, + ) + try: + assert not unchanged_watcher._needs_work(path) + assert unchanged_cursor.list_excluded() == [str(path)] + assert unchanged_cursor.list_retry_records() == [] + finally: + unchanged_watcher.stop() + + +def test_receipt_with_wrong_self_hash_is_rejected(tmp_path: Path) -> None: + receipt_path = tmp_path / "receipt.json" + body = {"schema": "test", "receipt_sha256": "placeholder"} + receipt_path.write_text(json.dumps(body), encoding="utf-8") + with pytest.raises(AssertionError, match="hash mismatch"): + verify_receipt(receipt_path) + + +def test_committed_candidate_receipt_is_self_hashed() -> None: + receipt = verify_receipt(REPO_ROOT / "docs/evidence/polylogue-excluded-cursor-live-proof-2026-08-06.json") + + assert receipt["schema"] == "polylogue.excluded-cursor-live-proof.v1" + assert receipt["execution"]["live_census"] == "not_run" + assert receipt["outcomes"] == { + "indexed": True, + "still_excluded": True, + "typed_terminal": True, + } diff --git a/tests/unit/sources/test_live_watcher_catchup_order.py b/tests/unit/sources/test_live_watcher_catchup_order.py index 63ab8c7515..5d79ad339c 100644 --- a/tests/unit/sources/test_live_watcher_catchup_order.py +++ b/tests/unit/sources/test_live_watcher_catchup_order.py @@ -41,6 +41,34 @@ def test_interleave_by_source_round_robins_families() -> None: assert claude_paths == sorted(claude_paths) +def test_interleave_by_source_prevents_large_family_starvation() -> None: + """A long Codex backlog cannot hide a smaller family from round one.""" + candidates = [_candidate("codex", f"/home/u/.codex/sessions/x/{index:02d}.jsonl") for index in range(20)] + [ + _candidate("codex", "/home/u/.codex/sessions/x/excluded-after-fingerprint-change.jsonl"), + _candidate("hermes", "/home/u/.hermes/sessions/retry.jsonl"), + ] + + ordered = live_watcher._interleave_by_source(candidates) + + assert {candidate.source_name for candidate in ordered[:2]} == {"codex", "hermes"} + assert [candidate.source_name for candidate in ordered[:2]] == ["codex", "hermes"] + + +def test_interleave_by_source_prioritizes_browser_capture_before_round_robin() -> None: + candidates = [ + _candidate("codex", "/home/u/.codex/sessions/x/codex.jsonl"), + _candidate("hermes", "/home/u/.hermes/sessions/retry.json"), + _candidate("browser-capture", "/home/u/.browser/capture-b.json"), + _candidate("browser-capture", "/home/u/.browser/capture-a.json"), + ] + + ordered = live_watcher._interleave_by_source(candidates) + + assert [candidate.source_name for candidate in ordered[:2]] == ["browser-capture", "browser-capture"] + assert [candidate.path.name for candidate in ordered[:2]] == ["capture-a.json", "capture-b.json"] + assert {candidate.source_name for candidate in ordered[2:4]} == {"codex", "hermes"} + + def test_interleave_by_source_empty_input_returns_empty() -> None: assert live_watcher._interleave_by_source([]) == []