diff --git a/prod/exec_unified/__init__.py b/prod/exec_unified/__init__.py index 2fd7ad0..3c2bc02 100644 --- a/prod/exec_unified/__init__.py +++ b/prod/exec_unified/__init__.py @@ -13,19 +13,25 @@ from .contract import ( ExecutionAdvice, ExecutionRequest, ProtectiveSpec, + SGrade, Side, UrgencyClass, ) +from .executor import ExecuteResult, ExecuteStatus, UnifiedExecutor from .router import ExecutionMethod, RoutingDecision, TtlDiscipline, decide __all__ = [ + "ExecuteResult", + "ExecuteStatus", "ExecutionAdvice", "ExecutionMethod", "ExecutionRequest", "ProtectiveSpec", "RoutingDecision", + "SGrade", "Side", "TtlDiscipline", + "UnifiedExecutor", "UrgencyClass", "decide", ] diff --git a/prod/exec_unified/_constants.py b/prod/exec_unified/_constants.py index 2e5b5fb..95270bc 100644 --- a/prod/exec_unified/_constants.py +++ b/prod/exec_unified/_constants.py @@ -29,6 +29,12 @@ ACQUIRE_MAX_CHASES: int = 0 # question: 2× vs ATR-scaled). PROVISIONAL default. PROTECTIVE_STOP_MULT_DEFAULT: Decimal = Decimal("2.0") +# Maker-quote lifetime for UNBOUNDED-TTL urgencies (HARVEST/ACQUIRE): how long a resting +# quote lives before the drive loop sweeps + re-evaluates it. Bounded quote lifetime + +# cancel/replace is the certified GTX discipline (spec §4-16). PROVISIONAL — calibrate vs +# L8 GTX fill latencies; measured in injected-clock seconds, never a hardcoded loop literal. +MAKER_QUOTE_TTL_S: float = 6.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 diff --git a/prod/exec_unified/contract.py b/prod/exec_unified/contract.py index 4158e13..c99308d 100644 --- a/prod/exec_unified/contract.py +++ b/prod/exec_unified/contract.py @@ -47,6 +47,24 @@ class UrgencyClass(Enum): UrgencyClass.HARVEST, UrgencyClass.ROTATE) +class SGrade(Enum): + """Size axis — orthogonal to the T-tier smartness ladder (spec §15.2). + + Computed at intake from Appendix E's power-law depth model (per-asset, stress-adjusted), + or supplied by the caller. Below T14 the layer executes S0/S1 directly and REFUSES S2+ + whole (the pain-fence) so nobody meets the large-size pain class by accident. + """ + + S0 = 0 # size << top-of-book depth (walk < 1 tick) — any tier executes directly + S1 = 1 # walks visible levels (expected walk >= spread) — T2+ prices the walk + S2 = 2 # exceeds book capacity in the urgency window — must be SLICED; refused below T14 + S3 = 3 # size IS the market (footprint moves MM behaviour) — MALKHUT/VIBRASS territory + + @property + def needs_slicing(self) -> bool: + return self in (SGrade.S2, SGrade.S3) + + @dataclass(frozen=True) class ProtectiveSpec: """Attach-a-dead-man's-stop geometry (spec §10). Mark-based, rides the entry order. @@ -113,6 +131,8 @@ class ExecutionRequest: protective: ProtectiveSpec | None = None advice: ExecutionAdvice | None = None deadline_ms: int | None = None # caller's patience budget (ROTATE/ACQUIRE) + s_grade: SGrade = SGrade.S0 # size regime (§15.2); default S0 = fits at touch + parent_request_id: str | None = None # slicer family linkage (§15.2 T*.s), telemetry only def __post_init__(self) -> None: if not self.request_id: diff --git a/prod/exec_unified/drive_loop.py b/prod/exec_unified/drive_loop.py index d049f04..46bdb44 100644 --- a/prod/exec_unified/drive_loop.py +++ b/prod/exec_unified/drive_loop.py @@ -36,6 +36,9 @@ OPENISH = frozenset({ }) # Stages that mean the position is gone. pink_direct.py:1106. CLOSED_STAGES = frozenset({"POSITION_CLOSED", "CLOSED", "TRADE_TERMINAL_WRITTEN", "IDLE"}) +# Stages that mean the venue rejected our order (post-only would cross, etc.). +# pink_direct.py:1000. A rejected maker quote still registers → resolves via the TTL path. +REJECTED_STAGES = frozenset({"ORDER_REJECTED", "EXIT_REJECTED"}) @dataclass(frozen=True) @@ -64,6 +67,12 @@ class ResubmitPlan: reduce_only: bool +# The plan the venue receives is the same shape whether it is an initial submit (executor, +# attempt 0) or a rebuilt retry/escalation (drive loop). "Resubmit" is the origin name; +# OrderPlan is the role name the executor uses. +OrderPlan = ResubmitPlan + + class ExecPort(Protocol): """The single seam to kernel + venue. Fakeable; the loop touches the world ONLY here.""" @@ -103,11 +112,21 @@ class DriveLoop: # ── classification ─────────────────────────────────────────────────────── def _is_resolved(self, wo: WorkingOrder, slot: SlotView) -> bool: - """Entry filled or exit done, per kernel truth (pink _entry_filled/_exit_done).""" + """Entry filled or exit done, per kernel truth (pink _entry_filled/_exit_done, L1099). + + ENTER: our entry filled — the slot carries OUR clientOrderId (the kernel tags the new + position with it; clientOrderId echo is the best-practice fill key, audit H5) and shows + size + an open stage. + + EXIT: the position is gone — size drained or a closed stage. **Deliberately SIZE-based, + not trade_id-based**: pink_direct.py:1105 could compare `slot_tid != wo.trade_id` only + because it REUSED the position's trade_id for the exit intent. This layer is agnostic + (contract §2) — the caller mints a fresh request_id for the exit and never hands us the + position's id — so an exit is "done" iff the position closed, which is the size signal. + """ 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) + return 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. diff --git a/prod/exec_unified/executor.py b/prod/exec_unified/executor.py new file mode 100644 index 0000000..0f138ea --- /dev/null +++ b/prod/exec_unified/executor.py @@ -0,0 +1,146 @@ +"""UnifiedExecutor — the initial-submit pipeline (spec §2, §6, §15.2). + +The two halves of the engine: + * ``UnifiedExecutor.execute(request, snapshot)`` — a NEW ExecutionRequest enters, is + graded/refused, routed (whether-maker), placed (where-in-book), and submitted. This file. + * ``DriveLoop`` — the working-order lifecycle after submit (TTL sweep, requote, escalate). + +This is the analogue of ``pink_direct.py``'s in-``step`` execution block (L1645-1700): +router-decides-HOW → placer-sets-WHERE → submit → register working. Pure orchestration over +the injected ``ExecPort``; it never sees capital/leverage/asset (contract §2). + +Composition rulings encoded here (not re-derivable, so read them): + * Placer declines (spread gate fails) → **cross only if the urgency crosses on expiry** + (PROTECT/HARVEST/ROTATE). ACQUIRE **abandons** — a missed entry is free (§4-1). + * The S2+ pain-fence (§15.2): below T14 an S2/S3 request is REFUSED WHOLE, loudly — never + silently sliced or dumped. Nobody meets the large-size pain class by accident. +""" +from __future__ import annotations + +import logging +from dataclasses import dataclass +from decimal import Decimal +from enum import Enum + +from . import _constants as K +from .contract import ExecutionRequest, SGrade, Side, UrgencyClass +from .drive_loop import ( + REJECTED_STAGES, + Action, + DriveLoop, + ExecutionMethod, + OrderPlan, + build_working_order, +) +from .placer import MarketSnapshot, pre_submit +from .router import RoutingDecision, TtlDiscipline, decide +from .working import WorkingOrder + +LOGGER = logging.getLogger(__name__) + + +class ExecuteStatus(Enum): + REFUSED_S2 = "refused_s2" # pain-fence: S2+ below T14, refused whole + CROSSED_TAKER = "crossed_taker" # crossed the spread now (CATASTROPHIC, or maker-impossible cross) + RESTING_MAKER = "resting_maker" # post-only quote resting; drive loop owns it now + REJECTED_RESTING = "rejected_resting" # post-only rejected; registered for instant TTL resolve + FILLED_ON_SUBMIT = "filled_on_submit" # order filled during the submit round-trip + ABANDONED = "abandoned" # ACQUIRE, maker impossible, no cross — a free miss + + +@dataclass(frozen=True) +class ExecuteResult: + status: ExecuteStatus + decision: RoutingDecision | None + plan: OrderPlan | None + working: WorkingOrder | None + reason: str + + +def _ttl_seconds(request: ExecutionRequest, decision: RoutingDecision) -> float: + """Working-quote lifetime before the drive loop re-evaluates it, from the urgency's TTL + discipline (spec §6). BOUNDED→its ms; DEADLINE→caller's deadline_ms; UNBOUNDED→the + certified GTX quote lifetime (§4-16).""" + if decision.ttl is TtlDiscipline.BOUNDED_MS and decision.max_ms is not None: + return decision.max_ms / 1000.0 + if decision.ttl is TtlDiscipline.DEADLINE: + return (request.deadline_ms or int(K.MAKER_QUOTE_TTL_S * 1000)) / 1000.0 + return K.MAKER_QUOTE_TTL_S # UNBOUNDED (HARVEST/ACQUIRE): bounded by chases, not clock + + +class UnifiedExecutor: + """One per exec instance. Shares its ``DriveLoop``'s port + registry.""" + + def __init__(self, drive_loop: DriveLoop, logger: logging.Logger = LOGGER) -> None: + self.dl = drive_loop + self.port = drive_loop.port + self.registry = drive_loop.registry + self.log = logger + + async def execute(self, request: ExecutionRequest, + snapshot: MarketSnapshot) -> ExecuteResult: + # ── 1. S2+ pain-fence (§15.2): refuse large size WHOLE, loudly, below T14 ── + if request.s_grade.needs_slicing: + reason = (f"S-grade {request.s_grade.name} exceeds book capacity — REFUSED whole " + f"below T14. Route through a slicer (TWAP/VWAP) or take a tier decision. " + f"Silently dumping this size is the pain class the fence exists to prevent.") + self.log.error("executor: %s (%s size=%s)", reason, request.asset, request.size) + return ExecuteResult(ExecuteStatus.REFUSED_S2, None, None, None, reason) + + decision = decide(request) + action = Action.ENTER if request.urgency is UrgencyClass.ACQUIRE else Action.EXIT + base = request.request_id + rid0 = request.client_order_id_core(0) # unique per attempt (audit H4) + reduce_only = action is Action.EXIT + + # ── 2. TAKER: cross now (CATASTROPHIC). No cleverness, ever. ── + if decision.method is ExecutionMethod.TAKER: + return await self._cross(request, decision, action, rid0, base, + reduce_only, reason=decision.rationale) + + # ── 3. MAKER: ask the placer WHERE (spread gate + touch, quantized). ── + placement = pre_submit(request, decision, snapshot) + if placement is None: + # Maker impossible right now (spread too wide). Compose by urgency: + if decision.cross_on_expiry: + return await self._cross(request, decision, action, rid0, base, + reduce_only, reason="spread_gate: maker impossible → cross") + # ACQUIRE: a missed entry is free — abandon, never cross (§4-1). + self.log.info("executor: %s ACQUIRE maker impossible (spread) → abandon", base) + return ExecuteResult(ExecuteStatus.ABANDONED, decision, None, None, + "spread_gate: ACQUIRE abandons, never crosses") + + # ── 4. Submit the maker quote, then classify (filled / resting / rejected). ── + maker_plan = OrderPlan( + request_id=rid0, base_request_id=base, asset=request.asset, side=request.side, + action=action, method=ExecutionMethod.MAKER, limit_price=placement.limit_price, + attempt=0, reduce_only=reduce_only, + ) + await self.port.submit(maker_plan) + wo = build_working_order( + request_id=rid0, base_request_id=base, asset=request.asset, side=request.side, + action=action, limit_price=placement.limit_price, decision=decision, + clock=self.port.clock(), ttl_s=_ttl_seconds(request, decision), attempt=0, + ) + rejected = self.port.slot_view().stage in REJECTED_STAGES + cls = self.dl.after_submit(wo, rejected=rejected) + if cls == "immediate_fill": + return ExecuteResult(ExecuteStatus.FILLED_ON_SUBMIT, decision, maker_plan, None, + "maker filled on submit") + if cls == "working_reject": + return ExecuteResult(ExecuteStatus.REJECTED_RESTING, decision, maker_plan, wo, + "post-only rejected — instant TTL resolve") + return ExecuteResult(ExecuteStatus.RESTING_MAKER, decision, maker_plan, wo, + "maker quote resting; drive loop owns it") + + async def _cross(self, request: ExecutionRequest, decision: RoutingDecision, + action: Action, rid0: str, base: str, reduce_only: bool, + *, reason: str) -> ExecuteResult: + plan = OrderPlan( + request_id=rid0, base_request_id=base, asset=request.asset, side=request.side, + action=action, method=ExecutionMethod.TAKER, limit_price=Decimal(0), + attempt=0, reduce_only=reduce_only, + ) + await self.port.submit(plan) + await self.port.pump() # a market cross settles fast; drain it before returning + return ExecuteResult(ExecuteStatus.CROSSED_TAKER, decision, plan, None, reason) diff --git a/prod/exec_unified/test_executor.py b/prod/exec_unified/test_executor.py new file mode 100644 index 0000000..32faa99 --- /dev/null +++ b/prod/exec_unified/test_executor.py @@ -0,0 +1,176 @@ +"""UnifiedExecutor composition tests — the initial-submit pipeline (spec §2/§6/§15.2). + +Each asserts a composition ruling that is NOT re-derivable from the parts: + * S2+ pain-fence refuses whole (mutation: let S2 through -> RED) + * placer-declines → cross iff urgency crosses; ACQUIRE abandons (mutation: cross ACQUIRE -> RED) + * TTL derived from the urgency's discipline + +Async via asyncio.run (no pytest-asyncio needed). Run: + /home/dolphin/siloqy_env/bin/python3 -m pytest prod/exec_unified/test_executor.py -q +""" +from __future__ import annotations + +import asyncio +from decimal import Decimal + +from prod.exec_unified.contract import ExecutionRequest, SGrade, Side, UrgencyClass +from prod.exec_unified.drive_loop import Action, DriveLoop, SlotView +from prod.exec_unified.executor import ExecuteStatus, UnifiedExecutor +from prod.exec_unified.placer import MarketSnapshot +from prod.exec_unified.router import ExecutionMethod +from prod.exec_unified.working import WorkingRegistry + +NARROW = MarketSnapshot(best_bid=Decimal("100"), best_ask=Decimal("100.05"), + spread_bps=Decimal("5"), tick=Decimal("0.01"), step=Decimal("0.001")) +WIDE = MarketSnapshot(best_bid=Decimal("100"), best_ask=Decimal("100.5"), + spread_bps=Decimal("20"), tick=Decimal("0.01"), step=Decimal("0.001")) +FLAT = SlotView(trade_id="", stage="ACKED", size=Decimal(0)) +# An exit is submitted against an OPEN position — the slot still shows it (size > 0) until +# the maker close fills. A flat slot at exit-submit time is nonsensical. +OPEN = SlotView(trade_id="pos-1", stage="POSITION_OPEN", size=Decimal("1")) + + +class FakeClock: + def __init__(self, t=1000.0): self.t = t + def __call__(self): return self.t + + +class FakePort: + def __init__(self, clock, slot=FLAT): + self._c = clock + self.slot = slot + self.submits = [] + self.pumps = 0 + def clock(self): return self._c() + def slot_view(self): return self.slot + def last_own_fill_at(self): return -1e9 + async def pump(self): self.pumps += 1 + async def cancel(self, wo): pass + async def open_positions(self): return [] + async def submit(self, plan): self.submits.append(plan) + + +def _exec(slot=FLAT): + clock = FakeClock() + port = FakePort(clock, slot) + reg = WorkingRegistry(clock) + return UnifiedExecutor(DriveLoop(port, reg)), port, reg + + +def _req(urgency, side=Side.BUY, **kw): + base = dict(request_id="req", asset="BTCUSDT", side=side, size=Decimal("0.001"), + urgency=urgency) + base.update(kw) + return ExecutionRequest(**base) + + +def _run(coro): + return asyncio.run(coro) + + +# ── S2+ pain-fence ─────────────────────────────────────────────────────────── + +def test_s2_refused_whole_no_submit(): + # Mutation: drop the fence -> this order would submit and RED here. + ex, port, _ = _exec() + r = _run(ex.execute(_req(UrgencyClass.ACQUIRE, s_grade=SGrade.S2), NARROW)) + assert r.status is ExecuteStatus.REFUSED_S2 + assert port.submits == [] # nothing hit the venue + assert "REFUSED" in r.reason + + +def test_s3_also_refused(): + ex, port, _ = _exec() + r = _run(ex.execute(_req(UrgencyClass.HARVEST, side=Side.SELL, s_grade=SGrade.S3), NARROW)) + assert r.status is ExecuteStatus.REFUSED_S2 + assert port.submits == [] + + +def test_s0_default_not_refused(): + ex, port, _ = _exec() + r = _run(ex.execute(_req(UrgencyClass.ACQUIRE), NARROW)) # default S0 + assert r.status is not ExecuteStatus.REFUSED_S2 + + +# ── routing → submit ───────────────────────────────────────────────────────── + +def test_catastrophic_crosses_taker(): + ex, port, _ = _exec() + r = _run(ex.execute(_req(UrgencyClass.CATASTROPHIC, side=Side.SELL), NARROW)) + assert r.status is ExecuteStatus.CROSSED_TAKER + assert len(port.submits) == 1 + assert port.submits[0].method is ExecutionMethod.TAKER + assert port.submits[0].reduce_only is True # exit reduces + assert port.submits[0].request_id == "req-0" # per-attempt id (H4) + + +def test_acquire_narrow_spread_rests_maker(): + ex, port, reg = _exec() + r = _run(ex.execute(_req(UrgencyClass.ACQUIRE), NARROW)) + assert r.status is ExecuteStatus.RESTING_MAKER + assert port.submits[0].method is ExecutionMethod.MAKER + assert port.submits[0].limit_price == Decimal("100") # placed at bid (touch), quantized + assert reg.working("req-0") is not None # drive loop now owns it + + +def test_acquire_wide_spread_abandons_never_crosses(): + # THE key composition ruling: a missed entry is free (§4-1). Mutation: cross here -> RED. + ex, port, _ = _exec() + r = _run(ex.execute(_req(UrgencyClass.ACQUIRE), WIDE)) + assert r.status is ExecuteStatus.ABANDONED + assert port.submits == [] # NOT crossed + + +def test_protect_wide_spread_crosses(): + # Contrast with ACQUIRE: PROTECT crosses when maker is impossible (cross_on_expiry). + ex, port, _ = _exec() + r = _run(ex.execute(_req(UrgencyClass.PROTECT, side=Side.SELL), WIDE)) + assert r.status is ExecuteStatus.CROSSED_TAKER + assert port.submits[0].method is ExecutionMethod.TAKER + + +def test_harvest_narrow_rests_maker(): + ex, port, reg = _exec(slot=OPEN) + r = _run(ex.execute(_req(UrgencyClass.HARVEST, side=Side.SELL), NARROW)) + assert r.status is ExecuteStatus.RESTING_MAKER + assert port.submits[0].limit_price == Decimal("100.05") # SELL at ask (touch) + + +def test_maker_filled_on_submit(): + ex, port, reg = _exec(slot=SlotView(trade_id="req-0", stage="POSITION_OPEN", size=Decimal("1"))) + r = _run(ex.execute(_req(UrgencyClass.ACQUIRE), NARROW)) + assert r.status is ExecuteStatus.FILLED_ON_SUBMIT + assert reg.working("req-0") is None # filled → not left working + + +def test_post_only_rejected_registers_for_instant_resolve(): + ex, port, reg = _exec(slot=SlotView(trade_id="", stage="ORDER_REJECTED", size=Decimal(0))) + r = _run(ex.execute(_req(UrgencyClass.ACQUIRE), NARROW)) + assert r.status is ExecuteStatus.REJECTED_RESTING + wo = reg.working("req-0") + assert wo is not None + assert reg.expired() == [wo] # deadline pulled to now + + +# ── TTL derivation from urgency discipline ─────────────────────────────────── + +def test_ttl_protect_bounded_2s(): + ex, port, reg = _exec(slot=OPEN) + _run(ex.execute(_req(UrgencyClass.PROTECT, side=Side.SELL), NARROW)) + wo = reg.working("req-0") + assert wo.deadline - wo.created == 2.0 # BOUNDED_MS = 2000 ms + + +def test_ttl_rotate_from_deadline_ms(): + ex, port, reg = _exec(slot=OPEN) + _run(ex.execute(_req(UrgencyClass.ROTATE, side=Side.SELL, deadline_ms=30000), NARROW)) + wo = reg.working("req-0") + assert wo.deadline - wo.created == 30.0 # caller's deadline + + +def test_ttl_acquire_unbounded_uses_quote_lifetime(): + from prod.exec_unified import _constants as K + ex, port, reg = _exec() + _run(ex.execute(_req(UrgencyClass.ACQUIRE), NARROW)) + wo = reg.working("req-0") + assert wo.deadline - wo.created == K.MAKER_QUOTE_TTL_S # bounded by chases, not clock