Files
sentiment-engine/prod/clean_arch/dita_v2/bingx_venue.py
Codex 3d75e4b59d dita_v2: unknown != flat (VenuePostAckError) + lossless telemetry lane
BUG CLASS (doc: BUGCLASS_INDETERMINATE_OUTCOME_20260713.md): an operation with an external
side effect has THREE outcomes — NOT_ATTEMPTED, ATTEMPTED_REFUSED, ATTEMPTED_INDETERMINATE
— and rollback is sound only for the first two. Collapsing the third into 'failed' is what
orphaned 6 live SHORTs: a post-ack TypeError reached rust_backend's 'except Exception ->
synthetic REJECTED -> FSM rollback', which asserted 'no order exists' about an order that
was already filled. Telemetry never had a veto; it hijacked the failure channel.

FIXES (no new seams, no re-architecture):
- venue.py: VenuePostAckError — typed channel meaning THE EFFECT EXISTS. Carries receipt.
- bingx_venue submit/submit_async: point-of-no-return fence. Post-ack bookkeeping failures
  raise VenuePostAckError instead of a bare exception.
- rust_backend (BOTH submit paths): catch VenuePostAckError FIRST -> no synthetic REJECT,
  no rollback. Slot stays working; E-feed FULL_FILL / reconcile settles the truth.

LOSSLESS TELEMETRY (HJ: 'DITAv2 exists precisely because seams dropped 40% of inputs'):
drop-oldest is data loss and is GONE. Exec path appends O(1) to an unbounded queue and
returns. A SEPARATE spiller thread (which never touches the plane, so a wedged plane cannot
starve it) parks the backlog above HWM into a durable append-only spool; the publisher
replays the spool when the plane recovers.
Proven: wedged-forever plane + 200k records -> 0 dropped, 195903 durable on disk, 4096 in
memory, 7.4 us/call on the exec path. Lossless AND memory-bounded. Healthy plane: 2000/2000.

STILL BROKEN, flagged to codex: the pre-ack branch rolls back on TIMEOUT — but a timeout is
the definition of INDETERMINATE (the order may have filled). Same bug class, older, live.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-13 10:31:39 +02:00

