Files
hms-mediaengine/tests/unit/test_cluster_routing.py
T
HMS MediaEngine Agent cc119d3751 Phase 2: State-Sync-Fundament – State-Store, Gruppen, Clock, Aktivierung
- ProjectStateStore (§6.4, §24.2): monotone State-/Projekt-Revisionen,
  Snapshot nach Verbindung, Delta mit neu/geaendert/geloescht,
  Szenenabruf als Zielzustand (Uebergang macht der Renderer, §18.1),
  Projekt-/Livezustand strikt getrennt
- GroupRouter (§6.3, §10.1): Ziele All/Node/Output/ServerGroup,
  Regeln selected/tag_query/all, Commit-Vorschau (§17.5)
- ClockEstimator (§6.4): RTT-Min-Filter (<=2x Min gegen Jitter),
  Offset-/Drift-Schaetzung (Drift erst ab 1 s Fenster belastbar),
  Showzeit -> lokale Node-Zeit
- ActivationCoordinator (§6.5): zeitgestempelte Preset-Aktivierung,
  200 ms Vorlauf, Arm/Execute/Ack-Kette, FAILED bei fehlender
  Arm-Bestaetigung (kein stilles Ueberbruecken)
- 35 neue Unit-Tests; Gesamtsuite 293 gruen, Ruff gruen
2026-09-11 01:31:26 +02:00

