CWM (103): core mechanics, exhaustive edge cases, numba, exchange mechanics Replay (118): exhaustive verification, microstructure, trajectory Training (190): asset classification, phase0 extensive, pipeline, exhaustive DSL (102): v2 syntax, expanded, new features ASEx (33): validate-before-mutate, single-writer Planner (48): MCTS, alternatives, hooks Counterparties (19): 9 adversarial agent policies Clock (30): event-driven reactor BingX (28): venue adapter IPC (8): Zinc SHM Storage (9): ClickHouse Risk (4): hard invariants State (17): frozen dataclass invariants Integration: E2E, concurrency, sync/async seams, hypothesis, fuzz, adversarial
230 lines
7.4 KiB
Python
230 lines
7.4 KiB
Python
"""
|
|
T19 UV Clock Host tests — event dispatch, staleness, edge-triggered BarFire.
|
|
"""
|
|
import time
|
|
import pytest
|
|
from malkhut.clock.events import (
|
|
BarFire, EventProvenance, EventType, ScanEvent,
|
|
StaleInput, TickEvent, TimerEvent,
|
|
)
|
|
from malkhut.clock.host import UVClock
|
|
from malkhut.clock.staleness import StalenessWatchdog
|
|
from malkhut.clock.deadnode import DeadNodeReaper
|
|
|
|
|
|
class TestEventProvenance:
|
|
def test_fresh_when_recent(self):
|
|
prov = EventProvenance(
|
|
scan_number=1, scan_ts=time.time_ns(),
|
|
ingest_ts=time.time_ns(), source="live",
|
|
)
|
|
assert prov.is_fresh
|
|
|
|
def test_age_increases(self):
|
|
prov = EventProvenance(
|
|
scan_number=1, scan_ts=time.time_ns() - 10_000_000_000,
|
|
ingest_ts=time.time_ns() - 10_000_000_000, source="live",
|
|
)
|
|
assert prov.age_ns > 0
|
|
|
|
def test_stale_threshold(self):
|
|
prov = EventProvenance(
|
|
scan_number=1, scan_ts=0, ingest_ts=0, source="live",
|
|
)
|
|
threshold = prov.stale_threshold_ns(5_850_000_000)
|
|
assert threshold == 8_775_000_000 # 5.85s * 1.5
|
|
|
|
|
|
class TestScanEvent:
|
|
def test_scan_event_type(self):
|
|
e = ScanEvent(scan_number=1, symbol="BTCUSDT")
|
|
assert e.event_type == EventType.SCAN
|
|
|
|
def test_scan_event_has_provenance(self):
|
|
e = ScanEvent(scan_number=1, symbol="BTCUSDT")
|
|
assert e.provenance is not None
|
|
assert e.provenance.scan_number == 1
|
|
|
|
|
|
class TestTickEvent:
|
|
def test_tick_event_type(self):
|
|
e = TickEvent(symbol="BTCUSDT", price=50000.0, bid=49999.0, ask=50001.0)
|
|
assert e.event_type == EventType.TICK
|
|
|
|
def test_tick_event_has_provenance(self):
|
|
e = TickEvent(symbol="BTCUSDT", price=50000.0, bid=49999.0, ask=50001.0)
|
|
assert e.provenance is not None
|
|
|
|
|
|
class TestBarFire:
|
|
def test_bar_fire_type(self):
|
|
e = BarFire(scan_number=1, symbol="BTCUSDT")
|
|
assert e.event_type == EventType.BAR_FIRE
|
|
|
|
def test_bar_fire_has_provenance(self):
|
|
e = BarFire(scan_number=1, symbol="BTCUSDT")
|
|
assert e.provenance is not None
|
|
assert e.provenance.scan_number == 1
|
|
|
|
|
|
class TestTimerEvent:
|
|
def test_timer_type(self):
|
|
e = TimerEvent(timer_id="watchdog")
|
|
assert e.event_type == EventType.TIMER
|
|
|
|
|
|
class TestStaleInput:
|
|
def test_stale_type(self):
|
|
e = StaleInput(source="scan", last_scan_ts=0, current_ts=10_000_000_000)
|
|
assert e.event_type == EventType.STALE_INPUT
|
|
|
|
|
|
class TestUVClock:
|
|
def test_subscribe_and_dispatch(self):
|
|
clock = UVClock()
|
|
received = []
|
|
clock.subscribe(EventType.SCAN, lambda e: received.append(e))
|
|
clock.emit_scan(1, "BTCUSDT", {"price": 50000})
|
|
assert len(received) == 1
|
|
assert received[0].scan_number == 1
|
|
|
|
def test_bar_fire_on_scan_advance(self):
|
|
clock = UVClock()
|
|
bar_fires = []
|
|
clock.subscribe(EventType.BAR_FIRE, lambda e: bar_fires.append(e))
|
|
clock.emit_scan(1, "BTCUSDT", {})
|
|
clock.emit_scan(2, "BTCUSDT", {})
|
|
assert len(bar_fires) == 2
|
|
|
|
def test_no_bar_fire_on_duplicate_scan(self):
|
|
clock = UVClock()
|
|
bar_fires = []
|
|
clock.subscribe(EventType.BAR_FIRE, lambda e: bar_fires.append(e))
|
|
clock.emit_scan(1, "BTCUSDT", {})
|
|
clock.emit_scan(1, "BTCUSDT", {}) # duplicate
|
|
assert len(bar_fires) == 1 # edge-triggered, not level
|
|
|
|
def test_tick_dispatch(self):
|
|
clock = UVClock()
|
|
ticks = []
|
|
clock.subscribe(EventType.TICK, lambda e: ticks.append(e))
|
|
clock.emit_tick("BTCUSDT", 50000.0, 49999.0, 50001.0)
|
|
assert len(ticks) == 1
|
|
assert ticks[0].price == 50000.0
|
|
|
|
def test_timer_dispatch(self):
|
|
clock = UVClock()
|
|
timers = []
|
|
clock.subscribe(EventType.TIMER, lambda e: timers.append(e))
|
|
clock.emit_timer("watchdog", {"check": True})
|
|
assert len(timers) == 1
|
|
|
|
def test_unsubscribe(self):
|
|
clock = UVClock()
|
|
received = []
|
|
handler = lambda e: received.append(e)
|
|
clock.subscribe(EventType.SCAN, handler)
|
|
clock.emit_scan(1, "BTCUSDT", {})
|
|
assert len(received) == 1
|
|
clock.unsubscribe(EventType.SCAN, handler)
|
|
clock.emit_scan(2, "BTCUSDT", {})
|
|
assert len(received) == 1
|
|
|
|
def test_multiple_subscribers(self):
|
|
clock = UVClock()
|
|
r1, r2 = [], []
|
|
clock.subscribe(EventType.SCAN, lambda e: r1.append(e))
|
|
clock.subscribe(EventType.SCAN, lambda e: r2.append(e))
|
|
clock.emit_scan(1, "BTCUSDT", {})
|
|
assert len(r1) == 1
|
|
assert len(r2) == 1
|
|
|
|
def test_scan_number_tracking(self):
|
|
clock = UVClock()
|
|
clock.emit_scan(5, "BTCUSDT", {})
|
|
assert clock.last_scan_number == 5
|
|
|
|
def test_event_count(self):
|
|
clock = UVClock()
|
|
clock.emit_scan(1, "BTCUSDT", {}) # scan + barfire = 2
|
|
clock.emit_tick("BTCUSDT", 50000.0, 49999.0, 50001.0) # tick = 1
|
|
assert clock.event_count == 3 # scan + barfire + tick
|
|
|
|
def test_staleness_check(self):
|
|
clock = UVClock(scan_cadence_ns=100_000_000) # 100ms cadence for test
|
|
# No events yet — should be stale
|
|
stale = clock.check_staleness()
|
|
# May or may not be stale depending on timing, but should not crash
|
|
|
|
def test_provenance_forwarded(self):
|
|
clock = UVClock()
|
|
received = []
|
|
clock.subscribe(EventType.SCAN, lambda e: received.append(e))
|
|
clock.emit_scan(1, "BTCUSDT", {})
|
|
assert received[0].provenance.source == "live"
|
|
|
|
def test_replay_source(self):
|
|
clock = UVClock()
|
|
received = []
|
|
clock.subscribe(EventType.SCAN, lambda e: received.append(e))
|
|
clock.emit_scan(1, "BTCUSDT", {}, source="replay")
|
|
assert received[0].provenance.source == "replay"
|
|
|
|
|
|
class TestStalenessWatchdog:
|
|
def test_heartbeat_resets_staleness(self):
|
|
w = StalenessWatchdog(cadence_ns=100_000_000)
|
|
prov = EventProvenance(
|
|
scan_number=1, scan_ts=time.time_ns(),
|
|
ingest_ts=time.time_ns(), source="live",
|
|
)
|
|
w.heartbeat(prov)
|
|
assert not w.is_stale
|
|
|
|
def test_stale_after_threshold(self):
|
|
w = StalenessWatchdog(cadence_ns=100) # 100ns cadence
|
|
prov = EventProvenance(
|
|
scan_number=1, scan_ts=time.time_ns() - 1000,
|
|
ingest_ts=time.time_ns() - 1000, source="live",
|
|
)
|
|
w.heartbeat(prov)
|
|
# Wait for staleness
|
|
import time as _time
|
|
_time.sleep(0.001)
|
|
stale = w.check()
|
|
assert stale is not None
|
|
|
|
def test_stale_count_increments(self):
|
|
w = StalenessWatchdog(cadence_ns=1)
|
|
prov = EventProvenance(
|
|
scan_number=1, scan_ts=0, ingest_ts=0, source="live",
|
|
)
|
|
w.heartbeat(prov)
|
|
w.check()
|
|
assert w.stale_count >= 1
|
|
|
|
def test_age_ns(self):
|
|
w = StalenessWatchdog()
|
|
prov = EventProvenance(
|
|
scan_number=1, scan_ts=time.time_ns(),
|
|
ingest_ts=time.time_ns(), source="live",
|
|
)
|
|
w.heartbeat(prov)
|
|
assert w.age_ns >= 0
|
|
|
|
|
|
class TestDeadNodeReaper:
|
|
def test_reaper_creates(self):
|
|
r = DeadNodeReaper()
|
|
assert r.reaped_count == 0
|
|
|
|
def test_sweep_no_orphans(self):
|
|
r = DeadNodeReaper(shm_path="/tmp")
|
|
removed = r.sweep()
|
|
assert isinstance(removed, list)
|
|
|
|
def test_sweep_nonexistent_path(self):
|
|
r = DeadNodeReaper(shm_path="/nonexistent_path_xyz")
|
|
removed = r.sweep()
|
|
assert removed == []
|