Files
sentiment-engine/MALKHUT/malkhut/clock/host.py
Codex 316d01079b malkhut(T7): UV Clock — event-driven reactor + staleness watchdog
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.
2026-07-11 10:36:30 +02:00

185 lines
5.9 KiB
Python

"""
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