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.
This commit is contained in:
7
MALKHUT/malkhut/clock/__init__.py
Normal file
7
MALKHUT/malkhut/clock/__init__.py
Normal file
@@ -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
|
||||||
94
MALKHUT/malkhut/clock/deadnode.py
Normal file
94
MALKHUT/malkhut/clock/deadnode.py
Normal file
@@ -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
|
||||||
168
MALKHUT/malkhut/clock/events.py
Normal file
168
MALKHUT/malkhut/clock/events.py
Normal file
@@ -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",
|
||||||
|
))
|
||||||
184
MALKHUT/malkhut/clock/host.py
Normal file
184
MALKHUT/malkhut/clock/host.py
Normal file
@@ -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
|
||||||
107
MALKHUT/malkhut/clock/staleness.py
Normal file
107
MALKHUT/malkhut/clock/staleness.py
Normal file
@@ -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
|
||||||
Reference in New Issue
Block a user