diff --git a/prod/clean_arch/dita_v2/bingx_venue.py b/prod/clean_arch/dita_v2/bingx_venue.py index 3ec8b9d..04e5dfc 100644 --- a/prod/clean_arch/dita_v2/bingx_venue.py +++ b/prod/clean_arch/dita_v2/bingx_venue.py @@ -298,6 +298,8 @@ class BingxVenueAdapter(VenueAdapter): retry_after_ms: int = 0, venue_order_status: str = "", venue_event_kind: str = "", + order_id: str = "", + client_order_id: str = "", details: dict[str, Any] | None = None, ) -> None: plane = self._telemetry_plane @@ -336,8 +338,10 @@ class BingxVenueAdapter(VenueAdapter): asset=asset, side=side, action=action, - order_id=str(getattr(order, "venue_order_id", "") or ""), - client_order_id=str(getattr(order, "venue_client_id", "") or ""), + 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, @@ -620,23 +624,36 @@ class BingxVenueAdapter(VenueAdapter): receipt = self._call_backend("submit_intent", legacy) events = self._events_from_submit(submitted, receipt, None, None) ack_row = dict(getattr(receipt, "raw_ack", {}) or {}) - self._publish_telemetry( - phase="submit:done", - status=str(getattr(receipt, "status", "") or _row_text(ack_row, "status", default="NEW")), - intent=intent, - endpoint="/openApi/swap/v2/trade/order", - method="POST", - message=_row_text(ack_row, "msg", "message", default=""), - order_id=_row_text(ack_row, "orderId", "orderID", default=str(getattr(receipt, "order_id", "") or "")), - client_order_id=_row_text(ack_row, "clientOrderID", "clientOrderId", default=str(getattr(receipt, "client_order_id", "") or intent.intent_id)), - 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, - }, - ) + # 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 + # the venue keeps the position (orphan). Observability must not un-do fills. + try: + self._publish_telemetry( + phase="submit:done", + status=str(getattr(receipt, "status", "") or _row_text(ack_row, "status", default="NEW")), + intent=intent, + endpoint="/openApi/swap/v2/trade/order", + method="POST", + message=_row_text(ack_row, "msg", "message", default=""), + order_id=_row_text(ack_row, "orderId", "orderID", default=str(getattr(receipt, "order_id", "") or "")), + client_order_id=_row_text(ack_row, "clientOrderID", "clientOrderId", default=str(getattr(receipt, "client_order_id", "") or intent.intent_id)), + 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 async def submit_async(self, intent: KernelIntent) -> List[VenueEvent]: @@ -665,23 +682,36 @@ class BingxVenueAdapter(VenueAdapter): receipt = await self.backend.submit_intent(legacy) events = self._events_from_submit(submitted, receipt, None, None) ack_row = dict(getattr(receipt, "raw_ack", {}) or {}) - self._publish_telemetry( - phase="submit:done", - status=str(getattr(receipt, "status", "") or _row_text(ack_row, "status", default="NEW")), - intent=intent, - endpoint="/openApi/swap/v2/trade/order", - method="POST", - message=_row_text(ack_row, "msg", "message", default=""), - order_id=_row_text(ack_row, "orderId", "orderID", default=str(getattr(receipt, "order_id", "") or "")), - client_order_id=_row_text(ack_row, "clientOrderID", "clientOrderId", default=str(getattr(receipt, "client_order_id", "") or intent.intent_id)), - 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, - }, - ) + # 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 + # the venue keeps the position (orphan). Observability must not un-do fills. + try: + self._publish_telemetry( + phase="submit:done", + status=str(getattr(receipt, "status", "") or _row_text(ack_row, "status", default="NEW")), + intent=intent, + endpoint="/openApi/swap/v2/trade/order", + method="POST", + message=_row_text(ack_row, "msg", "message", default=""), + order_id=_row_text(ack_row, "orderId", "orderID", default=str(getattr(receipt, "order_id", "") or "")), + client_order_id=_row_text(ack_row, "clientOrderID", "clientOrderId", default=str(getattr(receipt, "client_order_id", "") or intent.intent_id)), + 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 def _events_from_submit(self, intent: KernelIntent, receipt: Any, before, after) -> List[VenueEvent]: # noqa: ANN001