"""Venue adapter contracts for DITAv2.""" from __future__ import annotations from dataclasses import dataclass, field from datetime import datetime from typing import Any, AsyncIterator, Dict, List, Optional, Protocol from .contracts import ( KernelCommandType, KernelIntent, KernelEventKind, TradeSide, VenueEvent, VenueEventStatus, VenueOrder, VenueOrderStatus, ) from .exchange_event import ExchangeEvent class VenuePostAckError(Exception): """Raised when submit fails AFTER the venue has already accepted the order. THE POINT OF NO RETURN. A plain exception out of submit()/submit_async() is ambiguous: it may mean "the venue never got the order" (safe to roll the FSM back to IDLE) or "the venue filled it and a line below blew up" (rolling back is catastrophic — the venue keeps the position while the kernel believes it is flat). On 2026-07-13 a post-ack TypeError took the first interpretation and orphaned 6 live SHORTs. Post-ack failure is NOT failure. It is UNKNOWN — and unknown is never flat. Callers MUST NOT synthesise a REJECTED event for this error: leave the slot in its working state and let the E-feed FILL / reconcile settle the truth. """ def __init__(self, message: str, *, receipt: Any = None, events: Optional[List[VenueEvent]] = None): super().__init__(message) self.receipt = receipt self.events = events or [] class VenueAdapter(Protocol): """Abstract venue adapter used by the kernel.""" def submit(self, intent: KernelIntent) -> List[VenueEvent]: ... def cancel(self, order: VenueOrder, *, reason: str = "") -> List[VenueEvent]: ... def open_orders(self) -> List[VenueOrder]: ... def open_positions(self) -> List[Dict[str, Any]]: ... def reconcile(self) -> List[VenueEvent]: ... # ------------------------------------------------------------------ # Phase 2 — stream seam (spec G3) # ------------------------------------------------------------------ async def subscribe(self) -> AsyncIterator[ExchangeEvent]: """ Yield ExchangeEvent instances in arrival order. Implementations must handle reconnection, keepalive, and 24h rotation internally. The iterator never terminates normally — callers cancel it on shutdown. Both the WS and poll-failover paths implement this interface so the kernel layer is source-agnostic. """ ... # pragma: no cover async def account_snapshot(self) -> ExchangeEvent: """ Return a single ACCOUNT_UPDATE + POSITION_UPDATE merged event by calling the exchange REST API. Used for gap-backfill on reconnect and as the poll-failover path. """ ... # pragma: no cover