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:
Codex
2026-07-13 15:43:55 +02:00
parent d9b7e05531
commit bdc54fbeaa
5 changed files with 529 additions and 19 deletions

View File

@@ -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

View File

@@ -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

View File

@@ -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.

View 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"

View File

@@ -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."""