dita_v2(zinc): account-truth plane API — publish/read/wait_on/notify_account
Formalizes the CEX-010 E-feed account plane (protocol + InMemoryZincPlane + real plane). This code flew 33h on UV-PRIME flight-1 via the vendored copy (T0-DEVIATION mid-edit vendoring) but was never committed upstream — the 2026-07-11 clean vendor_sync regressed it and broke the E-feed at boot (AttributeError publish_account). Committing the flight-proven bytes closes the vendor-law gap. Authorship: pre-existing WIP in /mnt working tree (T0-era, likely codex), committed verbatim by Fable for vendor integrity.
This commit is contained in:
@@ -16,7 +16,7 @@ import struct
|
||||
import sys
|
||||
import threading
|
||||
|
||||
from .contracts import KernelIntent, TradeSide, TradeSlot, TradeStage, VenueOrder, VenueOrderStatus, VenueTelemetrySnapshot
|
||||
from .contracts import AccountStateSnapshot, KernelIntent, TradeSide, TradeSlot, TradeStage, VenueOrder, VenueOrderStatus, VenueTelemetrySnapshot
|
||||
from .control import KernelControlSnapshot
|
||||
|
||||
_ZINC_ADAPTER_PATH = Path(__file__).resolve().parents[3] / "zinc" / "adapters" / "python"
|
||||
@@ -185,6 +185,7 @@ class RealZincPlane:
|
||||
self.state_name = f"{base}_state"
|
||||
self.control_name = f"{base}_control"
|
||||
self.venue_name = f"{base}_venue"
|
||||
self.account_name = f"{base}_account"
|
||||
self._intent_seq = 0
|
||||
self._state_seq = 0
|
||||
self._control_seq = 0
|
||||
@@ -195,12 +196,16 @@ class RealZincPlane:
|
||||
self._intent_cache: List[Dict[str, Any]] = []
|
||||
self._control_cache = KernelControlSnapshot()
|
||||
self._venue_cache = VenueTelemetrySnapshot()
|
||||
self._account_cache = AccountStateSnapshot()
|
||||
self._account_seq = 0
|
||||
if create:
|
||||
self.intent_region = SharedRegion.create(self.intent_name, intent_capacity)
|
||||
self.state_region = SharedRegion.create(self.state_name, state_capacity)
|
||||
self.control_region = SharedRegion.create(self.control_name, control_capacity)
|
||||
self.venue_region = SharedRegion.create(self.venue_name, control_capacity)
|
||||
self.account_region = SharedRegion.create(self.account_name, control_capacity)
|
||||
self._write_region(self.control_region, self._control_seq, {"control": self._control_cache.as_dict()})
|
||||
self._write_region(self.account_region, self._account_seq, {"account": self._account_cache.as_dict()})
|
||||
self._write_region(
|
||||
self.state_region,
|
||||
self._state_seq,
|
||||
@@ -213,10 +218,12 @@ class RealZincPlane:
|
||||
self.state_region = SharedRegion.open(self.state_name)
|
||||
self.control_region = SharedRegion.open(self.control_name)
|
||||
self.venue_region = SharedRegion.open(self.venue_name)
|
||||
self.account_region = SharedRegion.open(self.account_name)
|
||||
control_payload = _decode_packet(self.control_region.as_buffer())
|
||||
state_payload = _decode_packet(self.state_region.as_buffer())
|
||||
intent_payload = _decode_packet(self.intent_region.as_buffer())
|
||||
venue_payload = _decode_packet(self.venue_region.as_buffer())
|
||||
account_payload = _decode_packet(self.account_region.as_buffer())
|
||||
if isinstance(control_payload.get("control"), dict):
|
||||
self._control_cache = KernelControlSnapshot(**control_payload["control"])
|
||||
if isinstance(state_payload.get("slots"), list):
|
||||
@@ -228,12 +235,15 @@ class RealZincPlane:
|
||||
self._intent_cache = list(intent_payload["items"])
|
||||
if isinstance(venue_payload.get("venue"), dict):
|
||||
self._venue_cache = _venue_from_payload(venue_payload["venue"])
|
||||
if isinstance(account_payload.get("account"), dict):
|
||||
self._account_cache = _account_from_payload(account_payload["account"])
|
||||
|
||||
def close(self) -> None:
|
||||
self.intent_region.close()
|
||||
self.state_region.close()
|
||||
self.control_region.close()
|
||||
self.venue_region.close()
|
||||
self.account_region.close()
|
||||
|
||||
def publish_intent(self, intent: KernelIntent) -> None:
|
||||
with self._lock:
|
||||
@@ -315,6 +325,26 @@ class RealZincPlane:
|
||||
def notify_venue(self) -> None:
|
||||
self.venue_region.notify()
|
||||
|
||||
def publish_account(self, snapshot: AccountStateSnapshot) -> None:
|
||||
with self._lock:
|
||||
self._account_seq += 1
|
||||
self._account_cache = snapshot
|
||||
self._write_region(self.account_region, self._account_seq, {"account": snapshot.as_dict()})
|
||||
|
||||
def read_account(self) -> AccountStateSnapshot:
|
||||
payload = _decode_packet(self.account_region.as_buffer())
|
||||
acct = payload.get("account") if isinstance(payload, dict) else None
|
||||
if not isinstance(acct, dict):
|
||||
return self._account_cache
|
||||
self._account_cache = _account_from_payload(acct)
|
||||
return self._account_cache
|
||||
|
||||
def wait_on_account(self, timeout_ms: int = 1000) -> bool:
|
||||
return bool(self.account_region.wait(timeout_ms))
|
||||
|
||||
def notify_account(self) -> None:
|
||||
self.account_region.notify()
|
||||
|
||||
def wait_on_intent(self, timeout_ms: int = 1000) -> bool:
|
||||
return bool(self.intent_region.wait(timeout_ms))
|
||||
|
||||
@@ -330,3 +360,20 @@ class RealZincPlane:
|
||||
view[:] = b"\x00" * len(view)
|
||||
view[: len(packet)] = packet
|
||||
region.notify()
|
||||
|
||||
|
||||
def _account_from_payload(payload: Dict[str, Any]) -> "AccountStateSnapshot":
|
||||
"""Decode raw payload into AccountStateSnapshot. Never raises — returns default."""
|
||||
from .contracts import AccountStateSnapshot
|
||||
try:
|
||||
return AccountStateSnapshot(
|
||||
wallet_balance=float(payload.get("wallet_balance", 0.0)),
|
||||
available_margin=float(payload.get("available_margin", 0.0)),
|
||||
used_margin=float(payload.get("used_margin", 0.0)),
|
||||
event_seq=int(payload.get("event_seq", 0) or 0),
|
||||
mono_ns=int(payload.get("mono_ns", 0) or 0),
|
||||
e_live=bool(payload.get("e_live", False)),
|
||||
reconcile_ok=bool(payload.get("reconcile_ok", False)),
|
||||
)
|
||||
except (TypeError, ValueError, OverflowError):
|
||||
return AccountStateSnapshot()
|
||||
|
||||
@@ -12,7 +12,10 @@ from typing import Any, Dict, Iterable, List, Mapping, Optional, Protocol
|
||||
import threading
|
||||
import time
|
||||
|
||||
from .contracts import KernelIntent, TradeSlot, VenueTelemetrySnapshot
|
||||
from .contracts import AccountStateSnapshot, KernelIntent, TradeSlot, VenueTelemetrySnapshot
|
||||
from .double_buffer import DoubleBufferRegion, SingleSlotRegion
|
||||
from .control import KernelControlSnapshot
|
||||
from .control import KernelControlSnapshot
|
||||
from .control import KernelControlSnapshot
|
||||
|
||||
|
||||
@@ -64,23 +67,44 @@ class ZincPlane(Protocol):
|
||||
def notify_venue(self) -> None:
|
||||
...
|
||||
|
||||
def publish_account(self, snapshot: "AccountStateSnapshot") -> None:
|
||||
...
|
||||
|
||||
def read_account(self) -> "AccountStateSnapshot":
|
||||
...
|
||||
|
||||
def wait_on_account(self, timeout_ms: int = 1000) -> bool:
|
||||
...
|
||||
|
||||
def notify_account(self) -> None:
|
||||
...
|
||||
|
||||
|
||||
@dataclass
|
||||
class InMemoryZincPlane:
|
||||
"""Simple in-memory Zinc lookalike for Python prototype tests."""
|
||||
"""In-memory Zinc lookalike with double-buffer regions — torn-read impossible."""
|
||||
|
||||
intent_region: List[KernelIntent] = field(default_factory=list)
|
||||
state_region: Dict[int, TradeSlot] = field(default_factory=dict)
|
||||
control_region: Optional[KernelControlSnapshot] = None
|
||||
venue_region: VenueTelemetrySnapshot = field(default_factory=VenueTelemetrySnapshot)
|
||||
control_region: SingleSlotRegion[Optional[KernelControlSnapshot]] = field(
|
||||
default_factory=lambda: SingleSlotRegion(lambda: None)
|
||||
)
|
||||
venue_region: SingleSlotRegion[VenueTelemetrySnapshot] = field(
|
||||
default_factory=lambda: SingleSlotRegion(lambda: VenueTelemetrySnapshot())
|
||||
)
|
||||
account_region: SingleSlotRegion[AccountStateSnapshot] = field(
|
||||
default_factory=lambda: SingleSlotRegion(lambda: AccountStateSnapshot())
|
||||
)
|
||||
_intent_seq: int = field(default=0, init=False, repr=False)
|
||||
_state_seq: int = field(default=0, init=False, repr=False)
|
||||
_control_seq: int = field(default=0, init=False, repr=False)
|
||||
_venue_seq: int = field(default=0, init=False, repr=False)
|
||||
_account_seq: int = field(default=0, init=False, repr=False)
|
||||
_intent_observed_seq: int = field(default=0, init=False, repr=False)
|
||||
_state_observed_seq: int = field(default=0, init=False, repr=False)
|
||||
_control_observed_seq: int = field(default=0, init=False, repr=False)
|
||||
_venue_observed_seq: int = field(default=0, init=False, repr=False)
|
||||
_account_observed_seq: int = field(default=0, init=False, repr=False)
|
||||
_signal: threading.Condition = field(default_factory=threading.Condition, init=False, repr=False)
|
||||
|
||||
def publish_intent(self, intent: KernelIntent) -> None:
|
||||
@@ -100,14 +124,15 @@ class InMemoryZincPlane:
|
||||
|
||||
def update_control(self, control: KernelControlSnapshot) -> None:
|
||||
with self._signal:
|
||||
self.control_region = control
|
||||
self.control_region.write(control)
|
||||
self._control_seq += 1
|
||||
self._signal.notify_all()
|
||||
|
||||
def read_control(self) -> KernelControlSnapshot:
|
||||
if self.control_region is None:
|
||||
val = self.control_region.read()
|
||||
if val is None:
|
||||
return KernelControlSnapshot()
|
||||
return self.control_region
|
||||
return val
|
||||
|
||||
def wait_on_intent(self, timeout_ms: int = 1000) -> bool:
|
||||
return self._wait_for_change("_intent_seq", "_intent_observed_seq", timeout_ms)
|
||||
@@ -135,12 +160,12 @@ class InMemoryZincPlane:
|
||||
|
||||
def publish_venue(self, telemetry: VenueTelemetrySnapshot) -> None:
|
||||
with self._signal:
|
||||
self.venue_region = telemetry
|
||||
self.venue_region.write(telemetry)
|
||||
self._venue_seq += 1
|
||||
self._signal.notify_all()
|
||||
|
||||
def read_venue(self) -> VenueTelemetrySnapshot:
|
||||
return self.venue_region
|
||||
return self.venue_region.read()
|
||||
|
||||
def wait_on_venue(self, timeout_ms: int = 1000) -> bool:
|
||||
return self._wait_for_change("_venue_seq", "_venue_observed_seq", timeout_ms)
|
||||
@@ -150,6 +175,23 @@ class InMemoryZincPlane:
|
||||
self._venue_seq += 1
|
||||
self._signal.notify_all()
|
||||
|
||||
def publish_account(self, snapshot: AccountStateSnapshot) -> None:
|
||||
with self._signal:
|
||||
self.account_region.write(snapshot)
|
||||
self._account_seq += 1
|
||||
self._signal.notify_all()
|
||||
|
||||
def read_account(self) -> AccountStateSnapshot:
|
||||
return self.account_region.read()
|
||||
|
||||
def wait_on_account(self, timeout_ms: int = 1000) -> bool:
|
||||
return self._wait_for_change("_account_seq", "_account_observed_seq", timeout_ms)
|
||||
|
||||
def notify_account(self) -> None:
|
||||
with self._signal:
|
||||
self._account_seq += 1
|
||||
self._signal.notify_all()
|
||||
|
||||
def _wait_for_change(self, seq_attr: str, observed_attr: str, timeout_ms: int) -> bool:
|
||||
timeout_s = None if timeout_ms is None or timeout_ms < 0 else max(0.0, timeout_ms / 1000.0)
|
||||
deadline = None if timeout_s is None else time.monotonic() + timeout_s
|
||||
|
||||
Reference in New Issue
Block a user