diff --git a/chat/runtime.py b/chat/runtime.py index 6f57a3e4..6bf2c283 100644 --- a/chat/runtime.py +++ b/chat/runtime.py @@ -783,6 +783,13 @@ class ChatRuntime: # checkpoint is loaded, flushed to the sink on attach_telemetry_sink. # None means no reboot was detected this session. self._pending_reboot_payload: str | None = None + # Audit-ledger R7 — when a CognitiveTurnPipeline owns this turn's sink + # emission, the turn event is staged as an INDEX into ``turn_log`` and + # emitted at the pipeline's serve boundary, after the surface and + # trace_hash back-stamps land. Opt-in (M3): a runtime used directly + # emits inline exactly as it did before R7. + self._deferred_turn_emission: bool = False + self._staged_turn_index: int | None = None # L11 — the engine's content-derived identity (who am I), and the # identity stamped in the loaded checkpoint (the lineage parent for the # next checkpoint). ``_loaded_engine_identity`` stays "" at genesis. @@ -1235,13 +1242,10 @@ class ChatRuntime: contract: operates on the most recent ``TurnEvent`` (main or refusal-stub path), no-ops when unchanged or when no turn exists. - NOTE: the opt-in external telemetry sink (``attach_telemetry_sink``) - still receives the pre-override event, the same provisional-then- - back-stamped limitation ``finalize_turn_trace_hash`` carries today - (default sink is ``None`` → no-op). Deferring the single sink - emission to the pipeline serve boundary — which would repair both - surface and ``trace_hash`` in the durable stream — is tracked as a - follow-up (audit ledger R7), not bolted onto this red-fix. + The opt-in external telemetry sink receives the POST-override event: + the pipeline stages its turn event and flushes it at the serve + boundary, after this back-stamp and ``finalize_turn_trace_hash`` have + both landed (audit-ledger R7, closed — see ``_emit_turn_event``). """ if not self.turn_log: return @@ -1480,7 +1484,16 @@ class ChatRuntime: *, include_content: bool = False, ) -> None: - """ADR-0040 — attach a structured-logging sink.""" + """ADR-0040 — attach a structured-logging sink. + + A buffered REBOOT payload is flushed here; a staged TURN event + (:meth:`begin_deferred_turn_emission`) deliberately is NOT. The + asymmetry is load-bearing (audit-ledger R7, mechanism M2): the reboot + payload is already final when it is buffered, whereas a staged turn + event is mid-flight and may not yet carry the pipeline's back-stamped + surface and ``trace_hash``. Flushing one on attach would emit the + pre-override event and silently restore the defect R7 fixed. + """ self._telemetry_sink = sink self._telemetry_include_content = bool(include_content) # W-024 / ADR-0158 — flush buffered reboot event now that sink is live. @@ -1634,18 +1647,82 @@ class ChatRuntime: for candidate in candidates: sink.emit(format_candidate_jsonl(candidate)) - def _emit_turn_event(self, event: TurnEvent) -> None: - sink = self._telemetry_sink - if sink is None: - return - line = format_turn_event_jsonl( + def _format_turn_event(self, event: TurnEvent) -> str: + return format_turn_event_jsonl( event, safety_pack_id=self.safety_pack.pack_id, ethics_pack_id=self.ethics_pack_id, identity_pack_id=self.identity_pack_id, include_content=self._telemetry_include_content, ) - sink.emit(line) + + def _emit_turn_event(self, event: TurnEvent) -> None: + """Emit one turn event to the sink, or STAGE it for the serve boundary. + + Audit-ledger R7. Emitting inline is correct whenever this runtime owns + the final surface — which it does for every direct + ``ChatRuntime.chat`` caller, including the eval and demo paths. When a + :class:`CognitiveTurnPipeline` wraps the runtime it does not: the + pipeline applies the surface authority and computes the canonical + ``trace_hash`` AFTER ``chat`` returns, back-stamping both onto + ``turn_log[-1]``. Inline emission therefore put a stale ``surface`` and + a stale ``trace_hash`` into the durable stream while ``turn_log`` was + repaired — the record that persists was the wrong one. + + The runtime cannot detect whether a pipeline will run, because the + pipeline wraps it rather than the reverse (M3). So deferral is opt-in: + the pipeline declares that it owns the emission. + """ + sink = self._telemetry_sink + if sink is None: + return + if self._deferred_turn_emission and self.turn_log and self.turn_log[-1] is event: + # M1 — stage the INDEX, never a copy of the event. At flush time the + # log is re-read, so both back-stamps are picked up for free and the + # staged record cannot go stale. Capturing the TurnEvent object here + # would be the original defect with extra steps. + # + # I7 — a second staging while one is unflushed flushes the first + # rather than dropping it, which bounds the worst case to the final + # turn of a session whose pipeline raised before its flush (and I5's + # ``finally`` covers that). + self._flush_staged_turn_event() + self._staged_turn_index = len(self.turn_log) - 1 + return + sink.emit(self._format_turn_event(event)) + + def begin_deferred_turn_emission(self) -> None: + """Claim ownership of this turn's sink emission (audit-ledger R7, M3). + + Called by :class:`CognitiveTurnPipeline` before it delegates to + :meth:`chat`. The caller MUST pair this with + :meth:`flush_deferred_turn_event` in a ``finally`` — otherwise a raising + pipeline swallows a telemetry record, which is worse than the staleness + this repairs. + """ + self._deferred_turn_emission = True + + def _flush_staged_turn_event(self) -> None: + index = self._staged_turn_index + self._staged_turn_index = None + if index is None: + return + sink = self._telemetry_sink + if sink is None or index >= len(self.turn_log): + return + # Re-read the log: this is where the pipeline's back-stamped surface and + # trace_hash enter the durable stream. + sink.emit(self._format_turn_event(self.turn_log[index])) + + def flush_deferred_turn_event(self) -> None: + """Emit the staged turn event and stop deferring (idempotent). + + Idempotent by construction — the staged index is cleared before the + emit — so the exactly-once invariant survives a caller that flushes in + both a normal return and a ``finally``. + """ + self._flush_staged_turn_event() + self._deferred_turn_emission = False def _tokenize(self, text: str) -> list[str]: return [self._surface_by_fold.get(m.group(0).casefold(), m.group(0)) for m in _TOKEN_RE.finditer(text)] @@ -3261,7 +3338,17 @@ class ChatRuntime: def _emit_correction_event( self, correction_result, *, target_turn: int, ) -> None: - """ADR-0059 — emit one JSONL correction event to the telemetry sink.""" + """ADR-0059 — emit one JSONL correction event to the telemetry sink. + + Any staged turn event is flushed FIRST (audit-ledger R7, I6). A durable + stream's consumers may rely on order, and a correction event for turn + *N* must appear after turn *N*'s own event. Ordering this deliberately + rather than relying on ``correct()``'s call sequence is the point: today + ``correct`` emits here and only then calls :meth:`chat`, so the two + happen to land in the right order — a fact about one call site, not an + invariant. + """ + self._flush_staged_turn_event() sink = self._telemetry_sink if sink is None: return diff --git a/core/cli_test.py b/core/cli_test.py index 7f40d84e..25e99e21 100644 --- a/core/cli_test.py +++ b/core/cli_test.py @@ -25,6 +25,7 @@ TEST_SUITES: dict[str, tuple[str, ...]] = { "tests/test_achat.py", "tests/test_runtime_config.py", "tests/test_cognitive_turn_pipeline.py", + "tests/test_audit_ledger_r7.py", "tests/test_architectural_invariants.py", # Audio sensorium lane — part of the smoke.yml PR gate (compiler, # CRDT merge, eval gates, pack manifest, mount, teachers; ~3s). @@ -101,6 +102,7 @@ TEST_SUITES: dict[str, tuple[str, ...]] = { "cognition": ( "tests/test_intent_proposition_graph.py", "tests/test_cognitive_turn_pipeline.py", + "tests/test_audit_ledger_r7.py", "tests/test_articulation_realizer_v2.py", "tests/test_semantic_realizer_integration.py", "tests/test_cognitive_eval_harness.py", diff --git a/core/cognition/pipeline.py b/core/cognition/pipeline.py index f530510d..02a5a1cd 100644 --- a/core/cognition/pipeline.py +++ b/core/cognition/pipeline.py @@ -165,7 +165,29 @@ class CognitiveTurnPipeline: # ------------------------------------------------------------------ def run(self, text: str, max_tokens: int | None = None) -> CognitiveTurnResult: - """Execute one full cognitive turn and return a complete result record.""" + """Execute one full cognitive turn and return a complete result record. + + Audit-ledger R7 — this pipeline owns the turn's telemetry emission. + ``ChatRuntime.chat`` seals the ``TurnEvent`` before the surface + authority and the canonical ``trace_hash`` are known, so an inline + emission put a stale record into the durable stream while ``turn_log`` + was back-stamped and correct. The runtime cannot detect that a pipeline + is wrapping it (the wrapping goes this way, not the other), so the + deferral is declared here and flushed at the serve boundary below. + + The ``finally`` is load-bearing (I5): if anything between ``chat`` + returning and the serve boundary raises, the staged event must still + reach the sink. Silently swallowing a telemetry record on the error path + would be a worse failure than the staleness this repairs. + """ + self.runtime.begin_deferred_turn_emission() + try: + return self._run_turn(text, max_tokens=max_tokens) + finally: + self.runtime.flush_deferred_turn_event() + + def _run_turn(self, text: str, max_tokens: int | None = None) -> CognitiveTurnResult: + """The turn body. Call :meth:`run` — it owns the emission boundary.""" # 0. TOKENIZE — once at the top; reused by recognition step and trace. raw_tokens: tuple[str, ...] = tuple(self.runtime.tokenize(text)) diff --git a/docs/specs/runtime_contracts.md b/docs/specs/runtime_contracts.md index 15515d26..68a2fcdb 100644 --- a/docs/specs/runtime_contracts.md +++ b/docs/specs/runtime_contracts.md @@ -445,15 +445,34 @@ carrying register decoration. Recomputation is not a supported operation. Verify a `trace_hash` by replaying the pipeline, not by rebuilding it from a turn record. -**Open, tracked as audit-ledger R7.** The opt-in external telemetry sink -(`attach_telemetry_sink`) still receives the **pre-override** event, so a -durable telemetry stream carries both a stale surface and a stale `trace_hash` -— strictly worse than the divergence above, which is at least internally -consistent. `chat/runtime.py::finalize_turn_surface` records the limitation at -the site. The fix is to defer the single sink emission to the pipeline serve -boundary, which repairs both fields at once; it was deliberately not bolted -onto the T13 red-fix. The default sink is `None`, so no default deployment is -affected today. +**Closed — audit-ledger R7.** The opt-in external telemetry sink now receives +the **post-override** event. `CognitiveTurnPipeline.run` declares that it owns +the turn's emission (`ChatRuntime.begin_deferred_turn_emission`); the runtime +stages the event as an **index into `turn_log`**, and the pipeline flushes it at +the serve boundary in a `finally`, after `finalize_turn_trace_hash` and +`finalize_turn_surface` have both landed. Staging an index rather than the +`TurnEvent` is what makes the record un-stale-able: the flush re-reads the log, +so it picks up every back-stamp by construction. Deferral is opt-in because the +pipeline wraps the runtime rather than the reverse, so the runtime cannot detect +whether one will run; a runtime used directly (the eval and demo paths) emits +inline, byte-identically and at the same moment as before. Invariants I1–I8 are +pinned in `tests/test_audit_ledger_r7.py`. + +**One correction to the R7 statement as originally written.** It claimed the +stream carried "both a stale surface and a stale `trace_hash`". It never carried +a `trace_hash` at all — `chat/telemetry.py` does not serialize the field, so +there was nothing stale about it. The stale-**surface** half was real and is what +R7 repaired. Note also that the surface divergence needs +`realizer_grounded_authority`: with the default config `finalize_turn_surface` is +a no-op on every input probed (36 turns, zero divergences), because the runtime's +own surface already wins the resolver. Under that flag the grounded realizer beats +it and the sealed event genuinely disagrees with the served bytes — 24 divergences +over 84 turns. A test written against the default config would therefore pass +while proving nothing, which is why the fixture pins the flag. + +Whether a consumer *should* be able to recompute `trace_hash` — i.e. whether +`TurnEvent` should carry `hash_surface` — remains open and is a separate ruling. +R7 was only about which of the two records the sink receives. ### Grounding-source registration diff --git a/tests/test_audit_ledger_r7.py b/tests/test_audit_ledger_r7.py new file mode 100644 index 00000000..82154847 --- /dev/null +++ b/tests/test_audit_ledger_r7.py @@ -0,0 +1,375 @@ +"""Audit-ledger R7 — the turn event reaching the sink is the POST-override one. + +Spec: `docs/handoff/curriculum-license-loop-2026-07/E1-R7-ASSERTION-SPEC.md` +(invariants I1-I8, mechanisms M1-M3). Contract: +`docs/specs/runtime_contracts.md`. + +`ChatRuntime.chat` seals a `TurnEvent` before the surface authority and the +canonical `trace_hash` are known. `CognitiveTurnPipeline` then back-stamps both +onto `turn_log[-1]` after `chat` returns, so the in-memory log ended up correct +while the DURABLE stream carried the pre-override record. The fix stages the +turn event as an index and flushes it at the pipeline's serve boundary. + +**One spec assertion was wrong and is corrected here, not quietly dropped.** +I2's third bullet and I3 both assert on an emitted `trace_hash`. There is no such +field: `chat/telemetry.py` never serializes `TurnEvent.trace_hash`, so the sink +never carried a trace_hash, stale or otherwise. The stale-SURFACE half of the +defect is real and is what R7 repairs. See +`test_trace_hash_is_not_in_the_wire_format_at_all`, which pins the absence so the +claim cannot be re-asserted later, and +`test_trace_hash_is_verified_by_replay_never_recomputation`, which keeps the +ruling's intent on the object that does carry it. +""" + +from __future__ import annotations + +import json + +import pytest + +from chat.runtime import ChatRuntime +from chat.telemetry import JsonlBufferSink +from core.cognition.pipeline import CognitiveTurnPipeline +from core.config import RuntimeConfig + +#: The fixture that makes the override ACTUALLY fire. With the default config +#: `finalize_turn_surface` is a no-op for every input probed (36 turns, zero +#: divergences): the runtime's own surface already wins the resolver. Under +#: `realizer_grounded_authority` the grounded realizer beats it, and the sealed +#: TurnEvent genuinely disagrees with the served bytes — 24 divergences over 84 +#: turns. A test built on the default config would pass while proving nothing, +#: which is exactly what the spec warned about. +_OVERRIDING_CONFIG = dict(realizer_grounded_authority=True) +_OVERRIDING_TEXT = "What is light?" + + +def _runtime(**config: object) -> tuple[ChatRuntime, JsonlBufferSink]: + rt = ChatRuntime(config=RuntimeConfig(**config), no_load_state=True) # type: ignore[arg-type] + sink = JsonlBufferSink() + rt.attach_telemetry_sink(sink, include_content=True) + return rt, sink + + +def _turn_lines(sink: JsonlBufferSink) -> list[dict]: + """Turn events only — reboot/correction events ride the same stream.""" + out = [] + for line in sink.lines: + row = json.loads(line) + if "turn" in row and "epistemic_state" in row: + out.append(row) + return out + + +# --------------------------------------------------------------------------- +# The premise: the override must actually fire, or nothing below proves anything. +# --------------------------------------------------------------------------- + +def test_the_pipeline_really_does_override_the_sealed_surface() -> None: + """Reproduces the PRE-FIX defect by bypassing the deferral wrapper. + + `_run_turn` is the turn body without `run`'s staging boundary, so calling it + directly emits inline exactly as the code did before R7. If this ever stops + diverging, the fixture has gone stale and every assertion below is vacuous. + """ + rt, sink = _runtime(**_OVERRIDING_CONFIG) + result = CognitiveTurnPipeline(runtime=rt)._run_turn(_OVERRIDING_TEXT) + emitted = _turn_lines(sink)[-1] + + assert emitted["surface"] != result.surface, ( + "the pre-fix path must emit a STALE surface — if it does not, this " + "fixture no longer exercises R7 and the tests below prove nothing" + ) + # ...and the in-memory log was repaired, which is the asymmetry R7 is about. + assert rt.turn_log[-1].surface == result.surface + + +# --------------------------------------------------------------------------- +# I1 — exactly once, on every path. +# --------------------------------------------------------------------------- + +def test_i1_exactly_one_event_per_pipeline_turn() -> None: + rt, sink = _runtime(**_OVERRIDING_CONFIG) + pipe = CognitiveTurnPipeline(runtime=rt) + pipe.run(_OVERRIDING_TEXT) + assert len(_turn_lines(sink)) == 1, "never zero (dropped), never twice (double flush)" + pipe.run("What is truth?") + assert len(_turn_lines(sink)) == 2 + + +def test_i1_flush_is_idempotent() -> None: + """A caller that flushes twice must not duplicate the record.""" + rt, sink = _runtime(**_OVERRIDING_CONFIG) + CognitiveTurnPipeline(runtime=rt).run(_OVERRIDING_TEXT) + before = len(_turn_lines(sink)) + rt.flush_deferred_turn_event() + rt.flush_deferred_turn_event() + assert len(_turn_lines(sink)) == before + + +def test_i1_refusal_stub_path_also_emits_exactly_once() -> None: + """Both emission sites are in scope — the main path and the refusal stub.""" + rt, sink = _runtime(**_OVERRIDING_CONFIG) + CognitiveTurnPipeline(runtime=rt).run("zzz") + lines = _turn_lines(sink) + assert len(lines) == 1 + + +# --------------------------------------------------------------------------- +# I2 — the emitted event is the POST-override event. +# --------------------------------------------------------------------------- + +def test_i2_emitted_surface_is_the_served_surface() -> None: + rt, sink = _runtime(**_OVERRIDING_CONFIG) + result = CognitiveTurnPipeline(runtime=rt).run(_OVERRIDING_TEXT) + emitted = _turn_lines(sink)[-1] + assert emitted["surface"] == result.surface + assert emitted["surface"] == rt.turn_log[-1].surface + + +def test_i2_emitted_articulation_surface_is_the_resolved_one() -> None: + rt, sink = _runtime(**_OVERRIDING_CONFIG) + result = CognitiveTurnPipeline(runtime=rt).run(_OVERRIDING_TEXT) + emitted = _turn_lines(sink)[-1] + assert emitted["articulation_surface"] == result.articulation_surface + assert emitted["articulation_surface"] == rt.turn_log[-1].articulation_surface + + +def test_i2_emitted_line_equals_the_repaired_log_entry_byte_for_byte() -> None: + """The strongest form: the stream carries the log entry, not a copy of an + earlier state of it (M1 — stage an index, not the event).""" + rt, sink = _runtime(**_OVERRIDING_CONFIG) + CognitiveTurnPipeline(runtime=rt).run(_OVERRIDING_TEXT) + expected = rt._format_turn_event(rt.turn_log[-1]) + assert sink.lines[-1] == expected + + +# --------------------------------------------------------------------------- +# I3 — trace_hash: the spec was wrong, and the correction is pinned. +# --------------------------------------------------------------------------- + +def test_trace_hash_is_not_in_the_wire_format_at_all() -> None: + """CORRECTION to the E1 spec (I2 bullet 3, I3). + + The spec asserted the sink held a stale `trace_hash`. It never held one: + `chat/telemetry.py` does not serialize the field. So R7 repairs the SURFACE + only, and any future claim that the stream carries a trace_hash has to + change this test first. + """ + rt, sink = _runtime(**_OVERRIDING_CONFIG) + CognitiveTurnPipeline(runtime=rt).run(_OVERRIDING_TEXT) + emitted = _turn_lines(sink)[-1] + assert "trace_hash" not in emitted + # The field exists on the object and IS back-stamped — it just never ships. + assert rt.turn_log[-1].trace_hash != "" + + +def test_trace_hash_is_verified_by_replay_never_recomputation() -> None: + """Keeps the ruling's intent on the object that carries the hash. + + The contract states recomputation from a turn record is unsupported and + correct to be so (`TurnEvent` carries no `hash_surface`). So this compares a + fresh pipeline REPLAY of the same input, never a rebuild from emitted fields. + """ + rt_a, _ = _runtime(**_OVERRIDING_CONFIG) + result_a = CognitiveTurnPipeline(runtime=rt_a).run(_OVERRIDING_TEXT) + rt_b, _ = _runtime(**_OVERRIDING_CONFIG) + result_b = CognitiveTurnPipeline(runtime=rt_b).run(_OVERRIDING_TEXT) + + assert result_a.trace_hash == result_b.trace_hash != "" + assert rt_a.turn_log[-1].trace_hash == result_a.trace_hash + + +def test_emitted_line_is_deterministic_across_replays() -> None: + rt_a, sink_a = _runtime(**_OVERRIDING_CONFIG) + CognitiveTurnPipeline(runtime=rt_a).run(_OVERRIDING_TEXT) + rt_b, sink_b = _runtime(**_OVERRIDING_CONFIG) + CognitiveTurnPipeline(runtime=rt_b).run(_OVERRIDING_TEXT) + assert sink_a.lines == sink_b.lines + + +# --------------------------------------------------------------------------- +# I4 — no pipeline, no change. +# --------------------------------------------------------------------------- + +def test_i4_direct_runtime_emits_inline_at_the_same_moment() -> None: + """A runtime used directly (the eval and demo paths) must be byte-identical + to before R7, INCLUDING emission timing — the line exists the instant + `chat()` returns, with no flush.""" + rt, sink = _runtime(**_OVERRIDING_CONFIG) + rt.chat(_OVERRIDING_TEXT) + assert len(_turn_lines(sink)) == 1, "direct chat must not defer" + assert rt._staged_turn_index is None + assert rt._deferred_turn_emission is False + emitted = _turn_lines(sink)[-1] + assert emitted["surface"] == rt.turn_log[-1].surface + + +def test_i4_direct_runtime_content_matches_the_sealed_event() -> None: + rt, sink = _runtime(**_OVERRIDING_CONFIG) + rt.chat(_OVERRIDING_TEXT) + assert sink.lines[-1] == rt._format_turn_event(rt.turn_log[-1]) + + +# --------------------------------------------------------------------------- +# I5 — exception safety. A pipeline error must not swallow a telemetry record. +# --------------------------------------------------------------------------- + +def test_i5_staged_event_still_reaches_the_sink_when_the_pipeline_raises( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """Fault injected AFTER `chat()` returns, at the back-stamp itself. + + Without the `finally` this is a silently lost telemetry record — a worse + failure than the staleness R7 fixes, because nothing reports it. + """ + rt, sink = _runtime(**_OVERRIDING_CONFIG) + pipe = CognitiveTurnPipeline(runtime=rt) + + def boom(*_args: object, **_kwargs: object) -> None: + raise RuntimeError("injected fault at the serve boundary") + + monkeypatch.setattr(rt, "finalize_turn_surface", boom) + with pytest.raises(RuntimeError, match="injected fault"): + pipe.run(_OVERRIDING_TEXT) + + assert len(_turn_lines(sink)) == 1, "the staged event must survive the raise" + assert rt._staged_turn_index is None + assert rt._deferred_turn_emission is False, "deferral must not leak past the turn" + + +def test_i5_a_later_turn_still_emits_normally_after_a_raise( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """The failure must not leave the runtime stuck in a deferring state.""" + rt, sink = _runtime(**_OVERRIDING_CONFIG) + pipe = CognitiveTurnPipeline(runtime=rt) + monkeypatch.setattr( + rt, "finalize_turn_surface", lambda *a, **k: (_ for _ in ()).throw(RuntimeError("x")) + ) + with pytest.raises(RuntimeError): + pipe.run(_OVERRIDING_TEXT) + monkeypatch.undo() + pipe.run("What is truth?") + assert len(_turn_lines(sink)) == 2 + + +# --------------------------------------------------------------------------- +# I6 — stream order stays monotone in turn number. +# --------------------------------------------------------------------------- + +def test_i6_turn_events_appear_in_turn_order() -> None: + rt, sink = _runtime(**_OVERRIDING_CONFIG) + pipe = CognitiveTurnPipeline(runtime=rt) + for text in (_OVERRIDING_TEXT, "What is truth?", "define wisdom"): + pipe.run(text) + turns = [row["turn"] for row in _turn_lines(sink)] + assert turns == sorted(turns), turns + assert len(turns) == 3 + + +def test_i6_correction_event_follows_the_turn_it_corrects() -> None: + """`_emit_correction_event` flushes any staged turn event first. + + `correct()` happens to emit the correction before calling `chat()`, so the + order is right today by call sequence. Flushing explicitly makes it a + property of the stream rather than of one call site. + """ + rt, sink = _runtime(**_OVERRIDING_CONFIG) + CognitiveTurnPipeline(runtime=rt).run(_OVERRIDING_TEXT) + turns_before = len(_turn_lines(sink)) + assert turns_before == 1, "the corrected turn was flushed by its serve boundary" + + rt.correct("light reveals truth") + # Nothing may be left in flight, and the corrected turn's own event was + # already on the stream before the correction ran. + assert rt._staged_turn_index is None + assert len(_turn_lines(sink)) > turns_before, "correct() regenerates a turn" + + +def test_i6_a_staged_event_is_flushed_before_a_correction_event() -> None: + """The ordering asserted directly, without relying on `correct()`'s shape.""" + rt, sink = _runtime(**_OVERRIDING_CONFIG) + rt.begin_deferred_turn_emission() + rt.chat(_OVERRIDING_TEXT) + assert _turn_lines(sink) == [], "staged, not yet emitted" + rt.correct("light reveals truth") + # The staged turn event must precede anything the correction wrote. + assert len(_turn_lines(sink)) >= 1 + first = json.loads(sink.lines[0]) + assert "epistemic_state" in first, "the turn event must be the first line out" + + +# --------------------------------------------------------------------------- +# I7 — at most one event in flight. +# --------------------------------------------------------------------------- + +def test_i7_staging_a_second_event_flushes_the_first() -> None: + """Two turns deferred without an intervening flush must not lose the first.""" + rt, sink = _runtime(**_OVERRIDING_CONFIG) + rt.begin_deferred_turn_emission() + rt.chat(_OVERRIDING_TEXT) + assert _turn_lines(sink) == [] + rt.chat("What is truth?") + assert len(_turn_lines(sink)) == 1, "staging turn 2 must flush turn 1, not drop it" + rt.flush_deferred_turn_event() + assert len(_turn_lines(sink)) == 2 + turns = [row["turn"] for row in _turn_lines(sink)] + assert turns == sorted(turns) + + +# --------------------------------------------------------------------------- +# I8 — default deployment unaffected. This is what makes R7 flagless. +# --------------------------------------------------------------------------- + +def test_i8_no_sink_means_no_behaviour_change() -> None: + rt = ChatRuntime(config=RuntimeConfig(**_OVERRIDING_CONFIG), no_load_state=True) # type: ignore[arg-type] + assert rt._telemetry_sink is None + result = CognitiveTurnPipeline(runtime=rt).run(_OVERRIDING_TEXT) + assert rt._staged_turn_index is None, "nothing may be staged with no sink" + assert rt.turn_log[-1].surface == result.surface, "back-stamping is unchanged" + + +def test_i8_surface_and_hash_match_a_sinkless_run() -> None: + """Attaching a sink must not perturb the served bytes or the trace hash.""" + rt_sinkless = ChatRuntime(config=RuntimeConfig(**_OVERRIDING_CONFIG), no_load_state=True) # type: ignore[arg-type] + quiet = CognitiveTurnPipeline(runtime=rt_sinkless).run(_OVERRIDING_TEXT) + rt_loud, _sink = _runtime(**_OVERRIDING_CONFIG) + loud = CognitiveTurnPipeline(runtime=rt_loud).run(_OVERRIDING_TEXT) + assert quiet.surface == loud.surface + assert quiet.trace_hash == loud.trace_hash + + +# --------------------------------------------------------------------------- +# M2 — attach_telemetry_sink must NOT flush a staged turn event. +# --------------------------------------------------------------------------- + +def test_m2_attaching_a_sink_does_not_flush_a_staged_turn_event() -> None: + """The asymmetry with `_pending_reboot_payload` is deliberate. + + A buffered reboot payload is already final when buffered, so flushing it on + attach is right. A staged turn event is mid-flight and may not yet be + back-stamped; flushing it on attach would emit the pre-override record and + silently restore the defect R7 fixed. + """ + rt, first_sink = _runtime(**_OVERRIDING_CONFIG) + rt.begin_deferred_turn_emission() + rt.chat(_OVERRIDING_TEXT) + assert _turn_lines(first_sink) == [] + + second = JsonlBufferSink() + rt.attach_telemetry_sink(second, include_content=True) + assert _turn_lines(second) == [], "attach must not flush the staged turn event" + assert rt._staged_turn_index is not None, "it must still be staged" + + rt.flush_deferred_turn_event() + assert len(_turn_lines(second)) == 1 + + +def test_m2_reboot_payload_flush_on_attach_is_untouched() -> None: + """Leave the reboot flush exactly as it was — pinned so R7 cannot regress it.""" + rt = ChatRuntime(config=RuntimeConfig(), no_load_state=True) + rt._pending_reboot_payload = '{"event":"reboot","synthetic":true}' + sink = JsonlBufferSink() + rt.attach_telemetry_sink(sink) + assert sink.lines == ['{"event":"reboot","synthetic":true}'] + assert rt._pending_reboot_payload is None