diff --git a/prod/exec_unified/dialect.py b/prod/exec_unified/dialect.py new file mode 100644 index 0000000..21a1ed6 --- /dev/null +++ b/prod/exec_unified/dialect.py @@ -0,0 +1,106 @@ +"""BingX venue dialect — the ONLY place venue-specific idioms live (spec §11, §4-14/15/17). + +The layer above speaks a normalized order language (Side, ExecutionMethod, Decimal size); +this module translates it into a BingX REST order payload at the boundary. Everything that is +"BingX law" is encoded here and nowhere else, so a second venue is a second dialect, not a +rewrite. + +Encoded findings (each a headstone): + * H4 / BI-1: clientOrderId unique per attempt, sanitized to the venue charset, length-capped + deterministically (BingX caps it) — a retry must not collide, FIX OrigClOrdID discipline. + * L7 (mm_ correction): POST_ONLY/IOC/FOK are BingX ``timeInForce`` values, NOT ``type`` — + the map that flattened them into ``type`` caused venue rejects. type = LIMIT|MARKET only. + * §4-15: quantize price to tick (conservatively by side, maker-safe) AND size to step + (floor — trade slightly less, never overshoot) BEFORE submit. + * §4-17: one-way flattening needs ``positionSide=BOTH``; MARKET carries no price/TIF; exits + are ``reduceOnly``. + +Pure stdlib + Decimal. No I/O, no dita_v2. +""" +from __future__ import annotations + +import hashlib +import re +from decimal import ROUND_CEILING, ROUND_FLOOR, Decimal + +from .contract import Side +from .router import ExecutionMethod + +# BingX clientOrderId: conservative cap + charset. BingX accepts alphanumerics plus a few +# separators; keep it tight and deterministic so a retry id never collides. +CLIENT_ORDER_ID_MAX = 40 +_CID_ALLOWED = re.compile(r"[^A-Za-z0-9_-]") +# Quote currencies, longest-first, for splitting the dash into a BingX wire symbol. +_QUOTES = ("USDT", "USDC", "BUSD", "USD") + + +def dash_symbol(asset: str) -> str: + """Canonical undashed ("BTCUSDT") → BingX wire symbol ("BTC-USDT"). Idempotent.""" + a = asset.upper() + if "-" in a: + return a + for q in _QUOTES: + if a.endswith(q) and len(a) > len(q): + return f"{a[:-len(q)]}-{q}" + return a + + +def client_order_id(request_id: str, attempt: int = 0, *, prefix: str = "u-") -> str: + """Mint a venue-legal clientOrderId, unique per attempt (audit H4). + + ``-`` sanitized to [A-Za-z0-9_-] and capped at + CLIENT_ORDER_ID_MAX. If it would overflow, truncate the core and append a deterministic + 8-hex digest of the FULL core so uniqueness (and per-attempt distinctness) survives. + """ + if attempt < 0: + raise ValueError(f"attempt must be >= 0, got {attempt}") + core = f"{request_id}-{attempt}" + cid = _CID_ALLOWED.sub("", f"{prefix}{core}") + if len(cid) <= CLIENT_ORDER_ID_MAX: + return cid + digest = hashlib.sha1(core.encode()).hexdigest()[:8] + keep = CLIENT_ORDER_ID_MAX - 9 # room for "-" + 8 hex + return f"{cid[:keep]}-{digest}" + + +def quantize_price(price: Decimal, tick: Decimal, side: Side) -> Decimal: + """Quantize a maker price to tick CONSERVATIVELY by side (never cross the touch): BUY floors + (never up into the ask), SELL ceils (never down into the bid). §4-15 + placer's fill-rate fix.""" + if tick <= 0: + raise ValueError(f"tick must be > 0, got {tick}") + rounding = ROUND_FLOOR if side is Side.BUY else ROUND_CEILING + return (price / tick).quantize(Decimal("1"), rounding=rounding) * tick + + +def quantize_size(size: Decimal, step: Decimal) -> Decimal: + """Floor size to the venue step — trade slightly less, NEVER overshoot (§4-15; L8: partials + cluster on low-price/high-qty symbols, XRP 31.446→31 etc.).""" + if step <= 0: + raise ValueError(f"step must be > 0, got {step}") + return (size / step).quantize(Decimal("1"), rounding=ROUND_FLOOR) * step + + +def build_order_payload(*, asset: str, side: Side, method: ExecutionMethod, + limit_price: Decimal, size: Decimal, reduce_only: bool, + client_order_id: str, tick: Decimal, step: Decimal, + position_side: str = "BOTH") -> dict: + """Normalized order → BingX REST payload. The single venue-idiom boundary.""" + qty = quantize_size(size, step) + if qty <= 0: + raise ValueError(f"quantized size is non-positive (size={size}, step={step}) — " + "below the venue minimum; caller must not submit dust") + payload = { + "symbol": dash_symbol(asset), + "side": side.value, # BUY | SELL (order direction) + "positionSide": position_side, # BOTH for one-way flatten (§4-17) + "quantity": str(qty), + "clientOrderID": client_order_id, + "reduceOnly": "true" if reduce_only else "false", + } + if method is ExecutionMethod.MAKER: + payload["type"] = "LIMIT" + payload["price"] = str(quantize_price(limit_price, tick, side)) + payload["timeInForce"] = "PostOnly" # L7: PostOnly is timeInForce, NOT type + else: + payload["type"] = "MARKET" # no price, no timeInForce on a market cross + return payload diff --git a/prod/exec_unified/drive_loop.py b/prod/exec_unified/drive_loop.py index 46bdb44..0523b5b 100644 --- a/prod/exec_unified/drive_loop.py +++ b/prod/exec_unified/drive_loop.py @@ -65,6 +65,7 @@ class ResubmitPlan: limit_price: Decimal # ignored when method is TAKER attempt: int reduce_only: bool + size: Decimal = Decimal("0") # base qty for the venue (threaded from the request) # The plan the venue receives is the same shape whether it is an initial submit (executor, @@ -252,18 +253,20 @@ class DriveLoop: 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, + size=wo.size, ) 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: + clock: float, ttl_s: float, size: Decimal = Decimal("0"), + 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, + size=size, attempt=attempt, meta={"max_reprices": decision.max_reprices, "cross_on_expiry": decision.cross_on_expiry}, ) diff --git a/prod/exec_unified/executor.py b/prod/exec_unified/executor.py index 0f138ea..cd9ec8a 100644 --- a/prod/exec_unified/executor.py +++ b/prod/exec_unified/executor.py @@ -114,13 +114,14 @@ class UnifiedExecutor: 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, + attempt=0, reduce_only=reduce_only, size=request.size, ) 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, + clock=self.port.clock(), ttl_s=_ttl_seconds(request, decision), + size=request.size, attempt=0, ) rejected = self.port.slot_view().stage in REJECTED_STAGES cls = self.dl.after_submit(wo, rejected=rejected) @@ -139,7 +140,7 @@ class UnifiedExecutor: 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, + attempt=0, reduce_only=reduce_only, size=request.size, ) await self.port.submit(plan) await self.port.pump() # a market cross settles fast; drain it before returning diff --git a/prod/exec_unified/friction.py b/prod/exec_unified/friction.py new file mode 100644 index 0000000..be1f990 --- /dev/null +++ b/prod/exec_unified/friction.py @@ -0,0 +1,212 @@ +"""Friction telemetry — "before-and-after or it didn't happen" (spec §12). + +This is the instrument that turns "smart exec is better" from a claim into a number. It +computes, per resolved order, the effective friction in basis points (adverse slippage vs +the caller's guideline price + realized fee) and the savings vs the naive-MARKET baseline a +T0 flight would have paid. Aggregated per urgency class, that IS the §15 rung-by-rung delta +and simultaneously MALKHUT's reward signal (§5-3) — one instrument, two consumers. + +TWO LAWS this module exists to honor, both learned the hard way: + + * §13 / §4-5 — **telemetry is a lossless SIDE-LANE and is NEVER allowed on the exec path.** + (b46ebd2: a telemetry write on the hot path is a latent stall/raise in the order loop.) + So ``FrictionJournal.record`` catches everything and CANNOT raise. A dead sink loses a + telemetry row; it never touches an order. The pure bps math *can* raise (bad reference + price is a programming error) — but only ever inside record()'s guard, off the exec path. + + * BingX commission SIGN — a NEGATIVE commission is a COST/debit (opposite Binance); + a positive one is a rebate. We paid for this: omp read −2.00 bps as a "rebate" when it + was a +2.00 bps cost. ``fee_cost_from_commission`` encodes the flip so no caller repeats + it. See [[bingx_maker_fee_and_commission_sign]]. + +Pure stdlib + Decimal. No I/O of its own; the CH sink is injected. +""" +from __future__ import annotations + +import logging +from dataclasses import dataclass, field +from decimal import Decimal +from typing import Callable + +from .contract import Side, UrgencyClass + +LOGGER = logging.getLogger(__name__) + +# Measured venue rates (VST, 2026-07-14) — [[bingx_maker_fee_and_commission_sign]]. +# Maker = 2.00 bps (refutes the 1 bp savings-table assumption); taker = 5.016 bps over 1,455 fills. +MAKER_FEE_BPS = Decimal("2.00") +TAKER_FEE_BPS = Decimal("5.016") +_BPS = Decimal("10000") + + +def fee_cost_from_commission(commission: Decimal) -> Decimal: + """Venue commission → signed COST in quote currency (positive = we paid). + + BingX reports commission NEGATIVE for a debit (a cost) and positive for a rebate — the + opposite of Binance. Cost is therefore ``-commission``: a −0.013 commission is a +0.013 + cost; a +0.005 rebate is a −0.005 cost (we earned it). The one-line flip that omp's Q1 + got backwards. THE mutation test: change this to ``return commission`` and the sign test + goes RED. + """ + return -commission + + +def adverse_slippage_bps(guideline_px: Decimal, fill_px: Decimal, side: Side) -> Decimal: + """Signed slippage of the fill vs the caller's guideline, in bps. POSITIVE = adverse + (we paid up on a BUY / sold down on a SELL); NEGATIVE = price improvement. + + The side flip is the whole point: a higher fill helps a seller but hurts a buyer, so the + same raw delta is friction for one side and improvement for the other.""" + if guideline_px <= 0: + raise ValueError(f"guideline_px must be > 0 to reference slippage, got {guideline_px}") + side_sign = Decimal(1) if side is Side.BUY else Decimal(-1) + return side_sign * (fill_px - guideline_px) / guideline_px * _BPS + + +def fee_bps(fee_cost: Decimal, notional: Decimal) -> Decimal: + """Realized fee as bps of notional. ``fee_cost`` is the SIGNED cost (see + fee_cost_from_commission) so a rebate lowers friction. POSITIVE = cost.""" + if notional <= 0: + raise ValueError(f"notional must be > 0 to reference fee bps, got {notional}") + return fee_cost / notional * _BPS + + +def effective_friction_bps(guideline_px: Decimal, fill_px: Decimal, side: Side, + fee_cost: Decimal, size: Decimal) -> Decimal: + """Total realized friction = adverse slippage + fee, in bps. The single number §12 asks + per order. POSITIVE = the order cost us that many bps vs a free fill at guideline.""" + notional = fill_px * size + return adverse_slippage_bps(guideline_px, fill_px, side) + fee_bps(fee_cost, notional) + + +def naive_baseline_bps(guideline_px: Decimal, touch_px: Decimal, side: Side, + taker_fee_bps: Decimal = TAKER_FEE_BPS) -> Decimal: + """What a T0 naive-MARKET cross at the touch would have cost, in bps: it crosses the + spread (guideline→touch, adverse by construction) and pays the taker fee. This is the + per-order baseline the smart fill is measured against (§12; the 2026-07-10 audit is the + frozen AGGREGATE baseline, this is its per-order model).""" + return adverse_slippage_bps(guideline_px, touch_px, side) + taker_fee_bps + + +def savings_vs_naive_bps(effective_bps: Decimal, baseline_bps: Decimal) -> Decimal: + """Smart-exec savings = baseline − effective. POSITIVE = we beat the naive MARKET; the + number that makes "$91–231K" a measurement instead of a hope (§4-20).""" + return baseline_bps - effective_bps + + +@dataclass(frozen=True) +class FrictionRecord: + """One resolved order's §12 telemetry row. Every field the journal DDL carries; the + derived bps are computed lazily so a malformed row can never poison construction.""" + + request_id: str + asset: str + side: Side + urgency: UrgencyClass + maker: bool # True = filled as maker (post-only rested), False = taker cross + guideline_px: Decimal # the caller's reference price at request time + placed_px: Decimal # where we actually rested/sent (0 for pure MARKET) + fill_px: Decimal # realized average fill + touch_px: Decimal # the touch a naive MARKET would have crossed (baseline ref) + size: Decimal + fee_cost: Decimal # SIGNED quote-currency cost (positive = paid; see fee_cost_from_commission) + chases: int = 0 # reprice attempts spent + queue_age_s: Decimal = Decimal(0) + advisor_id: str = "" # MALKHUT advice provenance ("" = none) + triage: str = "" # terminal disposition (filled/crossed/abandoned/…) + + @property + def notional(self) -> Decimal: + return self.fill_px * self.size + + def effective_bps(self) -> Decimal: + return effective_friction_bps(self.guideline_px, self.fill_px, self.side, + self.fee_cost, self.size) + + def baseline_bps(self) -> Decimal: + return naive_baseline_bps(self.guideline_px, self.touch_px, self.side) + + def savings_bps(self) -> Decimal: + return savings_vs_naive_bps(self.effective_bps(), self.baseline_bps()) + + +FrictionSink = Callable[[FrictionRecord], None] + + +class FrictionJournal: + """Side-lane friction telemetry (§4-5). ``record`` is fire-and-forget and CANNOT raise — + a failing sink loses the row and logs, never disturbs the order loop (§13, b46ebd2). + + The sink is injected: in prod it writes the §12 CH rows (``exec_smart_journal``); in + tests it's a list. The journal keeps an in-memory copy only when ``retain`` so aggregate + queries work without a database.""" + + def __init__(self, sink: FrictionSink | None = None, *, retain: bool = True, + logger: logging.Logger = LOGGER) -> None: + self._sink = sink + self._retain = retain + self.log = logger + self._rows: list[FrictionRecord] = [] + + def record(self, rec: FrictionRecord) -> None: + """Emit one row. Swallows ALL exceptions — telemetry never breaks execution.""" + if self._retain: + self._rows.append(rec) + if self._sink is None: + return + try: + self._sink(rec) + except Exception as exc: # noqa: BLE001 — side-lane law: never propagate + self.log.warning("friction.record: sink failed, row dropped: %r", exc) + + def rows(self) -> list[FrictionRecord]: + return list(self._rows) + + def friction_by_urgency(self) -> dict[UrgencyClass, Decimal]: + """Mean effective friction bps per urgency class — the §15 rung delta. Rows whose bps + can't be computed (bad reference px) are skipped, not fatal (side-lane).""" + buckets: dict[UrgencyClass, list[Decimal]] = {} + for r in self._rows: + try: + buckets.setdefault(r.urgency, []).append(r.effective_bps()) + except Exception: # noqa: BLE001 + continue + return {u: sum(v) / len(v) for u, v in buckets.items() if v} + + def savings_by_urgency(self) -> dict[UrgencyClass, Decimal]: + """Mean savings vs naive-MARKET per urgency — the before/after claim, per class.""" + buckets: dict[UrgencyClass, list[Decimal]] = {} + for r in self._rows: + try: + buckets.setdefault(r.urgency, []).append(r.savings_bps()) + except Exception: # noqa: BLE001 + continue + return {u: sum(v) / len(v) for u, v in buckets.items() if v} + + +# ── DDL — ships WITH the code (§12: the 404-storm lesson — tables that don't exist journal +# nothing). When this package syncs to /mnt, register EXEC_SMART_JOURNAL_DDL in the applier +# verify-set (prod/clickhouse/uv/apply_uv_ddl.py) so a missing table fails LOUD, not silent. +EXEC_SMART_JOURNAL_DDL = """\ +CREATE TABLE IF NOT EXISTS {db}.exec_smart_journal ( + ts DateTime64(3) DEFAULT now64(3), + request_id String, + asset String, + side Enum8('BUY' = 1, 'SELL' = 2), + urgency LowCardinality(String), + maker UInt8, + guideline_px Decimal(38, 12), + placed_px Decimal(38, 12), + fill_px Decimal(38, 12), + touch_px Decimal(38, 12), + size Decimal(38, 12), + fee_cost Decimal(38, 12), + effective_bps Decimal(18, 6), + baseline_bps Decimal(18, 6), + savings_bps Decimal(18, 6), + chases UInt16, + queue_age_s Decimal(18, 6), + advisor_id String, + triage LowCardinality(String) +) ENGINE = MergeTree ORDER BY (asset, ts) +""" diff --git a/prod/exec_unified/kernel_port.py b/prod/exec_unified/kernel_port.py new file mode 100644 index 0000000..919444f --- /dev/null +++ b/prod/exec_unified/kernel_port.py @@ -0,0 +1,177 @@ +"""KernelExecPort — the real ExecPort over the DITAv2 kernel + venue (spec §4, §7). + +This is the wiring that turns the pure engine (UnifiedExecutor + DriveLoop) into something +that drives a real venue. It implements the ``ExecPort`` seam by calling the DITAv2 +``ExecutionKernel`` exactly the way ``pink_direct.py`` does — process_intent_async, slot(), +on_venue_event, venue.reconcile/open_positions — WITHOUT importing or modifying the kernel. + +Two design choices that matter: + * The kernel's concrete intent type (``KernelIntent``, from the vendored dita_v2) is reached + through an INJECTED ``to_kernel_intent`` translator. So this module imports nothing from + dita_v2 — it stays testable against a fake kernel, and the real wiring supplies the real + translator (side/action/TIF mapping) from a module that may import dita_v2. Clean seam. + * Everything time/order-truth-based reads the kernel/venue, never ambient state. + +SEAM lessons preserved as referenced comments (zero-suppression — "we paid dearly for pink"): + * K≈E capital reconcile + enter-freeze live INSIDE the kernel (pink_direct.py:717-763); the + port only feeds venue events in via on_venue_event and never reimplements accounting. + * Fill ownership on a shared account (pink _fill_is_ours L588) is a venue-adapter concern; + the port trusts the kernel/venue's own-fill attribution and only latches the timestamp. +""" +from __future__ import annotations + +import inspect +import logging +import time +from dataclasses import dataclass +from decimal import Decimal +from typing import Any, Callable + +from .contract import Side +from .drive_loop import Action, ExecutionMethod, OrderPlan, SlotView +from .working import WorkingOrder + +LOGGER = logging.getLogger(__name__) + + +@dataclass(frozen=True) +class IntentSpec: + """Normalized, venue-agnostic intent. The injected translator turns this into the + kernel's concrete KernelIntent (side→LONG/SHORT, action→KernelCommandType, TIF→metadata). + Everything the kernel needs, nothing about the caller's world.""" + + request_id: str + base_request_id: str + asset: str + side: Side + action_kind: str # "ENTER" | "EXIT" | "CANCEL" + order_type: str # "MARKET" | "LIMIT" + limit_price: Decimal # 0 for MARKET + size: Decimal + tif: str # "PostOnly" | "GTC" + reduce_only: bool + reason: str + + +ToKernelIntent = Callable[[IntentSpec], Any] + + +class KernelExecPort: + """Implements the ExecPort protocol against a DITAv2 ExecutionKernel (duck-typed).""" + + def __init__(self, kernel: Any, to_kernel_intent: ToKernelIntent, *, + clock: Callable[[], float] = time.monotonic, + logger: logging.Logger = LOGGER) -> None: + self.kernel = kernel + self._to_ki = to_kernel_intent + self._clock = clock + self.log = logger + self._last_own_fill: float = 0.0 + + # ── read seams ─────────────────────────────────────────────────────────── + def clock(self) -> float: + return self._clock() + + def last_own_fill_at(self) -> float: + return self._last_own_fill + + def slot_view(self) -> SlotView: + """Kernel-truth snapshot of slot 0. Tolerant of partial mocks (pink _exec_slot_view L940).""" + try: + slot = self.kernel.slot(0) + except Exception: + return SlotView(trade_id="", stage="", size=Decimal(0)) + stage = getattr(slot, "fsm_state", None) + stage_name = getattr(stage, "value", None) or str(stage or "") + return SlotView( + trade_id=str(getattr(slot, "trade_id", "") or ""), + stage=str(stage_name), + size=_dec(getattr(slot, "size", 0)), + ) + + async def open_positions(self) -> list[dict]: + """Venue-truth flat probe (pink _exec_safe_to_requote L1059). Fail-safe is the + CALLER's job (drive_loop treats a raise as 'not provably flat').""" + venue = self.kernel.venue + rows = venue.open_positions() if hasattr(venue, "open_positions") else [] + if inspect.isawaitable(rows): + rows = await rows + return list(rows or []) + + # ── write seams ────────────────────────────────────────────────────────── + async def pump(self) -> None: + """Drain late venue events into the kernel (pink pump_venue_events L1228). The kernel + dedups (seen_event_ids) and owns capital settlement — we only feed events and latch + the own-fill timestamp so the requote hot-window can fail safe.""" + venue = self.kernel.venue + reconcile = getattr(venue, "reconcile", None) + if reconcile is None: + return + try: + events = reconcile() + if inspect.isawaitable(events): + events = await events + except Exception as exc: + self.log.warning("kernel_port.pump: venue reconcile failed: %r", exc) + return + for event in list(events or []): + try: + self.kernel.on_venue_event(event) + except Exception as exc: + self.log.warning("kernel_port.pump: on_venue_event failed: %r", exc) + continue + if _is_fill(event): + # REST reconcile lags WS fills by seconds — latch so the requote gate + # fails safe inside the hot window (pink_direct.py:648, :1055). + self._last_own_fill = self._clock() + + async def cancel(self, wo: WorkingOrder) -> None: + """Cancel a working quote via the kernel (idempotent on the venue; CANCEL_REJECT on a + filled order is harmless — pink _exec_cancel_working L1021).""" + spec = IntentSpec( + request_id=f"{wo.request_id}-cxl", base_request_id=wo.base_request_id, + asset=wo.asset, side=wo.side, action_kind="CANCEL", order_type="MARKET", + limit_price=Decimal(0), size=wo.size, tif="GTC", reduce_only=False, + reason="exec_unified:ttl_cancel", + ) + await self._process(spec) + + async def submit(self, plan: OrderPlan) -> None: + """Translate an OrderPlan → kernel intent and dispatch it. Exits cap their size to + the kernel's authoritative remaining position (pink _exit_intent_from_slot L1358) so a + stale caller size can neither strand nor overshoot a close.""" + size = plan.size + if plan.action is Action.EXIT: + slot_size = self.slot_view().size + if slot_size > 0: + size = min(plan.size, slot_size) if plan.size > 0 else slot_size + is_maker = plan.method is ExecutionMethod.MAKER + spec = IntentSpec( + request_id=plan.request_id, base_request_id=plan.base_request_id, + asset=plan.asset, side=plan.side, action_kind=plan.action.value, + order_type="LIMIT" if is_maker else "MARKET", + limit_price=plan.limit_price if is_maker else Decimal(0), + size=size, tif="PostOnly" if is_maker else "GTC", + reduce_only=plan.reduce_only, reason="exec_unified:submit", + ) + await self._process(spec) + + async def _process(self, spec: IntentSpec) -> Any: + ki = self._to_ki(spec) + result = self.kernel.process_intent_async(ki) + if inspect.isawaitable(result): + result = await result + return result + + +def _dec(v: object) -> Decimal: + try: + return Decimal(str(v or 0)) + except Exception: + return Decimal(0) + + +def _is_fill(event: Any) -> bool: + kind = getattr(getattr(event, "kind", None), "value", "") or str(getattr(event, "kind", "") or "") + status = getattr(getattr(event, "status", None), "value", "") or str(getattr(event, "status", "") or "") + return "FILL" in kind.upper() or status.upper() in ("FILLED", "PARTIALLY_FILLED") diff --git a/prod/exec_unified/test_dialect.py b/prod/exec_unified/test_dialect.py new file mode 100644 index 0000000..521ad40 --- /dev/null +++ b/prod/exec_unified/test_dialect.py @@ -0,0 +1,129 @@ +"""BingX dialect tests — venue idioms are law (spec §4-17). Each asserts a rule a real +venue reject taught us. Mutation notes inline. + +Run: /home/dolphin/siloqy_env/bin/python3 -m pytest prod/exec_unified/test_dialect.py -q +""" +from __future__ import annotations + +from decimal import Decimal + +import pytest + +from prod.exec_unified.contract import Side +from prod.exec_unified.dialect import ( + CLIENT_ORDER_ID_MAX, + build_order_payload, + client_order_id, + dash_symbol, + quantize_price, + quantize_size, +) +from prod.exec_unified.router import ExecutionMethod + + +# ── symbol dashing ─────────────────────────────────────────────────────────── + +@pytest.mark.parametrize("undashed,wire", [ + ("BTCUSDT", "BTC-USDT"), + ("ETHUSDT", "ETH-USDT"), + ("1000SHIBUSDT", "1000SHIB-USDT"), + ("USDCUSDT", "USDC-USDT"), + ("BTC-USDT", "BTC-USDT"), # already dashed → idempotent +]) +def test_dash_symbol(undashed, wire): + assert dash_symbol(undashed) == wire + + +# ── clientOrderId (audit H4 / BI-1) ────────────────────────────────────────── + +def test_client_order_id_format_and_per_attempt(): + assert client_order_id("req", 0) == "u-req-0" + assert client_order_id("req", 1) == "u-req-1" # retry gets a DIFFERENT id + assert client_order_id("req", 0) != client_order_id("req", 1) + + +def test_client_order_id_sanitizes_charset(): + cid = client_order_id("BTC/USDT:long #7", 0) # slashes, colon, space, hash + assert all(c.isalnum() or c in "_-" for c in cid) + + +def test_client_order_id_length_capped_but_still_unique_per_attempt(): + long_req = "x" * 200 + a, b = client_order_id(long_req, 0), client_order_id(long_req, 1) + assert len(a) <= CLIENT_ORDER_ID_MAX and len(b) <= CLIENT_ORDER_ID_MAX + assert a != b # digest preserves per-attempt distinctness + assert client_order_id(long_req, 0) == a # deterministic (no RNG) + + +def test_client_order_id_rejects_negative_attempt(): + with pytest.raises(ValueError): + client_order_id("req", -1) + + +# ── quantization (§4-15) ───────────────────────────────────────────────────── + +def test_quantize_price_conservative_by_side(): + # BUY floors (never up into the ask); SELL ceils (never down into the bid). + assert quantize_price(Decimal("100.017"), Decimal("0.01"), Side.BUY) == Decimal("100.01") + assert quantize_price(Decimal("100.011"), Decimal("0.01"), Side.SELL) == Decimal("100.02") + + +def test_quantize_size_floors_never_overshoots(): + assert quantize_size(Decimal("31.446"), Decimal("1")) == Decimal("31") + assert quantize_size(Decimal("0.00019"), Decimal("0.0001")) == Decimal("0.0001") + + +def test_quantize_guards_nonpositive(): + with pytest.raises(ValueError): + quantize_price(Decimal("100"), Decimal("0"), Side.BUY) + with pytest.raises(ValueError): + quantize_size(Decimal("1"), Decimal("0")) + + +# ── payload assembly (venue idioms) ────────────────────────────────────────── + +def _maker(): + return build_order_payload( + asset="BTCUSDT", side=Side.BUY, method=ExecutionMethod.MAKER, + limit_price=Decimal("64547.53"), size=Decimal("0.00019"), reduce_only=False, + client_order_id="u-req-0", tick=Decimal("0.1"), step=Decimal("0.0001"), + ) + + +def test_maker_payload_postonly_is_timeInForce_not_type(): + # THE mm_ L7 correction. Mutation: put "POST_ONLY" in type -> this goes RED. + p = _maker() + assert p["type"] == "LIMIT" # type is LIMIT|MARKET ONLY + assert p["timeInForce"] == "PostOnly" # PostOnly rides timeInForce + assert "POST_ONLY" not in p["type"] and "PostOnly" not in p["type"] + + +def test_maker_payload_quantized_and_shaped(): + p = _maker() + assert p["symbol"] == "BTC-USDT" + assert p["price"] == "64547.5" # floored to tick 0.1 (BUY) + assert p["quantity"] == "0.0001" # floored to step + assert p["positionSide"] == "BOTH" # one-way flatten idiom + assert p["reduceOnly"] == "false" + assert p["clientOrderID"] == "u-req-0" + + +def test_market_payload_has_no_price_or_tif(): + p = build_order_payload( + asset="ETHUSDT", side=Side.SELL, method=ExecutionMethod.TAKER, + limit_price=Decimal("0"), size=Decimal("0.5"), reduce_only=True, + client_order_id="u-req-0", tick=Decimal("0.01"), step=Decimal("0.001"), + ) + assert p["type"] == "MARKET" + assert "price" not in p and "timeInForce" not in p # a market cross carries neither + assert p["reduceOnly"] == "true" # exit reduces + + +def test_payload_refuses_dust_below_step(): + # size that floors to zero must not silently submit a 0-qty order. + with pytest.raises(ValueError): + build_order_payload( + asset="BTCUSDT", side=Side.BUY, method=ExecutionMethod.TAKER, + limit_price=Decimal("0"), size=Decimal("0.00005"), reduce_only=False, + client_order_id="u-req-0", tick=Decimal("0.1"), step=Decimal("0.0001"), + ) diff --git a/prod/exec_unified/test_friction.py b/prod/exec_unified/test_friction.py new file mode 100644 index 0000000..408a9ba --- /dev/null +++ b/prod/exec_unified/test_friction.py @@ -0,0 +1,157 @@ +"""Friction telemetry tests (§12). Each asserts a business invariant, not execution. +Two are load-bearing and mutation-annotated: the BingX commission SIGN and the side-lane +NEVER-RAISES law (b46ebd2). + +Run: /home/dolphin/siloqy_env/bin/python3 -m pytest prod/exec_unified/test_friction.py -q +""" +from __future__ import annotations + +from decimal import Decimal + +import pytest + +from prod.exec_unified.contract import Side, UrgencyClass +from prod.exec_unified.friction import ( + MAKER_FEE_BPS, + TAKER_FEE_BPS, + FrictionJournal, + FrictionRecord, + adverse_slippage_bps, + effective_friction_bps, + fee_bps, + fee_cost_from_commission, + naive_baseline_bps, + savings_vs_naive_bps, +) + + +# ── the commission SIGN (the omp headstone) ────────────────────────────────── + +def test_negative_commission_is_a_cost_not_a_rebate(): + # BingX: negative commission = DEBIT. omp read this as a rebate and was wrong. + # Mutation: `return commission` -> this goes RED. + assert fee_cost_from_commission(Decimal("-0.013")) == Decimal("0.013") # cost, positive + + +def test_positive_commission_is_a_rebate_negative_cost(): + assert fee_cost_from_commission(Decimal("0.005")) == Decimal("-0.005") # earned, negative cost + + +# ── adverse slippage: the side flip is the point ───────────────────────────── + +def test_buy_above_guideline_is_adverse_sell_is_improvement(): + # same raw +0.1% delta: a cost to the buyer, a gift to the seller. + buy = adverse_slippage_bps(Decimal("100"), Decimal("100.1"), Side.BUY) + sell = adverse_slippage_bps(Decimal("100"), Decimal("100.1"), Side.SELL) + assert buy == Decimal("10") # +10 bps adverse + assert sell == Decimal("-10") # −10 bps = price improvement + # Mutation: drop the side_sign flip -> sell would read +10 and this fails. + + +def test_buy_below_guideline_is_price_improvement(): + assert adverse_slippage_bps(Decimal("100"), Decimal("99.9"), Side.BUY) == Decimal("-10") + + +def test_slippage_rejects_nonpositive_reference(): + with pytest.raises(ValueError): + adverse_slippage_bps(Decimal("0"), Decimal("1"), Side.BUY) + + +# ── fee bps ────────────────────────────────────────────────────────────────── + +def test_fee_bps_of_notional(): + # 0.02 quote cost on 100 notional = 2 bps (the measured maker rate shape). + assert fee_bps(Decimal("0.02"), Decimal("100")) == Decimal("2") + + +def test_fee_bps_rejects_nonpositive_notional(): + with pytest.raises(ValueError): + fee_bps(Decimal("1"), Decimal("0")) + + +# ── effective friction = slippage + fee ────────────────────────────────────── + +def test_effective_friction_sums_slippage_and_fee(): + # BUY 1 unit, guideline 100, fill 100.1 (+10 bps), fee 0.020021 on 100.1 notional (~2 bps). + eff = effective_friction_bps(Decimal("100"), Decimal("100.1"), Side.BUY, + fee_cost=Decimal("0.020020"), size=Decimal("1")) + # 10 bps slippage + ~2 bps fee ≈ 12 bps + assert Decimal("11.9") < eff < Decimal("12.1") + + +# ── naive baseline + savings: the before/after ─────────────────────────────── + +def test_naive_baseline_crosses_spread_plus_taker(): + # BUY guideline 100, touch (ask) 100.05 -> 5 bps cross + 5.016 taker = ~10.016 bps. + base = naive_baseline_bps(Decimal("100"), Decimal("100.05"), Side.BUY) + assert base == Decimal("5") + TAKER_FEE_BPS + + +def test_savings_positive_when_smart_beats_naive(): + assert savings_vs_naive_bps(Decimal("4"), Decimal("10")) == Decimal("6") # saved 6 bps + assert savings_vs_naive_bps(Decimal("12"), Decimal("10")) == Decimal("-2") # worse than naive + + +def test_maker_rate_is_cheaper_than_taker(): + assert MAKER_FEE_BPS < TAKER_FEE_BPS # the whole reason T1 exists + + +# ── record derives its own bps ─────────────────────────────────────────────── + +def _rec(**kw): + base = dict(request_id="r", asset="BTCUSDT", side=Side.BUY, urgency=UrgencyClass.ACQUIRE, + maker=True, guideline_px=Decimal("100"), placed_px=Decimal("100"), + fill_px=Decimal("100"), touch_px=Decimal("100.05"), size=Decimal("1"), + fee_cost=Decimal("0.02")) + base.update(kw) + return FrictionRecord(**base) + + +def test_record_maker_fill_at_guideline_beats_naive(): + # maker filled exactly at guideline (0 slippage) paying 2 bps fee vs naive 5bps cross+5bps taker. + r = _rec() + assert r.effective_bps() == Decimal("2") # 0 slippage + 2 bps fee + assert r.savings_bps() > Decimal("7") # naive ~10.016 − 2 = ~8 bps saved + + +# ── THE side-lane law: telemetry NEVER raises on the exec path (b46ebd2) ────── + +def test_record_swallows_a_failing_sink(): + def boom(_): + raise RuntimeError("clickhouse down") + j = FrictionJournal(sink=boom) + j.record(_rec()) # must NOT raise — a dead sink loses a row, never an order. + # Mutation: remove the try/except in record() -> this raises and goes RED. + assert len(j.rows()) == 1 # retained locally even though the sink failed + + +def test_record_forwards_to_sink_when_healthy(): + seen = [] + j = FrictionJournal(sink=seen.append) + j.record(_rec(request_id="a")) + assert len(seen) == 1 and seen[0].request_id == "a" + + +# ── aggregates: the §15 rung delta ─────────────────────────────────────────── + +def test_friction_and_savings_by_urgency(): + j = FrictionJournal() + j.record(_rec(request_id="a", urgency=UrgencyClass.ACQUIRE, fill_px=Decimal("100"))) # 2 bps + j.record(_rec(request_id="b", urgency=UrgencyClass.ACQUIRE, fill_px=Decimal("100.1"), + fee_cost=Decimal("0.02"))) # ~12 bps + j.record(_rec(request_id="c", urgency=UrgencyClass.PROTECT, fill_px=Decimal("100"))) # 2 bps + fr = j.friction_by_urgency() + assert UrgencyClass.ACQUIRE in fr and UrgencyClass.PROTECT in fr + assert fr[UrgencyClass.PROTECT] == Decimal("2") + assert Decimal("6") < fr[UrgencyClass.ACQUIRE] < Decimal("8") # mean of ~2 and ~12 + sv = j.savings_by_urgency() + assert sv[UrgencyClass.PROTECT] > 0 # maker@guideline beats naive + + +def test_aggregate_skips_unbucketable_rows_without_dying(): + # a row with a poison reference px must not kill the whole aggregate (side-lane). + j = FrictionJournal() + j.record(_rec(request_id="ok")) + j.record(_rec(request_id="bad", guideline_px=Decimal("0"))) # effective_bps() will raise + fr = j.friction_by_urgency() # must still return the good one + assert fr[UrgencyClass.ACQUIRE] == Decimal("2") diff --git a/prod/exec_unified/test_kernel_port.py b/prod/exec_unified/test_kernel_port.py new file mode 100644 index 0000000..bbc456a --- /dev/null +++ b/prod/exec_unified/test_kernel_port.py @@ -0,0 +1,177 @@ +"""End-to-end: UnifiedExecutor → DriveLoop → KernelExecPort → (fake) DITAv2 kernel. + +Proves the pure engine drives a kernel-shaped backend correctly — translation (side/action/ +TIF/order_type), exit-size capping, maker-rest→fill, exit-escalation, taker-cross. The fake +kernel mimics ExecutionKernel's surface (slot()/process_intent_async/on_venue_event/venue). +The injected translator is identity, so the fake reads IntentSpec directly and we can assert +exactly what the port produced. + +Run: /home/dolphin/siloqy_env/bin/python3 -m pytest prod/exec_unified/test_kernel_port.py -q +""" +from __future__ import annotations + +import asyncio +from decimal import Decimal +from types import SimpleNamespace + +from prod.exec_unified.contract import ExecutionRequest, Side, UrgencyClass +from prod.exec_unified.drive_loop import Action, DriveLoop +from prod.exec_unified.executor import ExecuteStatus, UnifiedExecutor +from prod.exec_unified.kernel_port import IntentSpec, KernelExecPort +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")) + + +def _slot(trade_id, stage, size): + return SimpleNamespace(trade_id=trade_id, fsm_state=SimpleNamespace(value=stage), + size=size, asset="BTCUSDT") + + +class FakeVenue: + def __init__(self): + self.events = [] + self.positions = [] + def reconcile(self): + ev, self.events = self.events, [] + return ev + def open_positions(self): + return self.positions + + +class FakeKernel: + """Minimal ExecutionKernel surface. process_intent_async mutates slot 0 like the real + single-slot FSM: MARKET fills now; LIMIT rests (slot unchanged); CANCEL is a no-op.""" + def __init__(self, slot): + self._slot = slot + self.venue = FakeVenue() + self.max_slots = 1 + self.processed: list[IntentSpec] = [] + def slot(self, i): + return self._slot + async def process_intent_async(self, spec: IntentSpec): + self.processed.append(spec) + if spec.action_kind == "ENTER" and spec.order_type == "MARKET": + self._slot = _slot(spec.request_id, "POSITION_OPEN", spec.size) + elif spec.action_kind == "EXIT" and spec.order_type == "MARKET": + self._slot = _slot("", "IDLE", Decimal(0)) + return SimpleNamespace(accepted=True) + def on_venue_event(self, e): + fn = getattr(e, "apply", None) + if fn: + fn(self) + # test helpers + def force_open(self, tid, size="1"): + self._slot = _slot(tid, "POSITION_OPEN", Decimal(size)) + + +def _ident(spec): # local translator: kernel reads the normalized spec directly + return spec + + +class Clock: + def __init__(self, t=1000.0): self.t = t + def __call__(self): return self.t + + +def _stack(slot): + clock = Clock() + kernel = FakeKernel(slot) + port = KernelExecPort(kernel, _ident, clock=clock) + reg = WorkingRegistry(port.clock) + dl = DriveLoop(port, reg) + return UnifiedExecutor(dl), dl, kernel, reg + + +def _req(urgency, side=Side.BUY, size="0.001", **kw): + return ExecutionRequest(request_id="req", asset="BTCUSDT", side=side, + size=Decimal(size), urgency=urgency, **kw) + + +def _run(c): return asyncio.run(c) + + +# ── translation correctness ────────────────────────────────────────────────── + +def test_maker_entry_translated_to_limit_postonly(): + ex, dl, kernel, reg = _stack(_slot("", "IDLE", Decimal(0))) + r = _run(ex.execute(_req(UrgencyClass.ACQUIRE), NARROW)) + assert r.status is ExecuteStatus.RESTING_MAKER + spec = kernel.processed[-1] + assert spec.action_kind == "ENTER" + assert spec.order_type == "LIMIT" and spec.tif == "PostOnly" # maker + assert spec.side is Side.BUY + assert spec.limit_price == Decimal("100") and spec.size == Decimal("0.001") + assert reg.working("req-0") is not None + + +def test_catastrophic_translated_to_market_gtc(): + ex, dl, kernel, reg = _stack(_slot("pos", "POSITION_OPEN", Decimal("1"))) + r = _run(ex.execute(_req(UrgencyClass.CATASTROPHIC, side=Side.SELL, size="1"), NARROW)) + assert r.status is ExecuteStatus.CROSSED_TAKER + spec = kernel.processed[-1] + assert spec.action_kind == "EXIT" and spec.order_type == "MARKET" and spec.tif == "GTC" + assert spec.reduce_only is True + # MARKET exit closed the slot in the fake kernel + assert kernel.slot(0).size == Decimal(0) + + +# ── exit size capping (pink _exit_intent_from_slot) ────────────────────────── + +def test_exit_size_capped_to_slot(): + # Caller asks to close 10 but the slot holds 0.5 — the port caps to 0.5 (never overshoot). + ex, dl, kernel, reg = _stack(_slot("pos", "POSITION_OPEN", Decimal("0.5"))) + _run(ex.execute(_req(UrgencyClass.PROTECT, side=Side.SELL, size="10"), NARROW)) + spec = kernel.processed[-1] + assert spec.action_kind == "EXIT" + assert spec.size == Decimal("0.5") # capped to the authoritative remaining size + + +# ── full lifecycle: maker entry rests, then fills after TTL ────────────────── + +def test_entry_rests_then_fills_after_ttl(): + ex, dl, kernel, reg = _stack(_slot("", "IDLE", Decimal(0))) + r = _run(ex.execute(_req(UrgencyClass.ACQUIRE), NARROW)) + assert r.status is ExecuteStatus.RESTING_MAKER + wo = reg.working("req-0") + assert wo is not None + # the maker quote fills — kernel slot now shows our position (trade_id == our rid0) + kernel.force_open("req-0") + res = _run(dl.handle_expired(wo)) + from prod.exec_unified.drive_loop import Resolved + assert res is Resolved.FILLED_AFTER_TTL + assert reg.working("req-0") is None # cleared, not retried + + +# ── full lifecycle: maker exit never strands → escalates to MARKET ─────────── + +def test_maker_exit_escalates_to_market_and_closes(): + ex, dl, kernel, reg = _stack(_slot("pos", "POSITION_OPEN", Decimal("1"))) + r = _run(ex.execute(_req(UrgencyClass.PROTECT, side=Side.SELL, size="1"), NARROW)) + assert r.status is ExecuteStatus.RESTING_MAKER # maker exit resting + wo = reg.working("req-0") + # TTL fires, position still open → escalate to MARKET (never strand) + from prod.exec_unified.drive_loop import Resolved + res = _run(dl.handle_expired(wo)) + assert res is Resolved.ESCALATED_MARKET + assert kernel.processed[-1].order_type == "MARKET" # crossed + assert kernel.slot(0).size == Decimal(0) # position closed + + +# ── requote gate reads real venue positions ────────────────────────────────── + +def test_requote_gate_blocks_when_venue_holds_position(): + ex, dl, kernel, reg = _stack(_slot("", "IDLE", Decimal(0))) + # entry with a chase budget, so a miss would try to requote + _run(ex.execute(_req(UrgencyClass.ACQUIRE), NARROW)) + wo = reg.working("req-0") + wo.meta["max_reprices"] = 1 # allow a retry attempt + kernel.venue.positions = [{"positionAmt": "0.3"}] # venue is NOT flat + from prod.exec_unified.drive_loop import Resolved + res = _run(dl.handle_expired(wo)) + assert res is Resolved.REQUOTE_BLOCKED # port.open_positions() saw the position + # no new order submitted after the blocked requote + assert all(s.action_kind != "ENTER" or s.request_id == "req-0" for s in kernel.processed) diff --git a/prod/exec_unified/working.py b/prod/exec_unified/working.py index 06573b9..f8c6f27 100644 --- a/prod/exec_unified/working.py +++ b/prod/exec_unified/working.py @@ -39,6 +39,7 @@ class WorkingOrder: limit_price: Decimal deadline: float # in injected-clock units; <= now() means "resolve me" created: float + size: Decimal = Decimal("0") # base qty the caller asked for (venue needs it; exits cap to slot) attempt: int = 0 meta: dict = field(default_factory=dict)