diff --git a/ouroboros/budget_pause.py b/ouroboros/budget_pause.py index 8b4e62e30..54d24adfe 100644 --- a/ouroboros/budget_pause.py +++ b/ouroboros/budget_pause.py @@ -350,13 +350,22 @@ def observe_task_runs(root: Any, task_id: str, *, reason: str = "budget_resume_u # UNKNOWN, never "no open runs" (Astra run-a882315dbcd7 #2). if custody.custody_log_unreadable(pathlib.Path(root)): raise OSError("custody_log_unreadable") - runs = [run for run in custody.replay(pathlib.Path(root)).values() + from ouroboros.delegate_custody_memo import malformed_custody_lines_mentioning + + # ONE snapshot feeds both projections: a START_REQUESTED that becomes + # STARTED between two reads must land in one of them (Astra 6fe5 #1). + snapshot = list(custody.custody_rows(pathlib.Path(root))) + # An unparseable custody line naming this task may hide its request. + malformed = malformed_custody_lines_mentioning(pathlib.Path(root), mine) + if malformed is None or malformed: + raise OSError(f"custody_rows_incomplete:{'unknown' if malformed is None else malformed}") + runs = [run for run in custody.replay(pathlib.Path(root), rows=snapshot).values() if str(getattr(run, "task_id", "") or "") == mine and not getattr(run, "settled", True)] # A START_REQUESTED whose response was lost has no run id yet but may # be a live remote writer: unknown custody, never absence (#3). from ouroboros.delegate_pending import pending_invocations - pending = [row for row in pending_invocations(pathlib.Path(root)) + pending = [row for row in pending_invocations(pathlib.Path(root), rows=snapshot) if str(row.get("task_id") or "") == mine] pending_rows = [{"run_id": "", "invocation_id": str(row.get("invocation_id") or ""), "route": str(row.get("route") or ""), diff --git a/ouroboros/delegate_custody_memo.py b/ouroboros/delegate_custody_memo.py index 2c4464293..6c70bef31 100644 --- a/ouroboros/delegate_custody_memo.py +++ b/ouroboros/delegate_custody_memo.py @@ -94,6 +94,11 @@ class _ChainMemo: rows_view: Tuple[Dict[str, Any], ...] = () generation: int = 0 torn_archive_lines: int = 0 + # Custody-marked lines the fold could NOT parse (bounded). A START_REQUESTED + # joined onto a torn prefix lands here; readers that must prove absence + # consult it instead of trusting the silent skip (#1196, Astra 6fe5 #2). + malformed_marker_lines: List[bytes] = field(default_factory=list) + malformed_overflow: bool = False # (generation, folded state) for ``folded_state``; cloned on every return. state_cache: Optional[Tuple[int, Any]] = None @@ -228,6 +233,8 @@ def _fold_segment( break # a torn live tail completes on a later call # An archive never completes its torn tail: consume it, count it. memo.torn_archive_lines += 1 + if marker in raw: + _remember_malformed(memo, raw) consumed += len(raw) continue if inner: @@ -240,8 +247,12 @@ def _fold_segment( try: row = json.loads(raw.decode("utf-8", errors="replace")) except ValueError: + _remember_malformed(memo, raw) continue - if isinstance(row, dict) and str(row.get("type") or "").startswith(_custody()._ROW_MARKER): + if not isinstance(row, dict): + _remember_malformed(memo, raw) + continue + if str(row.get("type") or "").startswith(_custody()._ROW_MARKER): memo.rows.append(_compact_row(row, (stat.st_dev, stat.st_ino, offset, len(raw)))) try: after = os.fstat(handle.fileno()) @@ -251,6 +262,31 @@ def _fold_segment( prefix_sha256=hasher.hexdigest() if hasher is not None else "") +_MALFORMED_KEEP = 200 + + +def _remember_malformed(memo: _ChainMemo, raw: bytes) -> None: + if len(memo.malformed_marker_lines) >= _MALFORMED_KEEP: + memo.malformed_overflow = True + return + memo.malformed_marker_lines.append(bytes(raw[:65536])) + + +def malformed_custody_lines_mentioning(drive_root: Any, needle: str) -> Optional[int]: + """How many unparseable custody-marked lines mention ``needle``; None = unknown. + + None when the memo was bypassed (lenient read, nothing recorded) or the + bounded record overflowed: an absence proof cannot be built from that. + """ + key = _key(_custody().event_log_path(drive_root)) + token = str(needle or "").encode("utf-8") + with _lock_for(key): + memo, _rows = _refresh(drive_root) + if memo is None or memo.malformed_overflow: + return None + return sum(1 for raw in memo.malformed_marker_lines if token and token in raw) + + def _advance(memo: _ChainMemo, chain: List[Tuple[pathlib.Path, os.stat_result, bool]]) -> None: before = len(memo.rows) last_known = len(memo.segments) - 1 diff --git a/tests/test_budget_pause_astra_a882.py b/tests/test_budget_pause_astra_a882.py index ea3af0d9a..7ad97c9b8 100644 --- a/tests/test_budget_pause_astra_a882.py +++ b/tests/test_budget_pause_astra_a882.py @@ -305,3 +305,61 @@ def test_a_parallel_batch_raising_usage_accounting_error_waits_for_already_start assert raised_at - started_at >= 0.4 finally: b_done.wait(timeout=2.0) # never leak the worker thread past the test + + +# --- Astra run-6fe5bf761449 follow-ups: one snapshot, and malformed custody lines ---- + +def _write_event_log(root, lines): + from ouroboros import delegate_custody as custody + + path = custody.event_log_path(root) + path.parent.mkdir(parents=True, exist_ok=True) + path.write_bytes(b"".join(lines)) + return path + + +def test_both_custody_projections_come_from_one_snapshot(tmp_path, monkeypatch): + """Astra 6fe5 #1: replay and pending_invocations must fold the SAME rows, so a + START_REQUESTED that becomes STARTED between two reads is seen by one of them.""" + from ouroboros import budget_pause + from ouroboros import delegate_custody as custody + from ouroboros import delegate_pending + + snapshot = [{"type": custody.START_REQUESTED, "invocation_id": "inv-x", "task_id": "snap-1"}] + seen = {} + monkeypatch.setattr(custody, "custody_log_unreadable", lambda _root: False) + monkeypatch.setattr(custody, "custody_rows", lambda _root: tuple(snapshot)) + + def _replay(_root, rows=None): + seen["replay"] = rows + return {} + + def _pending(_root, rows=None): + seen["pending"] = rows + return [{"invocation_id": "inv-x", "task_id": "snap-1", "route": "codex"}] + + monkeypatch.setattr(custody, "replay", _replay) + monkeypatch.setattr(delegate_pending, "pending_invocations", _pending) + observed = budget_pause.observe_task_runs(tmp_path, "snap-1") + assert seen["replay"] is not None and seen["replay"] == seen["pending"] == snapshot + assert observed["coverage_basis"] == "pending_invocations_unbound" + + +def test_a_malformed_custody_line_naming_the_task_is_unknown_custody(tmp_path): + """Astra 6fe5 #2, on a REAL event log: a START_REQUESTED joined onto a torn + prefix is unparseable; the memo used to skip it silently and the observer then + proved "no open runs". Now it is an incomplete read, never absence.""" + import json + + from ouroboros import budget_pause + from ouroboros import delegate_custody as custody + + start = json.dumps({"type": custody.START_REQUESTED, "invocation_id": "inv-torn", + "task_id": "torn-1"}).encode() + _write_event_log(tmp_path, [b'{"type": "llm_round", "x": 1', start + b"\n"]) + observed = budget_pause.observe_task_runs(tmp_path, "torn-1") + assert observed["custody_read"] == "failed" + assert "custody_rows_incomplete" in observed["error"] + # Another task is not blocked by this task's torn line. + other = budget_pause.observe_task_runs(tmp_path, "someone-else") + assert other["custody_read"] == "ok" and other["coverage_basis"] == "no_open_runs" diff --git a/tests/test_budget_pause_exact.py b/tests/test_budget_pause_exact.py index 061fa0437..73ad2940a 100644 --- a/tests/test_budget_pause_exact.py +++ b/tests/test_budget_pause_exact.py @@ -475,7 +475,7 @@ def test_unreadable_external_custody_is_held_as_unknown_not_clean(tmp_path, monk assert observed["custody_read"] == "failed" and observed["runs"] == [] assert "boom" in observed["error"] # The supervisor-side twin shares the body: a replay that fails is the same typed fact. - monkeypatch.setattr(custody, "replay", lambda _r: (_ for _ in ()).throw(OSError("rows torn"))) + monkeypatch.setattr(custody, "replay", lambda _r, rows=None: (_ for _ in ()).throw(OSError("rows torn"))) grant_side = budget_pause.observe_task_runs(tmp_path, ctx.task_id) assert grant_side["custody_read"] == "failed" and "rows torn" in grant_side["error"]