diff --git a/CHANGELOG.md b/CHANGELOG.md index 57f63d3..c33d98a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,3 +1,7 @@ # Changelog ## 0.1.0 — initial public release + +- Quarantine every contributing record when a recognized credential spans same-lineage record boundaries, including zero-content gaps. +- Bound aggregate cross-record detector state with fail-closed, observable LRU eviction. +- Pair tool results with full, collision-resistant lineage and call identities. diff --git a/src/claude_code_memory/transcript.py b/src/claude_code_memory/transcript.py index 591dd37..8f55221 100644 --- a/src/claude_code_memory/transcript.py +++ b/src/claude_code_memory/transcript.py @@ -4,6 +4,7 @@ import hashlib import json +from collections import OrderedDict from collections.abc import Iterator from pathlib import Path from typing import Any @@ -11,6 +12,10 @@ from substrate_capture import normalize_message from substrate_capture.redaction import ( configured_secret_values, + credential_content_has_open_continuation, + credential_continuation_prefix, + credential_detector_overlap_chars, + credential_redaction_spans, iter_redacted_text_chunks, redact_text, ) @@ -21,6 +26,51 @@ MAX_MESSAGE_CHARS = MAX_BLOCK_BYTES _MAX_LINE_BYTES = 1024 * 1024 _SCAN_CHARS = 16 * 1024 +_MAX_OVERLAP_TOTAL_BYTES = 8 * 1024 * 1024 +_MAX_OVERLAP_LINEAGES = 64 +_TailSegments = list[tuple[int, int, list[dict[str, Any]]]] +_TailState = tuple[str, _TailSegments, bool, int] + + +class _LineageTailBudget: + """LRU-bounded raw overlap state with content-free eviction accounting.""" + + def __init__(self, *, max_bytes: int, max_lineages: int) -> None: + self.max_bytes = max_bytes + self.max_lineages = max_lineages + self.total_bytes = 0 + self.evicted_count = 0 + self._states: OrderedDict[str, _TailState] = OrderedDict() + + def get(self, key: str) -> _TailState | None: + state = self._states.get(key) + if state is not None: + self._states.move_to_end(key) + return state + + def preserve(self, key: str) -> None: + """Keep pending lexical context across a zero-visible-content record.""" + if key in self._states: + self._states.move_to_end(key) + + def set(self, key: str, state: _TailState) -> list[_TailState]: + previous = self._states.pop(key, None) + if previous is not None: + self.total_bytes -= previous[3] + self._states[key] = state + self.total_bytes += state[3] + evicted: list[_TailState] = [] + while ( + self.total_bytes > self.max_bytes + or len(self._states) > self.max_lineages + ): + _evicted_key, evicted_state = self._states.popitem(last=False) + self.total_bytes -= evicted_state[3] + self.evicted_count += 1 + evicted.append(evicted_state) + return evicted + + _TEXT_BLOCK_TYPES = frozenset({"text", "input_text", "output_text"}) _BINARY_KEYS = frozenset( { @@ -160,6 +210,24 @@ def _redacted_bound(value: Any, secrets: tuple[str, ...], maximum: int = 512) -> return redact_text(str(value or ""), secrets)[:maximum] +def _identity_digest(*values: Any) -> str: + """Return a collision-resistant key for full, unbounded capture identities.""" + digest = hashlib.sha256() + for value in values: + encoded = str(value or "").encode("utf-8") + digest.update(len(encoded).to_bytes(8, "big")) + digest.update(encoded) + return digest.hexdigest() + + +def _lineage_identity(record: dict[str, Any]) -> str: + return _identity_digest( + record.get("sessionId"), + record.get("agentId"), + record.get("isSidechain") is True, + ) + + def _coordinates( record: dict[str, Any], record_index: int, @@ -191,6 +259,7 @@ def _envelope( source_protocol: str = SOURCE_PROTOCOL_TRANSCRIPT, source_identity: str | None = None, attribution_reason_code: str | None = None, + pairing_tool_call_id: str | None = None, quarantine: bool = False, ) -> dict[str, Any]: safe = _safe_block(content, secrets) @@ -216,6 +285,14 @@ def _envelope( "observed_at": record.get("timestamp"), "platform_message_id": record.get("uuid"), **_coordinates(record, record_index, block_index, secrets), + # Internal full-identity keys are removed before normalization. Display + # bounds must never define trust or capture-stream membership. + "_pairing_lineage": _lineage_identity(record), + "_pairing_tool_call_id": ( + _identity_digest(pairing_tool_call_id) + if pairing_tool_call_id and not quarantine + else None + ), } if attribution_reason_code: message["attribution_reason_code"] = attribution_reason_code @@ -324,6 +401,7 @@ def _parse_record( source_protocol=SOURCE_PROTOCOL_TOOL, source_identity=name, attribution_reason_code=(None if call_id else "missing_tool_call_id"), + pairing_tool_call_id=raw_call_id, quarantine=block_quarantine, ) ) @@ -345,21 +423,16 @@ def _parse_record( tool_call_id=result_id, source_protocol=SOURCE_PROTOCOL_TOOL, attribution_reason_code="pairing_pending", + pairing_tool_call_id=raw_result_id, quarantine=block_quarantine, ) ) return parsed -def _lineage_key(message: dict[str, Any]) -> tuple[str, str, bool]: - ancestry = message.get("session_ancestry") - if not isinstance(ancestry, dict): - return ("", "", False) - return ( - str(ancestry.get("session_id") or ""), - str(ancestry.get("agent_id") or ""), - ancestry.get("is_sidechain") is True, - ) +def _lineage_key(message: dict[str, Any]) -> str: + value = message.get("_pairing_lineage") + return value if isinstance(value, str) else "" def _clear_tool_result_identity(message: dict[str, Any], reason: str) -> None: @@ -368,15 +441,16 @@ def _clear_tool_result_identity(message: dict[str, Any], reason: str) -> None: message["tool_call_id"] = None message["tool_name"] = None message["source_identity"] = None + message["_pairing_tool_call_id"] = None message["attribution_reason_code"] = reason def _pair_tool_results(messages: list[dict[str, Any]]) -> None: - calls: dict[tuple[tuple[str, str, bool], str], list[dict[str, Any]]] = {} + calls: dict[tuple[str, str], list[dict[str, Any]]] = {} global_calls: dict[str, list[dict[str, Any]]] = {} - results: dict[tuple[tuple[str, str, bool], str], list[dict[str, Any]]] = {} + results: dict[tuple[str, str], list[dict[str, Any]]] = {} for message in messages: - call_id = message.get("tool_call_id") + call_id = message.get("_pairing_tool_call_id") if not isinstance(call_id, str) or not call_id: continue key = (_lineage_key(message), call_id) @@ -389,7 +463,7 @@ def _pair_tool_results(messages: list[dict[str, Any]]) -> None: for message in messages: if message.get("role") != "tool_result": continue - call_id = message.get("tool_call_id") + call_id = message.get("_pairing_tool_call_id") if not isinstance(call_id, str) or not call_id: _clear_tool_result_identity(message, "missing_tool_call_id") continue @@ -416,6 +490,61 @@ def _pair_tool_results(messages: list[dict[str, Any]]) -> None: _clear_tool_result_identity(message, reason) +def _record_content(record: dict[str, Any]) -> str: + message = record.get("message") + if not isinstance(message, dict): + return "" + content = message.get("content") + if isinstance(content, str): + return content + if not isinstance(content, list): + return "" + return "".join( + piece + for piece in (_block_content(block) for block in content) + if piece is not None + ) + + +def _credential_crossing_spans( + previous_tail: str, current_content: str, secrets: tuple[str, ...] +) -> list[tuple[int, int]]: + """Locate full-suite credential matches completed across the next record.""" + if not previous_tail or not current_content: + return [] + boundary = len(previous_tail) + return [ + (start, end) + for start, end in credential_redaction_spans( + previous_tail + current_content, secrets + ) + if start < boundary < end + ] + + +def _quarantine_record( + messages: list[dict[str, Any]], *, code: str = "credential_detected" +) -> None: + for message in messages: + existing_codes = message.get("redaction_codes") + redaction_codes = ( + existing_codes + if isinstance(existing_codes, list) and "credential_detected" in existing_codes + else [code] + ) + message.update( + content="", + retained_bytes=0, + redaction_codes=redaction_codes, + truncated=False, + content_digest=hashlib.sha256(b"").hexdigest(), + tool_call_id=None, + tool_name=None, + source_identity=None, + _pairing_tool_call_id=None, + ) + + def _discard_line_remainder(stream: Any, chunk: bytes) -> None: while chunk and not chunk.endswith(b"\n"): chunk = stream.readline(_MAX_LINE_BYTES + 1) @@ -431,6 +560,14 @@ def read_messages(path: str | Path, *, include_sidechains: bool = True) -> list[ parsed: list[dict[str, Any]] = [] try: secrets = configured_secret_values() + overlap_chars = credential_detector_overlap_chars() + # Each lineage owns a detector-sized suffix. The LRU additionally caps + # aggregate raw state; eviction quarantines referenced records before + # dropping their context, and increments only a content-free counter. + stream_tails = _LineageTailBudget( + max_bytes=_MAX_OVERLAP_TOTAL_BYTES, + max_lineages=_MAX_OVERLAP_LINEAGES, + ) with Path(path).open("rb") as stream: record_index = 0 while len(parsed) < MAX_MESSAGES: @@ -455,14 +592,88 @@ def read_messages(path: str | Path, *, include_sidechains: bool = True) -> list[ if record.get("isSidechain") is True and not include_sidechains: record_index += 1 continue + current_messages: list[dict[str, Any]] = [] for message in _parse_record(record, record_index, secrets): if len(parsed) >= MAX_MESSAGES: break parsed.append(message) + current_messages.append(message) + current_content = _record_content(record) + stream_key = _lineage_identity(record) + previous_state = stream_tails.get(stream_key) + if not current_content: + stream_tails.preserve(stream_key) + record_index += 1 + continue + if previous_state is None: + previous_tail, previous_segments, previous_continuation = ( + "", + [], + False, + ) + else: + previous_tail, previous_segments, previous_continuation, _ = ( + previous_state + ) + combined = previous_tail + current_content + continuation_active = False + if previous_continuation: + continuation_chars, continuation_active = ( + credential_continuation_prefix(current_content) + ) + if continuation_chars: + _quarantine_record(current_messages) + crossing_spans = _credential_crossing_spans( + previous_tail, current_content, secrets + ) + for start, end in crossing_spans: + for segment_start, segment_end, segment_messages in previous_segments: + if segment_start < end and segment_end > start: + _quarantine_record(segment_messages) + _quarantine_record(current_messages) + current_detected = any( + "credential_detected" in message.get("redaction_codes", []) + for message in current_messages + ) + if ( + previous_continuation or crossing_spans or current_detected + ) and credential_content_has_open_continuation(combined, secrets): + continuation_active = True + boundary = len(previous_tail) + segments = previous_segments + [ + (boundary, len(combined), current_messages) + ] + cutoff = max(0, len(combined) - overlap_chars) + tail = combined[cutoff:] + tail_state: _TailState = ( + tail, + [ + ( + max(segment_start, cutoff) - cutoff, + segment_end - cutoff, + segment_messages, + ) + for segment_start, segment_end, segment_messages in segments + if segment_end > cutoff + ], + continuation_active, + len(tail.encode("utf-8")), + ) + for evicted_state in stream_tails.set(stream_key, tail_state): + seen_records: set[int] = set() + for _start, _end, evicted_messages in evicted_state[1]: + identity = id(evicted_messages) + if identity not in seen_records: + seen_records.add(identity) + _quarantine_record( + evicted_messages, code="overlap_state_evicted" + ) record_index += 1 _pair_tool_results(parsed) normalized: list[dict[str, Any]] = [] for index, message in enumerate(parsed): + message.pop("_pairing_lineage", None) + message.pop("_pairing_tool_call_id", None) item = normalize_message(message, index=index, secrets=secrets) if item is not None: normalized.append(item) diff --git a/src/substrate_capture/_vendor.json b/src/substrate_capture/_vendor.json index 10e0f44..a0bce37 100644 --- a/src/substrate_capture/_vendor.json +++ b/src/substrate_capture/_vendor.json @@ -8,7 +8,7 @@ "events.py": "fb1bb57e52a804f9a17a48b8657737d6c761c1db9c839abd7a2eb76e328bd6ff", "mcp.py": "50b2f046792ed7da20644c167a1bd38da464d2025ac1a0949599f2e5ac90097e", "py.typed": "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855", - "redaction.py": "e9bec198aa7ad018da359d2e9aa6df1dab717881bb41b1001348911b23e6439b", + "redaction.py": "40b57020c3957d1dda72e0b36fa81cf0b40598a968166d8fb1c91729b91016d3", "spool.py": "8ab56569e61da99f7fd8e1f539f4ecf3cdbf68cce5f50406212f5da4929aa977", "tools.py": "1e4887131c95d10e99550de4a3d4f38953661adbceb1fbec26bd5986552f5902" }, diff --git a/src/substrate_capture/redaction.py b/src/substrate_capture/redaction.py index f12b230..aaad325 100644 --- a/src/substrate_capture/redaction.py +++ b/src/substrate_capture/redaction.py @@ -20,12 +20,13 @@ _MAX_CONFIGURED_SECRETS = 64 _MAX_CONFIGURED_SECRET_BYTES = 64 * 1024 _MAX_CONFIGURED_SECRET_CHARS = 16 * 1024 -# The longest bounded fixed-pattern match is the orphan/private-key form: up to -# 128 lines of 4,096 base64 characters plus indentation and delimiters. Keep a -# conservative overlap so a match cannot straddle released text. The overlap -# is independent of total message size and stays comfortably below the import -# worker's 256 MiB RSS ceiling even for wide Unicode strings. -_STREAM_REDACTION_OVERLAP_CHARS = 640 * 1024 +# The longest bounded match is the orphan/private-key form. Each of its 128 +# lines can contain 32 leading spaces, the 10-character ``Proc-Type:`` label, +# 4,096 value characters, 32 trailing spaces, and CRLF. The longest footer is +# the 35-character encrypted-key form. Open-ended lexical rules use explicit +# continuation state instead of increasing this retained suffix. +_PRIVATE_KEY_ORPHAN_MAX_CHARS = 128 * (32 + 10 + 4_096 + 32 + 2) + 35 +_STREAM_REDACTION_OVERLAP_CHARS = _PRIVATE_KEY_ORPHAN_MAX_CHARS _STREAM_REDACTION_SCAN_CHARS = 256 * 1024 _STREAM_CONTINUATION_TERMINATORS = frozenset("\r\n\t ,;}]&#@'\"") _PRIORITY_SECRET_KEYS = ( @@ -266,6 +267,37 @@ def _redaction_spans(value: str, secrets: Sequence[str]) -> list[tuple[int, int] return merged +def credential_detector_overlap_chars() -> int: + """Return the shared bounded overlap required by the full detector suite.""" + return _STREAM_REDACTION_OVERLAP_CHARS + + +def credential_redaction_spans( + value: str, secrets: Sequence[str] = () +) -> list[tuple[int, int]]: + """Return raw spans matched by configured and pattern-based credential rules.""" + return _redaction_spans(value, secrets) + + +def credential_continuation_prefix(value: str) -> tuple[int, bool]: + """Return credential-continuation chars and whether continuation remains open.""" + for index, character in enumerate(value): + if character.isspace() or character in _STREAM_CONTINUATION_TERMINATORS: + return index, False + return len(value), True + + +def _span_can_continue( + value: str, start: int, end: int, secrets: Sequence[str] +) -> bool: + if end != len(value): + return False + return any( + extended_start <= start and extended_end > end + for extended_start, extended_end in _redaction_spans(value + "A", secrets) + ) + + class StreamingTextRedactor: """Bounded deterministic redaction for a repeatable streamed text source. @@ -353,6 +385,23 @@ def finish(self) -> Iterator[str]: yield rendered +def credential_content_has_open_continuation( + value: str, secrets: Sequence[str] = () +) -> bool: + """Return whether a detected credential lexically continues past ``value``.""" + if not value: + return False + redactor = StreamingTextRedactor(secrets) + for _rendered in redactor.feed(value): + pass + if redactor._suppress_continuation: + return True + return any( + _span_can_continue(redactor._buffer, start, end, secrets) + for start, end in _redaction_spans(redactor._buffer, secrets) + ) + + def iter_redacted_text_chunks( chunks: Iterable[str], secrets: Sequence[str] = (), diff --git a/tests/test_review185_regressions.py b/tests/test_review185_regressions.py index 2a46986..9ce7595 100644 --- a/tests/test_review185_regressions.py +++ b/tests/test_review185_regressions.py @@ -5,6 +5,7 @@ import multiprocessing import os import sys +import time from pathlib import Path import pytest @@ -61,6 +62,218 @@ def test_secret_split_across_inner_tool_result_text_blocks_must_not_survive(tmp_ assert message["content"] == "" +def test_secret_split_across_adjacent_records_in_same_stream_is_quarantined( + tmp_path, monkeypatch +): + monkeypatch.setenv("ANTHROPIC_API_KEY", SYNTHETIC_CREDENTIAL) + left, right = SYNTHETIC_CREDENTIAL[:13], SYNTHETIC_CREDENTIAL[13:] + records = [ + { + "type": "user", + "sessionId": "session", + "agentId": "side-A", + "isSidechain": True, + "message": { + "role": "user", + "content": [ + {"type": "tool_result", "tool_use_id": "left", "content": left} + ], + }, + }, + { + "type": "user", + "sessionId": "session", + "agentId": "side-B", + "isSidechain": True, + "message": { + "role": "user", + "content": [ + { + "type": "tool_result", + "tool_use_id": "unrelated", + "content": "safe", + } + ], + }, + }, + { + "type": "user", + "sessionId": "session", + "agentId": "side-A", + "isSidechain": True, + "message": { + "role": "user", + "content": [ + {"type": "tool_result", "tool_use_id": "right", "content": right} + ], + }, + }, + ] + path = tmp_path / "record-split.jsonl" + path.write_text("".join(json.dumps(record) + "\n" for record in records)) + + messages = read_messages(path) + side_a = [ + message + for message in messages + if message["session_ancestry"]["agent_id"] == "side-A" + ] + side_b = [ + message + for message in messages + if message["session_ancestry"]["agent_id"] == "side-B" + ] + assert len(side_a) == 2 and all(message["content"] == "" for message in side_a) + assert side_b[0]["content"] == "safe" + assert SYNTHETIC_CREDENTIAL not in "".join(message["content"] for message in messages) + + +def test_zero_visible_content_records_preserve_same_lineage_overlap(tmp_path, monkeypatch): + secret = "test-unicode-é🙂-record-secret-987654" + monkeypatch.setenv("ANTHROPIC_API_KEY", secret) + left, right = secret[:12], secret[12:] + media = [{"type": "image", "source": {"type": "base64", "data": "AAAA"}}] + records = [ + { + "type": "user", + "sessionId": "session", + "agentId": "side-A", + "isSidechain": True, + "message": { + "role": "user", + "content": [ + {"type": "tool_result", "tool_use_id": "left", "content": left} + ], + }, + }, + { + "type": "user", + "sessionId": "session", + "agentId": "side-A", + "isSidechain": True, + "message": {"role": "user", "content": ""}, + }, + { + "type": "user", + "sessionId": "session", + "agentId": "side-A", + "isSidechain": True, + "message": {"role": "user", "content": media}, + }, + { + "type": "user", + "sessionId": "session", + "agentId": "side-A", + "isSidechain": True, + "message": { + "role": "user", + "content": [ + {"type": "tool_result", "tool_use_id": "right", "content": right} + ], + }, + }, + ] + path = tmp_path / "zero-visible-gap.jsonl" + path.write_text("".join(json.dumps(record, ensure_ascii=False) + "\n" for record in records)) + + messages = read_messages(path) + assert len(messages) == 3 + assert messages[0]["content"] == "" and messages[-1]["content"] == "" + serialized = json.dumps(messages, ensure_ascii=False) + assert left not in serialized and right not in serialized + + monkeypatch.delenv("ANTHROPIC_API_KEY") + token = "ghp_" + "A" * 40 + provider_left, provider_right = token[:3], token[3:] + provider_records = [records[0], records[2], records[3]] + provider_records[0]["message"]["content"][0]["content"] = provider_left + provider_records[-1]["message"]["content"][0]["content"] = provider_right + provider_path = tmp_path / "media-only-gap.jsonl" + provider_path.write_text( + "".join(json.dumps(record) + "\n" for record in provider_records) + ) + provider_messages = read_messages(provider_path) + assert len(provider_messages) == 2 + assert all(message["content"] == "" for message in provider_messages) + + +@pytest.mark.parametrize( + ("split_points", "record_count"), + [ + ((3,), 2), + ((3, 17), 3), + ((2, 4, 18), 4), + ], +) +def test_provider_token_split_across_records_is_quarantined( + tmp_path, monkeypatch, split_points, record_count +): + monkeypatch.delenv("ANTHROPIC_API_KEY", raising=False) + token = "ghp_" + "A" * 40 + points = (0, *split_points, len(token)) + pieces = [token[start:end] for start, end in zip(points, points[1:])] + records = [ + { + "type": "user", + "sessionId": "session", + "agentId": "side-A", + "isSidechain": True, + "message": { + "role": "user", + "content": [ + { + "type": "tool_result", + "tool_use_id": f"result-{index}", + "content": piece, + } + ], + }, + } + for index, piece in enumerate(pieces) + ] + path = tmp_path / f"provider-{record_count}.jsonl" + path.write_text("".join(json.dumps(record) + "\n" for record in records)) + + messages = read_messages(path) + assert len(messages) == record_count + assert all(message["content"] == "" for message in messages) + assert all(message["redaction_codes"] == ["credential_detected"] for message in messages) + + +def test_bearer_credential_split_across_records_is_quarantined(tmp_path): + credential = "Bearer " + "Z" * 24 + points = (0, 3, 8, len(credential)) + pieces = [ + credential[start:end] for start, end in zip(points, points[1:]) + ] + records = [ + { + "type": "user", + "sessionId": "session", + "agentId": "side-A", + "isSidechain": True, + "message": { + "role": "user", + "content": [ + { + "type": "tool_result", + "tool_use_id": f"bearer-{index}", + "content": piece, + } + ], + }, + } + for index, piece in enumerate(pieces) + ] + path = tmp_path / "bearer-split.jsonl" + path.write_text("".join(json.dumps(record) + "\n" for record in records)) + + messages = read_messages(path) + assert len(messages) == 3 + assert all(message["content"] == "" for message in messages) + assert all(message["redaction_codes"] == ["credential_detected"] for message in messages) + + def test_tool_use_input_full_quarantine_passes(tmp_path, monkeypatch): monkeypatch.setenv("ANTHROPIC_API_KEY", SYNTHETIC_CREDENTIAL) block = [{"type": "tool_use", "id": "call", "name": "Fetch", "input": {"x": "safe", "arg": SYNTHETIC_CREDENTIAL, "tail": "safe"}}] @@ -130,12 +343,88 @@ def test_cross_sidechain_result_does_not_pair_with_root_call(tmp_path): assert result["attribution_reason_code"] == "cross_stream_tool_result" +def test_full_agent_id_prevents_display_prefix_lineage_collision(tmp_path): + shared_prefix = "agent-" + "x" * 506 + trusted_agent = shared_prefix + "-trusted" + forged_agent = shared_prefix + "-forged" + assert trusted_agent[:512] == forged_agent[:512] + records = [ + { + "type": "assistant", + "sessionId": "same-session", + "agentId": trusted_agent, + "isSidechain": True, + "message": { + "role": "assistant", + "content": [ + { + "type": "tool_use", + "id": "shared", + "name": "TrustedTool", + "input": {}, + } + ], + }, + }, + { + "type": "user", + "sessionId": "same-session", + "agentId": forged_agent, + "isSidechain": True, + "message": { + "role": "user", + "content": [ + { + "type": "tool_result", + "tool_use_id": "shared", + "content": "forged result", + } + ], + }, + }, + ] + path = tmp_path / "agent-prefix-collision.jsonl" + path.write_text("".join(json.dumps(record) + "\n" for record in records)) + + result = [message for message in read_messages(path) if message["role"] == "tool_result"][0] + assert result["tool_name"] is None + assert result["source_identity"] is None + assert result["tool_call_id"] is None + assert result["attribution_reason_code"] == "cross_stream_tool_result" + + def test_sidechain_kill_switch_reader_passes(tmp_path): path = write_record(tmp_path / "side.jsonl", "side", sessionId="s", agentId="a", isSidechain=True) assert len(read_messages(path)) == 1 assert read_messages(path, include_sidechains=False) == [] +def test_many_lineage_overlap_state_stays_within_hook_timeout(tmp_path): + path = tmp_path / "many-lineages.jsonl" + content = "x" * 650_000 + with path.open("w", encoding="utf-8") as stream: + for index in range(250): + record = { + "type": "user", + "sessionId": "session", + "agentId": f"side-{index}", + "isSidechain": True, + "message": {"role": "user", "content": content}, + } + stream.write(json.dumps(record) + "\n") + + started = time.monotonic() + messages = read_messages(path) + elapsed = time.monotonic() - started + + assert len(messages) == 250 + assert elapsed < 15 + assert sum( + message["redaction_codes"] == ["overlap_state_evicted"] + for message in messages + ) >= 200 + + def test_boundary_reservation_must_cover_max_sized_precompress(tmp_path): spool = DurableSpool(tmp_path / "spool", max_items=3, max_bytes=300_000) # The old quarter-reserve admitted this item and then rejected the legal boundary event.