fix(#1196): whole malformed custody lines (bounded, overflow=unknown); rows and integrity from one refresh (Astra 2bc1 #1/#2)

This commit is contained in:
Anton 2026-09-24 20:43:18 +03:00
parent e7b0510467
commit 4cee6aff01
3 changed files with 56 additions and 15 deletions

View file

@ -350,13 +350,14 @@ 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")
from ouroboros.delegate_custody_memo import malformed_custody_lines_mentioning
from ouroboros.delegate_custody_memo import custody_rows_with_integrity
# 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)
# STARTED between two reads must land in one of them (Astra 6fe5 #1),
# and its integrity is judged on that same read. An unparseable custody
# line naming this task may hide its request.
rows_read, malformed = custody_rows_with_integrity(pathlib.Path(root), mine)
snapshot = list(rows_read)
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()

View file

@ -99,6 +99,7 @@ class _ChainMemo:
# 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
malformed_bytes: int = 0
# (generation, folded state) for ``folded_state``; cloned on every return.
state_cache: Optional[Tuple[int, Any]] = None
@ -263,28 +264,36 @@ def _fold_segment(
_MALFORMED_KEEP = 200
_MALFORMED_BYTES_KEEP = 8 * 1024 * 1024
def _remember_malformed(memo: _ChainMemo, raw: bytes) -> None:
if len(memo.malformed_marker_lines) >= _MALFORMED_KEEP:
# WHOLE lines only: a truncated copy could drop the one id a reader needs
# (a START_REQUESTED joined after a huge torn prefix). Past the bound the
# record is unknown, never a shorter proof of absence (Astra 2bc1 #1).
if (len(memo.malformed_marker_lines) >= _MALFORMED_KEEP
or memo.malformed_bytes + len(raw) > _MALFORMED_BYTES_KEEP):
memo.malformed_overflow = True
return
memo.malformed_marker_lines.append(bytes(raw[:65536]))
memo.malformed_marker_lines.append(bytes(raw))
memo.malformed_bytes += len(raw)
def malformed_custody_lines_mentioning(drive_root: Any, needle: str) -> Optional[int]:
"""How many unparseable custody-marked lines mention ``needle``; None = unknown.
def custody_rows_with_integrity(drive_root: Any, needle: str) -> Tuple[Tuple[Dict[str, Any], ...], Optional[int]]:
"""ONE refresh: the rows AND how many unparseable custody lines mention ``needle``.
None when the memo was bypassed (lenient read, nothing recorded) or the
bounded record overflowed: an absence proof cannot be built from that.
The count is None when that same refresh bypassed the memo (lenient read,
nothing recorded) or the bounded record overflowed: no absence proof can be
built from it. Rows and integrity come from the same read, so a check can
never certify a different traversal than the one it judges (Astra 2bc1 #2).
"""
key = _key(_custody().event_log_path(drive_root))
token = str(needle or "").encode("utf-8")
with _lock_for(key):
memo, _rows = _refresh(drive_root)
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)
return rows, None
return rows, 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:

View file

@ -327,8 +327,11 @@ def test_both_custody_projections_come_from_one_snapshot(tmp_path, monkeypatch):
snapshot = [{"type": custody.START_REQUESTED, "invocation_id": "inv-x", "task_id": "snap-1"}]
seen = {}
from ouroboros import delegate_custody_memo
monkeypatch.setattr(custody, "custody_log_unreadable", lambda _root: False)
monkeypatch.setattr(custody, "custody_rows", lambda _root: tuple(snapshot))
monkeypatch.setattr(delegate_custody_memo, "custody_rows_with_integrity",
lambda _root, _needle: (tuple(snapshot), 0))
def _replay(_root, rows=None):
seen["replay"] = rows
@ -363,3 +366,31 @@ def test_a_malformed_custody_line_naming_the_task_is_unknown_custody(tmp_path):
# 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"
def test_a_huge_torn_prefix_cannot_truncate_the_task_id_away(tmp_path):
"""Astra 2bc1 #1: a >64 KiB torn unrelated prefix with this task's
START_REQUESTED joined after it must still count as this task's malformed line."""
import json
from ouroboros import budget_pause
from ouroboros import delegate_custody as custody
start = json.dumps({"type": custody.START_REQUESTED, "invocation_id": "inv-big",
"task_id": "big-1"}).encode()
prefix = b'{"type": "delegate_run_note", "blob": "' + b"x" * 200_000
_write_event_log(tmp_path, [prefix, start + b"\n"])
observed = budget_pause.observe_task_runs(tmp_path, "big-1")
assert observed["custody_read"] == "failed"
def test_a_bypassed_refresh_in_the_same_read_is_unknown(tmp_path, monkeypatch):
"""Astra 2bc1 #2: when the one refresh that produced the rows bypassed the memo
(lenient read), the observer cannot prove absence from those rows."""
from ouroboros import budget_pause
from ouroboros import delegate_custody_memo
_write_event_log(tmp_path, [b""])
monkeypatch.setattr(delegate_custody_memo, "_refresh", lambda _root: (None, ()))
observed = budget_pause.observe_task_runs(tmp_path, "bypass-1")
assert observed["custody_read"] == "failed" and "unknown" in observed["error"]