dita_v2(bingx_venue): telemetry OFF the exec path — bounded non-blocking lane
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 <noreply@anthropic.com>
This commit is contained in:
@@ -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,
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user