dita_v2(bingx_venue): restore _publish_telemetry order_id/client_order_id kwargs + make post-ack telemetry non-fatal

M1 re-vendor (89a1f1a) clobbered the T0-DEVIATION fix (20ad493): submit()/submit_async()
pass order_id=/client_order_id= post-ack, the signature dropped them. Every ENTER raised
TypeError AFTER the venue POST returned 200 OK -> rust_backend synthesised REJECTED ->
FSM rollback, while the venue kept the position. Flight-4: 8 bridged promotions, 6 orphan
SHORTs live on VST with a kernel that believes it is flat.

Two fixes:
1. signature accepts order_id/client_order_id again; snapshot prefers them (order is None
   on the submit path, so the ack row is the only id source).
2. post-ack telemetry wrapped: observability can never again veto an accepted order.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
Codex
2026-07-13 05:08:12 +02:00
parent 459215b7d8
commit b46ebd2f82

View File

@@ -298,6 +298,8 @@ class BingxVenueAdapter(VenueAdapter):
retry_after_ms: int = 0, retry_after_ms: int = 0,
venue_order_status: str = "", venue_order_status: str = "",
venue_event_kind: str = "", venue_event_kind: str = "",
order_id: str = "",
client_order_id: str = "",
details: dict[str, Any] | None = None, details: dict[str, Any] | None = None,
) -> None: ) -> None:
plane = self._telemetry_plane plane = self._telemetry_plane
@@ -336,8 +338,10 @@ class BingxVenueAdapter(VenueAdapter):
asset=asset, asset=asset,
side=side, side=side,
action=action, action=action,
order_id=str(getattr(order, "venue_order_id", "") or ""), order_id=str(order_id or getattr(order, "venue_order_id", "") or ""),
client_order_id=str(getattr(order, "venue_client_id", "") or ""), client_order_id=str(
client_order_id or getattr(order, "venue_client_id", "") or ""
),
venue_order_status=venue_order_status, venue_order_status=venue_order_status,
venue_event_kind=venue_event_kind, venue_event_kind=venue_event_kind,
message=message, message=message,
@@ -620,23 +624,36 @@ class BingxVenueAdapter(VenueAdapter):
receipt = self._call_backend("submit_intent", legacy) receipt = self._call_backend("submit_intent", legacy)
events = self._events_from_submit(submitted, receipt, None, None) events = self._events_from_submit(submitted, receipt, None, None)
ack_row = dict(getattr(receipt, "raw_ack", {}) or {}) ack_row = dict(getattr(receipt, "raw_ack", {}) or {})
self._publish_telemetry( # Post-ack: the order is LIVE at the venue. Telemetry is observability,
phase="submit:done", # never a veto — an exception escaping here reaches the caller's submit
status=str(getattr(receipt, "status", "") or _row_text(ack_row, "status", default="NEW")), # guard, which synthesises a REJECTED event and rolls the FSM back while
intent=intent, # the venue keeps the position (orphan). Observability must not un-do fills.
endpoint="/openApi/swap/v2/trade/order", try:
method="POST", self._publish_telemetry(
message=_row_text(ack_row, "msg", "message", default=""), phase="submit:done",
order_id=_row_text(ack_row, "orderId", "orderID", default=str(getattr(receipt, "order_id", "") or "")), status=str(getattr(receipt, "status", "") or _row_text(ack_row, "status", default="NEW")),
client_order_id=_row_text(ack_row, "clientOrderID", "clientOrderId", default=str(getattr(receipt, "client_order_id", "") or intent.intent_id)), intent=intent,
venue_order_status=str(getattr(receipt, "status", "") or _row_text(ack_row, "status", default="")), endpoint="/openApi/swap/v2/trade/order",
venue_event_kind=events[1].kind.value if len(events) > 1 else events[0].kind.value, method="POST",
details={ message=_row_text(ack_row, "msg", "message", default=""),
"filled_size": float((events[1].filled_size if len(events) > 1 else events[0].filled_size) or 0.0), order_id=_row_text(ack_row, "orderId", "orderID", default=str(getattr(receipt, "order_id", "") or "")),
"event_count": len(events), client_order_id=_row_text(ack_row, "clientOrderID", "clientOrderId", default=str(getattr(receipt, "client_order_id", "") or intent.intent_id)),
"asset": intent.asset, venue_order_status=str(getattr(receipt, "status", "") or _row_text(ack_row, "status", default="")),
}, venue_event_kind=events[1].kind.value if len(events) > 1 else events[0].kind.value,
) details={
"filled_size": float((events[1].filled_size if len(events) > 1 else events[0].filled_size) or 0.0),
"event_count": len(events),
"asset": intent.asset,
},
)
except Exception as exc: # pragma: no cover - defensive, see comment above
import logging as _log
_log.getLogger(__name__).critical(
"FIX(bingx_venue): post-ack telemetry failed for intent=%s asset=%s — "
"order is LIVE at venue, events returned unchanged: %s",
intent.intent_id, intent.asset, exc,
)
return events return events
async def submit_async(self, intent: KernelIntent) -> List[VenueEvent]: async def submit_async(self, intent: KernelIntent) -> List[VenueEvent]:
@@ -665,23 +682,36 @@ class BingxVenueAdapter(VenueAdapter):
receipt = await self.backend.submit_intent(legacy) receipt = await self.backend.submit_intent(legacy)
events = self._events_from_submit(submitted, receipt, None, None) events = self._events_from_submit(submitted, receipt, None, None)
ack_row = dict(getattr(receipt, "raw_ack", {}) or {}) ack_row = dict(getattr(receipt, "raw_ack", {}) or {})
self._publish_telemetry( # Post-ack: the order is LIVE at the venue. Telemetry is observability,
phase="submit:done", # never a veto — an exception escaping here reaches the caller's submit
status=str(getattr(receipt, "status", "") or _row_text(ack_row, "status", default="NEW")), # guard, which synthesises a REJECTED event and rolls the FSM back while
intent=intent, # the venue keeps the position (orphan). Observability must not un-do fills.
endpoint="/openApi/swap/v2/trade/order", try:
method="POST", self._publish_telemetry(
message=_row_text(ack_row, "msg", "message", default=""), phase="submit:done",
order_id=_row_text(ack_row, "orderId", "orderID", default=str(getattr(receipt, "order_id", "") or "")), status=str(getattr(receipt, "status", "") or _row_text(ack_row, "status", default="NEW")),
client_order_id=_row_text(ack_row, "clientOrderID", "clientOrderId", default=str(getattr(receipt, "client_order_id", "") or intent.intent_id)), intent=intent,
venue_order_status=str(getattr(receipt, "status", "") or _row_text(ack_row, "status", default="")), endpoint="/openApi/swap/v2/trade/order",
venue_event_kind=events[1].kind.value if len(events) > 1 else events[0].kind.value, method="POST",
details={ message=_row_text(ack_row, "msg", "message", default=""),
"filled_size": float((events[1].filled_size if len(events) > 1 else events[0].filled_size) or 0.0), order_id=_row_text(ack_row, "orderId", "orderID", default=str(getattr(receipt, "order_id", "") or "")),
"event_count": len(events), client_order_id=_row_text(ack_row, "clientOrderID", "clientOrderId", default=str(getattr(receipt, "client_order_id", "") or intent.intent_id)),
"asset": intent.asset, venue_order_status=str(getattr(receipt, "status", "") or _row_text(ack_row, "status", default="")),
}, venue_event_kind=events[1].kind.value if len(events) > 1 else events[0].kind.value,
) details={
"filled_size": float((events[1].filled_size if len(events) > 1 else events[0].filled_size) or 0.0),
"event_count": len(events),
"asset": intent.asset,
},
)
except Exception as exc: # pragma: no cover - defensive, see comment above
import logging as _log
_log.getLogger(__name__).critical(
"FIX(bingx_venue): post-ack telemetry failed for intent=%s asset=%s — "
"order is LIVE at venue, events returned unchanged: %s",
intent.intent_id, intent.asset, exc,
)
return events return events
def _events_from_submit(self, intent: KernelIntent, receipt: Any, before, after) -> List[VenueEvent]: # noqa: ANN001 def _events_from_submit(self, intent: KernelIntent, receipt: Any, before, after) -> List[VenueEvent]: # noqa: ANN001