exec(unified): sync local build — kernel_port + dialect §11 + friction §12 + size threading
Lands the /root/dev local increments (share was ENOSPC; now writable) into the canonical repo:
- kernel_port.py: real ExecPort over the DITAv2 kernel (duck-typed; injected to_kernel_intent).
- dialect.py §11: BingX boundary — clientOrderId(H4)/dash/quantize/payload (PostOnly=timeInForce).
- friction.py §12: effective-bps + naive-baseline savings; side-lane journal that can't raise
(b46ebd2); BingX commission-sign flip. DDL ships with code (register in applier verify-set).
- drive_loop/executor/working: Decimal size threaded through plan/working types (exit size cap).
All pure stdlib+Decimal, mutation-litmus RED on the two load-bearing asserts. 132 tests green.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
106
prod/exec_unified/dialect.py
Normal file
106
prod/exec_unified/dialect.py
Normal file
@@ -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).
|
||||||
|
|
||||||
|
``<prefix><request_id>-<attempt>`` 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
|
||||||
@@ -65,6 +65,7 @@ class ResubmitPlan:
|
|||||||
limit_price: Decimal # ignored when method is TAKER
|
limit_price: Decimal # ignored when method is TAKER
|
||||||
attempt: int
|
attempt: int
|
||||||
reduce_only: bool
|
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,
|
# 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,
|
request_id=f"{base}-{nxt}", base_request_id=base, asset=wo.asset,
|
||||||
side=wo.side, action=wo.action, method=method,
|
side=wo.side, action=wo.action, method=method,
|
||||||
limit_price=wo.limit_price, attempt=nxt, reduce_only=reduce_only,
|
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,
|
def build_working_order(*, request_id: str, base_request_id: str, asset: str, side: Side,
|
||||||
action: Action, limit_price: Decimal, decision: RoutingDecision,
|
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
|
"""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."""
|
(max_reprices / cross_on_expiry) so the loop never imports the brain's decision path."""
|
||||||
return WorkingOrder(
|
return WorkingOrder(
|
||||||
request_id=request_id, base_request_id=base_request_id, asset=asset, side=side,
|
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,
|
action=action, limit_price=limit_price, deadline=clock + ttl_s, created=clock,
|
||||||
attempt=attempt,
|
size=size, attempt=attempt,
|
||||||
meta={"max_reprices": decision.max_reprices,
|
meta={"max_reprices": decision.max_reprices,
|
||||||
"cross_on_expiry": decision.cross_on_expiry},
|
"cross_on_expiry": decision.cross_on_expiry},
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -114,13 +114,14 @@ class UnifiedExecutor:
|
|||||||
maker_plan = OrderPlan(
|
maker_plan = OrderPlan(
|
||||||
request_id=rid0, base_request_id=base, asset=request.asset, side=request.side,
|
request_id=rid0, base_request_id=base, asset=request.asset, side=request.side,
|
||||||
action=action, method=ExecutionMethod.MAKER, limit_price=placement.limit_price,
|
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)
|
await self.port.submit(maker_plan)
|
||||||
wo = build_working_order(
|
wo = build_working_order(
|
||||||
request_id=rid0, base_request_id=base, asset=request.asset, side=request.side,
|
request_id=rid0, base_request_id=base, asset=request.asset, side=request.side,
|
||||||
action=action, limit_price=placement.limit_price, decision=decision,
|
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
|
rejected = self.port.slot_view().stage in REJECTED_STAGES
|
||||||
cls = self.dl.after_submit(wo, rejected=rejected)
|
cls = self.dl.after_submit(wo, rejected=rejected)
|
||||||
@@ -139,7 +140,7 @@ class UnifiedExecutor:
|
|||||||
plan = OrderPlan(
|
plan = OrderPlan(
|
||||||
request_id=rid0, base_request_id=base, asset=request.asset, side=request.side,
|
request_id=rid0, base_request_id=base, asset=request.asset, side=request.side,
|
||||||
action=action, method=ExecutionMethod.TAKER, limit_price=Decimal(0),
|
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.submit(plan)
|
||||||
await self.port.pump() # a market cross settles fast; drain it before returning
|
await self.port.pump() # a market cross settles fast; drain it before returning
|
||||||
|
|||||||
212
prod/exec_unified/friction.py
Normal file
212
prod/exec_unified/friction.py
Normal file
@@ -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)
|
||||||
|
"""
|
||||||
177
prod/exec_unified/kernel_port.py
Normal file
177
prod/exec_unified/kernel_port.py
Normal file
@@ -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")
|
||||||
129
prod/exec_unified/test_dialect.py
Normal file
129
prod/exec_unified/test_dialect.py
Normal file
@@ -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"),
|
||||||
|
)
|
||||||
157
prod/exec_unified/test_friction.py
Normal file
157
prod/exec_unified/test_friction.py
Normal file
@@ -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")
|
||||||
177
prod/exec_unified/test_kernel_port.py
Normal file
177
prod/exec_unified/test_kernel_port.py
Normal file
@@ -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)
|
||||||
@@ -39,6 +39,7 @@ class WorkingOrder:
|
|||||||
limit_price: Decimal
|
limit_price: Decimal
|
||||||
deadline: float # in injected-clock units; <= now() means "resolve me"
|
deadline: float # in injected-clock units; <= now() means "resolve me"
|
||||||
created: float
|
created: float
|
||||||
|
size: Decimal = Decimal("0") # base qty the caller asked for (venue needs it; exits cap to slot)
|
||||||
attempt: int = 0
|
attempt: int = 0
|
||||||
meta: dict = field(default_factory=dict)
|
meta: dict = field(default_factory=dict)
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user