From df40e62104ec6c66e061e98bda09e4a584892529 Mon Sep 17 00:00:00 2001 From: Codex Date: Sat, 11 Jul 2026 22:46:33 +0200 Subject: [PATCH] =?UTF-8?q?dita=5Fv2(zinc):=20account-truth=20plane=20API?= =?UTF-8?q?=20=E2=80=94=20publish/read/wait=5Fon/notify=5Faccount?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- prod/clean_arch/dita_v2/real_zinc_plane.py | 49 +++++++++++++++++- prod/clean_arch/dita_v2/zinc_plane.py | 60 ++++++++++++++++++---- 2 files changed, 99 insertions(+), 10 deletions(-) diff --git a/prod/clean_arch/dita_v2/real_zinc_plane.py b/prod/clean_arch/dita_v2/real_zinc_plane.py index 94c9f82..e1f54c0 100644 --- a/prod/clean_arch/dita_v2/real_zinc_plane.py +++ b/prod/clean_arch/dita_v2/real_zinc_plane.py @@ -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() diff --git a/prod/clean_arch/dita_v2/zinc_plane.py b/prod/clean_arch/dita_v2/zinc_plane.py index aea04e6..0cb6292 100644 --- a/prod/clean_arch/dita_v2/zinc_plane.py +++ b/prod/clean_arch/dita_v2/zinc_plane.py @@ -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