259 lines
9.3 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Unit-Tests Zielrouting, Clock-Sync, zeitgestempelte Aktivierung
(PLAN.md §6.3, §6.4, §6.5)."""
from __future__ import annotations
import pytest
from hms_cluster import (
ActivationCoordinator,
ArmState,
ClockEstimator,
ClockSample,
GroupRouter,
GroupRule,
ServerGroup,
TargetKind,
)
# ---------- GroupRouter (§6.3 Bedienmodelle) ----------
@pytest.fixture()
def router() -> GroupRouter:
r = GroupRouter()
r.set_node_tags("node-a", frozenset({"stage-left"}))
r.set_node_tags("node-b", frozenset({"stage-right"}))
r.set_node_tags("node-c", frozenset({"stage-left", "stage-right"}))
r.set_node_outputs("node-a", ["out-1"])
r.set_node_outputs("node-c", ["out-2"])
return r
def test_resolve_all_targets(router: GroupRouter) -> None:
assert router.resolve(TargetKind.ALL) == frozenset({"node-a", "node-b", "node-c"})
def test_resolve_single_node(router: GroupRouter) -> None:
assert router.resolve(TargetKind.NODE, "node-a") == frozenset({"node-a"})
assert router.resolve(TargetKind.NODE, "unbekannt") == frozenset()
def test_resolve_by_output(router: GroupRouter) -> None:
assert router.resolve(TargetKind.OUTPUT, "out-2") == frozenset({"node-c"})
assert router.resolve(TargetKind.OUTPUT, "out-9") == frozenset()
def test_group_rule_selected(router: GroupRouter) -> None:
router.upsert_group(
ServerGroup(
id="g1",
name="Links",
rule=GroupRule.SELECTED,
node_ids=frozenset({"node-a", "unbekannt"}),
)
)
# unbekannte Mitglieder werden still gefiltert, bekannte bleiben
assert router.resolve(TargetKind.SERVER_GROUP, "g1") == frozenset({"node-a"})
def test_group_rule_tag_query(router: GroupRouter) -> None:
router.upsert_group(
ServerGroup(
id="g2",
name="Beide Bühnen",
rule=GroupRule.TAG_QUERY,
tags=frozenset({"stage-left"}),
)
)
assert router.resolve(TargetKind.SERVER_GROUP, "g2") == frozenset(
{"node-a", "node-c"}
)
def test_group_rule_all(router: GroupRouter) -> None:
router.upsert_group(ServerGroup(id="g3", name="Alle", rule=GroupRule.ALL))
assert router.resolve(TargetKind.SERVER_GROUP, "g3") == router.resolve(TargetKind.ALL)
def test_preview_matches_resolve(router: GroupRouter) -> None:
"""§17.5: Commit-Vorschau zeigt dieselben Ziele wie der Versand."""
router.upsert_group(
ServerGroup(id="g4", name="X", rule=GroupRule.SELECTED, node_ids=frozenset({"node-b"}))
)
assert router.preview_targets(TargetKind.SERVER_GROUP, "g4") == router.resolve(
TargetKind.SERVER_GROUP, "g4"
)
def test_unknown_group_resolves_empty(router: GroupRouter) -> None:
assert router.resolve(TargetKind.SERVER_GROUP, "gibts-nicht") == frozenset()
def test_remove_group(router: GroupRouter) -> None:
router.upsert_group(ServerGroup(id="g", name="Weg"))
router.remove_group("g")
assert router.get_group("g") is None
# ---------- ClockEstimator (§6.4 Clock Sync) ----------
def _sample(t0: int, rtt: int, offset: int) -> ClockSample:
"""Probe mit definiertem echtem Offset: node_time = t0 + rtt/2 + offset."""
node_time = t0 + rtt // 2 + offset
return ClockSample(t0_ns=t0, t1_ns=t0 + rtt, node_time_ns=node_time)
def test_clock_offset_estimated_from_low_rtt_samples() -> None:
"""Min-Filter: Proben mit Rauschen (hohe RTT) verschieben den Schätzer
nicht; die niedrigste RTT dominiert (§6.4)."""
est = ClockEstimator()
# echter Offset: +5 ms; einige Proben mit Jitter
est.feed(_sample(0, 1_000_000, 5_000_000)) # 1 ms RTT
est.feed(_sample(1_000_000, 50_000_000, 30_000_000)) # 50 ms RTT, Jitter
est.feed(_sample(2_000_000, 2_000_000, 5_500_000))
est.feed(_sample(3_000_000, 1_500_000, 4_800_000))
offset = est.offset_ns
assert offset is not None
assert 4_000_000 < offset < 6_000_000 # nahe am echten 5 ms
def test_clock_best_rtt_reported() -> None:
est = ClockEstimator()
est.feed(_sample(0, 20_000_000, 0))
est.feed(_sample(1, 5_000_000, 0))
assert est.rtt_ns == 5_000_000
def test_clock_no_data_returns_none() -> None:
est = ClockEstimator()
assert est.offset_ns is None
assert est.rtt_ns is None
assert est.drift_ppm is None
def test_clock_drift_estimated_over_time() -> None:
"""Drift: Offset wächst um 100 µs pro Sekunde = 100 ppm (§6.4)."""
est = ClockEstimator()
est.feed(_sample(t0=0, rtt=1_000_000, offset=0))
# 2 s später: 200 µs mehr Offset (200_000 ns) → 100 µs/s = 100 ppm
est.feed(_sample(t0=2_000_000_000, rtt=1_000_000, offset=200_000))
drift = est.drift_ppm
assert drift is not None
assert 80.0 <= drift <= 120.0 # ~100 ppm
def test_clock_drift_none_below_one_second_window() -> None:
"""Zu kurzes Fenster: Drift ist nicht belastbar (§6.4 Grenze)."""
est = ClockEstimator()
est.feed(_sample(0, 1_000_000, 0))
est.feed(_sample(100_000_000, 1_000_000, 50)) # nur 100 ms Abstand
assert est.drift_ppm is None
def test_clock_rejects_negative_rtt() -> None:
with pytest.raises(ValueError, match="negative RTT"):
ClockEstimator().feed(ClockSample(t0_ns=10, t1_ns=5, node_time_ns=0))
def test_clock_maps_show_time_to_node_time() -> None:
"""§29.1: Showzeit → lokale Monotonic über den Offset."""
est = ClockEstimator()
est.feed(_sample(0, 1_000_000, offset=10_000_000)) # +10 ms
mapped = est.map_show_time(1_000_000_000)
assert mapped is not None
assert mapped == 1_010_000_000 # Showzeit + Offset
# ohne Proben: keine Abbildung möglich
assert ClockEstimator().map_show_time(0) is None
def test_clock_sample_window_bounded() -> None:
"""Kein unbeschränkter Zustand (§33): Fenster bleibt begrenzt."""
est = ClockEstimator(max_samples=4)
for i in range(10):
est.feed(_sample(i * 1_000_000, 1_000_000, 0))
assert len(est._samples) <= 4
# ---------- ActivationCoordinator (§6.5) ----------
@pytest.fixture()
def two_node_setup():
router = GroupRouter()
router.set_node_tags("node-a", frozenset({"x"}))
router.set_node_tags("node-b", frozenset({"x"}))
coord = ActivationCoordinator(router)
return coord, router
def test_schedule_plans_with_lead_time(two_node_setup) -> None:
coord, _ = two_node_setup
planned = coord.schedule("scene-1", now_show_ns=1_000_000_000, lead_ns=200_000_000)
assert planned.execute_at_show_ns == 1_200_000_000 # §6.4: 100300 ms Vorlauf
assert planned.state is ArmState.CREATED
def test_due_requires_all_arms(two_node_setup) -> None:
"""§6.5: Execute nur, wenn ALLE Ziel-Nodes armed; sonst FAILED."""
coord, _ = two_node_setup
planned = coord.schedule("scene-1", now_show_ns=0, lead_ns=100)
coord.acknowledge_arm(planned.command_id, "node-a")
# Showzeit erreicht, aber node-b fehlt → FAILED, nicht still ausgeführt
due = coord.due(now_show_ns=1_000_000)
assert due == []
assert coord.get(planned.command_id).state is ArmState.FAILED
def test_due_executes_when_all_armed(two_node_setup) -> None:
coord, _ = two_node_setup
planned = coord.schedule("scene-1", now_show_ns=0, lead_ns=100)
coord.acknowledge_arm(planned.command_id, "node-a")
coord.acknowledge_arm(planned.command_id, "node-b")
due = coord.due(now_show_ns=1_000_000)
assert len(due) == 1 and due[0].command_id == planned.command_id
assert coord.get(planned.command_id).state is ArmState.ARMED
def test_execute_completes_when_all_nodes_report(two_node_setup) -> None:
coord, _ = two_node_setup
planned = coord.schedule("scene-1", now_show_ns=0, lead_ns=100)
coord.acknowledge_arm(planned.command_id, "node-a")
coord.acknowledge_arm(planned.command_id, "node-b")
coord.due(now_show_ns=1_000_000)
# beide Nodes melden ausgeführt (mit Ist-Zeit, §6.5)
coord.acknowledge_execute(planned.command_id, "node-a")
assert coord.get(planned.command_id).state is ArmState.ARMED # noch nicht komplett
coord.acknowledge_execute(planned.command_id, "node-b")
assert coord.get(planned.command_id).state is ArmState.EXECUTED
def test_expected_nodes_preview(two_node_setup) -> None:
"""§17.5: Vor dem Commit sichtbar, welche Nodes die Szene erhalten."""
coord, _ = two_node_setup
planned = coord.schedule("scene-1", now_show_ns=0, lead_ns=100)
assert coord.expected_nodes(planned) == frozenset({"node-a", "node-b"})
def test_acknowledge_unknown_command_rejected(two_node_setup) -> None:
coord, _ = two_node_setup
with pytest.raises(KeyError):
coord.acknowledge_arm("gibts-nicht", "node-a")
def test_due_ignores_already_handled(two_node_setup) -> None:
"""Erledigte/gescheiterte Aktivierungen werden nicht erneut geliefert."""
coord, _ = two_node_setup
planned = coord.schedule("scene-1", now_show_ns=0, lead_ns=100)
coord.acknowledge_arm(planned.command_id, "node-a")
coord.acknowledge_arm(planned.command_id, "node-b")
coord.due(now_show_ns=1_000_000) # erstmalig fällig → ARMED
second = coord.due(now_show_ns=2_000_000) # erneut aufgerufen: kein Duplikat
assert second == []
def test_due_before_showtime_returns_empty(two_node_setup) -> None:
coord, _ = two_node_setup
coord.schedule("scene-1", now_show_ns=1_000_000_000, lead_ns=200_000_000)
assert coord.due(now_show_ns=1_100_000_000) == [] # Showzeit noch nicht erreicht