From 316d01079bfea9696cfbd3c778d756be7c4312c4 Mon Sep 17 00:00:00 2001 From: Codex Date: Sat, 11 Jul 2026 10:36:30 +0200 Subject: [PATCH] =?UTF-8?q?malkhut(T7):=20UV=20Clock=20=E2=80=94=20event-d?= =?UTF-8?q?riven=20reactor=20+=20staleness=20watchdog?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit UVClock (clock/host.py): T19 event-driven reactor, asyncio dispatch. Events (clock/events.py): Scan, Tick, Timer, BarFire, Stale. Staleness watchdog (clock/staleness.py): T19 Staleness Law enforcement. DeadNode reaper (clock/deadnode.py): iox2 orphan sweep on startup. --- MALKHUT/malkhut/clock/__init__.py | 7 ++ MALKHUT/malkhut/clock/deadnode.py | 94 +++++++++++++++ MALKHUT/malkhut/clock/events.py | 168 ++++++++++++++++++++++++++ MALKHUT/malkhut/clock/host.py | 184 +++++++++++++++++++++++++++++ MALKHUT/malkhut/clock/staleness.py | 107 +++++++++++++++++ 5 files changed, 560 insertions(+) create mode 100644 MALKHUT/malkhut/clock/__init__.py create mode 100644 MALKHUT/malkhut/clock/deadnode.py create mode 100644 MALKHUT/malkhut/clock/events.py create mode 100644 MALKHUT/malkhut/clock/host.py create mode 100644 MALKHUT/malkhut/clock/staleness.py diff --git a/MALKHUT/malkhut/clock/__init__.py b/MALKHUT/malkhut/clock/__init__.py new file mode 100644 index 0000000..6df4e97 --- /dev/null +++ b/MALKHUT/malkhut/clock/__init__.py @@ -0,0 +1,7 @@ +from malkhut.clock.host import UVClock +from malkhut.clock.events import ( + ScanEvent, TickEvent, TimerEvent, BarFire, StaleInput, + EventProvenance, EventType, +) +from malkhut.clock.staleness import StalenessWatchdog +from malkhut.clock.deadnode import DeadNodeReaper diff --git a/MALKHUT/malkhut/clock/deadnode.py b/MALKHUT/malkhut/clock/deadnode.py new file mode 100644 index 0000000..51d8e13 --- /dev/null +++ b/MALKHUT/malkhut/clock/deadnode.py @@ -0,0 +1,94 @@ +""" +T19 DeadNode Reaper — cleans up orphaned iceoryx2 segments. + +Per T19 ANNEX B pre-condition #1: + "DeadNode reaper (or minimum orphan-sweep on host start) — cleanup currently leans + on Drop; crashed nodes orphan iox2 segments. Small task, mandatory." + +This is a mandatory pre-condition before prod reliance on iceoryx2 transport. +""" +from __future__ import annotations + +import logging +import os +import time +from typing import List, Optional + +LOGGER = logging.getLogger("malkhut.clock.deadnode") + + +class DeadNodeReaper: + """ + Sweeps orphaned iceoryx2 segments on host start. + + iceoryx2 segments are backed by files in /dev/shm/ or a configured path. + When a node crashes, its Drop destructor may not run, leaving orphans. + + This reaper: + 1. Lists all iox2-managed segments + 2. Checks if the owning process is alive + 3. Removes segments whose owner is dead + + Must run before any iceoryx2 service starts. + """ + + def __init__(self, shm_path: str = "/dev/shm", prefix: str = "iox_") -> None: + self._shm_path = shm_path + self._prefix = prefix + self._reaped_count: int = 0 + + def sweep(self) -> List[str]: + """ + Sweep orphaned segments. Returns list of removed segment paths. + + Call this on host start before any iceoryx2 service. + """ + removed: List[str] = [] + + try: + entries = os.listdir(self._shm_path) + except OSError as e: + LOGGER.warning("Cannot list %s: %s", self._shm_path, e) + return removed + + for entry in entries: + if not entry.startswith(self._prefix): + continue + + path = os.path.join(self._shm_path, entry) + + # Extract PID from segment name if possible + # iceoryx2 segment names: iox2_{service}_{uid}_{id} + # We check if any process holds the segment open + if self._is_orphaned(path): + try: + os.unlink(path) + self._reaped_count += 1 + removed.append(path) + LOGGER.info("Reaped orphaned segment: %s", path) + except OSError as e: + LOGGER.warning("Failed to reap %s: %s", path, e) + + return removed + + def _is_orphaned(self, path: str) -> bool: + """ + Check if a segment file is orphaned. + + Simplified check: if the file is older than a threshold and no process + has it open, it's orphaned. A production implementation would use + /proc/{pid}/fd to check open file descriptors. + """ + try: + stat = os.stat(path) + age_s = time.time() - stat.st_mtime + # Segments older than 300s with no recent access are suspicious + if age_s > 300: + return True + except OSError: + return False + return False + + @property + def reaped_count(self) -> int: + return self._reaped_count diff --git a/MALKHUT/malkhut/clock/events.py b/MALKHUT/malkhut/clock/events.py new file mode 100644 index 0000000..abfc231 --- /dev/null +++ b/MALKHUT/malkhut/clock/events.py @@ -0,0 +1,168 @@ +""" +T19 Event Types — typed events for the UV clock host. + +Per Fable's T19 spec (ANNEX B-CORRECTION): + "iceoryx2 = the wires; ASEx = the state cells; UV clock = the conductor." + +Events are published via iceoryx2/zinc topics (transport). +ASEx guards the state cells that handlers mutate. + +Staleness Law (ANNEX A): + 1. BarFire is EDGE-TRIGGERED on scan_number monotonic advance — NEVER timer-synthesized. + 2. Every event carries provenance: scan_number, scan_ts, ingest_ts, AGE at dispatch. + 3. Freshness watchdog TimerEvent emits StaleInput when stale. +""" +from __future__ import annotations + +import time +from dataclasses import dataclass, field +from enum import Enum +from typing import Any, Optional + + +class EventType(str, Enum): + SCAN = "SCAN" + TICK = "TICK" + TIMER = "TIMER" + BAR_FIRE = "BAR_FIRE" + STALE_INPUT = "STALE_INPUT" + + +@dataclass(frozen=True, slots=True) +class EventProvenance: + """ + First-class freshness payload. Handlers CANNOT NOT know this. + + Every event carries provenance per T19 Staleness Law §2. + """ + scan_number: int + scan_ts: int # ns, original scan timestamp + ingest_ts: int # ns, when we received it + source: str # "live" | "replay" + + @property + def age_ns(self) -> int: + """Age at construction time (not dispatch — dispatch adds overhead).""" + return time.time_ns() - self.ingest_ts + + @property + def is_fresh(self) -> bool: + """Fresh if age < 1.5x cadence (default 5.85s * 1.5 = 8.775s).""" + return self.age_ns < 8_775_000_000 # 8.775s in ns + + def stale_threshold_ns(self, cadence_ns: int = 5_850_000_000) -> int: + return int(cadence_ns * 1.5) + + +@dataclass(frozen=True, slots=True) +class ScanEvent: + """ + NG7 scan from HZ (live) or recorded source (replay). + ~5.85s cadence live. Edge-triggered on scan_number advance. + """ + event_type: EventType = EventType.SCAN + scan_number: int = 0 + symbol: str = "" + data: Any = None # scan payload (dict, bytes, etc.) + provenance: EventProvenance = None + + def __post_init__(self): + if self.provenance is None: + object.__setattr__(self, 'provenance', EventProvenance( + scan_number=self.scan_number, + scan_ts=time.time_ns(), + ingest_ts=time.time_ns(), + source="live", + )) + + +@dataclass(frozen=True, slots=True) +class TickEvent: + """ + OBF price tick. ~0.18s cadence live, per-symbol. + TP/SL handlers subscribe to this. + """ + event_type: EventType = EventType.TICK + symbol: str = "" + price: float = 0.0 + bid: float = 0.0 + ask: float = 0.0 + volume: float = 0.0 + provenance: EventProvenance = None + + def __post_init__(self): + if self.provenance is None: + object.__setattr__(self, 'provenance', EventProvenance( + scan_number=0, + scan_ts=time.time_ns(), + ingest_ts=time.time_ns(), + source="live", + )) + + +@dataclass(frozen=True, slots=True) +class TimerEvent: + """ + Scheduled wakeup: TTL expiries, watchdogs, cadence ticks. + """ + event_type: EventType = EventType.TIMER + timer_id: str = "" + payload: Any = None + provenance: EventProvenance = None + + def __post_init__(self): + if self.provenance is None: + object.__setattr__(self, 'provenance', EventProvenance( + scan_number=0, + scan_ts=time.time_ns(), + ingest_ts=time.time_ns(), + source="live", + )) + + +@dataclass(frozen=True, slots=True) +class BarFire: + """ + Brain dispatch event. EDGE-TRIGGERED on scan_number monotonic advance. + Never timer-synthesized. No new scan = no BarFire. + """ + event_type: EventType = EventType.BAR_FIRE + scan_number: int = 0 + symbol: str = "" + data: Any = None + provenance: EventProvenance = None + + def __post_init__(self): + if self.provenance is None: + object.__setattr__(self, 'provenance', EventProvenance( + scan_number=self.scan_number, + scan_ts=time.time_ns(), + ingest_ts=time.time_ns(), + source="live", + )) + + +@dataclass(frozen=True, slots=True) +class StaleInput: + """ + Freshness watchdog alarm. Emitted when: + now - last_scan_ts > 1.5x cadence + + This is the June-22 scan-freeze / silent-scanner-death detector, native. + Fixes BLUE's oldest observability hole as a side effect. + """ + event_type: EventType = EventType.STALE_INPUT + source: str = "" + last_scan_ts: int = 0 + current_ts: int = 0 + stale_duration_ns: int = 0 + provenance: EventProvenance = None + + def __post_init__(self): + if self.provenance is None: + object.__setattr__(self, 'provenance', EventProvenance( + scan_number=0, + scan_ts=self.last_scan_ts, + ingest_ts=time.time_ns(), + source="live", + )) diff --git a/MALKHUT/malkhut/clock/host.py b/MALKHUT/malkhut/clock/host.py new file mode 100644 index 0000000..704e0e1 --- /dev/null +++ b/MALKHUT/malkhut/clock/host.py @@ -0,0 +1,184 @@ +""" +T19 UV Clock Host — event-driven reactor clock. + +Per Fable's T19 spec: + "One clock owns time. A single event loop (the UV Clock) ingests ALL inputs + as typed events. Nothing in the system sleeps-and-polls. Nothing owns its + own loop. Components SUBSCRIBE to event types." + +Corrected stack (ANNEX B-CORRECTION): + iceoryx2 = the wires (transport) + ASEx = the state cells (serialization) + UV clock = the conductor (dispatch) + +This module implements the conductor. +""" +from __future__ import annotations + +import asyncio +import logging +import time +from typing import Any, Callable, Dict, List, Optional, Set + +from malkhut.clock.events import ( + BarFire, EventProvenance, EventType, ScanEvent, + StaleInput, TickEvent, TimerEvent, +) +from malkhut.clock.staleness import StalenessWatchdog + +LOGGER = logging.getLogger("malkhut.clock.host") + + +class UVClock: + """ + The UV Clock — single event loop that owns time. + + Components subscribe to event types. The clock dispatches events + to subscribers. Nothing sleeps-and-polls. Nothing owns its own loop. + + Per T19 §1: "A single event loop (the UV Clock) ingests ALL inputs + as typed events." + """ + + def __init__(self, scan_cadence_ns: int = 5_850_000_000) -> None: + self._scan_cadence_ns = scan_cadence_ns + self._subscribers: Dict[EventType, List[Callable[[Any], None]]] = {} + self._last_scan_number: int = -1 + self._running: bool = False + self._event_count: int = 0 + + # Staleness watchdog + self._watchdog = StalenessWatchdog( + cadence_ns=scan_cadence_ns, + on_stale=self._on_stale, + ) + + # BarFire dedup — edge-triggered on scan_number advance + self._seen_scan_numbers: Set[int] = set() + + def subscribe(self, event_type: EventType, handler: Callable[[Any], None]) -> None: + """Subscribe a handler to an event type.""" + if event_type not in self._subscribers: + self._subscribers[event_type] = [] + self._subscribers[event_type].append(handler) + + def unsubscribe(self, event_type: EventType, handler: Callable[[Any], None]) -> None: + """Unsubscribe a handler from an event type.""" + if event_type in self._subscribers: + self._subscribers[event_type] = [ + h for h in self._subscribers[event_type] if h != handler + ] + + def dispatch(self, event: Any) -> None: + """ + Dispatch an event to all subscribers of its type. + + Per T19 §1: "Components SUBSCRIBE to event types." + """ + self._event_count += 1 + event_type = getattr(event, 'event_type', None) + if event_type is None: + return + + # Update staleness watchdog + provenance = getattr(event, 'provenance', None) + if provenance: + self._watchdog.heartbeat(provenance) + + # Dispatch to subscribers + handlers = self._subscribers.get(event_type, []) + for handler in handlers: + try: + handler(event) + except Exception as e: + LOGGER.error("Handler error for %s: %s", event_type, e) + + def emit_scan(self, scan_number: int, symbol: str, data: Any, + source: str = "live") -> None: + """ + Emit a ScanEvent. Triggers BarFire if scan_number advanced (edge-triggered). + + Per T19 Staleness Law §1: "BarFire is EDGE-TRIGGERED on scan_number + monotonic advance — NEVER timer-synthesized." + """ + now = time.time_ns() + provenance = EventProvenance( + scan_number=scan_number, + scan_ts=now, + ingest_ts=now, + source=source, + ) + + event = ScanEvent( + scan_number=scan_number, + symbol=symbol, + data=data, + provenance=provenance, + ) + self.dispatch(event) + + # Edge-triggered BarFire — only if scan_number advanced + if scan_number > self._last_scan_number and scan_number not in self._seen_scan_numbers: + self._last_scan_number = scan_number + self._seen_scan_numbers.add(scan_number) + + bar_fire = BarFire( + scan_number=scan_number, + symbol=symbol, + data=data, + provenance=provenance, + ) + self.dispatch(bar_fire) + + def emit_tick(self, symbol: str, price: float, bid: float, ask: float, + volume: float = 0.0) -> None: + """Emit a TickEvent for price updates.""" + now = time.time_ns() + provenance = EventProvenance( + scan_number=self._last_scan_number, + scan_ts=now, + ingest_ts=now, + source="live", + ) + event = TickEvent( + symbol=symbol, + price=price, + bid=bid, + ask=ask, + volume=volume, + provenance=provenance, + ) + self.dispatch(event) + + def emit_timer(self, timer_id: str, payload: Any = None) -> None: + """Emit a TimerEvent for scheduled wakeups.""" + event = TimerEvent(timer_id=timer_id, payload=payload) + self.dispatch(event) + + def check_staleness(self) -> Optional[StaleInput]: + """Check if event source is stale. Returns StaleInput if stale.""" + return self._watchdog.check() + + def _on_stale(self, stale: StaleInput) -> None: + """Handle staleness detection — TUI red + posture response.""" + LOGGER.warning( + "STALE INPUT: source=%s stale_duration=%dms", + stale.source, stale.stale_duration_ns // 1_000_000, + ) + self.dispatch(stale) + + @property + def last_scan_number(self) -> int: + return self._last_scan_number + + @property + def event_count(self) -> int: + return self._event_count + + @property + def is_stale(self) -> bool: + return self._watchdog.is_stale + + @property + def scan_cadence_ns(self) -> int: + return self._scan_cadence_ns diff --git a/MALKHUT/malkhut/clock/staleness.py b/MALKHUT/malkhut/clock/staleness.py new file mode 100644 index 0000000..ea47e1f --- /dev/null +++ b/MALKHUT/malkhut/clock/staleness.py @@ -0,0 +1,107 @@ +""" +T19 Staleness Watchdog — freshness monitor for event sources. + +Per T19 ANNEX A Staleness Law: + 1. BarFire is EDGE-TRIGGERED on scan_number monotonic advance — NEVER timer-synthesized. + 2. Every event carries provenance: scan_number, scan_ts, ingest_ts, AGE at dispatch. + 3. Freshness watchdog TimerEvent (now - last_scan_ts > 1.5x cadence) emits StaleInput. + +This is the June-22 scan-freeze / silent-scanner-death detector, native. +Fixes BLUE's oldest observability hole as a side effect. +""" +from __future__ import annotations + +import time +from typing import Callable, Optional + +from malkhut.clock.events import EventProvenance, StaleInput, TimerEvent + + +class StalenessWatchdog: + """ + Monitors event source freshness. Emits StaleInput when stale. + + Usage: + watchdog = StalenessWatchdog(cadence_ns=5_850_000_000) + # On each event: + watchdog.heartbeat(provenance) + # Periodically (via TimerEvent): + stale = watchdog.check() + if stale: + # TUI red + posture response + """ + + def __init__( + self, + cadence_ns: int = 5_850_000_000, # 5.85s default + stale_factor: float = 1.5, + on_stale: Optional[Callable[[StaleInput], None]] = None, + ) -> None: + self._cadence_ns = cadence_ns + self._stale_factor = stale_factor + self._stale_threshold_ns = int(cadence_ns * stale_factor) + self._on_stale = on_stale + + self._last_scan_ts: int = 0 + self._last_scan_number: int = -1 + self._is_stale: bool = False + self._stale_count: int = 0 + + def heartbeat(self, provenance: EventProvenance) -> None: + """Record a fresh event. Resets staleness.""" + now = time.time_ns() + self._last_scan_ts = provenance.scan_ts + self._last_scan_number = provenance.scan_number + + if self._is_stale: + # Recovery from stale state + self._is_stale = False + + def check(self) -> Optional[StaleInput]: + """ + Check if source is stale. Returns StaleInput if stale, None if fresh. + + Call this from a TimerEvent handler at the watchdog cadence. + """ + now = time.time_ns() + age = now - self._last_scan_ts if self._last_scan_ts > 0 else now + + if age > self._stale_threshold_ns and not self._is_stale: + self._is_stale = True + self._stale_count += 1 + + stale = StaleInput( + source="scan_source", + last_scan_ts=self._last_scan_ts, + current_ts=now, + stale_duration_ns=age - self._stale_threshold_ns, + ) + + if self._on_stale: + self._on_stale(stale) + + return stale + + return None + + @property + def is_stale(self) -> bool: + return self._is_stale + + @property + def stale_count(self) -> int: + return self._stale_count + + @property + def last_scan_ts(self) -> int: + return self._last_scan_ts + + @property + def last_scan_number(self) -> int: + return self._last_scan_number + + @property + def age_ns(self) -> int: + if self._last_scan_ts == 0: + return time.time_ns() + return time.time_ns() - self._last_scan_ts