dita_v2(exec): INDETERMINATE submit is not REJECTED — unknown is never flat
A BingX read-timeout/reset/5xx after send means the answer was lost, not that the order failed. Classify every submit failure by what it PROVES: NOT_ATTEMPTED / REFUSED -> rollback sound; INDETERMINATE -> point-lookup our own clientOrderId (read-only, bounded, never a reconcile); unresolved stays UNKNOWN — no synthetic REJECT, no slot rollback, E-feed FILL settles truth. - prod/bingx/http.py: BingxHttpError.effect + order_may_exist, 9 raise sites tagged - adapters/bingx_direct.py: _lookup_own_order_by_client_id (never POSTs) - dita_v2/venue.py: VenueIndeterminateError(VenuePostAckError) — existing fences catch it - dita_v2/bingx_venue.py: both submit paths escalate INDETERMINATE receipts - 14 tests incl. kernel no-rollback invariant + genuine-REFUSED contrast Suite: 3416 passed, 19 skipped, 3 xfailed. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -28,7 +28,53 @@ from .urls import get_rest_base_urls
|
|||||||
|
|
||||||
|
|
||||||
class BingxHttpError(RuntimeError):
|
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)
|
@dataclass(frozen=True)
|
||||||
@@ -413,7 +459,8 @@ class BingxHttpClient:
|
|||||||
payload = dict(params)
|
payload = dict(params)
|
||||||
if signed:
|
if signed:
|
||||||
if not self._api_key or not self._secret_key:
|
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 = build_signed_params(
|
||||||
payload,
|
payload,
|
||||||
self._secret_key,
|
self._secret_key,
|
||||||
@@ -453,13 +500,15 @@ class BingxHttpClient:
|
|||||||
rate_limited=response.status_code == 429,
|
rate_limited=response.status_code == 429,
|
||||||
retry_after_ms=self._rate_limits.snapshot().rest_reset_ms,
|
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:
|
if base_index < len(self._base_urls) - 1:
|
||||||
continue
|
continue
|
||||||
if attempt < max_retries:
|
if attempt < max_retries:
|
||||||
await asyncio.sleep(delay)
|
await asyncio.sleep(delay)
|
||||||
break
|
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()
|
self._circuit_breaker.record_success()
|
||||||
if hostname is not None:
|
if hostname is not None:
|
||||||
await self._maybe_refresh_dns_cache(hostname)
|
await self._maybe_refresh_dns_cache(hostname)
|
||||||
@@ -487,7 +536,8 @@ class BingxHttpClient:
|
|||||||
except httpx.HTTPError as exc:
|
except httpx.HTTPError as exc:
|
||||||
self._logger.warning("HTTP error [att=%d url=%d %s]: %s — %s",
|
self._logger.warning("HTTP error [att=%d url=%d %s]: %s — %s",
|
||||||
attempt, base_index, method, type(exc).__name__, exc)
|
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 self._is_dns_resolution_error(exc):
|
||||||
if hostname is not None:
|
if hostname is not None:
|
||||||
cached_ips = self._dns_cache.resolve(hostname)
|
cached_ips = self._dns_cache.resolve(hostname)
|
||||||
@@ -528,12 +578,14 @@ class BingxHttpClient:
|
|||||||
return self._unwrap_response(data)
|
return self._unwrap_response(data)
|
||||||
return data
|
return data
|
||||||
except Exception as fallback_exc:
|
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:
|
if base_index < len(self._base_urls) - 1:
|
||||||
continue
|
continue
|
||||||
if last_error is not None:
|
if last_error is not None:
|
||||||
raise last_error
|
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()
|
delay = self._circuit_breaker.record_failure()
|
||||||
if base_index < len(self._base_urls) - 1:
|
if base_index < len(self._base_urls) - 1:
|
||||||
continue
|
continue
|
||||||
@@ -584,12 +636,14 @@ class BingxHttpClient:
|
|||||||
return self._unwrap_response(data)
|
return self._unwrap_response(data)
|
||||||
return data
|
return data
|
||||||
except Exception as fallback_exc:
|
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:
|
if base_index < len(self._base_urls) - 1:
|
||||||
continue
|
continue
|
||||||
if last_error is not None:
|
if last_error is not None:
|
||||||
raise last_error
|
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()
|
delay = self._circuit_breaker.record_failure()
|
||||||
if base_index < len(self._base_urls) - 1:
|
if base_index < len(self._base_urls) - 1:
|
||||||
continue
|
continue
|
||||||
@@ -625,7 +679,8 @@ class BingxHttpClient:
|
|||||||
def _unwrap_response(payload: dict[str, Any]) -> Any:
|
def _unwrap_response(payload: dict[str, Any]) -> Any:
|
||||||
code = int(payload.get("code", -1))
|
code = int(payload.get("code", -1))
|
||||||
if code != 0:
|
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")
|
return payload.get("data")
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
|
|||||||
@@ -578,6 +578,89 @@ class BingxDirectExecutionAdapter(ExecutionPort):
|
|||||||
# (Bybit, OKX, Binance) will have the same REST/WS split.
|
# (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:
|
async def submit_intent(self, intent: Intent) -> ExecutionReceipt:
|
||||||
symbol = self._instrument_venue_symbol(intent.asset)
|
symbol = self._instrument_venue_symbol(intent.asset)
|
||||||
if intent.action == DecisionAction.EXIT:
|
if intent.action == DecisionAction.EXIT:
|
||||||
@@ -707,6 +790,59 @@ class BingxDirectExecutionAdapter(ExecutionPort):
|
|||||||
0.0,
|
0.0,
|
||||||
)
|
)
|
||||||
except BingxHttpError as exc:
|
except BingxHttpError as exc:
|
||||||
|
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"
|
status = "RATE_LIMITED" if _is_rate_limited_error(exc) else "REJECTED"
|
||||||
ack_row = {
|
ack_row = {
|
||||||
"status": status,
|
"status": status,
|
||||||
@@ -714,9 +850,6 @@ class BingxDirectExecutionAdapter(ExecutionPort):
|
|||||||
"symbol": symbol,
|
"symbol": symbol,
|
||||||
"clientOrderId": client_order_id,
|
"clientOrderId": client_order_id,
|
||||||
}
|
}
|
||||||
fill_price = 0.0
|
|
||||||
ack = None
|
|
||||||
is_limit = False
|
|
||||||
|
|
||||||
# ── Gap 2: fee estimation (ESTIMATED_TAKER / ESTIMATED_MAKER) ────────
|
# ── Gap 2: fee estimation (ESTIMATED_TAKER / ESTIMATED_MAKER) ────────
|
||||||
# BingX REST ACK does not include commission. WS FILL_SETTLED will deliver
|
# 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
|
# 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
|
# target_size fallback applies only to accepted MARKET acks whose fill
|
||||||
# arrives later via WS (BingX REST ack omits executedQty on those).
|
# 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
|
fill_qty = 0.0
|
||||||
else:
|
else:
|
||||||
fill_qty = float(ack_row.get("executedQty") or ack_row.get("filledQty") or
|
fill_qty = float(ack_row.get("executedQty") or ack_row.get("filledQty") or
|
||||||
|
|||||||
@@ -45,6 +45,7 @@ from .utils import json_safe
|
|||||||
from .utils import safe_float
|
from .utils import safe_float
|
||||||
from .venue import VenueAdapter
|
from .venue import VenueAdapter
|
||||||
from .venue import VenuePostAckError
|
from .venue import VenuePostAckError
|
||||||
|
from .venue import VenueIndeterminateError
|
||||||
|
|
||||||
|
|
||||||
def _row_text(row: dict[str, Any], *keys: str, default: str = "") -> str:
|
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))
|
submitted = replace(intent, target_size=float(legacy.target_size))
|
||||||
# ── PRE-ACK ── safe to roll back if this raises.
|
# ── PRE-ACK ── safe to roll back if this raises.
|
||||||
receipt = self._call_backend("submit_intent", legacy)
|
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().
|
# ── POINT OF NO RETURN ── order is LIVE; see submit_async().
|
||||||
try:
|
try:
|
||||||
events = self._events_from_submit(submitted, receipt, None, None)
|
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.
|
# ── PRE-ACK ── an exception here means the venue never took the order.
|
||||||
# Safe for the caller to roll the FSM back.
|
# Safe for the caller to roll the FSM back.
|
||||||
receipt = await self.backend.submit_intent(legacy)
|
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.
|
# ── 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
|
# 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.
|
# means "no order exists", and the caller rolls back to flat on it.
|
||||||
|
|||||||
282
prod/clean_arch/dita_v2/test_indeterminate_submit.py
Normal file
282
prod/clean_arch/dita_v2/test_indeterminate_submit.py
Normal file
@@ -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"
|
||||||
@@ -40,6 +40,23 @@ class VenuePostAckError(Exception):
|
|||||||
self.events = events or []
|
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):
|
class VenueAdapter(Protocol):
|
||||||
"""Abstract venue adapter used by the kernel."""
|
"""Abstract venue adapter used by the kernel."""
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user