fix(#1196): one custody snapshot feeds replay and pending; unparseable custody lines naming the task make the read incomplete (Astra 6fe5 #1/#2)

This commit is contained in:
Anton 2026-09-24 20:38:31 +03:00
parent 6fd9eeccd4
commit e7b0510467
4 changed files with 107 additions and 4 deletions

View file

@ -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 ""),

View file

@ -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

View file

@ -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"

View file

@ -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"]