From 40e85ab8d10bf1d530da6922c0c8ecce4e9fa6f8 Mon Sep 17 00:00:00 2001 From: Codex Date: Mon, 13 Jul 2026 08:58:52 +0200 Subject: [PATCH] =?UTF-8?q?dita=5Fv2(bingx=5Fvenue):=20telemetry=20OFF=20t?= =?UTF-8?q?he=20exec=20path=20=E2=80=94=20bounded=20non-blocking=20lane?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Standing steer (HJ, months): non-blocking behaviour + explicit execution priority. Telemetry is not on the ladder at all, yet it sat ON the order path, inline and able to raise — which is how it orphaned 6 live positions. Exec path now does exactly one thing: pack primitives, append to a bounded ring, return. No plane write, no snapshot construction, no I/O, no lock, no raise. Snapshot build + plane publish move to a daemon drain lane. Ring is deque(maxlen=4096): drops OLDEST on full and counts the drops — telemetry loss is always preferable to exec backpressure, but it stays observable. Measured: 1000 exec-path calls against a 2s-BLOCKING plane = 3.65 ms total (3.65 us/call). Inline, that was 2000 s of stall. Wedged lane + 50k pushes -> ring pinned at 4096, 45903 counted drops, exec path never backpressured. Healthy plane still gets 5/5. Co-Authored-By: Claude Opus 4.8 --- prod/clean_arch/dita_v2/bingx_venue.py | 132 ++++++++++++++++++------- 1 file changed, 97 insertions(+), 35 deletions(-) diff --git a/prod/clean_arch/dita_v2/bingx_venue.py b/prod/clean_arch/dita_v2/bingx_venue.py index 37f21f2..5d6f8ba 100644 --- a/prod/clean_arch/dita_v2/bingx_venue.py +++ b/prod/clean_arch/dita_v2/bingx_venue.py @@ -13,6 +13,7 @@ import inspect import itertools import re import threading +from collections import deque from datetime import datetime, timezone from dataclasses import replace from typing import Any, Iterable, List, Optional @@ -282,8 +283,64 @@ 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. + # + # 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 + + 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_wake = threading.Event() + self._tel_dropped = 0 + thread = threading.Thread( + target=self._drain_telemetry, name="bingx-telemetry-lane", daemon=True + ) + self._tel_thread = thread + thread.start() + + def _drain_telemetry(self) -> None: + """Low-priority lane: build snapshots + publish. Never touches the exec path.""" + while True: + self._tel_wake.wait(timeout=0.5) + self._tel_wake.clear() + ring = getattr(self, "_tel_ring", None) + if ring is None: + continue + while True: + try: + fields = ring.popleft() + except IndexError: + break + try: + plane = self._telemetry_plane + publish = ( + getattr(plane, "publish_venue", None) if plane is not None else None + ) + if publish is None: + continue + 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", + fields.get("phase", "?"), exc, + ) + def set_telemetry_plane(self, zinc_plane: Any | None) -> None: self._telemetry_plane = zinc_plane + if zinc_plane is not None: + self._start_telemetry_lane() def _publish_telemetry( self, @@ -302,20 +359,19 @@ class BingxVenueAdapter(VenueAdapter): client_order_id: str = "", details: dict[str, Any] | None = None, ) -> None: - # TOTAL: this method cannot raise. It runs INLINE on the order path — - # submit()/submit_async() call it immediately after the venue has already - # accepted the order. An exception escaping here reaches the caller's - # submit guard, which synthesises a REJECTED event and rolls the FSM back - # while the venue keeps the position (orphan; 6 of them on 2026-07-13). + # 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. # - # The guard used to cover only publish(), leaving the attribute/coercion - # block below it unprotected — which is how a signature skew orphaned live - # positions past a defence that existed precisely to prevent it. Everything - # that can raise now lives inside the try. Observability never vetoes a fill. + # 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: - plane = self._telemetry_plane - publish = getattr(plane, "publish_venue", None) if plane is not None else None - if publish is None: + ring = getattr(self, "_tel_ring", None) + if ring is None or self._telemetry_plane is None: return slot_id = 0 trade_id = "" @@ -335,36 +391,42 @@ 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 - publish( - VenueTelemetrySnapshot( - phase=phase, - status=status, - venue="bingx", - endpoint=endpoint, - method=method, - intent_id=intent_id, - trade_id=trade_id, - slot_id=slot_id, - asset=asset, - side=side, - action=action, - order_id=str(order_id or getattr(order, "venue_order_id", "") or ""), - client_order_id=str( + 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 + # Timestamp is stamped HERE (event time), not on the drain lane. + ring.append( + { + "phase": phase, + "status": status, + "venue": "bingx", + "endpoint": endpoint, + "method": method, + "intent_id": intent_id, + "trade_id": trade_id, + "slot_id": slot_id, + "asset": asset, + "side": side, + "action": action, + "order_id": str(order_id or getattr(order, "venue_order_id", "") or ""), + "client_order_id": str( client_order_id or getattr(order, "venue_client_id", "") or "" ), - venue_order_status=venue_order_status, - venue_event_kind=venue_event_kind, - message=message, - retry_after_ms=int(retry_after_ms or 0), - timestamp=datetime.now(timezone.utc), - details=dict(details or {}), - ) + "venue_order_status": venue_order_status, + "venue_event_kind": venue_event_kind, + "message": message, + "retry_after_ms": int(retry_after_ms or 0), + "timestamp": datetime.now(timezone.utc), + "details": dict(details or {}), + } ) + self._tel_wake.set() except Exception as exc: import logging as _log _log.getLogger(__name__).warning( - "FIX(bingx_venue): telemetry publish failed (phase=%s) — swallowed, " + "FIX(bingx_venue): telemetry enqueue failed (phase=%s) — swallowed, " "order path unaffected: %s", phase, exc, )