From 69e92679c749cd6e80b6827b59c5cce2995a14a6 Mon Sep 17 00:00:00 2001 From: Layne Penney Date: Thu, 20 Aug 2026 11:00:56 -0500 Subject: [PATCH] fix(events): repair torn outbox lines and report unreadable ones Fix 4 of the torn-line sweep, and the two-part contract it was ruled as. TERMINATOR REPAIR. emit() appended with "a" and wrote json + "\n". A previous write that died between write() and fsync() leaves a last line with no terminator, so the next append GLUES two records into one. The damage runs FORWARD from the tear: the torn record and THE NEXT HEALTHY APPEND fuse into one unparseable line, while the record before the tear is untouched. A torn write therefore costs that record and the next one written after it, permanently, because every later append builds on the glued line. emit() now probes the last byte under the existing write lock and heals the seam first. THE COUNT. Both readers skipped unusable lines in silence. For the channel bridge an unreadable line is a message that never reaches a channel. It is reported on EVERY read, not once: the cursor filter applies only to lines that parse, so a line with no usable seq can never be advanced past, wherever it sits -- position is irrelevant and mid-file lines repeat exactly as trailing ones do. Deliberate, not incidental, and its cost is named in the residuals. read_events_detailed() now returns the events AND the lines it could not read, from the SAME read. read_events() stays list-shaped for the ELEVEN call sites that index and len() it -- all tests, and ZERO production callers remain once the bridge moves to read_events_detailed(), so the wrapper is a test-compatibility surface rather than a load-bearing API. An earlier draft of this message said seventeen; that was a substring artifact counting a def, two prose mentions inside strings, and four hits on an unrelated _read_events helper in another test file. The bridge reports on stderr so stdout stays parseable. Reported from the read path ONLY. _current_seq() runs once per emit inside the write lock, where a count would be per-APPEND and aimed at whoever happened to be writing. Its docstring says so, because an omission and a decision look identical in code. Also guarded, each with witnesses: the decode moved to bytes-per-line so a single invalid byte cannot escape from outside every guard; the parse guard's exception tuple is now derived from what json.loads can raise rather than from what had been seen -- the old (JSONDecodeError, TypeError) caught syntax errors but missed RecursionError, which is not a ValueError, so deep nesting escaped; structure is checked separately, since a valid JSON array parses and is still not an event; seq values are type-validated, since bool is an int subclass and a float becomes inf and serializes as Infinity; and a corrupt cursor no longer bricks reads. THREE THINGS THIS FIX GOT WRONG FIRST, none found by its own witnesses: - An early version swallowed OSError in the line iterator. That is exactly the OSError-to-zero fallback an earlier fix removed on purpose: swallow it and _current_seq returns 0, then emit allocates sequence numbers that duplicate live ones. A pre-existing test stood guard and caught it. Content errors are data and get skipped-and-counted; an I/O error is not knowing what the file holds and must fail closed. The reader tolerates FileNotFoundError only, because rotation renames the file, and propagates every other OSError. - MalformedLine and EventRead were dataclasses. This module is loaded out-of-tree by spawned workers via spec_from_file_location + exec_module, which does not register it in sys.modules, and dataclass resolves field types through exactly that. Every worker died at import. Now plain classes, with a witness that fails in 0.06s naming the cause instead of after a 10s timeout that reads like flakiness. - The default report echoed raw line content to stderr, verbatim, including an API-key-shaped string in a measured probe. stderr is copied into CI logs and transcripts. Redacting was rejected: matching secret-shaped patterns in arbitrary bytes is a denylist over untrusted input, the same defect the parse guard exists to avoid. The excerpt stays on the data object; the default report prints ordinal and reason, both structural, with show_content opt-in. One existing test's monkeypatch target moved read_text -> read_bytes because the read verb changed; its contract and every assertion are unchanged, and it would otherwise have gone green by missing its target. Both it and the new witness now undo the patch before reading the file back, since a verification must not travel the path the test deliberately sabotaged. A reviewer measured both of those statements against an earlier version of this message and of the description, where the glue direction was backwards and the count was claimed to be reported once. The suite was green while the prose said the opposite of the code, because nothing pinned that behavior. Three witnesses now do, including the mid-file case, which shows the rule is broader than the trailing-line framing that surfaced it. A third review then caught a further claim, in the description only: that a torn last line self-heals at the next emit. That is true of one tear and false of the other, and the bare word "torn" hides the difference. An UNTERMINATED record -- complete, only its newline lost -- DOES recover fully once the seam is healed. A TRUNCATED record -- the write stopped mid-record -- never does: its bytes were never written, so it stays unreadable and every read reports it until the file is repaired. A FOURTH review then caught the opposite overcorrection, which a later draft of this message had made -- claiming a torn record never becomes readable -- and the unterminated-tear control disproves it. Both cases are now named separately and the bare word is not used to claim anything about recovery. That prose was wrong only because every tear fixture until then dropped the trailing newline and left the record COMPLETE -- the lucky case. Three witnesses now cover the realistic tear, with the UNTERMINATED tear kept as a discriminating control so the finding cannot be mistaken for the norm. What the repair buys, stated exactly: it cannot recover a record that was never fully written, and it saves the NEXT one. Correcting the message and the description was not enough, and a second review caught that: the same backwards claim was still in the production repair comment and in the test module docstring. Fixing the two cited surfaces and stopping there left the code asserting one direction while the description asserted the other -- an artifact contradicting itself, which is worse than being uniformly wrong. A sweep for the whole class rather than the cited instances found a third live instance neither review named: the reported-once claim, still standing in that same docstring. Those particular corrections changed comments only and left the executable syntax tree identical, verified by parsing both revisions and comparing with docstrings stripped. A FIFTH review then caught the same overclaim surviving in the test module comment -- "A TORN RECORD IS NEVER REPAIRED BY A LATER EMIT" -- which the unterminated control disproves. The class query that was supposed to have swept it missed it because the query was CASE-SENSITIVE and the comment is upper case: a false negative from my own instrument, of the kind noted one revision earlier and then committed in the next query. That round also renamed the control from "benign" to "unterminated" so the tests and the prose use one vocabulary, so unlike the previous round the syntax tree DID change here and the suite was re-run rather than reasoned about. A SIXTH review, from the other reviewer, found three defects and all three were PROSE -- the code it gated was clean. The first is the one worth the round: the reported-once claim was still live in the channel-bridge comment, in the bridge hunk of ALL FIVE versions, and every class sweep I ran missed it because it PARAPHRASES the sentence the earlier reviews cited rather than repeating it. A query built from a cited instance finds COPIES, not paraphrases; the class is the CLAIM, and the only instrument that finds a paraphrase is reading the hunk. The second: the parse-guard sentence read as if the new tuple were open-ended. It is enumerated -- (ValueError, RecursionError). What changed is where the enumeration comes FROM: what json.loads can raise, rather than what had been seen. The third: "seventeen call sites" was an artifact of a substring query that counted a def, two prose mentions inside strings, and four hits on an unrelated _read_events helper in another test file. The real count is eleven, all tests, zero production callers. BOTH reviewers co-signed seventeen at v1 from that same defective query, which is why it needed re-deriving rather than re-reading -- and the class sweep for it then found three more copies of the count in docstrings that no review had cited. A SEVENTH review found two more, and one of them was wrong in the CODE rather than in the prose about it. The bridge comment said an unterminated record is reported and then heals at the next emit. Measured false: an unterminated record whose bytes are COMPLETE is never reported at all, because _iter_outbox() splits on b"\n" and a complete final chunk parses on the spot -- before any emit, repair or no repair. The sentence implied a window of unreadability that does not exist. The gap that let it survive is now pinned: every tear fixture in this file emitted AFTER tearing, so nothing had ever asked what the READER alone does with a torn file. W11 is that pair -- the unterminated case reporting nothing across two reads, the truncated case reporting one as its discriminating control. The reviewer's own probe became the witness. Its two siblings, the W10 header and the description's table, are scoped to tear -> emit -> read and are true there, so they are deliberately unchanged. The second finding was stale suite totals, and the mechanism differs from the one proposed: the +3 is not the earlier merge, which added 21 tests that are inside both baselines, but a test directory the narrow invocation does not collect. Both figures were correct measurements of different scopes and the defect was publishing one without naming which -- which is also how the three mutation rows above went stale, unasked about by any review. RESIDUALS, named rather than presented as clean: this is a third copy of these primitives, and consolidating them is outside this fix's scope; the same content-echo property exists in the prototype reporter merged earlier; unbounded line length remains a read-layer resource limit and is undefended; and a TRUNCATED record -- an UNTERMINATED one is readable throughout and never reaches the report at all -- is permanently unreadable and reported on every read, so a consumer polling in a loop warns every cycle until the file is repaired -- suppressing that would need persisted already-reported state and would hide a re-occurring fault, so it is the lesser cost but it is a cost. Evidence: 37 witnesses, 0 before. Nine mutation rows, RE-MEASURED for this version over the two events spec files (71 tests), each killing witnesses whose failure TYPE matches the mutation, restores hash-verified with the unmutated pair as a control: decode 11, exception tuple 1, terminator repair 5, seq validation 4, count silenced 19, reporter silenced 1, OSError swallowed 3, FileNotFoundError widened 1, cursor guard 1. Three of those rows were STALE -- published as 8, 3 and 12 where the true figures are 11, 5 and 19 -- because they were measured early and never re-derived while five review rounds added witnesses to the files they count. Full suite at BOTH scopes, base pinned by SHA, since omitting the scope is what made the last version's totals unverifiable. On origin/dev@93675677 in an isolated worktree: pytest gr2/tests 1032 passed / 5 failed, pytest gr2 1035 / 5. On this head: 1069 / 5 and 1072 / 5. Delta exactly +37 in both, matching the witness count, failure sets identical in both directions. The scopes differ by gr2/gr2_overlay/tests/test_overlay_refs_namespace.py -- 3 tests, measured passing on both sides, which the narrow invocation does not collect. Lint unchanged: production 6 = 6, edited test file 23 = 23, new file 0. Premium boundary: grip is OSS; this is local file mechanics over opaque paths and carries no identity, org, or policy content. Co-Authored-By: Claude --- gr2/python_cli/channel_bridge.py | 26 +- gr2/python_cli/events.py | 290 ++++++++++++-- gr2/tests/test_events.py | 20 +- gr2/tests/test_events_torn_line.py | 592 +++++++++++++++++++++++++++++ 4 files changed, 896 insertions(+), 32 deletions(-) create mode 100644 gr2/tests/test_events_torn_line.py diff --git a/gr2/python_cli/channel_bridge.py b/gr2/python_cli/channel_bridge.py index a4c08e08..505f4b49 100644 --- a/gr2/python_cli/channel_bridge.py +++ b/gr2/python_cli/channel_bridge.py @@ -2,7 +2,7 @@ Translates outbox events into channel messages per the mapping table in HOOK-EVENT-CONTRACT.md section 8. Uses cursor-based consumption from -events.read_events(). +events.read_events_detailed(), which also reports lines it could not read. The bridge is a pure function layer: format_event() maps an event dict to a message string (or None), and run_bridge() orchestrates cursor reads and @@ -14,7 +14,7 @@ from pathlib import Path from typing import Callable -from .events import read_events +from .events import read_events_detailed, warn_unreadable _CONSUMER_NAME = "channel_bridge" @@ -103,9 +103,27 @@ def run_bridge( The post_fn receives formatted message strings; the caller decides how to deliver them (recall_channel, print, log, etc.). """ - events = read_events(workspace_root, _CONSUMER_NAME) + read = read_events_detailed(workspace_root, _CONSUMER_NAME) + # An unreadable outbox line is a message that never reaches a channel, and + # the bridge is the only place that says so. + # + # It is reported on EVERY read, not once: the cursor filter compares seq + # against the cursor and only lines that PARSE have a usable seq, so an + # unreadable line can never be advanced past, wherever it sits in the file + # (witnessed terminal AND mid-file). An UNTERMINATED record never reaches + # this report at all: the reader splits on b"\n", so a COMPLETE record whose + # only loss was its terminator is the final chunk and parses on the spot -- + # before any emit, repair or no repair (W11, and its truncated control). + # What seam repair protects is the NEXT append, which would otherwise be + # glued onto it. So everything reported here is a TRUNCATED record, and for + # those the report is permanent: a bridge polling in a loop warns every + # cycle until someone repairs the file. That cost is real and is recorded + # as a residual. + # + # stderr keeps stdout parseable for callers that consume the posted count. + warn_unreadable(read) posted = 0 - for event in events: + for event in read.events: msg = format_event(event) if msg is not None: post_fn(msg) diff --git a/gr2/python_cli/events.py b/gr2/python_cli/events.py index 0793108e..826ba967 100644 --- a/gr2/python_cli/events.py +++ b/gr2/python_cli/events.py @@ -13,6 +13,7 @@ import os import sys import time +from collections.abc import Iterator from contextlib import contextmanager from datetime import datetime, timezone from enum import Enum @@ -120,21 +121,170 @@ def _event_write_lock(outbox: Path): fcntl.flock(lock_file.fileno(), fcntl.LOCK_UN) +# DELIBERATELY PLAIN CLASSES, NOT @dataclass, AND THE REASON IS LOAD-BEARING. +# +# This module is loaded out-of-tree by spawned workers via +# importlib.util.spec_from_file_location() + exec_module(), which does NOT +# register the module in sys.modules. @dataclass resolves field types through +# sys.modules[cls.__module__].__dict__, so under that loader it raises +# AttributeError: 'NoneType' object has no attribute '__dict__' AT IMPORT -- +# every worker dies before running a line of its own. +# +# Measured: adding @dataclass here killed both writers in the concurrent-emit +# integrity test before either reached sequence allocation. The failure was +# visible an hour earlier in an ad-hoc probe and was dismissed as a loader +# artifact; it was a portability constraint on this file. Anything importable by +# a spawned worker must import without sys.modules registration. +# test_events_torn_line.py carries a witness for exactly this. + + +class MalformedLine: + """One line the reader could not turn into an event, and why.""" + + __slots__ = ("ordinal", "reason", "excerpt") + + def __init__(self, ordinal: int, reason: str, excerpt: str) -> None: + self.ordinal = ordinal + self.reason = reason + self.excerpt = excerpt + + def __repr__(self) -> str: # pragma: no cover - debugging aid + return f"MalformedLine(ordinal={self.ordinal!r}, reason={self.reason!r})" + + def __eq__(self, other: object) -> bool: + if not isinstance(other, MalformedLine): + return NotImplemented + return (self.ordinal, self.reason, self.excerpt) == ( + other.ordinal, + other.reason, + other.excerpt, + ) + + +class EventRead: + """Events read, AND the lines that could not be read. + + Both halves come from the SAME read. A second pass to "check health" would + describe a different moment, and the outbox is appended to concurrently. + """ + + __slots__ = ("events", "malformed") + + def __init__( + self, + events: list[dict[str, object]], + malformed: tuple[MalformedLine, ...], + ) -> None: + self.events = events + self.malformed = malformed + + def __repr__(self) -> str: # pragma: no cover - debugging aid + return f"EventRead(events={len(self.events)}, malformed={len(self.malformed)})" + + +def _decode_line(raw: bytes) -> tuple[str | None, str]: + """Decode one line, or say why it could not be decoded. + + The DECODE is the acquisition, not a transformation of an existing string: + reading the outbox in text mode manufactures the line as text outside every + guard, so a single invalid byte escapes as UnicodeDecodeError from a place + no parse guard can reach. Reading bytes and decoding per line puts the one + operation that can fail on file CONTENT inside the funnel. + """ + try: + return raw.decode("utf-8"), "" + except UnicodeDecodeError as exc: + return None, str(exc) + + +def _read_object(line: str) -> tuple[dict[str, object] | None, str]: + """Parse one line into a JSON object, or say why it is not one. + + This IS an enumerated tuple -- what changed is how the entries were chosen. + The previous guard, (JSONDecodeError, TypeError), was a list of what had been + SEEN: it caught syntax errors, and missed RecursionError, which json.loads + raises on deeply nested input and which is NOT a ValueError, so deep nesting + escaped a reader whose contract is to not raise on file content. + + This tuple is derived from what the OPERATION can raise: ValueError as the + base class covering JSONDecodeError and any other value error, plus + RecursionError, the one thing json.loads raises that ValueError does not + cover. Structure -- is the result an object? -- is then checked separately + below, because a valid JSON array parses cleanly and is still not an event. + """ + try: + obj = json.loads(line) + except (ValueError, RecursionError) as exc: + return None, str(exc) + if not isinstance(obj, dict): + return None, f"line is a {type(obj).__name__}, not an object" + return obj, "" + + +def _read_seq(obj: dict[str, object]) -> tuple[int | None, str]: + """The event's sequence number, or why it cannot be used as one. + + bool is an int subclass, so `isinstance(True, int)` is True and `max(0, True)` + quietly yields 1. A float slips through comparison and arithmetic and then + fails at serialization: 1e999 is valid JSON input, becomes inf, and + json.dumps writes `Infinity`, which is not valid JSON output. Validate the + type here rather than discovering it downstream. + """ + seq = obj.get("seq") + if isinstance(seq, bool) or not isinstance(seq, int): + return None, f"seq is {type(seq).__name__}, not an integer" + return seq, "" + + +def _iter_outbox(outbox: Path) -> Iterator[tuple[dict[str, object] | None, str, bytes]]: + """Yield (object, reason, raw) per line, never raising on file CONTENT. + + Splitting bytes on b"\n" rather than iterating text: see _decode_line. + """ + # DELIBERATELY CATCHES NOTHING. An I/O failure is not a malformed line: it + # is not knowing what the file contains, and the caller's correct response + # differs. _current_seq() is on the WRITE path, where swallowing this + # returns 0 and emit() then allocates a sequence number that duplicates + # existing ones -- silent event-log corruption, which is why an earlier fix + # removed exactly this OSError-to-zero fallback and left a test standing + # guard over it. Reintroducing it here was caught by that test and not by + # any witness of mine. + blob = outbox.read_bytes() + for raw in blob.split(b"\n"): + if not raw.strip(): + continue + line, reason = _decode_line(raw) + if line is None: + yield None, reason, raw + continue + obj, reason = _read_object(line) + yield obj, reason, raw + + +def _excerpt(raw: bytes) -> str: + return raw[:120].decode("utf-8", "replace") + + def _current_seq(outbox: Path) -> int: + """Highest sequence number in the outbox. + + DELIBERATELY REPORTS NO COUNT, and that is a decision rather than an + omission. This runs once per emit, inside the write lock, so a count here + would be a per-APPEND number reported to whoever happened to be writing -- + the wrong altitude and the wrong audience. Unreadable lines are surfaced + once, from read_events_detailed(), where the number is per-READ and reaches + a consumer that can act on it. + """ if not outbox.exists(): return 0 - text = outbox.read_text() last_seq = 0 - for line in text.strip().split("\n"): - line = line.strip() - if not line: + for obj, _reason, _raw in _iter_outbox(outbox): + if obj is None: continue - try: - obj = json.loads(line) - if isinstance(obj, dict) and "seq" in obj: - last_seq = max(last_seq, obj["seq"]) - except (json.JSONDecodeError, TypeError): + seq, _seq_reason = _read_seq(obj) + if seq is None: continue + last_seq = max(last_seq, seq) return last_seq @@ -189,8 +339,25 @@ def emit( event["agent_id"] = agent_id event.update(payload) - with outbox.open("a") as event_file: - event_file.write(json.dumps(event, separators=(",", ":")) + "\n") + # TERMINATOR REPAIR. A previous write that died between write() and + # fsync() leaves a last line with no "\n". Appending onto that GLUES + # two records into one line, and the damage runs FORWARD from the + # tear: the torn record and THE NEXT HEALTHY APPEND fuse into one + # unparseable line, while the record before the tear is untouched. + # So a torn write costs that record and the next one written after + # it, permanently, because every later append builds on the glued + # line. Probe the last byte and heal the seam before writing. + # (Direction measured, not reasoned: see the glue witness in + # test_events_torn_line.py. An earlier version of this comment + # stated it backwards.) + with outbox.open("a+b") as event_file: + event_file.seek(0, os.SEEK_END) + if event_file.tell() > 0: + event_file.seek(-1, os.SEEK_END) + if event_file.read(1) != b"\n": + event_file.write(b"\n") + payload_bytes = json.dumps(event, separators=(",", ":")).encode("utf-8") + event_file.write(payload_bytes + b"\n") event_file.flush() os.fsync(event_file.fileno()) except Exception as exc: @@ -228,27 +395,52 @@ def emit_after_outcome( ) -def read_events(workspace_root: Path, consumer: str) -> list[dict[str, object]]: +def read_events_detailed(workspace_root: Path, consumer: str) -> EventRead: + """New events for `consumer`, AND the lines that could not be read. + + This is the primitive; read_events() is the list-shaped wrapper kept for the + eleven existing call sites, all of them tests. The count is a RETURN VALUE rather than + hidden state: these are module-level functions with no instance to hang + health on, and a module-level accumulator would be wrong under concurrent + readers, which the outbox explicitly has. + + An unreadable line here is a LOST EVENT -- for the channel bridge it is a + message that never reaches a channel -- so silence is the failure, not the + safe default. + """ outbox = _outbox_path(workspace_root) if not outbox.exists(): - return [] + return EventRead([], ()) cursor = _load_cursor(workspace_root, consumer) last_seq = cursor.get("last_seq", 0) + if isinstance(last_seq, bool) or not isinstance(last_seq, int): + # A hand-edited or truncated cursor must not brick every future read. + last_seq = 0 events: list[dict[str, object]] = [] - text = outbox.read_text() - for line in text.strip().split("\n"): - line = line.strip() - if not line: - continue - try: - obj = json.loads(line) - except json.JSONDecodeError: + malformed: list[MalformedLine] = [] + try: + lines = list(_iter_outbox(outbox)) + except FileNotFoundError: + # ONLY this one, and only on the read path: _maybe_rotate() renames the + # outbox, so a reader can legitimately lose the file between exists() + # and the read. The events are not gone, they are in an archive. Any + # OTHER OSError is a real failure and propagates -- a reader that + # swallows EIO reports "no new events" forever. + return EventRead([], ()) + for ordinal, (obj, reason, raw) in enumerate(lines, start=1): + if obj is None: + malformed.append(MalformedLine(ordinal, reason, _excerpt(raw))) continue - if not isinstance(obj, dict): + seq, seq_reason = _read_seq(obj) + if seq is None: + # An event whose seq is unusable cannot be ordered against the + # cursor. Comparing it anyway is how "x" <= 0 raises TypeError from + # inside a reader whose job is to not raise on file content. + malformed.append(MalformedLine(ordinal, seq_reason, _excerpt(raw))) continue - if obj.get("seq", 0) <= last_seq: + if seq <= last_seq: continue events.append(obj) @@ -265,7 +457,57 @@ def read_events(workspace_root: Path, consumer: str) -> list[dict[str, object]]: }, ) - return events + return EventRead(events, tuple(malformed)) + + +def read_events(workspace_root: Path, consumer: str) -> list[dict[str, object]]: + """New events for `consumer`, DISCARDING the unreadable-line report. + + Kept list-shaped because eleven call sites index and len() the result, all + of them tests -- no production caller remains once the bridge moves to + read_events_detailed(), so this is a test-compatibility surface. + It discards information, which is exactly the silence this sweep exists to + remove -- so if you are writing a NEW consumer, call read_events_detailed() + and say something about `malformed`. The one production consumer, the + channel bridge, does. + """ + return read_events_detailed(workspace_root, consumer).events + + +def warn_unreadable(read: EventRead, stream=None, *, show_content: bool = False) -> bool: + """Report unreadable outbox lines on stderr. True when any were reported. + + stderr, never stdout: a --json consumer's stdout must stay parseable, and a + warning written there turns a health report into a parse failure. + + THE EXCERPT IS NOT PRINTED BY DEFAULT, and that is a deliberate call rather + than lost fidelity. `ordinal` and `reason` are STRUCTURAL -- a line number + and a parser complaint -- and carry no payload. The excerpt is CONTENT, and + printing it moves bytes out of a file the operator already owns into places + that get copied: CI logs, terminal scrollback, transcripts pasted into a + chat. Measured: a malformed line containing an API-key-shaped string echoed + that string verbatim. + + Redacting it instead was rejected. Matching "secret-looking" patterns in + arbitrary bytes is a denylist over untrusted input, which leaks by + construction -- the same defect shape this module's parse guard exists to + avoid. So the excerpt stays on MalformedLine, where a caller that needs it + can ask, and stays out of the default report. Pass show_content=True to + include it when you are debugging a specific file and know what is in it. + """ + if not read.malformed: + return False + out = stream if stream is not None else sys.stderr + count = len(read.malformed) + plural = "" if count == 1 else "s" + print( + f"warning: skipped {count} unreadable line{plural} in the event outbox", + file=out, + ) + for bad in read.malformed: + detail = f": {bad.excerpt}" if show_content else "" + print(f" line {bad.ordinal}: {bad.reason}{detail}", file=out) + return True def _load_cursor(workspace_root: Path, consumer: str) -> dict[str, object]: diff --git a/gr2/tests/test_events.py b/gr2/tests/test_events.py index eb78ede9..3eab6758 100644 --- a/gr2/tests/test_events.py +++ b/gr2/tests/test_events.py @@ -487,19 +487,27 @@ def test_emit_fails_closed_on_write_failure(self, workspace: Path): def test_emit_fails_closed_when_existing_sequence_cannot_be_read( self, workspace: Path, monkeypatch: pytest.MonkeyPatch ): - """Restoring _current_seq's OSError-to-zero fallback must turn this RED.""" + """Restoring _current_seq's OSError-to-zero fallback must turn this RED. + + MONKEYPATCH TARGET UPDATED read_text -> read_bytes (torn-line sweep fix 4). + The outbox is now read as BYTES so a single invalid byte cannot escape as + UnicodeDecodeError from outside every guard. This test's CONTRACT and every + assertion below are unchanged; only the read verb it intercepts moved, because + patching read_text no longer intercepts anything and the test would have gone + green by missing its target rather than by the behavior holding. + """ from gr2.python_cli.events import EventEmitError, EventType, _outbox_path, emit outbox = _outbox_path(workspace) outbox.write_text('{"seq":41}\n') - original_read_text = Path.read_text + original_read_bytes = Path.read_bytes def fail_for_outbox(path: Path, *args, **kwargs): if path == outbox: raise OSError("forced sequence-read failure") - return original_read_text(path, *args, **kwargs) + return original_read_bytes(path, *args, **kwargs) - monkeypatch.setattr(Path, "read_text", fail_for_outbox) + monkeypatch.setattr(Path, "read_bytes", fail_for_outbox) with pytest.raises(EventEmitError) as exc_info: emit( event_type=EventType.LANE_ENTERED, @@ -510,6 +518,10 @@ def fail_for_outbox(path: Path, *args, **kwargs): ) assert isinstance(exc_info.value.__cause__, OSError) + # Undo before the read-back: the patch now targets read_bytes, which is + # also how this assertion reads the file. A verification must not travel + # the path the test deliberately sabotaged. + monkeypatch.undo() assert outbox.read_bytes() == b'{"seq":41}\n' def test_emit_creates_events_dir_if_missing(self, workspace: Path): diff --git a/gr2/tests/test_events_torn_line.py b/gr2/tests/test_events_torn_line.py new file mode 100644 index 00000000..cc613074 --- /dev/null +++ b/gr2/tests/test_events_torn_line.py @@ -0,0 +1,592 @@ +"""Witnesses for the event outbox's torn-line contract and its unreadable-line count. + +Fix 4 of the torn-line sweep. Two parts, and they are separate claims: + + 1. TERMINATOR REPAIR -- an append onto a file whose last line lost its "\n" + must not GLUE two records together. The damage runs FORWARD from the tear: + the torn record and THE NEXT HEALTHY APPEND fuse into one unparseable + line, while the record before the tear is untouched. So a torn write costs + that record and the next one written after it, permanently, because every + later append builds on the glued line. + + 2. THE COUNT -- a line the reader cannot use is a LOST EVENT, and for the + channel bridge it is a message that never reaches a channel. It is + reported on EVERY read, not once: the cursor filter applies only to lines + that PARSE, so a line with no usable seq can never be advanced past, + wherever it sits -- position is irrelevant, and mid-file lines repeat + exactly as trailing ones do. It is reported from read_events_detailed() + only: _current_seq() runs once per emit inside the write lock, where a + count would be per-APPEND and aimed at whoever happened to be writing. + +Both statements above were WRONG in an earlier version of this file, in the +same words, while every test here passed. Prose is not covered by the suite +unless something pins the behavior it describes, so the witnesses at the end of +this file exist to pin exactly these two claims: which record the glue destroys, +and how long an unreadable line keeps being reported. +""" +from __future__ import annotations + +import io +import json +from pathlib import Path + +import pytest +from gr2.python_cli.channel_bridge import run_bridge +from gr2.python_cli.events import ( + EventType, + _current_seq, + _outbox_path, + emit, + read_events, + read_events_detailed, + warn_unreadable, +) + + +@pytest.fixture +def workspace(tmp_path: Path) -> Path: + (tmp_path / ".grip" / "events").mkdir(parents=True) + return tmp_path + + +def _emit(workspace: Path, **payload) -> None: + emit( + EventType.PROPAGATION_RECEIPT, + workspace, + actor="witness", + owner_unit="unit", + payload=payload or {"note": "n"}, + ) + + +def _lines(workspace: Path) -> list[bytes]: + raw = _outbox_path(workspace).read_bytes() + return [ln for ln in raw.split(b"\n") if ln.strip()] + + +# --- W1: terminator repair -------------------------------------------------- + +def test_append_onto_torn_line_does_not_glue_records(workspace: Path): + _emit(workspace, note="first") + outbox = _outbox_path(workspace) + # Tear the file exactly as a process killed mid-write leaves it. + blob = outbox.read_bytes() + assert blob.endswith(b"\n") + outbox.write_bytes(blob[:-1]) + + _emit(workspace, note="second") + + lines = _lines(workspace) + assert len(lines) == 2, "the append glued two records into one line" + for ln in lines: + json.loads(ln) # both must still parse + + +def test_torn_line_repair_survives_a_second_tear(workspace: Path): + """One repair is not a fix if the next tear re-breaks it.""" + _emit(workspace, note="a") + outbox = _outbox_path(workspace) + outbox.write_bytes(outbox.read_bytes()[:-1]) + _emit(workspace, note="b") + outbox.write_bytes(outbox.read_bytes()[:-1]) + _emit(workspace, note="c") + + lines = _lines(workspace) + assert len(lines) == 3 + notes = [json.loads(ln)["note"] for ln in lines] + assert notes == ["a", "b", "c"] + + +def test_no_spurious_terminator_on_a_fresh_outbox(workspace: Path): + """The repair must not write a leading blank line into an empty file.""" + _emit(workspace, note="only") + assert _outbox_path(workspace).read_bytes().count(b"\n") == 1 + + +def test_torn_line_does_not_lose_the_earlier_event_to_a_reader(workspace: Path): + """The point of the repair, stated as the consumer sees it.""" + _emit(workspace, note="earlier") + outbox = _outbox_path(workspace) + outbox.write_bytes(outbox.read_bytes()[:-1]) + _emit(workspace, note="later") + + read = read_events_detailed(workspace, "c") + notes = [e.get("note") for e in read.events] + assert notes == ["earlier", "later"] + assert read.malformed == () + + +# --- W2: undecodable bytes -------------------------------------------------- + +def test_invalid_utf8_does_not_brick_emit(workspace: Path): + """A single bad byte must not make the outbox permanently unwritable. + + _current_seq() runs inside emit(); if it raises, EVERY future emit fails. + That is worse than a crash -- it is a permanent denial with no recovery + path short of deleting the file. + """ + outbox = _outbox_path(workspace) + outbox.write_bytes(b'{"seq":1,"note":"ok"}\n\xff\xfe not utf-8\n') + _emit(workspace, note="after") # must not raise + assert any(b"after" in ln for ln in _lines(workspace)) + + +def test_invalid_utf8_is_counted_not_raised(workspace: Path): + outbox = _outbox_path(workspace) + outbox.write_bytes(b'{"seq":1,"note":"ok"}\n\xff\xfe\n{"seq":2,"note":"two"}\n') + read = read_events_detailed(workspace, "c") + assert [e["note"] for e in read.events] == ["ok", "two"] + assert len(read.malformed) == 1 + assert "utf-8" in read.malformed[0].reason.lower() + + +def test_seq_allocation_survives_an_undecodable_line(workspace: Path): + outbox = _outbox_path(workspace) + outbox.write_bytes(b'{"seq":7,"note":"ok"}\n\xff\n') + assert _current_seq(outbox) == 7 + + +# --- W3: hostile but syntactically valid content ---------------------------- + +@pytest.mark.parametrize( + "raw,why", + [ + (b'{"seq":1e999,"note":"inf"}', "float seq becomes inf and serializes as Infinity"), + (b'{"seq":"1","note":"str"}', "string seq cannot be ordered against the cursor"), + (b'{"seq":true,"note":"bool"}', "bool is an int subclass and max() accepts it silently"), + (b'{"seq":null,"note":"none"}', "null seq"), + (b'[1,2,3]', "a JSON array is not an event"), + (b'"just a string"', "a JSON string is not an event"), + (b'{"seq":1,', "truncated object"), + (b'not json at all', "not JSON"), + ], +) +def test_hostile_line_is_counted_and_never_raises(workspace: Path, raw: bytes, why: str): + outbox = _outbox_path(workspace) + outbox.write_bytes(b'{"seq":1,"note":"good"}\n' + raw + b"\n") + read = read_events_detailed(workspace, "c") + assert [e["note"] for e in read.events] == ["good"], why + assert len(read.malformed) == 1, why + assert _current_seq(outbox) == 1, why + + +def test_deeply_nested_json_is_counted_not_raised(workspace: Path): + """RecursionError is raised by json.loads and is NOT a ValueError. + + The previous guard caught (JSONDecodeError, TypeError) and would have let + this escape -- a denylist over untrusted input leaks by construction. + """ + outbox = _outbox_path(workspace) + bomb = b"[" * 20000 + b"]" * 20000 + outbox.write_bytes(b'{"seq":1,"note":"good"}\n' + bomb + b"\n") + read = read_events_detailed(workspace, "c") + assert [e["note"] for e in read.events] == ["good"] + assert len(read.malformed) == 1 + + +def test_emit_still_allocates_a_usable_seq_beside_hostile_lines(workspace: Path): + outbox = _outbox_path(workspace) + outbox.write_bytes(b'{"seq":3,"note":"good"}\n{"seq":1e999,"note":"inf"}\n') + _emit(workspace, note="next") + last = json.loads(_lines(workspace)[-1]) + assert last["seq"] == 4, "an unusable seq must not poison allocation" + json.dumps(last) # must remain serializable -- inf would emit `Infinity` + + +# --- W4: the count reaches a consumer --------------------------------------- + +def test_warn_unreadable_writes_to_the_given_stream(workspace: Path): + outbox = _outbox_path(workspace) + outbox.write_bytes(b'{"seq":1,"note":"ok"}\n\xff\n') + read = read_events_detailed(workspace, "c") + buf = io.StringIO() + assert warn_unreadable(read, stream=buf) is True + assert "skipped 1 unreadable line" in buf.getvalue() + + +def test_warn_unreadable_is_silent_when_clean(workspace: Path): + """The negative case. Without it, a reporter that always warns passes.""" + _emit(workspace, note="fine") + read = read_events_detailed(workspace, "c") + buf = io.StringIO() + assert warn_unreadable(read, stream=buf) is False + assert buf.getvalue() == "" + + +def test_bridge_reports_on_stderr_and_still_posts_good_events(workspace: Path, capsys): + outbox = _outbox_path(workspace) + outbox.write_bytes( + b'{"seq":1,"type":"propagation.receipt","note":"ok"}\n\xff\n' + ) + posted: list[str] = [] + count = run_bridge(workspace, post_fn=posted.append) + captured = capsys.readouterr() + assert "unreadable" in captured.err, "the count never reached a consumer" + assert captured.out == "", "a health warning on stdout breaks --json callers" + assert count == len(posted) + + +def test_bridge_is_silent_on_a_clean_outbox(workspace: Path, capsys): + _emit(workspace, note="clean") + run_bridge(workspace, post_fn=lambda _m: None) + assert "unreadable" not in capsys.readouterr().err + + +# --- W5: the cursor is untrusted input too ---------------------------------- + +def test_a_corrupt_cursor_does_not_brick_reads(workspace: Path): + """The cursor is a file on disk; a truncated write makes last_seq a non-int. + + Comparing seq against it would raise from inside a reader whose contract is + to not raise on file content. + """ + _emit(workspace, note="one") + cursors = workspace / ".grip" / "events" / "cursors" + cursors.mkdir(parents=True, exist_ok=True) + (cursors / "c.json").write_text(json.dumps({"consumer": "c", "last_seq": "not-an-int"})) + read = read_events_detailed(workspace, "c") + assert [e["note"] for e in read.events] == ["one"] + + +# --- the wrapper keeps its shape ------------------------------------------- + +def test_read_events_wrapper_still_returns_a_plain_list(workspace: Path): + """Eleven call sites index and len() this. The shape is the contract. + + All of them are tests: no production caller of read_events() remains once the + bridge moves to read_events_detailed(). The wrapper is a test-compatibility + surface, not a load-bearing API. + """ + _emit(workspace, note="a") + events = read_events(workspace, "c") + assert isinstance(events, list) + assert len(events) == 1 + assert events[0]["note"] == "a" + + +# --- W6: an I/O failure is NOT a malformed line ----------------------------- +# +# This boundary was NOT found by any witness above. It was found by an existing +# test in test_events.py whose docstring stands guard over it, after the first +# version of this fix swallowed OSError inside the line iterator. Content errors +# are data and get skipped-and-counted; an I/O error is not knowing what the file +# holds, and on the WRITE path that becomes a duplicate sequence number. + +def test_emit_fails_closed_when_the_outbox_cannot_be_read(workspace: Path, monkeypatch): + """Swallowing this returns seq 0, and emit then reuses live sequence numbers.""" + from gr2.python_cli.events import EventEmitError + + outbox = _outbox_path(workspace) + outbox.write_bytes(b'{"seq":41}\n') + original = Path.read_bytes + + def boom(path: Path, *a, **k): + if path == outbox: + raise OSError("forced sequence-read failure") + return original(path, *a, **k) + + monkeypatch.setattr(Path, "read_bytes", boom) + with pytest.raises(EventEmitError) as exc: + _emit(workspace, note="must not land") + assert isinstance(exc.value.__cause__, OSError) + # The read-back must not travel the sabotaged path -- my first version of + # this witness verified the file using the very method it had broken, and + # failed for that reason rather than for anything about emit(). + monkeypatch.undo() + assert outbox.read_bytes() == b'{"seq":41}\n', "a failed emit must not mutate the log" + + +def test_reader_tolerates_the_outbox_vanishing_mid_read(workspace: Path, monkeypatch): + """_maybe_rotate() renames the outbox, so this race is real and benign.""" + outbox = _outbox_path(workspace) + _emit(workspace, note="a") + original = Path.read_bytes + + def vanish(path: Path, *a, **k): + if path == outbox: + raise FileNotFoundError("rotated away") + return original(path, *a, **k) + + monkeypatch.setattr(Path, "read_bytes", vanish) + read = read_events_detailed(workspace, "c") + assert read.events == [] + assert read.malformed == () + + +def test_reader_propagates_a_real_io_error(workspace: Path, monkeypatch): + """The discriminating control for the case above. + + Without this, catching FileNotFoundError could widen to OSError and a reader + would report 'no new events' forever while the disk was failing. + """ + outbox = _outbox_path(workspace) + _emit(workspace, note="a") + original = Path.read_bytes + + def eio(path: Path, *a, **k): + if path == outbox: + raise OSError("EIO") + return original(path, *a, **k) + + monkeypatch.setattr(Path, "read_bytes", eio) + with pytest.raises(OSError): + read_events_detailed(workspace, "c") + + +# --- W7: the report must not amplify file content into logs ----------------- + +def test_default_report_does_not_echo_line_content(workspace: Path): + """stderr gets copied -- into CI logs, scrollback, pasted transcripts. + + Found by asking what this new output path could carry, not by a failing + test. `reason` is structural; the excerpt is content, and it is available on + the data object for callers that genuinely need it. + """ + outbox = _outbox_path(workspace) + outbox.write_bytes(b'{"seq":1,"note":"ok"}\n{"token":"sk-NOTREAL-abcdef","broken\n') + read = read_events_detailed(workspace, "c") + buf = io.StringIO() + warn_unreadable(read, stream=buf) + assert "NOTREAL" not in buf.getvalue() + assert "line 2" in buf.getvalue(), "the line must still be identified" + assert "NOTREAL" in read.malformed[0].excerpt, "fidelity kept on the data object" + + +def test_show_content_opt_in_still_works(workspace: Path): + """The discriminating control: without it, an always-redacting reporter passes.""" + outbox = _outbox_path(workspace) + outbox.write_bytes(b'{"seq":1,"note":"ok"}\n{"token":"sk-NOTREAL-abcdef","broken\n') + read = read_events_detailed(workspace, "c") + buf = io.StringIO() + warn_unreadable(read, stream=buf, show_content=True) + assert "NOTREAL" in buf.getvalue() + + +# --- W8: this module must import WITHOUT sys.modules registration ----------- + +def test_events_module_loads_under_an_out_of_tree_loader(tmp_path: Path): + """Spawned workers load events.py by path, without registering it. + + importlib.util.spec_from_file_location() + exec_module() leaves the module + OUT of sys.modules, and @dataclass resolves its field types through + sys.modules[cls.__module__].__dict__ -- so a dataclass here raises + AttributeError AT IMPORT and every worker dies before running its own code. + + That is not hypothetical: adding @dataclass to this module killed both + writers in the concurrent-emit integrity test before either reached + sequence allocation. This witness exists so the next person to reach for a + dataclass here finds out in one second instead of in a concurrency test + whose failure looks like flakiness. + """ + import importlib.util + + events_path = Path(__file__).resolve().parents[1] / "python_cli" / "events.py" + spec = importlib.util.spec_from_file_location("events_out_of_tree_probe", events_path) + assert spec is not None and spec.loader is not None + module = importlib.util.module_from_spec(spec) + # DELIBERATELY NOT registered in sys.modules -- that is the whole point. + spec.loader.exec_module(module) + + assert hasattr(module, "read_events_detailed") + assert module.MalformedLine(1, "why", "x").ordinal == 1 + assert module.EventRead([], ()).malformed == () + + +# --- W9: how long an unreadable line keeps being reported -------------------- +# +# These exist because a REVIEWER measured this and the PR body claimed the +# opposite, with the whole suite green. Nothing pinned the real behavior, so +# prose and code were free to disagree. The reviewer found the TERMINAL case; +# measuring the mid-file case showed the same thing, so the rule is simpler and +# broader than the finding that produced it. + +def test_unreadable_line_is_reported_on_every_read_terminal(workspace: Path): + """A trailing unreadable line is counted again on the next read.""" + _emit(workspace, note="good") + outbox = _outbox_path(workspace) + outbox.write_bytes(outbox.read_bytes() + b"\xff not utf-8\n") + + first = read_events_detailed(workspace, "c") + second = read_events_detailed(workspace, "c") + assert [e["note"] for e in first.events] == ["good"] + assert len(first.malformed) == 1 + assert second.events == [], "the good event was consumed, as expected" + assert len(second.malformed) == 1, "the unreadable line is reported again" + + +def test_unreadable_line_is_reported_on_every_read_midfile(workspace: Path): + """MID-FILE too, which is the part the terminal framing would hide. + + The cursor filter `seq <= last_seq` only applies to lines that PARSE. A line + with no usable seq can never be advanced past, wherever it sits -- so + position is irrelevant and "trailing" is not the distinguishing property. + """ + outbox = _outbox_path(workspace) + outbox.write_bytes(b'{"seq":1,"note":"a"}\n\xff\n{"seq":2,"note":"b"}\n') + + first = read_events_detailed(workspace, "c") + second = read_events_detailed(workspace, "c") + assert [e["note"] for e in first.events] == ["a", "b"] + assert len(first.malformed) == 1 + assert second.events == [], "both good events were consumed" + assert len(second.malformed) == 1, "a mid-file unreadable line repeats too" + + +def test_glue_destroys_the_next_append_not_the_preceding_record(workspace: Path): + """The measured direction of the damage, which the body stated backwards. + + Tear after record B, then append C WITHOUT the repair. A is untouched; B and + C fuse into one unparseable line. The damage runs FORWARD from the tear. + """ + _emit(workspace, note="A") + _emit(workspace, note="B") + outbox = _outbox_path(workspace) + outbox.write_bytes(outbox.read_bytes()[:-1]) # tear after B + # bypass emit() to reproduce the pre-fix append + outbox.write_bytes( + outbox.read_bytes() + json.dumps({"seq": 3, "note": "C"}).encode("utf-8") + b"\n" + ) + + lines = _lines(workspace) + assert len(lines) == 2, "B and C fused into one line" + assert json.loads(lines[0])["note"] == "A", "the record BEFORE the tear survives" + with pytest.raises(ValueError): + json.loads(lines[1]) # B+C, unparseable + + +# --- W10: the TRUNCATED tear, which no fixture ABOVE this line exercises ----- +# +# "TORN" IS TWO DIFFERENT FAILURES WITH OPPOSITE OUTCOMES, AND USING THE BARE +# WORD IS WHAT MADE FOUR SEPARATE DESCRIPTIONS OF THIS CODE WRONG: +# +# UNTERMINATED -- the record is COMPLETE and only its "\n" was lost. Seam +# repair ends the line and it parses. RECOVERS FULLY, zero malformed. +# TRUNCATED -- the write stopped mid-record. Seam repair stops the next append +# being glued on, and the record itself NEVER becomes readable, because its +# bytes were never written. +# +# Every terminator fixture ABOVE this line uses read_bytes()[:-1] -- the +# UNTERMINATED case, and the lucky one. The fixtures below are the truncated +# case. That gap is why prose describing this code kept being wrong: the only +# tear ever exercised was the one that recovers. +# +# Measured: unterminated -> 3 events, 0 malformed. Truncated -> 2 events, 1 +# malformed, and still 1 on the second and third read. +# +# What the repair actually buys, stated so it holds for BOTH cases: it cannot +# recover a record that was never fully written, and it saves the NEXT one. +# Without repair you lose two records; with it you lose the one the crash +# truncated -- and if the crash truncated nothing, you lose none. + +def _tear_mid_record(outbox: Path) -> None: + """Truncate the last record mid-JSON, as a partial write leaves it.""" + lines = outbox.read_bytes().split(b"\n") + body = [ln for ln in lines if ln.strip()] + body[-1] = body[-1][:20] + outbox.write_bytes(b"\n".join(body)) # no trailing newline either + + +def test_realistic_tear_saves_the_next_append_but_not_the_torn_record(workspace: Path): + _emit(workspace, note="A") + _emit(workspace, note="B") + outbox = _outbox_path(workspace) + _tear_mid_record(outbox) + + _emit(workspace, note="C") + + read = read_events_detailed(workspace, "c") + notes = [e.get("note") for e in read.events] + assert notes == ["A", "C"], "the NEXT append survives; the truncated record cannot" + assert len(read.malformed) == 1, "the truncated record is unreadable, permanently" + + +def test_a_truncated_record_never_heals(workspace: Path): + """A TRUNCATED record is reported on every read, forever; no emit repairs it. + + THE SEAM HEALS; A TRUNCATED RECORD DOES NOT. The next append is no longer + glued on, which is the whole of what the repair achieves for this case -- + the truncated record's bytes were never written and nothing can reconstruct + them. An UNTERMINATED record is the other case entirely and does recover; + its witness is directly below and the two are each other's control. + """ + _emit(workspace, note="A") + _emit(workspace, note="B") + outbox = _outbox_path(workspace) + _tear_mid_record(outbox) + _emit(workspace, note="C") + + counts = [len(read_events_detailed(workspace, "c").malformed) for _ in range(3)] + assert counts == [1, 1, 1], f"expected a permanent malformed line, got {counts}" + + +def test_unterminated_tear_DOES_fully_recover(workspace: Path): + """The discriminating control, and it has already earned its keep twice. + + Without it the two witnesses above would also pass against an implementation + that simply never recovers anything, and the truncated case would read as the + norm. It then caught its own author: a draft of the PR body claimed a torn + record never becomes readable, and THIS witness disproves that. An + UNTERMINATED record -- complete, only its newline lost -- must come back with + zero malformed lines. + """ + _emit(workspace, note="A") + _emit(workspace, note="B") + outbox = _outbox_path(workspace) + outbox.write_bytes(outbox.read_bytes()[:-1]) # newline only + + _emit(workspace, note="C") + + read = read_events_detailed(workspace, "c") + assert [e.get("note") for e in read.events] == ["A", "B", "C"] + assert read.malformed == () + + +# --- W11: the READER against a torn file, with NO intervening emit ---------- +# +# Every tear fixture above emits after tearing, so until this pair nothing ever +# asked what the READER alone does with a torn file. That gap is exactly how a +# false sentence survived into the bridge comment: it said an unterminated +# record is reported and then heals at the next emit, implying a window in +# which it is unreadable. There is no such window. _iter_outbox() splits on +# b"\n", so a COMPLETE record whose only loss was its terminator is simply the +# final chunk, and it parses -- repair or no repair, emit or no emit. +# +# The truncated case below is the discriminating control: without it this pair +# would also pass against a reader that never reported anything at all. + + +def test_unterminated_record_is_never_reported_even_without_an_emit(workspace: Path): + """A complete record missing only its "\n" is readable immediately. + + Not "recovers at the next emit" -- never unreadable in the first place. + Read twice, because a claim of transience would show as a difference + between the reads. + """ + _emit(workspace, note="A") + _emit(workspace, note="B") + outbox = _outbox_path(workspace) + outbox.write_bytes(outbox.read_bytes()[:-1]) # newline only; no emit follows + + for i in (1, 2): + read = read_events_detailed(workspace, f"c{i}") + assert [e.get("note") for e in read.events] == ["A", "B"], f"read {i}" + assert read.malformed == (), f"read {i}: expected nothing reported, got {read.malformed}" + + +def test_truncated_record_IS_reported_without_an_emit(workspace: Path): + """The discriminating control for the witness above. + + Byte-identical setup except the tear cuts mid-record instead of at the + terminator. If this one also came back empty, the witness above would be + measuring a reader that reports nothing rather than a record that is + readable. + """ + _emit(workspace, note="A") + _emit(workspace, note="B") + outbox = _outbox_path(workspace) + _tear_mid_record(outbox) # no emit follows + + for i in (1, 2): + read = read_events_detailed(workspace, f"c{i}") + assert [e.get("note") for e in read.events] == ["A"], f"read {i}" + assert len(read.malformed) == 1, f"read {i}: expected the truncated record reported"