1263 lines
60 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""DITAv2 BingX venue adapter.
This is a thin normalization layer over the existing direct BingX execution
surface. It converts BingX REST/account/order payloads into DITAv2
``VenueEvent`` / ``VenueOrder`` objects without reimplementing exchange logic.
"""
from __future__ import annotations
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
from prod.clean_arch.dita import DecisionAction as LegacyDecisionAction
from prod.clean_arch.dita import Intent as LegacyIntent
from prod.clean_arch.dita import TradeSide as LegacyTradeSide
from prod.bingx.http import BingxHttpError
from .contracts import (
KernelCommandType,
KernelEventKind,
KernelIntent,
TradeSide,
VenueTelemetrySnapshot,
VenueEvent,
VenueEventStatus,
VenueOrder,
VenueOrderStatus,
)
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:
for key in keys:
value = row.get(key)
if value is None:
continue
text = str(value)
if text:
return text
return default
def _row_float(row: dict[str, Any], *keys: str, default: float = 0.0) -> float:
for key in keys:
try:
value = float(row.get(key) or 0.0)
except Exception:
continue
if value == value and value not in (float("inf"), float("-inf")) and value != 0.0:
return value
return default
def _normalize_status(status: str) -> str:
return str(status or "").strip().upper()
def _trade_side_from_row(row: dict[str, Any], *, fallback: TradeSide = TradeSide.FLAT) -> TradeSide:
side_raw = _row_text(row, "side", "positionSide", default="").upper()
signed_qty = _row_float(row, "positionAmt", "positionQty", "positionSize", "quantity", "pa", default=0.0)
if side_raw in {"BUY", "LONG"}:
return TradeSide.LONG
if side_raw in {"SELL", "SHORT"}:
return TradeSide.SHORT
if signed_qty < 0:
return TradeSide.SHORT
if signed_qty > 0:
return TradeSide.LONG
return fallback
def _http_error_status(exc_msg: str) -> str:
"""Map a BingxHttpError message to a venue status string.
HTTP 429 and 5xx are transient → RATE_LIMITED so the slot can retry.
4xx (non-429) are genuine client-side rejections → REJECTED.
Transport / DNS / circuit-breaker errors (no HTTP prefix) are transient.
"""
m = exc_msg.upper()
if "HTTP 429" in m:
return "RATE_LIMITED"
for code in ("500", "501", "502", "503", "504"):
if f"HTTP {code}" in m:
return "RATE_LIMITED"
if "HTTP 4" in m:
return "REJECTED"
return "RATE_LIMITED"
def _venue_event_status_from_row(status: str) -> VenueEventStatus:
normalized = _normalize_status(status)
if normalized in {"NEW", "ACKED", "PENDING", "CREATED"}:
return VenueEventStatus.ACKED
if normalized in {"RATE_LIMITED", "THROTTLED"}:
return VenueEventStatus.RATE_LIMITED
if normalized in {"PARTIALLY_FILLED", "PARTIAL_FILL"}:
return VenueEventStatus.PARTIALLY_FILLED
if normalized in {"FILLED", "FULL_FILL"}:
return VenueEventStatus.FILLED
if normalized in {"CANCELED", "CANCELLED", "EXPIRED"}:
return VenueEventStatus.CANCELED
if normalized in {"REJECTED", "FAILED"}:
return VenueEventStatus.REJECTED
if normalized in {"CANCEL_REJECTED", "CANCEL_REJECT"}:
return VenueEventStatus.CANCELED_REJECTED
return VenueEventStatus.ACKED
def _venue_order_status_from_row(status: str) -> VenueOrderStatus:
normalized = _normalize_status(status)
if normalized in {"NEW", "ACKED", "PENDING", "CREATED"}:
return VenueOrderStatus.NEW
if normalized in {"RATE_LIMITED", "THROTTLED"}:
return VenueOrderStatus.NEW
if normalized in {"PARTIALLY_FILLED", "PARTIAL_FILL"}:
return VenueOrderStatus.PARTIALLY_FILLED
if normalized in {"FILLED", "FULL_FILL"}:
return VenueOrderStatus.FILLED
if normalized in {"CANCELED", "CANCELLED", "EXPIRED"}:
return VenueOrderStatus.CANCELED
if normalized in {"REJECTED", "FAILED"}:
return VenueOrderStatus.REJECTED
return VenueOrderStatus.NEW
def _position_qty(row: dict[str, Any]) -> float:
qty = _row_float(row, "positionAmt", "positionQty", "positionSize", "quantity", "pa", default=0.0)
if qty != 0.0:
return abs(qty)
return abs(_row_float(row, "executedQty", "filledQty", "z", default=0.0))
def _position_price(row: dict[str, Any]) -> float:
return _row_float(row, "entryPrice", "avgPrice", "avgEntryPrice", "ep", "ap", "price", "lastFillPrice", "tradePrice")
def _mapping_for_snapshot(rows: Iterable[dict[str, Any]]) -> dict[str, dict[str, Any]]:
mapping: dict[str, dict[str, Any]] = {}
for row in rows:
client_id = _row_text(row, "clientOrderID", "clientOrderId", default="")
order_id = _row_text(row, "orderId", "orderID", "id", default="")
key = client_id or order_id
if key:
mapping[key] = dict(row)
if order_id and order_id not in mapping:
mapping[order_id] = dict(row)
return mapping
def _venue_order_from_row(
row: dict[str, Any],
*,
internal_trade_id: str = "",
fallback_side: TradeSide = TradeSide.FLAT,
) -> VenueOrder:
side = _trade_side_from_row(row, fallback=fallback_side)
client_id = _row_text(row, "clientOrderID", "clientOrderId", default="")
order_id = _row_text(row, "orderId", "orderID", "id", default="")
intended = _row_float(row, "origQty", "quantity", "q", "positionAmt", "positionQty", default=0.0)
if intended <= 0:
intended = _position_qty(row)
return VenueOrder(
internal_trade_id=internal_trade_id or client_id or order_id,
venue_order_id=order_id,
venue_client_id=client_id,
side=side,
intended_size=abs(float(intended or 0.0)),
filled_size=abs(_row_float(row, "executedQty", "filledQty", "z", "lastFilledQty", default=0.0)),
average_fill_price=_position_price(row),
status=_venue_order_status_from_row(_row_text(row, "status", "X", default="NEW")),
metadata={"raw": dict(row)},
)
def _event_id(seq: itertools.count) -> str:
return f"EV-{next(seq):08d}"
def _rate_limit_retry_after_ms(row: dict[str, Any]) -> int:
raw_retry = row.get("retryAfter") or row.get("retry_after_ms") or row.get("retryAfterMs")
if raw_retry is None:
msg = _row_text(row, "msg", "message", default="")
match = re.search(r"unblocked after (\d+)", msg)
if match:
try:
ts = int(match.group(1))
now_ms = int(datetime.now(timezone.utc).timestamp() * 1000)
return max(0, ts - now_ms)
except Exception:
return 0
return 0
try:
return max(0, int(float(raw_retry)))
except Exception:
return 0
class BingxVenueAdapter(VenueAdapter):
"""Normalizes BingX execution responses into DITAv2 venue events."""
# Shared thread-pool executor reused across all adapter instances and
# all calls. Threads are created once and recycled, eliminating the
# per-call creation/destruction overhead of the old pattern.
_EXECUTOR: concurrent.futures.ThreadPoolExecutor | None = None
_EXECUTOR_LOCK: threading.Lock = threading.Lock()
@classmethod
def _get_executor(cls) -> concurrent.futures.ThreadPoolExecutor:
if cls._EXECUTOR is None:
with cls._EXECUTOR_LOCK:
if cls._EXECUTOR is None:
# max_workers=3 so three concurrent HTTP calls (balance,
# positions, openOrders) can proceed simultaneously without
# serialising on the pool.
cls._EXECUTOR = concurrent.futures.ThreadPoolExecutor(
max_workers=3,
thread_name_prefix="bingx_adapter",
)
return cls._EXECUTOR
def __init__(self, backend: Any | None = None, *, config: Any | None = None, zinc_plane: Any | None = None) -> None:
if backend is None:
if config is None:
raise ValueError("BingxVenueAdapter requires a backend or config")
from prod.clean_arch.adapters.bingx_direct import BingxDirectExecutionAdapter
backend = BingxDirectExecutionAdapter(config)
self.backend = backend
self._telemetry_plane = zinc_plane
self._event_seq = itertools.count(1)
# Thread-safe snapshot cache — reads from a snapshot may arrive from
# the kernel thread while _backend_snapshot writes from the pool thread.
self._snap_lock = threading.Lock()
self._last_snapshot = None
self._snapshot_ready = threading.Event()
self._snapshot_ready.set() # initially ready (no pending write)
# Maximum seconds to wait for a single backend HTTP call. BingX REST
# round-trips are ~0.5–2 s in normal conditions; 30 s is generous enough
# to survive transient slowdowns without hanging the process forever (O5).
_BACKEND_TIMEOUT_S: float = 30.0
def _run(self, result: Any) -> Any:
if inspect.isawaitable(result):
try:
asyncio.get_running_loop()
except RuntimeError:
return asyncio.run(result)
# Inside a running event loop: submit to the shared singleton
# executor so threads are reused across calls.
pool = self._get_executor()
try:
return pool.submit(asyncio.run, result).result(timeout=self._BACKEND_TIMEOUT_S)
except TimeoutError as exc:
raise TimeoutError(
f"BingX backend call exceeded {self._BACKEND_TIMEOUT_S}s timeout"
) from exc
return result
def close(self) -> None:
"""V2: release the class-level thread-pool and any backend HTTP session."""
executor = self.__class__._EXECUTOR
if executor is not None:
with self.__class__._EXECUTOR_LOCK:
if self.__class__._EXECUTOR is executor:
self.__class__._EXECUTOR = None
executor.shutdown(wait=False)
_maybe_close_backend = getattr(self.backend, "close", None)
if _maybe_close_backend is not None:
try:
_maybe_close_backend()
except Exception:
pass
# ── telemetry lane — NON-BLOCKING **AND** LOSSLESS ───────────────────────
# Two standing principles, both non-negotiable, and they do NOT conflict:
#
# 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_q: Any = deque() # UNBOUNDED: the exec path never drops
self._tel_wake = threading.Event()
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: spill under pressure, publish, replay. Never drops."""
while True:
self._tel_wake.wait(timeout=0.5)
self._tel_wake.clear()
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 = q.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:
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) — "
"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
if zinc_plane is not None:
self._start_telemetry_lane()
def _publish_telemetry(
self,
*,
phase: str,
status: str,
intent: KernelIntent | None = None,
order: VenueOrder | None = None,
endpoint: str = "",
method: str = "",
message: str = "",
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:
# 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:
q = getattr(self, "_tel_q", None)
if q is None or self._telemetry_plane is None:
return
slot_id = 0
trade_id = ""
asset = ""
side = TradeSide.FLAT
action = ""
intent_id = ""
if intent is not None:
slot_id = int(getattr(intent, "slot_id", 0) or 0)
trade_id = str(getattr(intent, "trade_id", "") or "")
asset = str(getattr(intent, "asset", "") or "")
side = getattr(intent, "side", TradeSide.FLAT)
action = str(getattr(intent, "action", "") or "")
intent_id = str(getattr(intent, "intent_id", "") or "")
if order is not None:
slot_id = int(order.metadata.get("slot_id", slot_id) or slot_id)
trade_id = str(order.internal_trade_id or trade_id)
asset = str(order.metadata.get("asset") or asset)
side = order.side or side
# 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.
q.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 {}),
}
)
self._tel_wake.set()
except Exception as exc:
import logging as _log
_log.getLogger(__name__).warning(
"FIX(bingx_venue): telemetry enqueue failed (phase=%s) — swallowed, "
"order path unaffected: %s",
phase, exc,
)
def _call_backend(self, method_name: str, *args: Any, **kwargs: Any) -> Any:
method = getattr(self.backend, method_name, None)
if method is None:
raise AttributeError(f"backend has no method {method_name}")
return self._run(method(*args, **kwargs))
def _backend_snapshot(self, *, include_history: bool = False, timeout_ms: float = 5000.0):
"""Fetch a fresh snapshot from the backend and cache it thread-safely.
Design (industry best-practice reader-writer pattern):
- A caller that needs a fresh snapshot *waits* on ``_snapshot_ready``
before reading, so it never sees a stale partial write.
- While a snapshot fetch is in-flight, the lock is cleared; concurrent
callers block on ``_snapshot_ready`` with a timeout. If the fetch
succeeds in time they get the fresh snapshot; if it times out they
fall back to ``_last_snapshot`` (an eventually-consistent design —
stale data that *was* consistent is safer than no data).
- The write is guarded by ``_snap_lock`` so concurrent writes are
serialised and ``_last_snapshot`` is never partially assigned.
"""
if not self._snapshot_ready.wait(timeout=timeout_ms / 1000.0):
# Timeout waiting for a previous snapshot write — return the
# last-known-good snapshot rather than blocking the caller.
with self._snap_lock:
return self._last_snapshot
self._snapshot_ready.clear()
try:
snapshot = self._call_backend("refresh_state", None, include_history=include_history)
except Exception:
self._snapshot_ready.set()
raise
with self._snap_lock:
self._last_snapshot = snapshot
self._snapshot_ready.set()
return snapshot
@staticmethod
def _legacy_intent(intent: KernelIntent, *, kernel: Any | None = None) -> LegacyIntent:
action = LegacyDecisionAction.ENTER if intent.action == KernelCommandType.ENTER else LegacyDecisionAction.EXIT
side = LegacyTradeSide.SHORT if intent.side == TradeSide.SHORT else LegacyTradeSide.LONG
metadata = dict(intent.metadata)
metadata["_order_type"] = getattr(intent, "order_type", "MARKET")
metadata["_limit_price"] = float(getattr(intent, "limit_price", 0.0) or 0.0)
target_size = float(intent.target_size)
if intent.action == KernelCommandType.EXIT and kernel is not None:
try:
slot = kernel.slot(int(intent.slot_id))
active = slot.active_exit_order
if active is not None and str(slot.trade_id) == str(intent.trade_id):
target_size = float(active.intended_size or target_size)
metadata["exit_leg_index"] = int(slot.active_leg_index or 0)
metadata["exit_leg_ratio"] = float(slot.next_exit_ratio())
except Exception:
pass
return LegacyIntent(
timestamp=intent.timestamp,
trade_id=intent.trade_id,
decision_id=intent.intent_id,
asset=intent.asset,
action=action,
side=side,
reason=intent.reason,
target_size=target_size,
leverage=float(intent.leverage),
reference_price=float(intent.reference_price),
confidence=1.0,
bars_held=0,
exit_leg_ratios=tuple(intent.exit_leg_ratios or (1.0,)),
metadata=metadata,
)
async def connect(self) -> bool:
"""Async connect — awaits backend.connect() in the caller's event loop.
The old sync path called self._run(backend.connect()) which spawns a
thread-pool asyncio.run(), creating the httpx AsyncClient in a temporary
loop that immediately closes. Every subsequent request then raises
"Event loop is closed" or "bound to a different event loop".
Awaiting directly here fixes that: the client is always created in the
main running loop.
"""
conn_fn = getattr(self.backend, "connect", None)
if conn_fn is not None:
result = conn_fn()
if inspect.isawaitable(result):
await result
# backend.connect() already called refresh_state() — no second fetch needed
return True
async def cancel_async(self, order: VenueOrder, *, reason: str = "") -> List[VenueEvent]:
"""Async cancel — runs in the caller's event loop, no thread-pool deadlock.
The sync cancel() path goes through _call_backend → _run → thread-pool →
asyncio.run() in a new thread. The aiohttp session is bound to the main
event loop, so using it from a different loop deadlocks — same bug that
was fixed for submit via submit_async. This version awaits backend.cancel()
directly in the caller's (main) event loop.
"""
self._publish_telemetry(
phase="cancel:start",
status="REQUESTED",
order=order,
endpoint="/openApi/swap/v2/trade/order",
method="DELETE",
message=reason,
details={"asset": str(order.metadata.get("asset") or "")},
)
cancel_fn = getattr(self.backend, "cancel", None)
if cancel_fn is not None:
response = await cancel_fn(order, reason=reason)
else:
response = None
events = self._events_from_cancel(order, response, None, None, reason=reason)
return events
def cancel(self, order: VenueOrder, *, reason: str = "") -> List[VenueEvent]:
# _events_from_cancel never reads before/after — snapshots are dead weight.
# NOTE: if backend.cancel is async (BingxDirectExecutionAdapter), this sync
# path goes through the thread-pool and will deadlock in a running event loop.
# Use cancel_async() from async contexts (process_intent_async already does).
self._publish_telemetry(
phase="cancel:start",
status="REQUESTED",
order=order,
endpoint="/openApi/swap/v2/trade/order",
method="DELETE",
message=reason,
details={"asset": str(order.metadata.get("asset") or "")},
)
response = None
if hasattr(self.backend, "cancel"):
response = self._call_backend("cancel", order, reason=reason)
else:
client = getattr(self.backend, "_client", None)
instrument_symbol = ""
if hasattr(self.backend, "_instrument_venue_symbol"):
asset = str(order.metadata.get("asset") or "")
if not asset:
slot_id = int(order.metadata.get("slot_id", 0) or 0)
if hasattr(self, "_kernel_ref") and self._kernel_ref is not None:
try:
asset = self._kernel_ref.slot(slot_id).asset
except Exception:
pass
if not asset:
asset = str(order.metadata.get("asset") or "")
instrument_symbol = str(self.backend._instrument_venue_symbol(asset)) if asset else ""
if client is None or not instrument_symbol:
raise RuntimeError("backend does not expose a cancel surface")
params = {"symbol": instrument_symbol}
if order.venue_order_id:
params["orderId"] = order.venue_order_id
else:
params["clientOrderId"] = order.venue_client_id
try:
response = self._run(client.signed_delete("/openApi/swap/v2/trade/order", params))
except BingxHttpError as exc:
# W10: map HTTP error class to status — 429/5xx are transient, 4xx are real rejections
response = {"status": _http_error_status(str(exc)), "msg": str(exc), "orderId": order.venue_order_id, "clientOrderId": order.venue_client_id}
events = self._events_from_cancel(order, response, None, None, reason=reason)
return events
def open_orders(self) -> List[VenueOrder]:
# Use backend._state (populated by await backend.connect()) rather than
# _backend_snapshot() → _call_backend() → _run() → pool.submit(asyncio.run)
# which spawns a temporary event loop, creates the httpx AsyncClient inside
# it, then closes that loop — every subsequent HTTP call then raises
# "Event loop is closed" or "asyncio.locks.Event bound to a different loop".
backend_state = getattr(self.backend, "_state", None)
if backend_state is not None:
return [_venue_order_from_row(row) for row in (backend_state.open_orders or [])]
snapshot = self._backend_snapshot(include_history=False)
return [_venue_order_from_row(row) for row in (snapshot.open_orders or [])]
def open_positions(self) -> List[dict[str, Any]]:
# Same rationale as open_orders(): prefer cached backend._state to avoid
# the thread-pool asyncio.run() path that corrupts the httpx session.
backend_state = getattr(self.backend, "_state", None)
if backend_state is not None:
return [dict(row) for row in (backend_state.open_positions or {}).values()]
snapshot = self._backend_snapshot(include_history=False)
return [dict(row) for row in (snapshot.open_positions or {}).values()]
async def reconcile(self) -> List[VenueEvent]: # type: ignore[override]
"""Fetch open-order state from BingX and return any pending VenueEvents.
WHY ASYNC: the old sync version called _backend_snapshot() → _call_backend()
→ _run() → pool.submit(asyncio.run, coro).result(timeout=30s). That spawned
a *new* event loop in a thread-pool thread. The BingxHttpClient (aiohttp
session) is bound to the *main* event loop — using it from a different loop
silently deadlocks. BingX VST responds in ~500ms; the deadlock made every
reconcile call block the main event loop for the full 30s timeout.
FIX: declare async, call backend.refresh_state() directly with await so it
runs in the *caller's* (main) event loop where the session lives.
pump_venue_events() already has `if inspect.isawaitable(events): await events`
— zero caller changes required.
include_history=False: all_orders/all_fills require a symbol (symbol=None
skips them anyway), so include_history=True was fetching nothing extra.
"""
# FILL VISIBILITY (2026-06-10): when the kernel slot owns an asset,
# fetch symbol-scoped history (all_orders + all_fills) so a maker
# entry that FILLED — and therefore left openOrders — reaches the FSM
# as a FULL_FILL event. With symbol=None the snapshot skips history
# entirely: the FSM stayed fill-blind (slot size 0 in ENTRY_WORKING),
# the DecisionEngine saw "no position", and re-entered → the live
# double-entries at 15:20 and 17:24 UTC.
self._publish_telemetry(
phase="reconcile:start",
status="REQUESTED",
endpoint="/openApi/swap/v2/trade/openOrders",
method="GET",
details={"include_history": True},
)
recon_symbol = None
kernel = getattr(self, "_kernel_ref", None)
if kernel is not None:
try:
slot = kernel.slot(0)
if not slot.is_free() and getattr(slot, "asset", ""):
recon_symbol = str(slot.asset)
except Exception:
recon_symbol = None
try:
snapshot = await self.backend.refresh_state(
recon_symbol, include_history=recon_symbol is not None
)
except Exception as exc:
import logging as _log
_log.getLogger(__name__).warning("reconcile: refresh_state failed: %s", exc)
self._publish_telemetry(
phase="reconcile:error",
status="ERROR",
endpoint="/openApi/swap/v2/trade/openOrders",
method="GET",
message=str(exc),
details={"symbol": recon_symbol or ""},
)
return []
self._publish_telemetry(
phase="reconcile:done",
status="OK",
endpoint="/openApi/swap/v2/trade/openOrders",
method="GET",
details={
"symbol": recon_symbol or "",
"open_orders": len(getattr(snapshot, "open_orders", []) or []),
"positions": len(getattr(snapshot, "open_positions", {}) or {}),
"fills": len(getattr(snapshot, "all_fills", []) or []),
},
)
return self._events_from_snapshot(snapshot)
def submit(self, intent: KernelIntent) -> List[VenueEvent]:
# Snapshots dropped: receipt executedQty fields take precedence (same as submit_async)
self._publish_telemetry(
phase="submit:start",
status="REQUESTED",
intent=intent,
endpoint="/openApi/swap/v2/trade/order",
method="POST",
details={"action": intent.action.value, "order_type": str(getattr(intent, "order_type", "MARKET") or "MARKET")},
)
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)
# ── 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
# 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]:
"""Async submit — runs in the caller's event loop, no thread-pool deadlock.
The sync submit() calls _backend_snapshot() × 2 + submit_intent via
_run() → asyncio.run() in a thread-pool → new event loop → aiohttp
session (main-loop-bound) deadlocks → 30s timeout on every ENTER/EXIT.
This version awaits the backend directly. The before/after snapshots
are omitted: fill size comes from the receipt's executedQty field, and
the WS account stream delivers FULL_FILL events independently.
Passing None for snapshots makes _filled_size_from_snapshots return 0.0
(a safe fallback; the receipt fields take precedence).
"""
self._publish_telemetry(
phase="submit:start",
status="REQUESTED",
intent=intent,
endpoint="/openApi/swap/v2/trade/order",
method="POST",
details={"action": intent.action.value, "order_type": str(getattr(intent, "order_type", "MARKET") or "MARKET")},
)
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)
# ── 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
# 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
ack_row = dict(getattr(receipt, "raw_ack", {}) or {})
status = _normalize_status(getattr(receipt, "status", "") or _row_text(ack_row, "status", default="NEW"))
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))
if status in {"RATE_LIMITED", "THROTTLED"}:
return [
VenueEvent(
timestamp=getattr(receipt, "timestamp", datetime.now(timezone.utc)),
event_id=_event_id(self._event_seq),
trade_id=intent.trade_id,
slot_id=intent.slot_id,
kind=KernelEventKind.RATE_LIMITED,
status=VenueEventStatus.RATE_LIMITED,
venue_order_id=order_id,
venue_client_id=client_order_id,
side=intent.side,
asset=intent.asset,
price=safe_float(getattr(receipt, "price", 0.0), 0.0),
size=float(intent.target_size or 0.0),
filled_size=0.0,
remaining_size=float(intent.target_size or 0.0),
reason=_row_text(ack_row, "msg", "message", default="BINGX_RATE_LIMITED"),
raw_payload=ack_row or json_safe(receipt),
metadata={"intent_id": intent.intent_id, "action": intent.action.value, "retry_after_ms": _rate_limit_retry_after_ms(ack_row)},
)
]
base_event = VenueEvent(
timestamp=getattr(receipt, "timestamp", datetime.now(timezone.utc)),
event_id=_event_id(self._event_seq),
trade_id=intent.trade_id,
slot_id=intent.slot_id,
kind=KernelEventKind.ORDER_ACK,
status=VenueEventStatus.ACKED,
venue_order_id=order_id,
venue_client_id=client_order_id,
side=intent.side,
asset=intent.asset,
price=safe_float(getattr(receipt, "price", 0.0), 0.0),
size=float(intent.target_size or 0.0),
filled_size=0.0,
remaining_size=float(intent.target_size or 0.0),
reason="",
raw_payload=ack_row or json_safe(receipt),
metadata={"intent_id": intent.intent_id, "action": intent.action.value},
)
if status in {"REJECTED", "FAILED"}:
return [
VenueEvent(
**{**base_event.__dict__, "event_id": _event_id(self._event_seq), "kind": KernelEventKind.ORDER_REJECT, "status": VenueEventStatus.REJECTED, "reason": _row_text(ack_row, "msg", "message", default="BINGX_ORDER_REJECTED")},
)
]
# Extract friction fields annotated by submit_intent (Gap 1/2/3).
fee_estimated = float(ack_row.get("_fee_estimated") or 0.0)
fee_source = str(ack_row.get("_fee_source") or "")
is_maker_est = bool(ack_row.get("_is_maker_est", False))
mark_at_submit = float(ack_row.get("_mark_at_submit") or 0.0)
slippage_bps = float(ack_row.get("_slippage_bps") or 0.0)
exchange_ts = int(ack_row.get("_exchange_ts") or 0)
events = [base_event]
fill_status = _venue_event_status_from_row(status)
filled_size = _row_float(ack_row, "executedQty", "cumFilledQty", "filledQty", "lastFilledQty", default=0.0)
snapshot_fill_size = self._filled_size_from_snapshots(before, after, intent.asset)
if filled_size <= 0:
filled_size = snapshot_fill_size
emit_fill = fill_status in {VenueEventStatus.PARTIALLY_FILLED, VenueEventStatus.FILLED} or snapshot_fill_size > 0.0
if emit_fill:
if filled_size <= 0:
filled_size = float(intent.target_size or 0.0)
remaining_size = max(0.0, float(intent.target_size or 0.0) - float(filled_size))
fill_kind = KernelEventKind.FULL_FILL if fill_status == VenueEventStatus.FILLED or remaining_size <= 1e-12 else KernelEventKind.PARTIAL_FILL
events.append(
VenueEvent(
timestamp=base_event.timestamp,
event_id=_event_id(self._event_seq),
trade_id=intent.trade_id,
slot_id=intent.slot_id,
kind=fill_kind,
status=VenueEventStatus.FILLED if fill_kind == KernelEventKind.FULL_FILL else VenueEventStatus.PARTIALLY_FILLED,
venue_order_id=order_id,
venue_client_id=client_order_id,
side=intent.side,
asset=intent.asset,
# FILL price must be a TRUE fill price (avgPrice/lastFillPrice).
# Never fall back to the order's nominal "price" or the submit
# receipt price: for BingX MARKET orders that is the protective
# bound (±20-25% from mark) — it poisoned realized PnL on every
# market fill (FET −$5,990 mis-book, 2026-06-11). 0.0 = unknown;
# the kernel refuses to compute PnL from a missing price.
price=safe_float(_row_float(ack_row, "avgPrice", "ap", "lastFillPrice", "L", default=0.0), 0.0),
size=float(intent.target_size or 0.0),
filled_size=float(filled_size),
remaining_size=float(remaining_size),
reason="",
raw_payload=ack_row or json_safe(receipt),
metadata={"intent_id": intent.intent_id, "action": intent.action.value},
# Gap 1/2/3: fee + friction fields populated from submit_intent annotations.
# fee_source="ESTIMATED_*" until WS FILL_SETTLED updates it to "WS_SETTLED".
fee=fee_estimated,
fee_asset="USDT",
fee_source=fee_source,
is_maker=is_maker_est,
exchange_ts=exchange_ts,
slippage_bps=slippage_bps,
mark_at_submit=mark_at_submit,
)
)
return events
def _events_from_cancel(self, order: VenueOrder, response: Any, before, after, *, reason: str = "") -> List[VenueEvent]: # noqa: ANN001
raw = response if isinstance(response, dict) else {}
status = _normalize_status(_row_text(raw, "status", default="CANCELED"))
if status in {"RATE_LIMITED", "THROTTLED"}:
self._publish_telemetry(
phase="cancel:done",
status=status,
order=order,
endpoint="/openApi/swap/v2/trade/order",
method="DELETE",
message=reason or _row_text(raw, "msg", "message", default="BINGX_RATE_LIMITED"),
venue_order_status=VenueEventStatus.RATE_LIMITED.value,
venue_event_kind=KernelEventKind.RATE_LIMITED.value,
retry_after_ms=_rate_limit_retry_after_ms(raw),
details={"order_status": status, "asset": str(order.metadata.get("asset") or "")},
)
return [
VenueEvent(
timestamp=datetime.now(timezone.utc),
event_id=_event_id(self._event_seq),
trade_id=order.internal_trade_id or order.venue_client_id,
slot_id=int(order.metadata.get("slot_id", 0) or 0),
kind=KernelEventKind.RATE_LIMITED,
status=VenueEventStatus.RATE_LIMITED,
venue_order_id=order.venue_order_id,
venue_client_id=order.venue_client_id,
side=order.side,
asset=str(order.metadata.get("asset") or ""),
price=safe_float(_row_float(raw, "avgPrice", "ap", "price", "lastFillPrice", default=order.average_fill_price), 0.0),
size=float(order.intended_size or 0.0),
filled_size=float(order.filled_size or 0.0),
remaining_size=float(order.remaining_size),
reason=reason or _row_text(raw, "msg", "message", default="BINGX_RATE_LIMITED"),
raw_payload=raw or {"orderId": order.venue_order_id, "clientOrderId": order.venue_client_id, "status": status or "RATE_LIMITED"},
metadata={**dict(order.metadata), "retry_after_ms": _rate_limit_retry_after_ms(raw)},
)
]
event_status = _venue_event_status_from_row(status)
kind = KernelEventKind.CANCEL_ACK if event_status == VenueEventStatus.CANCELED else KernelEventKind.CANCEL_REJECT
if event_status == VenueEventStatus.CANCELED_REJECTED:
kind = KernelEventKind.CANCEL_REJECT
self._publish_telemetry(
phase="cancel:done",
status=status or event_status.value,
order=order,
endpoint="/openApi/swap/v2/trade/order",
method="DELETE",
message=reason or _row_text(raw, "msg", "message", default=""),
venue_order_status=event_status.value,
venue_event_kind=kind.value,
retry_after_ms=_rate_limit_retry_after_ms(raw),
details={"order_status": status, "asset": str(order.metadata.get("asset") or "")},
)
return [
VenueEvent(
timestamp=datetime.now(timezone.utc),
event_id=_event_id(self._event_seq),
trade_id=order.internal_trade_id or order.venue_client_id,
slot_id=int(order.metadata.get("slot_id", 0) or 0),
kind=kind,
status=event_status,
venue_order_id=order.venue_order_id,
venue_client_id=order.venue_client_id,
side=order.side,
asset=str(order.metadata.get("asset") or ""),
price=safe_float(_row_float(raw, "avgPrice", "ap", "price", "lastFillPrice", default=order.average_fill_price), 0.0),
size=float(order.intended_size or 0.0),
filled_size=float(order.filled_size or 0.0),
remaining_size=float(order.remaining_size),
reason=reason or _row_text(raw, "msg", "message", default="BINGX_CANCEL_ACK" if kind == KernelEventKind.CANCEL_ACK else "BINGX_CANCEL_REJECT"),
raw_payload=raw or {"orderId": order.venue_order_id, "clientOrderId": order.venue_client_id, "status": status or event_status.value},
metadata=dict(order.metadata),
)
]
def _events_from_snapshot(self, snapshot: Any) -> List[VenueEvent]: # noqa: ANN001
events: list[VenueEvent] = []
seen: set[tuple[str, str, str]] = set()
for row in getattr(snapshot, "open_orders", []) or []:
if not isinstance(row, dict):
continue
event = self._event_from_row(row, slot_id=0)
key = (event.venue_client_id, event.venue_order_id, event.kind.value)
if key not in seen:
seen.add(key)
events.append(event)
for row in getattr(snapshot, "all_orders", []) or []:
if not isinstance(row, dict):
continue
event = self._event_from_row(row, slot_id=0)
key = (event.venue_client_id, event.venue_order_id, event.kind.value)
if key not in seen:
seen.add(key)
events.append(event)
for row in getattr(snapshot, "all_fills", []) or []:
if not isinstance(row, dict):
continue
event = self._fill_event_from_row(row)
key = (event.venue_client_id, event.venue_order_id, event.kind.value)
if key not in seen:
seen.add(key)
events.append(event)
return events
def _event_from_row(self, row: dict[str, Any], *, slot_id: int) -> VenueEvent:
status = _normalize_status(_row_text(row, "status", "X", default="NEW"))
event_status = _venue_event_status_from_row(status)
kind = {
VenueEventStatus.ACKED: KernelEventKind.ORDER_ACK,
VenueEventStatus.PARTIALLY_FILLED: KernelEventKind.PARTIAL_FILL,
VenueEventStatus.FILLED: KernelEventKind.FULL_FILL,
VenueEventStatus.CANCELED: KernelEventKind.CANCEL_ACK,
VenueEventStatus.REJECTED: KernelEventKind.ORDER_REJECT,
VenueEventStatus.CANCELED_REJECTED: KernelEventKind.CANCEL_REJECT,
VenueEventStatus.RATE_LIMITED: KernelEventKind.RATE_LIMITED,
}.get(event_status, KernelEventKind.ORDER_ACK)
size = _row_float(row, "origQty", "quantity", "q", "positionAmt", default=0.0)
filled = _row_float(row, "executedQty", "cumFilledQty", "filledQty", "z", "lastFilledQty", default=0.0)
if filled <= 0.0 and kind in {KernelEventKind.PARTIAL_FILL, KernelEventKind.FULL_FILL}:
filled = size
# For FILL events only true fill-price fields qualify; the nominal
# "price" is the MARKET bound price on BingX and must never feed PnL.
# Non-fill events (ACK/CANCEL/REJECT) may keep it as informational.
if kind in {KernelEventKind.PARTIAL_FILL, KernelEventKind.FULL_FILL}:
row_price = _row_float(row, "avgPrice", "ap", "lastFillPrice", "L", default=0.0)
else:
row_price = _row_float(row, "avgPrice", "ap", "price", "lastFillPrice", default=0.0)
return VenueEvent(
timestamp=datetime.now(timezone.utc),
event_id=_event_id(self._event_seq),
trade_id=_row_text(row, "tradeId", "trade_id", default=_row_text(row, "clientOrderId", "clientOrderID", default="")),
slot_id=slot_id,
kind=kind,
status=event_status,
venue_order_id=_row_text(row, "orderId", "orderID", "id", default=""),
venue_client_id=_row_text(row, "clientOrderID", "clientOrderId", "c", default=""),
side=_trade_side_from_row(row),
asset=_row_text(row, "symbol", default=""),
price=safe_float(row_price, 0.0),
size=abs(float(size or 0.0)),
filled_size=abs(float(filled or 0.0)),
remaining_size=max(0.0, abs(float(size or 0.0)) - abs(float(filled or 0.0))),
reason=_row_text(row, "msg", "message", default=""),
raw_payload=dict(row),
metadata={"source": "bingx"},
)
def _fill_event_from_row(self, row: dict[str, Any]) -> VenueEvent:
status = _normalize_status(_row_text(row, "status", "X", default="FILLED"))
event_status = _venue_event_status_from_row(status)
kind = KernelEventKind.FULL_FILL if event_status == VenueEventStatus.FILLED else KernelEventKind.PARTIAL_FILL
return VenueEvent(
timestamp=datetime.now(timezone.utc),
event_id=_event_id(self._event_seq),
trade_id=_row_text(row, "tradeId", "trade_id", default=_row_text(row, "clientOrderId", "clientOrderID", default="")),
slot_id=0,
kind=kind,
status=event_status,
venue_order_id=_row_text(row, "orderId", "orderID", "id", default=""),
venue_client_id=_row_text(row, "clientOrderID", "clientOrderId", "c", default=""),
side=_trade_side_from_row(row),
asset=_row_text(row, "symbol", default=""),
# True fill-price fields only — nominal "price" excluded (MARKET
# bound-price poisoning; see _events_from_submit note).
price=safe_float(_row_float(row, "lastFillPrice", "L", "avgPrice", "ap", default=0.0), 0.0),
size=abs(_row_float(row, "executedQty", "z", "lastFilledQty", default=0.0)),
filled_size=abs(_row_float(row, "lastFilledQty", "l", "z", default=0.0)),
remaining_size=max(0.0, abs(_row_float(row, "executedQty", "z", "lastFilledQty", default=0.0)) - abs(_row_float(row, "lastFilledQty", "l", "z", default=0.0))),
reason=_row_text(row, "msg", "message", default=""),
raw_payload=dict(row),
metadata={"source": "bingx"},
)
@staticmethod
def _filled_size_from_snapshots(before: Any, after: Any, asset: str) -> float: # noqa: ANN001
def _lookup(snapshot: Any) -> float:
positions = getattr(snapshot, "open_positions", {}) or {}
for key, row in positions.items():
symbol = _row_text(row, "symbol", default=str(key))
if symbol.replace("-", "").replace("_", "").upper() == asset.replace("-", "").replace("_", "").upper():
return _position_qty(row)
return 0.0
before_qty = _lookup(before)
after_qty = _lookup(after)
diff = abs(before_qty - after_qty)
return diff