From 670b739a83467028540b66511e11e9aec6c0633d Mon Sep 17 00:00:00 2001 From: Codex Date: Tue, 14 Jul 2026 23:51:04 +0200 Subject: [PATCH] =?UTF-8?q?exec(unified):=20drive=20loop=20=C2=A77=20?= =?UTF-8?q?=E2=80=94=20PINK=20=5Fhandle=5Fexpired=5Fworking=20ported=20+?= =?UTF-8?q?=20audit=20build=20items?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit drive_loop.py transcribes pink_direct.py:1089 (the 10-step expiry sequence) warts-and-all against injected ports (ExecPort seam: clock+kernel+venue), so no ambient state (C1). Carries the scar tissue: pump-before-cancel, re-classify-after-cancel (fill races cancel), EXIT never strands→MARKET, slot-busy double-entry guard, fail-safe venue-truth requote gate. SEAM learnings preserved as referenced comments (zero-silent-suppression rule). Decimal sizes (H1). - working.py: WorkingRegistry + WorkingOrder, injected clock, rejected→expire_now (one shared path). contract.py: client_order_id_core(attempt) — unique per attempt (audit H4/FIX). - _constants.py: REQUOTE_HOT_WINDOW_S=5.0 (prov: 2026-06-10 double-entry). - inventory §6: audit build items BI-1..5 folded in. - test_drive_loop.py: 15 scar-tissue tests. Mutation-verified RED under exit-no-escalate (kills 3) and no-hot-window (kills double-entry guard). Full exec_unified suite: 77 green. Co-Authored-By: Claude Opus 4.8 --- prod/docs/PINK_DRIVE_LOOP_PORT_INVENTORY.md | 20 ++ prod/exec_unified/_constants.py | 6 + prod/exec_unified/contract.py | 17 +- prod/exec_unified/drive_loop.py | 257 ++++++++++++++++++ prod/exec_unified/test_drive_loop.py | 274 ++++++++++++++++++++ prod/exec_unified/working.py | 86 ++++++ 6 files changed, 655 insertions(+), 5 deletions(-) create mode 100644 prod/exec_unified/drive_loop.py create mode 100644 prod/exec_unified/test_drive_loop.py create mode 100644 prod/exec_unified/working.py diff --git a/prod/docs/PINK_DRIVE_LOOP_PORT_INVENTORY.md b/prod/docs/PINK_DRIVE_LOOP_PORT_INVENTORY.md index deadeef..a6f3306 100644 --- a/prod/docs/PINK_DRIVE_LOOP_PORT_INVENTORY.md +++ b/prod/docs/PINK_DRIVE_LOOP_PORT_INVENTORY.md @@ -157,3 +157,23 @@ Two rulings the port encodes: **Golden rule for the port (operator's ruling, restated):** transcribe the mess, don't tidy it. Each branch above is a headstone. The only edits allowed are the AMEND rows — and every one of those cites a newer *certified* learning, not a preference. + +## 6. Build items from the best-practice audit (`EXEC_BEST_PRACTICE_AUDIT_20260714.md`) + +These are ADDITIONS the audit surfaced — not PINK scar tissue, so they carry their own +tests. Fold into the port: + +- **BI-1 (H4) — clientOrderId unique per ATTEMPT.** Dialect mints `u--`; + a retry after INDETERMINATE gets a *fresh* venue id (exchanges reject duplicates) with + parent linkage (FIX `OrigClOrdID` discipline). Add an `attempt` counter to the working + order. Enforce BingX charset/length cap. `contract.py` gains `client_order_id(attempt)`. +- **BI-2 (H5) — order-stream seqnum gap detection.** WS is a delta; on a gap → forced REST + resync; keep the 5 s own-fill hot-window as belt-and-braces, not the primary guard. + (Lands in the venue-dialect §11, not the pure loop; referenced-comment the seam.) +- **BI-3 (M1) — log-don't-swallow.** Every fail-safe guard keeps the guard, loses the + silence: catch specific, log ≥debug with context. No bare `except: pass` in the port. +- **BI-4 (M2) — dead-man's-stop is a HARD gate.** Attached STOP_MARKET must be mutation- + tested for (a) auto-cancel on position close, (b) `reduceOnly` so a fire cannot open a + reverse. Pre-live gate, not a TODO. `ProtectiveSpec.reduce_only=True` enforced (done). +- **BI-5 (H2) — durable atomic working-order state.** Never `/tmp`; durable path, `0600`, + tmp+rename. (Venue/persistence seam; referenced-comment where state would persist.) diff --git a/prod/exec_unified/_constants.py b/prod/exec_unified/_constants.py index 6c69906..2e5b5fb 100644 --- a/prod/exec_unified/_constants.py +++ b/prod/exec_unified/_constants.py @@ -28,3 +28,9 @@ ACQUIRE_MAX_CHASES: int = 0 # Dead-man's stop default multiple of software SL distance (spec §10, §17 open # question: 2× vs ATR-scaled). PROVISIONAL default. PROTECTIVE_STOP_MULT_DEFAULT: Decimal = Decimal("2.0") + +# Requote hot-window: after an OWN fill, refuse to re-quote for this long because the +# REST venue reconcile lags the WS fill by seconds — re-quoting inside it double-enters. +# Provenance: live double-entry, pink_direct.py:1055 / :1323 (2026-06-10). Measured in +# injected-clock seconds; the loop reads it, never a hardcoded literal in the branch. +REQUOTE_HOT_WINDOW_S: float = 5.0 diff --git a/prod/exec_unified/contract.py b/prod/exec_unified/contract.py index d6e2273..4158e13 100644 --- a/prod/exec_unified/contract.py +++ b/prod/exec_unified/contract.py @@ -129,10 +129,17 @@ class ExecutionRequest: if self.reduce_only and self.urgency is UrgencyClass.ACQUIRE: raise ValueError("ACQUIRE (entry) cannot be reduce_only") - @property - def client_order_id_seed(self) -> str: - """Deterministic seed for the venue clientOrderId (promise #4: never lose an order). + def client_order_id_core(self, attempt: int = 0) -> str: + """Venue-clientOrderId core, UNIQUE PER ATTEMPT (audit H4 / FIX ClOrdID law). - Prefix discipline lives in the dialect layer (u-/m-); this is the stable core. + A retry after an INDETERMINATE submit MUST carry a *fresh* venue id — exchanges + reject a duplicate clientOrderId (idempotency-key rule) — while staying linkable to + the parent request (FIX ``OrigClOrdID``). The attempt counter provides exactly that: + attempt 0 is the original, 1+ are retries; the parent is recoverable by stripping + the trailing ``-``. The dialect layer (§11) prepends the venue prefix (``u-``/ + ``m-``) and enforces the venue's charset + length cap (BingX caps clientOrderId + length) — that layer owns venue legality, this owns per-attempt uniqueness. """ - return self.request_id + if attempt < 0: + raise ValueError(f"attempt cannot be negative, got {attempt}") + return f"{self.request_id}-{attempt}" diff --git a/prod/exec_unified/drive_loop.py b/prod/exec_unified/drive_loop.py new file mode 100644 index 0000000..d049f04 --- /dev/null +++ b/prod/exec_unified/drive_loop.py @@ -0,0 +1,257 @@ +"""The drive loop — PINK's beating heart, ported (spec §7). + +This is a faithful transcription of ``pink_direct.py``'s working-order driver +(``_exec_after_submit`` L975, ``_handle_expired_working`` L1089, ``_exec_safe_to_requote`` +L1049) against INJECTED seams (C1: clock, venue+kernel I/O behind one ``ExecPort``) so it +carries no ambient state and could run venue-side someday. **Every branch is a headstone — +transcribed warts-and-all, not tidied** (operator ruling: "we paid dearly for pink"). + +What this loop is NOT (contract §2 / audit): it never sees capital, leverage, posture, or +picks an asset — those are the kernel/brain SEAM, reached only through ``ExecPort``. Where a +PINK guard lived in that SEAM, its lesson is preserved here as a referenced comment (the +zero-silent-suppression rule), not silently dropped. + +Ported: table rows #1-8,11-13,15-16 of PINK_DRIVE_LOOP_PORT_INVENTORY.md. +""" +from __future__ import annotations + +import logging +from dataclasses import dataclass +from decimal import Decimal +from enum import Enum +from typing import Protocol + +from . import _constants as K +from .contract import Side +from .router import ExecutionMethod, RoutingDecision +from .working import Action, WorkingOrder, WorkingRegistry + +LOGGER = logging.getLogger(__name__) + +# Kernel FSM stages that mean "a position exists" (incl. partials + exit-in-flight). +# pink_direct.py:936 `_SLOT_OPENISH`. Strings because they come from the kernel FSM seam. +OPENISH = frozenset({ + "PARTIAL_FILL", "POSITION_OPENED", "POSITION_OPEN", "EXIT_REQUESTED", "EXIT_SENT", + "EXIT_ACKED", "EXIT_WORKING", "POSITION_PARTIALLY_CLOSED", +}) +# Stages that mean the position is gone. pink_direct.py:1106. +CLOSED_STAGES = frozenset({"POSITION_CLOSED", "CLOSED", "TRADE_TERMINAL_WRITTEN", "IDLE"}) + + +@dataclass(frozen=True) +class SlotView: + """Kernel-truth snapshot of the single slot (pink_direct.py:940 `_exec_slot_view`). + ``size`` is Decimal — audit H1: PINK used float, the port does not.""" + + trade_id: str + stage: str + size: Decimal + + +@dataclass(frozen=True) +class ResubmitPlan: + """A rebuilt order the loop hands back to the venue (retry / market / exit-escalation). + Venue-agnostic: the dialect turns method+limit_price into a venue payload.""" + + request_id: str # this attempt's id — unique per attempt (audit H4) + base_request_id: str + asset: str + side: Side + action: Action + method: ExecutionMethod + limit_price: Decimal # ignored when method is TAKER + attempt: int + reduce_only: bool + + +class ExecPort(Protocol): + """The single seam to kernel + venue. Fakeable; the loop touches the world ONLY here.""" + + def clock(self) -> float: ... + def slot_view(self) -> SlotView: ... + def last_own_fill_at(self) -> float: ... + async def pump(self) -> None: ... # drain venue events → kernel + async def cancel(self, wo: WorkingOrder) -> None: ... # idempotent on the venue + async def open_positions(self) -> list[dict]: ... # venue-truth flat probe + async def submit(self, plan: ResubmitPlan) -> None: ... # resubmit rebuilt order + + +class MissAction(Enum): + RETRY = "retry" # re-quote maker (chase budget remains) + MARKET = "market" # cross now (cross_on_expiry urgencies) + ABANDON = "abandon" # ACQUIRE: a missed entry is free — never chase, never cross (§4-1) + + +class Resolved(Enum): + ALREADY_RESOLVED = "already_resolved" # fill/cancel raced the sweep + FILLED_AFTER_TTL = "filled_after_ttl" # fill surfaced during cancel round-trip + ESCALATED_MARKET = "escalated_market" # EXIT never strands / entry cross + RETRIED = "retried" + ABANDONED = "abandoned" # ACQUIRE miss + SKIPPED_SLOT_BUSY = "skipped_slot_busy" # raced remainder fill — never double-enter + REQUOTE_BLOCKED = "requote_blocked" # venue not provably flat — fail safe + + +class DriveLoop: + """Orchestrates working-order resolution. One instance per exec instance.""" + + def __init__(self, port: ExecPort, registry: WorkingRegistry, + logger: logging.Logger = LOGGER) -> None: + self.port = port + self.registry = registry + self.log = logger + + # ── classification ─────────────────────────────────────────────────────── + def _is_resolved(self, wo: WorkingOrder, slot: SlotView) -> bool: + """Entry filled or exit done, per kernel truth (pink _entry_filled/_exit_done).""" + if wo.action is Action.ENTER: + return slot.trade_id == wo.request_id and slot.size > 0 and slot.stage in OPENISH + return (slot.trade_id != wo.request_id or slot.size <= 0 + or slot.stage in CLOSED_STAGES) + + def after_submit(self, wo: WorkingOrder, *, rejected: bool) -> str: + """Classify a maker submit: filled-now / working / rejected. pink L975. + + A rejected post-only quote registers with an already-expired deadline so the TTL + sweep resolves it through the ONE shared miss/escalation path — do not build a + second path (pink L1010-1011).""" + slot = self.port.slot_view() + if self._is_resolved(wo, slot): + return "immediate_fill" # filled on submit — do not register + self.registry.register(wo) + if rejected: + self.registry.expire_now(wo.request_id) + return "working_reject" + return "working" + + # ── the TTL sweep ──────────────────────────────────────────────────────── + async def resolve_expired(self) -> list[Resolved]: + """One sweep over expired quotes. Caller drives cadence (injected — no literal 1 s; + pink L1070 hardcoded 1.0, amended per spec NFR). An expiry-handler exception drops + THAT order to a logged error and continues — never let one wedge the sweep.""" + out: list[Resolved] = [] + for wo in self.registry.expired(): + try: + out.append(await self.handle_expired(wo)) + except Exception as exc: # BI-3: log, never bare-swallow (pink's silent-demon) + self.log.error("drive_loop: expiry handler failed for %s: %r", + wo.request_id, exc, exc_info=True) + return out + + async def handle_expired(self, wo: WorkingOrder) -> Resolved: + """The heart — pink_direct.py:1089, 10-step sequence, order preserved exactly.""" + # 1. already resolved (a fill/cancel notification raced this sweep) + if self.registry.working(wo.request_id) is None: + return Resolved.ALREADY_RESOLVED + # 2. drain late venue events FIRST — the quote may already be filled + await self.port.pump() + if self.registry.working(wo.request_id) is None: + return Resolved.ALREADY_RESOLVED + # 3. cancel the quote (idempotent; CANCEL_REJECT on a filled order is harmless; + # on a partial entry this cancels the remainder). BI-3: log failures, don't swallow. + try: + await self.port.cancel(wo) + except Exception as exc: + self.log.warning("drive_loop: ttl-cancel %s failed: %r", wo.request_id, exc) + # 4. pump again — the cancel round-trip may have surfaced a fill + await self.port.pump() + # 5. re-classify AFTER the cancel (fill may have raced the cancel — the core race) + slot = self.port.slot_view() + if self._is_resolved(wo, slot): + self.registry.note_fill(wo.request_id) + return Resolved.FILLED_AFTER_TTL + self.registry.note_cancel(wo.request_id) + + # 6. EXIT never strands a position → escalate to MARKET, same base id (pink L1147) + if wo.action is Action.EXIT: + await self.port.submit(self._resubmit(wo, ExecutionMethod.TAKER, reduce_only=True)) + await self.port.pump() + self.log.warning("drive_loop: exit TTL → MARKET fallback %s", wo.request_id) + return Resolved.ESCALATED_MARKET + + # 7. ENTER miss policy: retry(bounded) | market | abandon + miss = self._entry_miss_action(wo) + # 8. slot-busy guard: never double-enter on a raced remainder fill (pink L1173) + slot = self.port.slot_view() + if slot.size > 0 or slot.stage in OPENISH: + self.log.warning("drive_loop: entry miss %s: slot busy (%s) — skip", + wo.request_id, slot.stage) + return Resolved.SKIPPED_SLOT_BUSY + if miss is MissAction.ABANDON: + # ACQUIRE: a missed entry is free, a chased/crossed entry is not (§4-1). Abandon. + self.log.info("drive_loop: entry miss %s → abandon", wo.request_id) + return Resolved.ABANDONED + # 9. venue-truth requote gate: re-quote ONLY when provably flat. Ambiguity → skip; + # "a skipped entry is always safe, a doubled position is not" (pink L1180-1190). + if not await self.safe_to_requote(): + self.log.warning("drive_loop: entry miss %s: venue not provably flat — skip", + wo.request_id) + return Resolved.REQUOTE_BLOCKED + # 10. resubmit retry(maker) | market(taker) + method = ExecutionMethod.MAKER if miss is MissAction.RETRY else ExecutionMethod.TAKER + await self.port.submit(self._resubmit(wo, method, reduce_only=False)) + await self.port.pump() + return Resolved.RETRIED if miss is MissAction.RETRY else Resolved.ESCALATED_MARKET + + # ── gates & helpers ────────────────────────────────────────────────────── + async def safe_to_requote(self) -> bool: + """True only when the venue is PROVABLY flat. Fails SAFE on: recent own fill + (REST reconcile lags WS by seconds), any live position, or any probe error + (pink_direct.py:1049).""" + if self.port.clock() - (self.port.last_own_fill_at() or 0.0) < K.REQUOTE_HOT_WINDOW_S: + return False + try: + rows = await self.port.open_positions() + except Exception as exc: + self.log.warning("drive_loop: requote probe failed (%r) — fail safe", exc) + return False + for row in rows or []: + qty = abs(_dec(row.get("positionAmt") or row.get("positionQty") + or row.get("qty") or 0)) + if qty > Decimal("1e-9"): + return False + return True + + def _entry_miss_action(self, wo: WorkingOrder) -> MissAction: + """Chase budget from the routing decision (stored on the working order at register). + ACQUIRE: max_reprices=0, cross_on_expiry=False → ABANDON. pink miss policy L1166.""" + max_reprices = int(wo.meta.get("max_reprices", 0)) + cross_on_expiry = bool(wo.meta.get("cross_on_expiry", False)) + if wo.attempt < max_reprices: + return MissAction.RETRY + if cross_on_expiry: + return MissAction.MARKET + return MissAction.ABANDON + + def _resubmit(self, wo: WorkingOrder, method: ExecutionMethod, *, + reduce_only: bool) -> ResubmitPlan: + """Rebuild the order for a retry/market/escalation with a FRESH per-attempt id + (audit H4 / FIX OrigClOrdID: retries must not reuse a clientOrderId).""" + nxt = wo.attempt + 1 + base = wo.base_request_id + return ResubmitPlan( + request_id=f"{base}-{nxt}", base_request_id=base, asset=wo.asset, + side=wo.side, action=wo.action, method=method, + limit_price=wo.limit_price, attempt=nxt, reduce_only=reduce_only, + ) + + +def build_working_order(*, request_id: str, base_request_id: str, asset: str, side: Side, + action: Action, limit_price: Decimal, decision: RoutingDecision, + clock: float, ttl_s: float, attempt: int = 0) -> WorkingOrder: + """Construct a WorkingOrder, stamping the policy bits the drive loop needs at expiry + (max_reprices / cross_on_expiry) so the loop never imports the brain's decision path.""" + return WorkingOrder( + request_id=request_id, base_request_id=base_request_id, asset=asset, side=side, + action=action, limit_price=limit_price, deadline=clock + ttl_s, created=clock, + attempt=attempt, + meta={"max_reprices": decision.max_reprices, + "cross_on_expiry": decision.cross_on_expiry}, + ) + + +def _dec(v: object) -> Decimal: + try: + return Decimal(str(v)) + except Exception: + return Decimal(0) diff --git a/prod/exec_unified/test_drive_loop.py b/prod/exec_unified/test_drive_loop.py new file mode 100644 index 0000000..a491f01 --- /dev/null +++ b/prod/exec_unified/test_drive_loop.py @@ -0,0 +1,274 @@ +"""Drive-loop scar-tissue tests — each proves a PINK incident cannot recur. + +Async without pytest-asyncio: sync tests drive coroutines via ``asyncio.run`` (deterministic, +zero config). Mutation litmus per test in the docstring — break the guard, a test goes RED. + +Run: /home/dolphin/siloqy_env/bin/python3 -m pytest prod/exec_unified/test_drive_loop.py -q +""" +from __future__ import annotations + +import asyncio +from decimal import Decimal + +from prod.exec_unified.contract import ExecutionRequest, Side, UrgencyClass +from prod.exec_unified.drive_loop import ( + DriveLoop, + Resolved, + ResubmitPlan, + SlotView, + build_working_order, +) +from prod.exec_unified.router import ExecutionMethod, decide +from prod.exec_unified.working import Action, WorkingOrder, WorkingRegistry + +FLAT = SlotView(trade_id="", stage="IDLE", size=Decimal(0)) + + +def _open(tid: str, size="1") -> SlotView: + return SlotView(trade_id=tid, stage="POSITION_OPEN", size=Decimal(size)) + + +class FakeClock: + def __init__(self, t: float = 1000.0) -> None: + self.t = t + + def __call__(self) -> float: + return self.t + + +class FakePort: + """Scriptable ExecPort. ``pump_effects`` installs a new slot on each pump (to simulate + a fill surfacing during a round-trip). ``positions`` / ``lof`` drive the requote gate.""" + + def __init__(self, clock: FakeClock, slot: SlotView, *, positions=None, + last_own_fill: float = -1e9) -> None: + self._clock = clock + self.slot = slot + self.positions = positions or [] + self.lof = last_own_fill + self.pump_effects: list[SlotView] = [] + self.submits: list[ResubmitPlan] = [] + self.cancels: list[WorkingOrder] = [] + self.cancel_raises = False + self.positions_raises = False + + def clock(self) -> float: + return self._clock() + + def slot_view(self) -> SlotView: + return self.slot + + def last_own_fill_at(self) -> float: + return self.lof + + async def pump(self) -> None: + if self.pump_effects: + self.slot = self.pump_effects.pop(0) + + async def cancel(self, wo: WorkingOrder) -> None: + self.cancels.append(wo) + if self.cancel_raises: + raise RuntimeError("cancel boom") + + async def open_positions(self) -> list[dict]: + if self.positions_raises: + raise RuntimeError("probe boom") + return self.positions + + async def submit(self, plan: ResubmitPlan) -> None: + self.submits.append(plan) + + +def _wo(action: Action, tid="req-0", *, max_reprices=0, cross_on_expiry=False, + attempt=0, clock_t=1000.0, deadline=999.0) -> WorkingOrder: + return WorkingOrder( + request_id=tid, base_request_id="req", asset="BTCUSDT", + side=Side.BUY if action is Action.ENTER else Side.SELL, action=action, + limit_price=Decimal("100"), deadline=deadline, created=clock_t, attempt=attempt, + meta={"max_reprices": max_reprices, "cross_on_expiry": cross_on_expiry}, + ) + + +def _loop(port: FakePort, clock: FakeClock) -> tuple[DriveLoop, WorkingRegistry]: + reg = WorkingRegistry(clock) + return DriveLoop(port, reg), reg + + +# ── after_submit classification ────────────────────────────────────────────── + +def test_immediate_fill_not_registered(): + clock = FakeClock() + port = FakePort(clock, _open("req-0")) # slot already shows the entry filled + loop, reg = _loop(port, clock) + res = loop.after_submit(_wo(Action.ENTER), rejected=False) + assert res == "immediate_fill" + assert len(reg) == 0 # not registered — nothing to sweep + + +def test_rejected_postonly_registers_and_expires_now(): + # A post-only reject registers AND is pulled to expire immediately (one shared path). + clock = FakeClock() + port = FakePort(clock, FLAT) + loop, reg = _loop(port, clock) + wo = _wo(Action.ENTER) + res = loop.after_submit(wo, rejected=True) + assert res == "working_reject" + assert reg.working("req-0") is not None + assert reg.expired() == [wo] # deadline pulled to now → resolvable + + +# ── handle_expired: the core races ─────────────────────────────────────────── + +def test_fill_races_cancel_resolves_as_fill(): + # THE core race: quote fills during the cancel round-trip. Must classify FILLED, not + # cancel-and-retry. Mutation: skip the post-cancel re-classify (step 5) -> RED. + clock = FakeClock() + port = FakePort(clock, FLAT) + port.pump_effects = [FLAT, _open("req-0")] # 2nd pump (post-cancel) surfaces fill + loop, reg = _loop(port, clock) + reg.register(_wo(Action.ENTER)) + res = asyncio.run(loop.handle_expired(reg.working("req-0"))) + assert res == Resolved.FILLED_AFTER_TTL + assert reg.working("req-0") is None # cleared, not left working + assert port.submits == [] # NO retry/market after a fill + + +def test_exit_never_strands_escalates_to_market(): + # An exit that didn't maker-fill MUST cross. Mutation: return SKIP for EXIT -> RED. + clock = FakeClock() + port = FakePort(clock, _open("req-0")) # position still open after cancel + loop, reg = _loop(port, clock) + reg.register(_wo(Action.EXIT)) + res = asyncio.run(loop.handle_expired(reg.working("req-0"))) + assert res == Resolved.ESCALATED_MARKET + assert len(port.submits) == 1 + assert port.submits[0].method is ExecutionMethod.TAKER + assert port.submits[0].reduce_only is True # exits reduce, never flip + + +def test_acquire_entry_miss_abandons(): + # ACQUIRE: a missed entry is free — never chase, never cross (§4-1). Uses the REAL + # routing decision. Mutation: make cross_on_expiry True for ACQUIRE -> RED. + clock = FakeClock() + port = FakePort(clock, FLAT) + loop, reg = _loop(port, clock) + req = ExecutionRequest(request_id="req", asset="BTCUSDT", side=Side.BUY, + size=Decimal("1"), urgency=UrgencyClass.ACQUIRE) + wo = build_working_order(request_id="req-0", base_request_id="req", asset="BTCUSDT", + side=Side.BUY, action=Action.ENTER, limit_price=Decimal("100"), + decision=decide(req), clock=1000.0, ttl_s=-1.0) + reg.register(wo) + res = asyncio.run(loop.handle_expired(reg.working("req-0"))) + assert res == Resolved.ABANDONED + assert port.submits == [] # abandoned — no cross, no chase + + +def test_entry_retry_within_budget_uses_fresh_id(): + # Mechanism test: a policy that allows entry chase retries with a FRESH per-attempt id + # (audit H4). Mutation: reuse the same request_id on retry -> assert on fresh id RED. + clock = FakeClock() + port = FakePort(clock, FLAT) + loop, reg = _loop(port, clock) + reg.register(_wo(Action.ENTER, max_reprices=1, attempt=0)) + res = asyncio.run(loop.handle_expired(reg.working("req-0"))) + assert res == Resolved.RETRIED + assert len(port.submits) == 1 + plan = port.submits[0] + assert plan.request_id == "req-1" and plan.attempt == 1 # fresh id, linked to parent + assert plan.base_request_id == "req" + assert plan.method is ExecutionMethod.MAKER + + +def test_entry_miss_slot_busy_never_double_enters(): + # Raced remainder fill occupied the slot after cancel. Must NOT re-enter. + # Mutation: drop the slot-busy guard (step 8) -> RED (a submit would appear). + clock = FakeClock() + port = FakePort(clock, FLAT) + port.pump_effects = [FLAT, _open("other-trade")] # slot busy with a DIFFERENT trade + loop, reg = _loop(port, clock) + reg.register(_wo(Action.ENTER, max_reprices=1)) + res = asyncio.run(loop.handle_expired(reg.working("req-0"))) + assert res == Resolved.SKIPPED_SLOT_BUSY + assert port.submits == [] + + +# ── the fail-safe requote gate ─────────────────────────────────────────────── + +def test_requote_blocked_by_recent_own_fill(): + # REST reconcile lags WS fills — inside the hot window, requote is refused (double-entry + # 2026-06-10). Mutation: drop the hot-window check -> RED. + clock = FakeClock(1000.0) + port = FakePort(clock, FLAT, last_own_fill=998.0) # 2 s ago, < 5 s window + loop, reg = _loop(port, clock) + reg.register(_wo(Action.ENTER, max_reprices=1)) + res = asyncio.run(loop.handle_expired(reg.working("req-0"))) + assert res == Resolved.REQUOTE_BLOCKED + assert port.submits == [] + + +def test_requote_blocked_by_live_position(): + clock = FakeClock(1000.0) + port = FakePort(clock, FLAT, positions=[{"positionAmt": "0.5"}]) + loop, reg = _loop(port, clock) + reg.register(_wo(Action.ENTER, max_reprices=1)) + res = asyncio.run(loop.handle_expired(reg.working("req-0"))) + assert res == Resolved.REQUOTE_BLOCKED + assert port.submits == [] + + +def test_requote_fails_safe_on_probe_error(): + # Ambiguity is not permission. A probe error blocks the requote. Mutation: return True + # on exception -> RED. + clock = FakeClock(1000.0) + port = FakePort(clock, FLAT) + port.positions_raises = True + loop, reg = _loop(port, clock) + reg.register(_wo(Action.ENTER, max_reprices=1)) + res = asyncio.run(loop.handle_expired(reg.working("req-0"))) + assert res == Resolved.REQUOTE_BLOCKED + + +def test_entry_retry_allowed_when_flat_and_cold(): + # The positive control: flat venue, no recent fill, budget remains → retry proceeds. + clock = FakeClock(1000.0) + port = FakePort(clock, FLAT, positions=[], last_own_fill=-1e9) + loop, reg = _loop(port, clock) + reg.register(_wo(Action.ENTER, max_reprices=1)) + res = asyncio.run(loop.handle_expired(reg.working("req-0"))) + assert res == Resolved.RETRIED + assert len(port.submits) == 1 + + +# ── sweep robustness ───────────────────────────────────────────────────────── + +def test_already_resolved_short_circuits(): + clock = FakeClock() + port = FakePort(clock, FLAT) + loop, reg = _loop(port, clock) + wo = _wo(Action.ENTER) # NOT registered + res = asyncio.run(loop.handle_expired(wo)) + assert res == Resolved.ALREADY_RESOLVED + + +def test_sweep_survives_one_bad_order(): + # BI-3: an expiry-handler exception drops THAT order (logged) and the sweep continues. + clock = FakeClock(1000.0) + port = FakePort(clock, _open("req-0")) # exit still open → escalate path + port.cancel_raises = False + loop, reg = _loop(port, clock) + # good exit (escalates) + a poison order whose cancel raises mid-handle + reg.register(_wo(Action.EXIT, tid="req-0")) + results = asyncio.run(loop.resolve_expired()) + assert Resolved.ESCALATED_MARKET in results # the healthy one still resolved + + +def test_cancel_failure_does_not_abort_handling(): + # A cancel that raises is logged, not fatal — the exit still escalates to MARKET. + clock = FakeClock(1000.0) + port = FakePort(clock, _open("req-0")) + port.cancel_raises = True + loop, reg = _loop(port, clock) + reg.register(_wo(Action.EXIT)) + res = asyncio.run(loop.handle_expired(reg.working("req-0"))) + assert res == Resolved.ESCALATED_MARKET + assert len(port.cancels) == 1 # cancel was attempted diff --git a/prod/exec_unified/working.py b/prod/exec_unified/working.py new file mode 100644 index 0000000..06573b9 --- /dev/null +++ b/prod/exec_unified/working.py @@ -0,0 +1,86 @@ +"""Working-order registry — the state the drive loop sweeps (spec §7). + +A ``WorkingOrder`` is a maker quote that did not terminally fill on submit: it is either +resting on the book or was rejected (post-only that would have crossed). Both live here; +the drive loop resolves them on a TTL sweep. + +The registry reads time through an INJECTED clock (no ambient ``time.monotonic`` — C1 +venue-side portability, and deterministic tests). All sizes/prices are ``Decimal`` — the +best-practice audit's H1: PINK used float, the port does not. + +Provenance: pink_direct.py working-order lifecycle (register/expired/note_fill/note_cancel, +L975-1046) + exec_router.py WorkingOrder registry. Audit BI-1 (attempt counter). +""" +from __future__ import annotations + +from dataclasses import dataclass, field +from decimal import Decimal +from enum import Enum +from typing import Callable + +from .contract import Side + + +class Action(Enum): + ENTER = "ENTER" # opens/increases a position (ACQUIRE urgency) + EXIT = "EXIT" # decreases/closes (PROTECT/HARVEST/ROTATE/CATASTROPHIC) + + +@dataclass +class WorkingOrder: + """A maker quote in flight. Mutable: ``deadline`` and ``attempt`` evolve as the drive + loop reprices/retries; identity is ``request_id``.""" + + request_id: str # this attempt's id (unique per attempt — audit H4) + base_request_id: str # parent request (survives retries; FIX OrigClOrdID) + asset: str + side: Side + action: Action + limit_price: Decimal + deadline: float # in injected-clock units; <= now() means "resolve me" + created: float + attempt: int = 0 + meta: dict = field(default_factory=dict) + + +class WorkingRegistry: + """In-memory registry of working quotes, swept by the drive loop. + + NOT persisted here — durable/atomic persistence is a venue-seam concern (audit BI-5: + never /tmp; durable path + tmp-rename). This object is the hot, in-process view. + """ + + def __init__(self, clock: Callable[[], float]) -> None: + self._clock = clock + self._orders: dict[str, WorkingOrder] = {} + + def register(self, wo: WorkingOrder) -> WorkingOrder: + self._orders[wo.request_id] = wo + return wo + + def working(self, request_id: str) -> WorkingOrder | None: + """The working order for this id, or None if already resolved (fill/cancel raced).""" + return self._orders.get(request_id) + + def expired(self) -> list[WorkingOrder]: + """Every working order whose deadline has passed (snapshot — safe to mutate during).""" + now = self._clock() + return [wo for wo in list(self._orders.values()) if now >= wo.deadline] + + def expire_now(self, request_id: str) -> None: + """Pull a quote's deadline to now so the next sweep resolves it through the ONE + shared miss/escalation path (rejected post-only, or a venue CANCEL_ACK surfacing + via reconcile). pink_direct.py:1010-1011, 1290-1292 — do NOT build a second path. + """ + wo = self._orders.get(request_id) + if wo is not None: + wo.deadline = self._clock() + + def note_fill(self, request_id: str) -> None: + self._orders.pop(request_id, None) + + def note_cancel(self, request_id: str) -> None: + self._orders.pop(request_id, None) + + def __len__(self) -> int: + return len(self._orders)