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