SmartExecBridge: a drop-in for PromotionBridge.try_promote that rests entries as PostOnly GTX
makers (via exec_unified.router policy) instead of always paying the taker cross. DORMANT —
selected only by UV_SMART_EXEC=1 (default off); nothing imports it yet, so the running flight
is untouched. T1 = smallest bug surface (spec §15):
- ENTER -> maker (ACQUIRE): PostOnly LIMIT @ touch; unfilled by next scan tick -> CANCEL/abandon
('a missed entry is free' §4-1). No chase, no cross, no requote race.
- EXIT -> MARKET unchanged (never strand; zero new exit risk in v1).
- the 6s scan tick IS the drive clock: sweep-stale-then-place each promote.
Reuses router.decide + friction side-lane (never blocks a promote); fail-soft (never raises
into the scan loop). 11 tests, 2 mutation litmus RED (entry-maker policy; stale-sweep).
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
227 lines
11 KiB
Python
227 lines
11 KiB
Python
"""Smart-exec bolt for the PRIME flight (T1 — "better orders, usable right now").
|
|
|
|
Bolts the unified exec layer's *whether-maker* policy (prod.exec_unified.router) onto the
|
|
PRIME promotion path so entries try to REST as maker (PostOnly GTX) instead of always paying
|
|
the taker cross. It is a drop-in alternative to the naive one-shot MARKET promote in
|
|
``promotion.PromotionBridge`` — selected by the ``UV_SMART_EXEC`` env flag, DORMANT by default
|
|
so the running flight is unaffected until an operator flips it.
|
|
|
|
T1 policy (deliberately the smallest nonzero bug surface — spec §15 T1):
|
|
* ENTER → maker (ACQUIRE): PostOnly LIMIT at the reference touch. If it does not fill by the
|
|
next scan tick it is CANCELLED and abandoned — "a missed entry is free" (§4-1). No chase,
|
|
no cross, no requote race. The one new state is "resting GTX with a TTL", the one new
|
|
transition is "TTL → cancel", and both are read off the kernel slot, not invented here.
|
|
* EXIT → MARKET (unchanged naive behaviour): exits never strand, and adding maker-then-cross
|
|
to the exit leg is T2+; v1 takes ZERO new risk on the exit side.
|
|
* The 6 s scan tick IS the drive clock: each promote first sweeps a stale resting entry
|
|
(cancel if the slot still shows it working past its TTL), THEN places the new one.
|
|
|
|
WHAT THIS DOES NOT DO (honest negative constraints, §13):
|
|
* It never changes side/size/asset (reads them off the intent the promotion layer built).
|
|
* It never crosses on a missed ENTER (abandons). It never delays or vetoes an EXIT.
|
|
* It writes friction telemetry on a SIDE LANE only (never blocks a promote).
|
|
|
|
Reuses: prod.exec_unified.router.decide (policy), .friction (telemetry). dita_v2 imported
|
|
lazily so importing this module never drags the kernel in.
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
import os
|
|
import time
|
|
from dataclasses import dataclass, field
|
|
from decimal import Decimal
|
|
from typing import Any, Optional
|
|
|
|
from prod.exec_unified.contract import ExecutionRequest, Side, UrgencyClass
|
|
from prod.exec_unified.friction import FrictionJournal, FrictionRecord
|
|
from prod.exec_unified.router import ExecutionMethod, decide
|
|
|
|
LOGGER = logging.getLogger("uv.smart_bridge")
|
|
|
|
SMART_EXEC_ENV = "UV_SMART_EXEC"
|
|
# Resting entry TTL. The scan tick (~6 s) is the coarse clock; anything older than this at the
|
|
# next tick is treated as expired and swept. Kept >= one tick so a fresh quote gets a full tick.
|
|
DEFAULT_ENTRY_TTL_S = 6.0
|
|
|
|
|
|
def smart_exec_enabled() -> bool:
|
|
"""True iff the operator has opted this flight into smart exec. Off ⇒ caller uses the
|
|
naive promote. Re-read every call (no caching) so it can be toggled without a code change."""
|
|
return os.environ.get(SMART_EXEC_ENV, "0").strip() == "1"
|
|
|
|
|
|
def _side_from_intent(intent: Any) -> Side:
|
|
"""Kernel intent carries POSITION side (LONG/SHORT); the ORDER side we quote is BUY to open
|
|
a long / sell to open a short (this path is entries only). Verified vs bingx_venue:627."""
|
|
val = getattr(getattr(intent, "side", None), "value", "") or str(getattr(intent, "side", ""))
|
|
return Side.BUY if val.upper() == "LONG" else Side.SELL
|
|
|
|
|
|
def _is_enter(intent: Any) -> bool:
|
|
val = getattr(getattr(intent, "action", None), "value", "") or str(getattr(intent, "action", ""))
|
|
return val.upper() == "ENTER"
|
|
|
|
|
|
def urgency_for(intent: Any) -> UrgencyClass:
|
|
"""T1 urgency assignment: entries are patient makers that abandon; everything else (exits)
|
|
crosses now. This is the ONLY place the T1 policy chooses a class."""
|
|
return UrgencyClass.ACQUIRE if _is_enter(intent) else UrgencyClass.CATASTROPHIC
|
|
|
|
|
|
@dataclass
|
|
class _Resting:
|
|
"""One tracked resting maker (single slot). Enough to sweep it at the next tick."""
|
|
intent_id: str
|
|
trade_id: str
|
|
placed_at: float
|
|
reference_price: float
|
|
|
|
|
|
@dataclass
|
|
class SmartExecBridge:
|
|
"""Drop-in for PromotionBridge.try_promote with T1 maker-first entries.
|
|
|
|
Same gate (injected), same fail-soft contract (never raises out of try_promote). Holds the
|
|
single-slot resting-entry state so it can sweep a stale quote before placing a new one."""
|
|
|
|
gate: Any # PromotionGate (two-man rule) — reused verbatim
|
|
kernel: Any # DITAv2 ExecutionKernel (.process_intent / .slot)
|
|
friction: FrictionJournal = field(default_factory=FrictionJournal)
|
|
entry_ttl_s: float = DEFAULT_ENTRY_TTL_S
|
|
clock: Any = time.monotonic
|
|
stats: Any = None # optional shared _Stats from promotion (duck-typed)
|
|
_resting: Optional[_Resting] = None
|
|
|
|
# ── slot truth (tolerant of partial kernels) ─────────────────────────────
|
|
def _slot_stage_size(self) -> tuple[str, Decimal, str]:
|
|
try:
|
|
slot = self.kernel.slot(0)
|
|
except Exception:
|
|
return "", Decimal(0), ""
|
|
stage = getattr(getattr(slot, "fsm_state", None), "value", None) or str(getattr(slot, "fsm_state", "") or "")
|
|
size = Decimal(str(getattr(slot, "size", 0) or 0))
|
|
tid = str(getattr(slot, "trade_id", "") or "")
|
|
return str(stage), size, tid
|
|
|
|
def _entry_still_working(self, resting: _Resting) -> bool:
|
|
stage, size, tid = self._slot_stage_size()
|
|
# unfilled entry ⇒ still in a working/pre-fill stage AND no position size yet, same trade
|
|
working = stage in ("ENTRY_WORKING", "ORDER_REQUESTED", "ORDER_SENT", "ORDER_ACKED", "IDLE")
|
|
return working and size <= 0 and (tid == resting.trade_id or tid == "")
|
|
|
|
def _sweep_stale_entry(self) -> None:
|
|
"""Cancel a resting maker entry that has outlived its TTL and still hasn't filled. Pump
|
|
happens inside the kernel on the next process_intent; we only decide to cancel."""
|
|
r = self._resting
|
|
if r is None:
|
|
return
|
|
expired = (self.clock() - r.placed_at) >= self.entry_ttl_s
|
|
if not expired:
|
|
return
|
|
if self._entry_still_working(r):
|
|
try:
|
|
self.kernel.process_intent(self._cancel_intent(r))
|
|
LOGGER.info("SMART: swept stale resting entry %s (ttl %.1fs)", r.intent_id, self.entry_ttl_s)
|
|
except Exception as exc: # fail-soft: a failed cancel must not block
|
|
LOGGER.warning("SMART: sweep cancel failed for %s: %r", r.intent_id, exc)
|
|
self._resting = None
|
|
|
|
def _cancel_intent(self, r: _Resting) -> Any:
|
|
from prod.clean_arch.dita_v2.contracts import KernelCommandType, KernelIntent, TradeSide
|
|
from datetime import datetime, timezone
|
|
return KernelIntent(
|
|
timestamp=datetime.now(timezone.utc), intent_id=f"{r.intent_id}-cxl",
|
|
trade_id=r.trade_id, slot_id=0, asset="", side=TradeSide.FLAT,
|
|
action=KernelCommandType.CANCEL, reference_price=0.0, target_size=0.0, leverage=1.0,
|
|
reason="uv_smart:ttl_abandon", metadata={"smart_exec": True},
|
|
)
|
|
|
|
def _as_maker(self, intent: Any) -> Any:
|
|
"""Return a maker-shaped copy of the promotion-built intent: LIMIT @ touch, PostOnly.
|
|
Everything else (side/size/asset/ids) is preserved untouched."""
|
|
import dataclasses
|
|
meta = dict(getattr(intent, "metadata", {}) or {})
|
|
meta.update({"smart_exec": True, "tif": "PostOnly", "uv_exec_tier": "T1"})
|
|
return dataclasses.replace(
|
|
intent, order_type="LIMIT",
|
|
limit_price=float(getattr(intent, "reference_price", 0.0) or 0.0),
|
|
metadata=meta,
|
|
)
|
|
|
|
def try_promote(self, *, decision: dict, ctx: Any = None,
|
|
build_intent: Any = None) -> bool:
|
|
"""Gated, fail-soft T1 promote. ``build_intent`` is the promotion layer's
|
|
``build_kernel_intent_from_decision`` (injected so this module never imports flight code)."""
|
|
try:
|
|
if not self.gate.is_active():
|
|
self._bump("skipped")
|
|
return False
|
|
if not decision.get("has_entry") and not decision.get("is_exit"):
|
|
# nothing actionable this tick, but still sweep a stale resting quote
|
|
self._sweep_stale_entry()
|
|
self._bump("skipped")
|
|
return False
|
|
|
|
intent = build_intent(decision=decision, ctx=ctx, reason="uv_smart")
|
|
self._sweep_stale_entry() # pump-before-place: clear a stale quote first
|
|
|
|
u = urgency_for(intent)
|
|
req = self._request_from(intent, u)
|
|
routing = decide(req)
|
|
maker = routing.method is ExecutionMethod.MAKER
|
|
|
|
send = self._as_maker(intent) if maker else intent # taker path = the naive intent verbatim
|
|
outcome = self.kernel.process_intent(send)
|
|
dc = getattr(getattr(outcome, "diagnostic_code", ""), "value", None) or str(
|
|
getattr(outcome, "diagnostic_code", ""))
|
|
accepted = dc == "OK"
|
|
|
|
if maker and accepted:
|
|
self._resting = _Resting(
|
|
intent_id=str(getattr(send, "intent_id", "")),
|
|
trade_id=str(getattr(send, "trade_id", "")),
|
|
placed_at=self.clock(),
|
|
reference_price=float(getattr(send, "reference_price", 0.0) or 0.0),
|
|
)
|
|
self._emit_friction(send, req, maker)
|
|
self._bump("accepted" if accepted else "rejected")
|
|
LOGGER.info("SMART PROMOTE: %s method=%s -> %s (asset=%s qty=%.4f)",
|
|
getattr(send, "intent_id", "?"), routing.method.value, dc,
|
|
getattr(send, "asset", "?"), float(getattr(send, "target_size", 0.0) or 0.0))
|
|
return accepted
|
|
except Exception as exc: # never raise into the scan loop
|
|
self._bump("errors")
|
|
LOGGER.error("SMART PROMOTE ERROR: %s: %s", type(exc).__name__, exc)
|
|
return False
|
|
|
|
def _request_from(self, intent: Any, urgency: UrgencyClass) -> ExecutionRequest:
|
|
return ExecutionRequest(
|
|
request_id=str(getattr(intent, "intent_id", "") or "uv"),
|
|
asset=str(getattr(intent, "asset", "") or "BTCUSDT"),
|
|
side=_side_from_intent(intent),
|
|
size=Decimal(str(getattr(intent, "target_size", 0) or 0)),
|
|
urgency=urgency,
|
|
)
|
|
|
|
def _emit_friction(self, intent: Any, req: ExecutionRequest, maker: bool) -> None:
|
|
ref = Decimal(str(getattr(intent, "reference_price", 0) or 0))
|
|
if ref <= 0:
|
|
return # no reference ⇒ nothing to measure (side-lane)
|
|
self.friction.record(FrictionRecord(
|
|
request_id=req.request_id, asset=req.asset, side=req.side, urgency=req.urgency,
|
|
maker=maker, guideline_px=ref, placed_px=ref, fill_px=ref, touch_px=ref,
|
|
size=req.size, fee_cost=Decimal(0), triage="placed",
|
|
))
|
|
|
|
def _bump(self, field_name: str) -> None:
|
|
s = self.stats
|
|
if s is None:
|
|
return
|
|
try:
|
|
setattr(s, field_name, getattr(s, field_name, 0) + 1)
|
|
if field_name in ("accepted", "rejected"):
|
|
s.processed = getattr(s, "processed", 0) + 1
|
|
except Exception:
|
|
pass
|