diff --git a/prod/bingx/http.py b/prod/bingx/http.py index c8ff5e6..f4a84b2 100644 --- a/prod/bingx/http.py +++ b/prod/bingx/http.py @@ -28,7 +28,53 @@ from .urls import get_rest_base_urls class BingxHttpError(RuntimeError): - pass + """A BingX HTTP failure, carrying WHETHER THE ORDER MAY EXIST. + + An operation with a side effect has three outcomes, not two, and rollback is + sound for only two of them: + + NOT_ATTEMPTED — the request never left us (no creds, connect refused, DNS). + Provably no order. Safe to roll back. + REFUSED — the venue answered and rejected it (business error code, + 4xx). Provably no order. Safe to roll back. + INDETERMINATE — the request was sent and we did not get a usable answer + (read timeout, connection reset, 5xx, unparseable body). + THE ORDER MAY BE LIVE. Rolling back here asserts "no order + exists" on exactly the evidence that cannot support it, and + orphans a real position. Reconcile by clientOrderId instead. + + Default is INDETERMINATE on purpose: when we do not know, we must not claim + absence. Unknown is not flat. + """ + + NOT_ATTEMPTED = "NOT_ATTEMPTED" + REFUSED = "REFUSED" + INDETERMINATE = "INDETERMINATE" + + def __init__(self, message: str, *, effect: str = "INDETERMINATE") -> None: + super().__init__(message) + self.effect = effect + + @property + def order_may_exist(self) -> bool: + return self.effect == self.INDETERMINATE + + +def _effect_of_httpx_error(exc: Exception) -> str: + """Classify a transport failure by whether the request could have been acted on.""" + # Never left the box -> provably no order. + if isinstance(exc, (httpx.ConnectError, httpx.ConnectTimeout)): + return BingxHttpError.NOT_ATTEMPTED + # Sent, but the answer was lost (read/write/pool timeout, reset, protocol + # error). The venue may well have matched it. UNKNOWN. + return BingxHttpError.INDETERMINATE + + +def _effect_of_status(status_code: int) -> str: + """4xx = the venue judged and refused. 5xx = it may have acted before failing.""" + if 400 <= status_code < 500: + return BingxHttpError.REFUSED + return BingxHttpError.INDETERMINATE @dataclass(frozen=True) @@ -413,7 +459,8 @@ class BingxHttpClient: payload = dict(params) if signed: if not self._api_key or not self._secret_key: - raise BingxHttpError("BingX API credentials are required for signed requests") + raise BingxHttpError("BingX API credentials are required for signed requests", + effect=BingxHttpError.NOT_ATTEMPTED) payload = build_signed_params( payload, self._secret_key, @@ -453,13 +500,15 @@ class BingxHttpClient: rate_limited=response.status_code == 429, retry_after_ms=self._rate_limits.snapshot().rest_reset_ms, ) - last_error = BingxHttpError(f"HTTP {response.status_code}: {text or response.reason_phrase}") + last_error = BingxHttpError(f"HTTP {response.status_code}: {text or response.reason_phrase}", + effect=_effect_of_status(response.status_code)) if base_index < len(self._base_urls) - 1: continue if attempt < max_retries: await asyncio.sleep(delay) break - raise BingxHttpError(f"HTTP {response.status_code}: {text or response.reason_phrase}") + raise BingxHttpError(f"HTTP {response.status_code}: {text or response.reason_phrase}", + effect=_effect_of_status(response.status_code)) self._circuit_breaker.record_success() if hostname is not None: await self._maybe_refresh_dns_cache(hostname) @@ -487,7 +536,8 @@ class BingxHttpClient: except httpx.HTTPError as exc: self._logger.warning("HTTP error [att=%d url=%d %s]: %s — %s", attempt, base_index, method, type(exc).__name__, exc) - last_error = BingxHttpError(f"{method} {path} failed: {exc}") + last_error = BingxHttpError(f"{method} {path} failed: {exc}", + effect=_effect_of_httpx_error(exc)) if self._is_dns_resolution_error(exc): if hostname is not None: cached_ips = self._dns_cache.resolve(hostname) @@ -528,12 +578,14 @@ class BingxHttpClient: return self._unwrap_response(data) return data except Exception as fallback_exc: - last_error = BingxHttpError(f"{method} {path} failed: {fallback_exc}") + last_error = BingxHttpError(f"{method} {path} failed: {fallback_exc}", + effect=_effect_of_httpx_error(fallback_exc)) if base_index < len(self._base_urls) - 1: continue if last_error is not None: raise last_error - raise BingxHttpError(f"{method} {path} failed: {exc}") + raise BingxHttpError(f"{method} {path} failed: {exc}", + effect=_effect_of_httpx_error(exc)) delay = self._circuit_breaker.record_failure() if base_index < len(self._base_urls) - 1: continue @@ -584,12 +636,14 @@ class BingxHttpClient: return self._unwrap_response(data) return data except Exception as fallback_exc: - last_error = BingxHttpError(f"{method} {path} failed: {fallback_exc}") + last_error = BingxHttpError(f"{method} {path} failed: {fallback_exc}", + effect=_effect_of_httpx_error(fallback_exc)) if base_index < len(self._base_urls) - 1: continue if last_error is not None: raise last_error - raise BingxHttpError(f"{method} {path} failed: {exc}") + raise BingxHttpError(f"{method} {path} failed: {exc}", + effect=_effect_of_httpx_error(exc)) delay = self._circuit_breaker.record_failure() if base_index < len(self._base_urls) - 1: continue @@ -625,7 +679,8 @@ class BingxHttpClient: def _unwrap_response(payload: dict[str, Any]) -> Any: code = int(payload.get("code", -1)) if code != 0: - raise BingxHttpError(payload.get("msg", f"BingX error code {code}")) + raise BingxHttpError(payload.get("msg", f"BingX error code {code}"), + effect=BingxHttpError.REFUSED) return payload.get("data") @staticmethod diff --git a/prod/clean_arch/adapters/bingx_direct.py b/prod/clean_arch/adapters/bingx_direct.py index f4d3f52..8706f54 100644 --- a/prod/clean_arch/adapters/bingx_direct.py +++ b/prod/clean_arch/adapters/bingx_direct.py @@ -578,6 +578,89 @@ class BingxDirectExecutionAdapter(ExecutionPort): # (Bybit, OKX, Binance) will have the same REST/WS split. # ───────────────────────────────────────────────────────────────────────── + async def _lookup_own_order_by_client_id( + self, symbol: str, client_order_id: str, *, attempts: int = 3 + ) -> Any: + """Did THE ONE ORDER I JUST SENT land? A read-only point lookup. Nothing else. + + *** THIS IS NOT A RECONCILE. DO NOT GROW IT INTO ONE. *** + + PINK's reconcile was a periodic venue->kernel STATE SYNC: it pulled broad venue + state and wrote it back into slots, which produced ghost positions, reseed + loops and double entries. That pattern is forbidden here. Hard boundaries, + and any future edit that crosses one is a bug: + + - Scope : exactly ONE order, identified by an idempotency key WE chose + before sending. Never "all orders", never "open positions". + - Side effects : NONE. It is a GET. It never cancels, never re-POSTs, never + writes a slot, never adopts unknown venue state. + - Re-submitting: NEVER. Re-POSTing an indeterminate order is exactly how you + get a double fill. We only ever ask. + - Bounded : `attempts` tries with backoff, then it gives up and SAYS SO. + It is allowed to answer "I don't know". It is never allowed + to guess. + + Why it must exist: a read timeout / reset / 5xx means the request WAS SENT and + the answer was lost. The order may be live. The only alternative to asking is + guessing, and guessing "rejected" is what orphans a real position. + + Returns: + dict -> the venue's row for OUR order. It exists. Adopt it as truth. + "ABSENT" -> the venue authoritatively has no such clientOrderId. Only now is + a rollback sound. + None -> truth not established. The caller MUST NOT assume flat. + """ + delay = 0.5 + for attempt in range(1, attempts + 1): + try: + resp = await self._client.signed_get( + "/openApi/swap/v2/trade/order", + {"symbol": symbol, "clientOrderID": client_order_id}, + ) + row = dict(unwrap_order_payload(resp)) if isinstance(resp, dict) else {} + handle = row.get("orderId") or row.get("orderID") or row.get("order_id") + if handle: + LOGGER.critical( + "FIX(bingx_direct): INDETERMINATE submit RECONCILED — order IS LIVE " + "at venue (clientOrderId=%s orderId=%s status=%s). Adopting venue truth " + "instead of fabricating a REJECT.", + client_order_id, handle, row.get("status"), + ) + return row + # Venue answered, and it has no such order. + LOGGER.warning( + "FIX(bingx_direct): INDETERMINATE submit reconciled — venue reports NO order " + "for clientOrderId=%s. Rollback is sound.", client_order_id, + ) + return "ABSENT" + except BingxHttpError as exc: + # A query that the venue REFUSES (e.g. "order not exist") is an + # authoritative answer: the order is not there. + if getattr(exc, "effect", "") == BingxHttpError.REFUSED: + LOGGER.warning( + "FIX(bingx_direct): venue refused the lookup for clientOrderId=%s (%s) " + "— treating as ABSENT, rollback is sound.", client_order_id, exc, + ) + return "ABSENT" + LOGGER.error( + "FIX(bingx_direct): reconcile attempt %d/%d for clientOrderId=%s failed: %s", + attempt, attempts, client_order_id, exc, + ) + except Exception as exc: # transport died mid-lookup + LOGGER.error( + "FIX(bingx_direct): reconcile attempt %d/%d for clientOrderId=%s errored: %s", + attempt, attempts, client_order_id, exc, + ) + if attempt < attempts: + await asyncio.sleep(delay) + delay *= 2 + LOGGER.critical( + "FIX(bingx_direct): COULD NOT ESTABLISH ORDER TRUTH for clientOrderId=%s after %d " + "attempts. The order MAY BE LIVE. Refusing to claim it was rejected — the kernel " + "must NOT roll back to flat.", client_order_id, attempts, + ) + return None + async def submit_intent(self, intent: Intent) -> ExecutionReceipt: symbol = self._instrument_venue_symbol(intent.asset) if intent.action == DecisionAction.EXIT: @@ -707,16 +790,66 @@ class BingxDirectExecutionAdapter(ExecutionPort): 0.0, ) except BingxHttpError as exc: - status = "RATE_LIMITED" if _is_rate_limited_error(exc) else "REJECTED" - ack_row = { - "status": status, - "msg": str(exc), - "symbol": symbol, - "clientOrderId": client_order_id, - } - fill_price = 0.0 ack = None is_limit = False + fill_price = 0.0 + effect = getattr(exc, "effect", BingxHttpError.INDETERMINATE) + + if effect == BingxHttpError.INDETERMINATE and not _is_rate_limited_error(exc): + # THE ORDER MAY BE LIVE. Before this branch existed, a BingX read + # timeout / 5xx was reported to the kernel as a confident REJECTED + # with fill_qty=0 — the kernel rolled the slot back to flat while the + # venue happily kept the position. That is the orphan-maker. + # + # We do not guess. We ask the venue about OUR OWN clientOrderId. + LOGGER.critical( + "FIX(bingx_direct): submit outcome INDETERMINATE (%s) symbol=%s " + "clientOrderId=%s — order MAY be live. Looking it up; NOT assuming " + "rejection.", exc, symbol, client_order_id, + ) + found = await self._lookup_own_order_by_client_id(symbol, client_order_id) + if isinstance(found, dict): + # It exists. The venue is the authority — adopt its row as the ack. + ack_row = dict(found) + status = str(ack_row.get("status") or "ACKED") + for key in ("avgPrice", "avgFilledPrice", "price", "lastFillPrice"): + try: + value = float(ack_row.get(key) or 0.0) + except Exception: + value = 0.0 + if value > 0: + fill_price = value + break + elif found == "ABSENT": + # The venue authoritatively has no such order. Rollback is sound. + status = "REJECTED" + ack_row = { + "status": status, + "msg": f"venue confirms no such order after indeterminate submit: {exc}", + "symbol": symbol, + "clientOrderId": client_order_id, + } + else: + # Truth not established. We refuse to claim rejection. The venue + # layer turns this into a VenueIndeterminateError so the kernel + # does NOT roll back to flat; the E-feed FILL settles it if it filled. + status = "INDETERMINATE" + ack_row = { + "status": status, + "msg": f"order state UNKNOWN after indeterminate submit: {exc}", + "symbol": symbol, + "clientOrderId": client_order_id, + } + else: + # NOT_ATTEMPTED / REFUSED / rate-limited: the venue provably does not + # have this order. A rejection here is the truth, not a guess. + status = "RATE_LIMITED" if _is_rate_limited_error(exc) else "REJECTED" + ack_row = { + "status": status, + "msg": str(exc), + "symbol": symbol, + "clientOrderId": client_order_id, + } # ── Gap 2: fee estimation (ESTIMATED_TAKER / ESTIMATED_MAKER) ──────── # BingX REST ACK does not include commission. WS FILL_SETTLED will deliver @@ -725,7 +858,10 @@ class BingxDirectExecutionAdapter(ExecutionPort): # FIX 2026-07-11 (bingx_direct): never fabricate a fill for a rejected order — the # target_size fallback applies only to accepted MARKET acks whose fill # arrives later via WS (BingX REST ack omits executedQty on those). - if status in ("REJECTED", "RATE_LIMITED"): + if status in ("REJECTED", "RATE_LIMITED", "INDETERMINATE"): + # INDETERMINATE: we do not know that it filled, so we must not invent a + # fill size from target_size. We also do not know that it DIDN'T — that is + # why the status is not REJECTED. The venue layer escalates from here. fill_qty = 0.0 else: fill_qty = float(ack_row.get("executedQty") or ack_row.get("filledQty") or diff --git a/prod/clean_arch/dita_v2/bingx_venue.py b/prod/clean_arch/dita_v2/bingx_venue.py index db64642..4dac5eb 100644 --- a/prod/clean_arch/dita_v2/bingx_venue.py +++ b/prod/clean_arch/dita_v2/bingx_venue.py @@ -45,6 +45,7 @@ from .utils import json_safe from .utils import safe_float from .venue import VenueAdapter from .venue import VenuePostAckError +from .venue import VenueIndeterminateError def _row_text(row: dict[str, Any], *keys: str, default: str = "") -> str: @@ -853,6 +854,14 @@ class BingxVenueAdapter(VenueAdapter): 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) + # ── UNKNOWN OUTCOME ── see submit_async(): the order may be live, so the + # kernel must NOT roll back to flat on it. + if str(getattr(receipt, "status", "") or "").upper() == "INDETERMINATE": + raise VenueIndeterminateError( + f"submit outcome UNKNOWN for intent={intent.intent_id} asset={intent.asset} " + f"— order may be live at the venue; refusing to claim rejection", + receipt=receipt, + ) # ── POINT OF NO RETURN ── order is LIVE; see submit_async(). try: events = self._events_from_submit(submitted, receipt, None, None) @@ -920,6 +929,17 @@ class BingxVenueAdapter(VenueAdapter): # ── 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) + # ── UNKNOWN OUTCOME ── the submit timed out / 5xx'd and the adapter could + # not establish the truth even after asking the venue about our own + # clientOrderId. The order MAY be live. Escalate as indeterminate so the + # kernel does NOT roll the slot back to flat (VenueIndeterminateError is a + # VenuePostAckError, so the existing no-rollback fences catch it). + if str(getattr(receipt, "status", "") or "").upper() == "INDETERMINATE": + raise VenueIndeterminateError( + f"submit outcome UNKNOWN for intent={intent.intent_id} asset={intent.asset} " + f"— order may be live at the venue; refusing to claim rejection", + receipt=receipt, + ) # ── 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. diff --git a/prod/clean_arch/dita_v2/test_indeterminate_submit.py b/prod/clean_arch/dita_v2/test_indeterminate_submit.py new file mode 100644 index 0000000..b3ea4e9 --- /dev/null +++ b/prod/clean_arch/dita_v2/test_indeterminate_submit.py @@ -0,0 +1,282 @@ +"""INDETERMINATE submit outcomes: "not done" != "tried and failed" != "unknown". + +A BingX read timeout / reset / 5xx means the request WAS SENT and the answer was +lost — the order may be live. Before this class of fix, `bingx_direct` reported such +a failure to the kernel as a confident REJECTED with fill_qty=0, and the kernel rolled +the slot back to flat while the venue kept the position. That is an orphan, and it is +the same disease as the 2026-07-13 post-ack skew, one layer deeper. + +The rule these tests protect: + + Rollback is sound ONLY when the effect is provably absent. + NOT_ATTEMPTED (never sent) -> rollback OK + REFUSED (venue said no) -> rollback OK + INDETERMINATE (answer lost) -> NEVER roll back. Ask the venue about our own + clientOrderId; if truth cannot be established, + stay UNKNOWN and let the E-feed FILL settle it. + +The contrast tests matter as much as the recovery tests: a fix that refuses to roll +back on a *genuine* rejection would strand every slot forever. +""" + +from __future__ import annotations + +import asyncio +import itertools +import threading +from datetime import datetime, timezone +from types import SimpleNamespace + +import httpx +import pytest + +from prod.bingx.http import ( + BingxHttpError, + _effect_of_httpx_error, + _effect_of_status, +) +from prod.clean_arch.adapters.bingx_direct import BingxDirectExecutionAdapter +from prod.clean_arch.dita_v2.bingx_venue import BingxVenueAdapter +from prod.clean_arch.dita_v2.contracts import ( + KernelCommandType, + KernelEventKind, + KernelIntent, + TradeSide, + TradeStage, +) +from prod.clean_arch.dita_v2.rust_backend import ExecutionKernel +from prod.clean_arch.dita_v2.venue import VenueIndeterminateError, VenuePostAckError + + +# ── 1. Classification: what does the failure actually PROVE? ───────────────── + + +def test_read_timeout_is_indeterminate_not_rejection(): + """The killer case: the request was sent, the answer was lost.""" + assert _effect_of_httpx_error(httpx.ReadTimeout("timed out")) == BingxHttpError.INDETERMINATE + + +def test_connect_failures_are_not_attempted(): + """Never left the box -> provably no order -> rollback is sound.""" + assert _effect_of_httpx_error(httpx.ConnectError("refused")) == BingxHttpError.NOT_ATTEMPTED + assert _effect_of_httpx_error(httpx.ConnectTimeout("no route")) == BingxHttpError.NOT_ATTEMPTED + + +def test_5xx_is_indeterminate_and_4xx_is_refused(): + assert _effect_of_status(500) == BingxHttpError.INDETERMINATE + assert _effect_of_status(503) == BingxHttpError.INDETERMINATE + assert _effect_of_status(400) == BingxHttpError.REFUSED + assert _effect_of_status(429) == BingxHttpError.REFUSED + + +def test_unknown_failure_defaults_to_indeterminate(): + """When we do not know, we must NOT claim absence.""" + assert BingxHttpError("mystery").effect == BingxHttpError.INDETERMINATE + assert BingxHttpError("mystery").order_may_exist is True + assert BingxHttpError("nope", effect=BingxHttpError.REFUSED).order_may_exist is False + + +# ── 2. The point lookup: one order, by our own key, read-only ──────────────── + + +class _LookupClient: + """Records every call so we can prove the lookup NEVER re-POSTs an order.""" + + def __init__(self, get_results): + self._get_results = list(get_results) + self.gets: list[tuple[str, dict]] = [] + self.posts: list[tuple[str, dict]] = [] + + async def signed_get(self, path, params=None): + self.gets.append((path, dict(params or {}))) + result = self._get_results.pop(0) + if isinstance(result, Exception): + raise result + return result + + async def signed_post(self, path, params=None, **_kw): # pragma: no cover - must never run + self.posts.append((path, dict(params or {}))) + raise AssertionError("lookup must NEVER submit an order") + + +def _adapter(client) -> BingxDirectExecutionAdapter: + adapter = BingxDirectExecutionAdapter.__new__(BingxDirectExecutionAdapter) + adapter._client = client + return adapter + + +def _order_row(**over): + """Shape as signed_get actually delivers it: the envelope is already unwrapped + (BingxHttpClient._unwrap_response returns payload["data"]), so a query order + lookup yields {"order": {...}}.""" + row = {"orderId": "venue-1", "clientOrderId": "cid-1", "status": "FILLED", + "avgPrice": "100.0", "executedQty": "10.0"} + row.update(over) + return {"order": row} + + +def test_lookup_finds_the_order_and_never_posts(): + client = _LookupClient([_order_row()]) + row = asyncio.run(_adapter(client)._lookup_own_order_by_client_id("TRX-USDT", "cid-1")) + assert isinstance(row, dict) and row["orderId"] == "venue-1" + assert client.posts == [], "lookup must never re-submit — that is how you double-fill" + assert client.gets[0][1]["clientOrderID"] == "cid-1", "must query by OUR idempotency key" + + +def test_lookup_reports_absent_when_venue_has_no_such_order(): + """Venue answered authoritatively: no order. ONLY now is rollback sound.""" + client = _LookupClient([{"code": 0, "data": {}}]) + assert asyncio.run( + _adapter(client)._lookup_own_order_by_client_id("TRX-USDT", "cid-1") + ) == "ABSENT" + + +def test_lookup_treats_venue_refusal_as_absent(): + client = _LookupClient([BingxHttpError("order not exist", effect=BingxHttpError.REFUSED)]) + assert asyncio.run( + _adapter(client)._lookup_own_order_by_client_id("TRX-USDT", "cid-1") + ) == "ABSENT" + + +def test_lookup_returns_none_when_truth_cannot_be_established(): + """It is allowed to say "I don't know". It is never allowed to guess.""" + client = _LookupClient([httpx.ReadTimeout("t")] * 3) + result = asyncio.run( + _adapter(client)._lookup_own_order_by_client_id("TRX-USDT", "cid-1", attempts=3) + ) + assert result is None, "unresolved must be None — never 'ABSENT', never a fabricated row" + assert len(client.gets) == 3, "bounded: exactly `attempts` tries, then give up" + assert client.posts == [] + + +# ── 3. The venue layer escalates UNKNOWN instead of claiming rejection ─────── + + +def _intent(trade_id="ind-1") -> KernelIntent: + return KernelIntent( + timestamp=datetime.now(timezone.utc), + intent_id=f"intent-{trade_id}", trade_id=trade_id, slot_id=0, + asset="TRX-USDT", action=KernelCommandType.ENTER, side=TradeSide.SHORT, + reason="indeterminate-regression", target_size=10.0, leverage=1.0, + reference_price=100.0, exit_leg_ratios=(1.0,), metadata={}, + ) + + +def _receipt(status): + return SimpleNamespace( + status=status, order_id="", client_order_id="cid-1", price=0.0, quantity=0.0, + timestamp=datetime.now(timezone.utc), + raw_ack={"status": status, "clientOrderId": "cid-1"}, + ) + + +class _Backend: + def __init__(self, status, is_async=False): + self._status = status + self._is_async = is_async + self.submit_count = 0 + + def submit_intent(self, _legacy): + self.submit_count += 1 + return _receipt(self._status) + + async def submit_intent_async(self, _legacy): # pragma: no cover - shim + return self.submit_intent(_legacy) + + +class _AsyncBackend(_Backend): + async def submit_intent(self, _legacy): # type: ignore[override] + self.submit_count += 1 + return _receipt(self._status) + + +def _venue(backend) -> BingxVenueAdapter: + venue = BingxVenueAdapter.__new__(BingxVenueAdapter) + venue.backend = backend + venue._event_seq = itertools.count(1) + venue._snap_lock = threading.Lock() + venue._snapshot_ready = threading.Event() + venue._snapshot_ready.set() + venue._last_snapshot = None + venue._telemetry_plane = None + return venue + + +def test_sync_venue_escalates_indeterminate_receipt(monkeypatch): + venue = _venue(_Backend("INDETERMINATE")) + monkeypatch.setattr(venue, "_publish_telemetry", lambda **_f: None) + with pytest.raises(VenueIndeterminateError): + venue.submit(_intent()) + + +def test_async_venue_escalates_indeterminate_receipt(monkeypatch): + venue = _venue(_AsyncBackend("INDETERMINATE")) + monkeypatch.setattr(venue, "_publish_telemetry", lambda **_f: None) + with pytest.raises(VenueIndeterminateError): + asyncio.run(venue.submit_async(_intent())) + + +def test_indeterminate_is_a_post_ack_error_so_existing_fences_catch_it(): + """Subclassing is load-bearing: the kernel's no-rollback fences key on the base.""" + assert issubclass(VenueIndeterminateError, VenuePostAckError) + + +# ── 4. THE INVARIANT: the kernel must never conclude "flat" on UNKNOWN ─────── + + +def test_kernel_does_not_roll_back_slot_on_indeterminate_submit(monkeypatch): + """kernel_believes_flat => venue_is_flat. An UNKNOWN order may be live, so the + slot must NOT return to IDLE and NO synthetic REJECT may be emitted.""" + venue = _venue(_Backend("INDETERMINATE")) + monkeypatch.setattr(venue, "_publish_telemetry", lambda **_f: None) + + with ExecutionKernel(max_slots=1, venue=venue) as kernel: + outcome = kernel.process_intent(_intent(trade_id="ind-sync")) + slot = kernel._get_slot(0) + + kinds = [e.kind for e in outcome.emitted_events] + assert KernelEventKind.ORDER_REJECT not in kinds, ( + "a synthetic REJECT on an UNKNOWN order is the orphan-maker" + ) + assert slot.fsm_state is not TradeStage.IDLE, ( + "slot rolled back to IDLE while the order may be LIVE at the venue" + ) + + +def test_kernel_does_not_roll_back_slot_on_indeterminate_submit_async(monkeypatch): + venue = _venue(_AsyncBackend("INDETERMINATE")) + monkeypatch.setattr(venue, "_publish_telemetry", lambda **_f: None) + + async def exercise(): + with ExecutionKernel(max_slots=1, venue=venue) as kernel: + outcome = await kernel.process_intent_async(_intent(trade_id="ind-async")) + return outcome, kernel._get_slot(0) + + outcome, slot = asyncio.run(exercise()) + kinds = [e.kind for e in outcome.emitted_events] + assert KernelEventKind.ORDER_REJECT not in kinds + assert slot.fsm_state is not TradeStage.IDLE + + +# ── 5. CONTRAST: a genuine rejection MUST still roll back ──────────────────── +# Over-correcting here would strand every slot in ORDER_REQUESTED forever, which is +# its own outage. Proving the safe path still works is part of the fix. + + +def test_genuine_rejection_still_rolls_the_slot_back(monkeypatch): + class _RejectingBackend: + def submit_intent(self, _legacy): + raise BingxHttpError("insufficient margin", effect=BingxHttpError.REFUSED) + + venue = _venue(_RejectingBackend()) + monkeypatch.setattr(venue, "_publish_telemetry", lambda **_f: None) + + with ExecutionKernel(max_slots=1, venue=venue) as kernel: + outcome = kernel.process_intent(_intent(trade_id="refused")) + slot = kernel._get_slot(0) + + kinds = [e.kind for e in outcome.emitted_events] + assert KernelEventKind.ORDER_REJECT in kinds, ( + "a PROVEN refusal must still produce a REJECT — otherwise the slot strands" + ) + assert slot.fsm_state is TradeStage.IDLE, "provably-no-order must free the slot" diff --git a/prod/clean_arch/dita_v2/venue.py b/prod/clean_arch/dita_v2/venue.py index 15e16c9..7bc7ce6 100644 --- a/prod/clean_arch/dita_v2/venue.py +++ b/prod/clean_arch/dita_v2/venue.py @@ -40,6 +40,23 @@ class VenuePostAckError(Exception): self.events = events or [] +class VenueIndeterminateError(VenuePostAckError): + """The submit outcome is UNKNOWN: the request was sent, the answer was lost. + + A read timeout / connection reset / 5xx means BingX may have matched the order. + The adapter has already asked the venue about our clientOrderId and could not + establish the truth, so we are left with genuine uncertainty. + + Deliberately a subclass of VenuePostAckError, because it demands the SAME + response: the effect may exist, therefore DO NOT roll the slot back to flat and + DO NOT synthesise a REJECT. Leave the slot working and let the E-feed FILL / + the account stream settle it. Callers that already fence VenuePostAckError get + this behaviour for free. + + Unknown is not flat. It is not failed either. It is unknown. + """ + + class VenueAdapter(Protocol): """Abstract venue adapter used by the kernel."""