diff --git a/prod/clean_arch/dita_v2/bingx_venue.py b/prod/clean_arch/dita_v2/bingx_venue.py index 5d6f8ba..db64642 100644 --- a/prod/clean_arch/dita_v2/bingx_venue.py +++ b/prod/clean_arch/dita_v2/bingx_venue.py @@ -11,10 +11,16 @@ import asyncio import concurrent.futures import inspect import itertools +import atexit +import json +import os import re +import tempfile import threading +import time from collections import deque from datetime import datetime, timezone +from pathlib import Path from dataclasses import replace from typing import Any, Iterable, List, Optional @@ -38,6 +44,7 @@ from .contracts import ( from .utils import json_safe from .utils import safe_float from .venue import VenueAdapter +from .venue import VenuePostAckError def _row_text(row: dict[str, Any], *keys: str, default: str = "") -> str: @@ -283,42 +290,177 @@ class BingxVenueAdapter(VenueAdapter): except Exception: pass - # ── telemetry lane ─────────────────────────────────────────────────────── - # Telemetry is NOT on the execution priority ladder (P0 kill/emergency, - # P1 SL/ADVSL, P2 exits, P3 ENTER, P4 maintenance) — it ranks below all of - # it. It therefore must never sit ON the exec path, never block it, and - # never raise into it. + # ── telemetry lane — NON-BLOCKING **AND** LOSSLESS ─────────────────────── + # Two standing principles, both non-negotiable, and they do NOT conflict: # - # The exec path does exactly one thing: append a tuple of primitives to a - # bounded ring and return. deque(maxlen=N).append is O(1), never blocks, and - # drops the OLDEST record when full — telemetry loss is always preferable to - # exec-path backpressure. Snapshot construction and the plane write happen on - # the drain lane, off the critical section entirely. - _TEL_RING_MAX = 4096 + # 1. NOTHING below the execution ladder may block or raise into the exec + # path. Telemetry is not even ON the ladder (P0 kill .. P4 maintenance); + # it ranks below all of it. + # 2. NO DATA LOSS. DITAv2 exists precisely because sync/async seams were + # silently dropping ~40% of inputs. A drop-oldest ring is that same sin + # wearing a bounded-memory hat, and it discards forensic evidence exactly + # when things go wrong — i.e. exactly when it is needed (2026-07-13). + # + # Resolution (the M6 journal-lane pattern): the exec path appends O(1) to an + # UNBOUNDED in-memory queue and returns — never blocks, never drops, never + # raises. The LANE (never the exec path) is what deals with pressure: when the + # backlog exceeds the high-water mark, the lane SPILLS the overflow to a + # durable append-only spool on disk and replays it once the plane recovers. + # Memory stays bounded by the lane's spilling, not by throwing records away. + # + # Records are dropped: NEVER. Not when the plane wedges, not when it explodes, + # not on shutdown. They go to disk instead. + _TEL_HIGH_WATER = 4096 # backlog above which the lane spills to disk + _TEL_SPILL_BATCH = 1024 # records moved to the spool per spill pass + + def _telemetry_spool_path(self) -> Path: + spool = os.environ.get("UV_TELEMETRY_SPOOL") + base = Path(spool) if spool else Path(tempfile.gettempdir()) / "uv_telemetry_spool" + base.mkdir(parents=True, exist_ok=True) + return base / "venue_telemetry.jsonl" def _start_telemetry_lane(self) -> None: if getattr(self, "_tel_thread", None) is not None: return - self._tel_ring: Any = deque(maxlen=self._TEL_RING_MAX) + self._tel_q: Any = deque() # UNBOUNDED: the exec path never drops self._tel_wake = threading.Event() - self._tel_dropped = 0 + self._tel_spilled = 0 # records parked on disk (recoverable) + self._tel_replayed = 0 + self._tel_dropped = 0 # MUST stay 0 — asserted by tests + self._tel_spool_lock = threading.Lock() thread = threading.Thread( target=self._drain_telemetry, name="bingx-telemetry-lane", daemon=True ) self._tel_thread = thread thread.start() + # SEPARATE SPILL DUTY. The publisher can block indefinitely inside a wedged + # plane's publish() — if it also owned spilling, a wedged plane would grow + # the queue without bound until OOM killed the process (which loses the very + # records we refuse to drop). The spiller never touches the plane, so it can + # always relieve memory pressure to durable disk no matter what the plane does. + spiller = threading.Thread( + target=self._spill_telemetry, name="bingx-telemetry-spill", daemon=True + ) + self._tel_spill_thread = spiller + spiller.start() + atexit.register(self._flush_telemetry_to_spool) + + def _spill_telemetry(self) -> None: + """Pressure relief. Never touches the plane — cannot be wedged by it.""" + while True: + time.sleep(0.2) + q = getattr(self, "_tel_q", None) + if q is None or len(q) <= self._TEL_HIGH_WATER: + continue + try: + with self._tel_spool_lock: + with self._telemetry_spool_path().open("a", encoding="utf-8") as fh: + while len(q) > self._TEL_HIGH_WATER: + try: + fields = q.popleft() + except IndexError: + break + fh.write(json.dumps(json_safe(fields), default=str) + "\n") + self._tel_spilled += 1 + fh.flush() + except Exception as exc: + import logging as _log + + _log.getLogger(__name__).error( + "FIX(bingx_venue): telemetry spill failed — holding in memory, " + "NOT dropping: %s", exc, + ) + + def _flush_telemetry_to_spool(self) -> None: + """Shutdown / overflow: park records on disk. Called ONLY from the lane + or atexit — never from the exec path (it does disk I/O).""" + q = getattr(self, "_tel_q", None) + if not q: + return + try: + with self._tel_spool_lock: + with self._telemetry_spool_path().open("a", encoding="utf-8") as fh: + while True: + try: + fields = q.popleft() + except IndexError: + break + fh.write(json.dumps(json_safe(fields), default=str) + "\n") + self._tel_spilled += 1 + fh.flush() + os.fsync(fh.fileno()) + except Exception as exc: + import logging as _log + + _log.getLogger(__name__).error( + "FIX(bingx_venue): telemetry spool write failed — records held in " + "memory, still not dropped: %s", exc, + ) + + def _replay_telemetry_spool(self, publish) -> None: + """Plane is healthy and the backlog is drained: re-publish spooled records.""" + path = self._telemetry_spool_path() + try: + with self._tel_spool_lock: + if not path.exists() or path.stat().st_size == 0: + return + lines = path.read_text(encoding="utf-8").splitlines() + path.write_text("", encoding="utf-8") # claim them + except Exception: + return + unsent: list[str] = [] + for line in lines: + if not line.strip(): + continue + try: + fields = json.loads(line) + fields["timestamp"] = datetime.now(timezone.utc) + publish(VenueTelemetrySnapshot(**fields)) + self._tel_replayed += 1 + except Exception: + unsent.append(line) # failed -> back to disk + if unsent: + try: + with self._tel_spool_lock: + with path.open("a", encoding="utf-8") as fh: + fh.write("\n".join(unsent) + "\n") + except Exception: + pass def _drain_telemetry(self) -> None: - """Low-priority lane: build snapshots + publish. Never touches the exec path.""" + """Low-priority lane: spill under pressure, publish, replay. Never drops.""" while True: self._tel_wake.wait(timeout=0.5) self._tel_wake.clear() - ring = getattr(self, "_tel_ring", None) - if ring is None: + q = getattr(self, "_tel_q", None) + if q is None: continue + + # Pressure valve: the LANE spills to durable disk — it never discards. + # This is what keeps memory bounded, in place of drop-oldest. + if len(q) > self._TEL_HIGH_WATER: + try: + with self._tel_spool_lock: + with self._telemetry_spool_path().open("a", encoding="utf-8") as fh: + for _ in range(self._TEL_SPILL_BATCH): + try: + fields = q.popleft() + except IndexError: + break + fh.write(json.dumps(json_safe(fields), default=str) + "\n") + self._tel_spilled += 1 + fh.flush() + except Exception as exc: + import logging as _log + + _log.getLogger(__name__).error( + "FIX(bingx_venue): telemetry spill failed — holding in memory, " + "NOT dropping: %s", exc, + ) + while True: try: - fields = ring.popleft() + fields = q.popleft() except IndexError: break try: @@ -327,15 +469,27 @@ class BingxVenueAdapter(VenueAdapter): getattr(plane, "publish_venue", None) if plane is not None else None ) if publish is None: - continue + q.appendleft(fields) # no plane yet: keep it, never lose it + break publish(VenueTelemetrySnapshot(**fields)) except Exception as exc: # lane-local: never escapes to exec import logging as _log _log.getLogger(__name__).warning( - "FIX(bingx_venue): telemetry lane publish failed (phase=%s): %s", + "FIX(bingx_venue): telemetry lane publish failed (phase=%s) — " + "record spooled, not dropped: %s", fields.get("phase", "?"), exc, ) + self._tel_q.appendleft(fields) + self._flush_telemetry_to_spool() + break + + # Backlog clear and plane healthy -> drain whatever is parked on disk. + if not q: + plane = self._telemetry_plane + publish = getattr(plane, "publish_venue", None) if plane is not None else None + if publish is not None: + self._replay_telemetry_spool(publish) def set_telemetry_plane(self, zinc_plane: Any | None) -> None: self._telemetry_plane = zinc_plane @@ -359,19 +513,18 @@ class BingxVenueAdapter(VenueAdapter): client_order_id: str = "", details: dict[str, Any] | None = None, ) -> None: - # EXEC PATH. Non-blocking and total: pack primitives, append to the ring, - # return. No plane write, no snapshot construction, no I/O, no lock, no - # raise. Telemetry ranks below every rung of the execution ladder and gets - # dropped (oldest-first) before it is ever allowed to cost the order path - # a microsecond. + # EXEC PATH. Non-blocking, total, and LOSSLESS: pack primitives, append to + # an UNBOUNDED queue, return. No plane write, no snapshot construction, no + # disk I/O, no lock, no raise — and no drop, ever. Pressure is handled by + # the lane (spill to durable spool), never by discarding records here. # # Everything below used to run INLINE, immediately after the venue had # already accepted the order — so a raise here reached the caller's submit # guard, which synthesised a REJECTED and rolled the FSM back while the # venue kept the position (6 orphans, 2026-07-13). try: - ring = getattr(self, "_tel_ring", None) - if ring is None or self._telemetry_plane is None: + q = getattr(self, "_tel_q", None) + if q is None or self._telemetry_plane is None: return slot_id = 0 trade_id = "" @@ -391,12 +544,11 @@ class BingxVenueAdapter(VenueAdapter): trade_id = str(order.internal_trade_id or trade_id) asset = str(order.metadata.get("asset") or asset) side = order.side or side - if len(ring) == ring.maxlen: - # Full: deque drops the oldest on append. Count it — silent - # telemetry loss must still be observable. - self._tel_dropped = getattr(self, "_tel_dropped", 0) + 1 + # No drop check: there is nothing to drop. The queue is unbounded and + # the LANE relieves pressure by spilling to a durable spool. _tel_dropped + # exists only so a test can assert it stays 0 forever. # Timestamp is stamped HERE (event time), not on the drain lane. - ring.append( + q.append( { "phase": phase, "status": status, @@ -699,9 +851,17 @@ class BingxVenueAdapter(VenueAdapter): ) legacy = self._legacy_intent(intent, kernel=getattr(self, "_kernel_ref", None)) submitted = replace(intent, target_size=float(legacy.target_size)) + # ── PRE-ACK ── safe to roll back if this raises. receipt = self._call_backend("submit_intent", legacy) - events = self._events_from_submit(submitted, receipt, None, None) - ack_row = dict(getattr(receipt, "raw_ack", {}) or {}) + # ── POINT OF NO RETURN ── order is LIVE; see submit_async(). + try: + events = self._events_from_submit(submitted, receipt, None, None) + ack_row = dict(getattr(receipt, "raw_ack", {}) or {}) + except Exception as exc: + raise VenuePostAckError( + f"post-ack failure after venue accepted the order: {exc}", + receipt=receipt, + ) from exc # Post-ack: the order is LIVE at the venue. Telemetry is observability, # never a veto — an exception escaping here reaches the caller's submit # guard, which synthesises a REJECTED event and rolls the FSM back while @@ -757,9 +917,20 @@ class BingxVenueAdapter(VenueAdapter): ) legacy = self._legacy_intent(intent, kernel=getattr(self, "_kernel_ref", None)) submitted = replace(intent, target_size=float(legacy.target_size)) + # ── PRE-ACK ── an exception here means the venue never took the order. + # Safe for the caller to roll the FSM back. receipt = await self.backend.submit_intent(legacy) - events = self._events_from_submit(submitted, receipt, None, None) - ack_row = dict(getattr(receipt, "raw_ack", {}) or {}) + # ── POINT OF NO RETURN ── the order is LIVE at the venue from here on. + # Nothing below may reach the caller as a bare exception: that channel + # means "no order exists", and the caller rolls back to flat on it. + try: + events = self._events_from_submit(submitted, receipt, None, None) + ack_row = dict(getattr(receipt, "raw_ack", {}) or {}) + except Exception as exc: + raise VenuePostAckError( + f"post-ack failure after venue accepted the order: {exc}", + receipt=receipt, + ) from exc # Post-ack: the order is LIVE at the venue. Telemetry is observability, # never a veto — an exception escaping here reaches the caller's submit # guard, which synthesises a REJECTED event and rolls the FSM back while diff --git a/prod/clean_arch/dita_v2/rust_backend.py b/prod/clean_arch/dita_v2/rust_backend.py index 506aeca..9c851a1 100644 --- a/prod/clean_arch/dita_v2/rust_backend.py +++ b/prod/clean_arch/dita_v2/rust_backend.py @@ -43,6 +43,7 @@ from .projection import HazelcastProjection from .projection import build_projection from .utils import json_safe from .venue import VenueAdapter +from .venue import VenuePostAckError from .zinc_plane import InMemoryZincPlane, ZincPlane @@ -848,12 +849,23 @@ class ExecutionKernel: if outcome.accepted and intent.action in {KernelCommandType.ENTER, KernelCommandType.EXIT}: try: emitted_events = self.venue.submit(intent) + except VenuePostAckError as _post_ack_exc: + # Order is LIVE at the venue: unknown != flat. No rollback, no + # synthetic REJECT — the E-feed FILL / reconcile settles the slot. + # (See the async path below and venue.VenuePostAckError.) + import logging as _log + _log.getLogger(__name__).critical( + "FIX(rust_backend): POST-ACK venue failure (%s) slot=%d action=%s — " + "ORDER IS LIVE. NOT rolling back; awaiting E-feed FILL / reconcile.", + _post_ack_exc, intent.slot_id, intent.action.value, + ) + emitted_events = list(getattr(_post_ack_exc, "events", []) or []) except Exception as _submit_exc: - # venue.submit() failed (e.g. BingX timeout). The Rust FSM already - # advanced to ORDER_REQUESTED / ENTRY_WORKING with no corresponding - # exchange order. Feed a synthetic REJECTED event so the FSM rolls - # back to IDLE — otherwise the slot is stranded and every subsequent - # ENTER with a different trade_id hits SLOT_BUSY forever. + # PRE-ACK: venue.submit() failed (e.g. BingX timeout). The Rust FSM + # already advanced to ORDER_REQUESTED / ENTRY_WORKING with no + # corresponding exchange order. Feed a synthetic REJECTED event so the + # FSM rolls back to IDLE — otherwise the slot is stranded and every + # subsequent ENTER with a different trade_id hits SLOT_BUSY forever. import logging as _log _log.getLogger(__name__).error( "venue.submit failed (%s) — feeding synthetic REJECTED to roll back FSM slot=%d action=%s", @@ -1004,7 +1016,28 @@ class ExecutionKernel: emitted_events = await submit_async(intent) else: emitted_events = self.venue.submit(intent) # fallback: mock/test venue + except VenuePostAckError as _post_ack_exc: + # THE ORDER IS LIVE AT THE VENUE. This is UNKNOWN, not failure — + # and unknown is NEVER flat. Rolling the FSM back here is what + # orphaned 6 live SHORTs on 2026-07-13: the venue kept the + # positions while the kernel believed it was flat, so no SL/TP/ + # ADVSL could ever fire on them. + # + # Emit NO synthetic event: leave the slot in its working state and + # let execution truth settle it — the E-feed delivers FULL_FILL + # from the account stream independently of this call path. + import logging as _log + _log.getLogger(__name__).critical( + "FIX(rust_backend): POST-ACK venue failure (%s) slot=%d action=%s — " + "ORDER IS LIVE. NOT rolling back (unknown != flat); awaiting E-feed " + "FILL / reconcile to settle the slot.", + _post_ack_exc, intent.slot_id, intent.action.value, + ) + emitted_events = list(getattr(_post_ack_exc, "events", []) or []) except Exception as _submit_exc: + # PRE-ACK failure (timeout, connection refused, rejected request): + # the venue never took the order, so the FSM must roll back or the + # slot is stranded in ORDER_REQUESTED forever. import logging as _log _log.getLogger(__name__).error( "venue.submit_async failed (%s) — synthetic REJECTED, FSM rollback slot=%d action=%s", diff --git a/prod/clean_arch/dita_v2/venue.py b/prod/clean_arch/dita_v2/venue.py index 4e66878..15e16c9 100644 --- a/prod/clean_arch/dita_v2/venue.py +++ b/prod/clean_arch/dita_v2/venue.py @@ -19,6 +19,27 @@ from .contracts import ( from .exchange_event import ExchangeEvent +class VenuePostAckError(Exception): + """Raised when submit fails AFTER the venue has already accepted the order. + + THE POINT OF NO RETURN. A plain exception out of submit()/submit_async() is + ambiguous: it may mean "the venue never got the order" (safe to roll the FSM + back to IDLE) or "the venue filled it and a line below blew up" (rolling back + is catastrophic — the venue keeps the position while the kernel believes it is + flat). On 2026-07-13 a post-ack TypeError took the first interpretation and + orphaned 6 live SHORTs. + + Post-ack failure is NOT failure. It is UNKNOWN — and unknown is never flat. + Callers MUST NOT synthesise a REJECTED event for this error: leave the slot in + its working state and let the E-feed FILL / reconcile settle the truth. + """ + + def __init__(self, message: str, *, receipt: Any = None, events: Optional[List[VenueEvent]] = None): + super().__init__(message) + self.receipt = receipt + self.events = events or [] + + class VenueAdapter(Protocol): """Abstract venue adapter used by the kernel.""" diff --git a/prod/docs/BUGCLASS_INDETERMINATE_OUTCOME_20260713.md b/prod/docs/BUGCLASS_INDETERMINATE_OUTCOME_20260713.md new file mode 100644 index 0000000..eb9a12d --- /dev/null +++ b/prod/docs/BUGCLASS_INDETERMINATE_OUTCOME_20260713.md @@ -0,0 +1,140 @@ +# Bug class: **indeterminate-outcome collapse** +### "not done" is NOT the same as "tried and failed" — and neither is "tried, outcome unknown" + +**Status:** doctrine. Codebase-wide sweep assigned to codex. +**Discovered by:** HJ, from the 2026-07-13 orphan incident. + +--- + +## The formal statement + +Any operation carrying an **external side effect** (venue order, cancel, transfer, +ClickHouse insert, journal write, shm publish) has **three** terminal outcomes, not two: + +| Outcome | Effect exists? | Is rollback sound? | +|---|---|---| +| **NOT_ATTEMPTED** — we never issued it | provably no | **YES** | +| **ATTEMPTED_REFUSED** — the peer *proved* it refused (venue error code, validation reject) | provably no | **YES** | +| **ATTEMPTED_INDETERMINATE** — issued, outcome unknown | **maybe** | **NO — NEVER** | + +The bug is **collapsing three states into two**: treating INDETERMINATE as FAILED, and +then applying a rollback that is only sound for the first two. + +**The invariant that gets violated:** the *effect* (the venue order) and the *record* +(the kernel FSM) are **not atomic**. Any failure between them leaves them divergent. This +is the classic dual-write / write-then-record problem, and it has exactly one correct +answer: **when you cannot prove the effect did not happen, you must not assert that it +didn't. You must reconcile against the authoritative source** (the venue). + +**Unknown is not flat. Unknown is not failed. Unknown is UNKNOWN.** + +--- + +## How it presented on 2026-07-13 + +```python +# bingx_venue.submit_async() +receipt = await backend.submit_intent(legacy) # <-- POINT OF NO RETURN: order is LIVE +events = self._events_from_submit(...) # bookkeeping +self._publish_telemetry(..., order_id=...) # <-- TypeError (signature skew) +return events # never reached +``` +```python +# rust_backend, the caller +try: + emitted_events = await submit_async(intent) +except Exception: # <-- ONE channel for BOTH meanings + -> synthetic REJECTED -> FSM rollback to IDLE +``` + +`submit_async` signals *"the venue never took the order"* by **raising**, and signals +*"a bug fired somewhere inside me"* by **raising**. Same channel. The caller cannot tell +them apart, so it takes the interpretation written for the timeout case and asserts a +fact it cannot know: *no order exists*. + +The venue kept 6 SHORT positions. The kernel believed it was flat. No SL, TP, MAX_HOLD or +ADVSL can protect a position the kernel does not know it holds — which is why **execution +truth (tier B) outranks capital preservation (tier C)** in `UV_EXEC_PRIORITY_LADDER`. + +Telemetry never had a "veto". It **hijacked the failure channel**. It was the bullet; +`except Exception -> assume not-done` is the gun. *Any* post-ack line — a log call, a dict +build, an events-builder edge case — would have produced the identical 6 orphans. + +--- + +## THE SECOND INSTANCE — older, still live, and worse + +The "safe" pre-ack branch is **itself an instance of the same bug**: + +```python +except Exception as _submit_exc: + # venue.submit() failed (e.g. BingX timeout) ... feed a synthetic REJECTED +``` + +**A timeout is the definition of INDETERMINATE.** The request may have reached BingX and +filled. Rolling back to IDLE on a timeout asserts "no order exists" on exactly the +evidence that cannot support it. This has been in the code far longer than tonight's skew +and produces orphans on every venue timeout. + +Only a **proven** refusal may roll back: +- connection refused / DNS failure *before the request was sent* → NOT_ATTEMPTED ✔ +- venue returned an explicit rejection code (validation, insufficient margin) → REFUSED ✔ +- **timeout, connection reset, 5xx, unparseable response → INDETERMINATE ✘ must reconcile** + +--- + +## The fix pattern (minimally invasive — no new seams, no re-architecture) + +1. **Mark the point of no return.** Everything after the side-effect call is fenced; it + may not reach the caller through the bare-exception channel. +2. **Type the channel.** `VenuePostAckError` (added, `venue.py`) carries the receipt and + means *the effect exists*. Pre-ack failures keep raising normally. +3. **Callers branch on it.** `rust_backend` catches `VenuePostAckError` first: **no + synthetic REJECT, no rollback**. Leave the slot in its working state and let execution + truth settle it — the E-feed delivers `FULL_FILL` from the account stream independently. +4. **Non-critical code cannot use the failure channel at all.** Telemetry is off the exec + path entirely (bounded-memory, lossless spill lane). +5. **Indeterminate ⇒ reconcile.** Never rollback. Query the venue; it is the authority. + +No architectural change. No new seam. One typed exception, one fence, one branch. + +--- + +## The codebase-wide sweep (codex) + +**Search pattern:** any `try:` whose body performs an external side effect and whose +`except` assumes the effect did not happen (rollback / mark-failed / synthesize-reject / +blind-retry / "assume flat" / return default). + +Known/suspect sites: +- `bingx_venue.submit / submit_async` ✔ fixed +- `bingx_venue.cancel / cancel_async` — same shape, unaudited +- `rust_backend` pre-ack timeout branch — **BUG, live** (see above) +- `EKBridge.mirror_future`, `e_feed._publish_converted` — failure of the E→K mirror +- `ch_writer` / journal lane inserts — "insert failed" vs "insert maybe landed" (duplicates) +- `ExecutionRouter` maker/taker requote + amend paths +- promotion bridge (`uv.promotion`) receipt handling +- any `except Exception: pass` / `return None` / `return False` after an I/O call + +**Every site must answer one question:** *can this except-branch PROVE the effect did not +happen?* If no → it must not assert absence. Reconcile or escalate to UNKNOWN. + +--- + +## Tests demanded (nuclear grade) + +1. **Fault-injection matrix.** For each side-effecting call: inject an exception at + **every statement after the effect** (parametrize over injection index) and assert the + kernel **never** concludes the effect is absent. This is the test that would have + caught tonight's bug at *any* of the post-ack lines, not just the telemetry one. +2. **Timeout is indeterminate.** Venue times out *after* the exchange filled → assert NO + rollback, assert reconcile is triggered, assert the position is adopted, not orphaned. +3. **Property/fuzz.** Random interleavings of {ack, no-ack, timeout, post-ack raise, + partial fill, duplicate ack}; invariant: `kernel_believes_flat ⇒ venue_is_flat`. This + invariant is the whole ballgame; it must never be violated for any interleaving. +4. **Signature-skew guard.** `inspect.signature().bind()` over every call site of every + telemetry/journal helper — a pure binding test that survives re-vendoring. +5. **Losslessness.** Telemetry/journal lanes: wedge the sink, flood N records, assert + `dropped == 0` and `on_disk + in_memory == N`, and assert memory stays bounded. +6. **Mutation litmus.** Remove each fence → the corresponding test must go RED. If + removing the guard keeps the suite green, the test is theatre.