From cc119d3751845e5bb21415cfa1dbc3d9fd2b08e9 Mon Sep 17 00:00:00 2001 From: HMS MediaEngine Agent Date: Fri, 11 Sep 2026 01:31:26 +0200 Subject: [PATCH] =?UTF-8?q?Phase=202:=20State-Sync-Fundament=20=E2=80=93?= =?UTF-8?q?=20State-Store,=20Gruppen,=20Clock,=20Aktivierung?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 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 --- STATUS.md | 27 +- packages/cluster/hms_cluster/__init__.py | 21 +- packages/cluster/hms_cluster/activation.py | 148 ++++++++++ packages/cluster/hms_cluster/clock.py | 117 ++++++++ packages/cluster/hms_cluster/groups.py | 116 ++++++++ .../hms_persistence/state_store.py | 165 +++++++++++ tests/unit/test_cluster_routing.py | 258 ++++++++++++++++++ tests/unit/test_state_store.py | 112 ++++++++ 8 files changed, 958 insertions(+), 6 deletions(-) create mode 100644 packages/cluster/hms_cluster/activation.py create mode 100644 packages/cluster/hms_cluster/clock.py create mode 100644 packages/cluster/hms_cluster/groups.py create mode 100644 packages/persistence/hms_persistence/state_store.py create mode 100644 tests/unit/test_cluster_routing.py create mode 100644 tests/unit/test_state_store.py diff --git a/STATUS.md b/STATUS.md index 99fd800..a2975f1 100644 --- a/STATUS.md +++ b/STATUS.md @@ -37,7 +37,7 @@ Phase 0 (abgeschlossen, soweit ohne Hardware möglich): - [x] Beispielplugins (Passthrough, Gaussian Blur mit 3 AQ-Varianten), Tools - [x] 121 Unit-/Integrationstests grün, Ruff grün (Belege: TEST_REPORT.md) -Phase 1 (im Bau): +Phase 1 (plattformneutral abgeschlossen, ADR-0008): - [x] IPC-Verbindung Control Core ↔ Renderer: Handshake (hello/welcome), vollständiger Snapshot nach Verbindung, Deltas mit monotoner Revision, @@ -66,11 +66,30 @@ Phase 1 (im Bau): - [ ] mDNS-Echtnetz-Betrieb mit zeroconf auf Zielsystemen (Modell fertig; Multicast-Test gehört zu Gate 1, ADR-0009) +Phase 2 (im Bau): + +- [x] Domänenmodell: Project/Composition/Layer/Source/EffectInstance/ + MediaAsset/OutputSurface/PresetScene mit Validierungen (§10.1) +- [x] Medien-Engine: PlaybackController (Play/Pause/Stop/Retrigger, Loop/ + Once/Ping-Pong, In/Out, ±4x Speed, Ende-Ereignis), PreloadSlot für + atomaren Clipwechsel, MediaLibrary mit Duplikaterkennung und Bank- + Slots, ContentManifest mit SHA-256 und Chunk-Hashes (§12, §13, §6.4) +- [x] Project-State-Store mit monotonen Revisionen, Snapshot/Delta + (neu/geändert/gelöscht), Szenenaktivierung als Zielzustand, + Projekt-/Livezustand getrennt (§6.4, §24.2, §18.1) +- [x] Servergruppen + Zielrouting: All/Node/Output/Group, Zielregeln + selected/tag_query/all, Commit-Vorschau (§6.3, §10.1, §17.5) +- [x] Clock-Offset-/Drift-Messung: RTT-Min-Filter (≤2× Min), Drift erst + ab 1 s Fenster, Showzeit→lokale-Zeit-Abbildung (§6.4) +- [x] zeitgestempelte Preset-Aktivierung: Vorlauf 200 ms, Arm/Execute/ + Ack, FAILED bei fehlender Arm-Bestätigung (§6.5) +- [ ] Renderer-Anbindung: IPC-Deltas aus dem State-Store in den Renderer + ## Nächste drei Aufgaben -1. Phase-2-Domänenmodell: Composition/Layer/MediaAsset in hms_domain (§10.1) -2. Renderer-IPC-Handshake an Control Core anbinden (Snapshot/Delta-Fluss) -3. mDNS-Echtnetz mit zeroconf auf Zielsystemen + Gate-1-LAN-Test vorbereiten +1. Renderer-Anbindung: IPC-Deltas aus dem State-Store (§6.4) +2. Phase-3-Vorbereitung: Plugin-Loader/Lifecycle und Backend-Adapter (§14) +3. Echtes mDNS mit zeroconf auf Zielsystemen (Gate-1-LAN-Test) ## Ausstehende Hardware-Validierung (ADR-0008 Testplan) diff --git a/packages/cluster/hms_cluster/__init__.py b/packages/cluster/hms_cluster/__init__.py index 0340ed9..7cd933f 100644 --- a/packages/cluster/hms_cluster/__init__.py +++ b/packages/cluster/hms_cluster/__init__.py @@ -1,6 +1,12 @@ -"""hms_cluster – Clusterprotokoll, Node-Registry, Paarung, Discovery -(PLAN.md §6.3, §6.5; ADR-0009, ADR-0010).""" +"""hms_cluster – Clusterprotokoll, Registry, Paarung, Discovery, Gruppen, +Clock-Sync, zeitgestempelte Aktivierung (§6.3–§6.5; ADR-0009/0010).""" +from hms_cluster.activation import ( + ActivationCoordinator, + ArmState, + PlannedActivation, +) +from hms_cluster.clock import ClockEstimator, ClockSample, now_monotonic_ns from hms_cluster.discovery import ( DISCOVERY_PROTOCOL_VERSION, SERVICE_TYPE, @@ -8,6 +14,7 @@ from hms_cluster.discovery import ( ServiceInfo, capability_digest, ) +from hms_cluster.groups import GroupRouter, GroupRule, ServerGroup, TargetKind from hms_cluster.message import ClusterMessage, CommandStatus, CommandTracker from hms_cluster.pairing import ( PairingPin, @@ -51,4 +58,14 @@ __all__ = [ "NodeEntry", "NodeHealth", "NodeRegistry", + "GroupRouter", + "GroupRule", + "ServerGroup", + "TargetKind", + "ClockEstimator", + "ClockSample", + "now_monotonic_ns", + "ActivationCoordinator", + "ArmState", + "PlannedActivation", ] diff --git a/packages/cluster/hms_cluster/activation.py b/packages/cluster/hms_cluster/activation.py new file mode 100644 index 0000000..75e77b7 --- /dev/null +++ b/packages/cluster/hms_cluster/activation.py @@ -0,0 +1,148 @@ +"""Zeitgestempelte Preset-Aktivierung über Gruppen (§6.5, §6.4). + +Koordinierter Show-Modus (§6.3): Der Coordinator verteilt zeitgestempelte +Preset-/Parameterkommandos mit Preload/Arm/Ack und execute_at typischerweise +100–300 ms im Voraus. Zwei-Phasen-Aktivierung (§6.5): vollständig ins +Staging übertragen und validieren, danach atomar auf dieselbe Revision +schalten. + +execute_at ist eine Showzeit (Coordinator-Zeitbasis); jeder Node bildet sie +über seinen Clock-Offset auf die lokale monotone Zeit ab (§6.4). +""" + +from __future__ import annotations + +import uuid +from dataclasses import dataclass, field +from enum import StrEnum + +from hms_cluster.clock import now_monotonic_ns +from hms_cluster.groups import GroupRouter, TargetKind + +DEFAULT_LEAD_NS = 200_000_000 # 200 ms Vorlauf (§6.4: typisch 100–300 ms) + + +class ArmState(StrEnum): + """Phasen eines geplanten Preset-Starts (§6.5).""" + + CREATED = "created" + ARMED = "armed" + EXECUTED = "executed" + FAILED = "failed" + + +@dataclass +class PlannedActivation: + """Ein zeitgestempelter Gruppen-Command (§6.5).""" + + command_id: str + scene_id: str + target_kind: TargetKind + target_id: str | None + execute_at_show_ns: int # Coordinator-Showzeit + created_ns: int + state: ArmState = ArmState.CREATED + armed_nodes: frozenset[str] = field(default_factory=frozenset) + executed_nodes: frozenset[str] = field(default_factory=frozenset) + + @classmethod + def plan( + cls, + scene_id: str, + target_kind: TargetKind, + target_id: str | None, + now_show_ns: int, + lead_ns: int = DEFAULT_LEAD_NS, + ) -> PlannedActivation: + """Plant die Ausführung `lead_ns` im Voraus (§6.4).""" + return cls( + command_id=str(uuid.uuid4()), + scene_id=scene_id, + target_kind=target_kind, + target_id=target_id, + execute_at_show_ns=now_show_ns + lead_ns, + created_ns=now_monotonic_ns(), + ) + + +class ActivationCoordinator: + """Koordiniert Preload/Arm/Execute über eine Node-Gruppe (§6.5). + + - schedule(): plant die Aktivierung mit Vorlauf + - acknowledge_arm(): Node bestätigt arm (Preflight grün) + - acknowledge_execute(): Node hat ausgeführt (mit Ist-Zeit) + - due(): Aktivierungen, deren Showzeit erreicht ist (pro Tick) + + Fehlerfälle (§6.5): fehlt eine Arm-Bestätigung, blockiert die + Ausführung standardmäßig; das wird als FAILED sichtbar dokumentiert + und nicht still überbrückt. + """ + + def __init__(self, router: GroupRouter) -> None: + self._router = router + self._planned: dict[str, PlannedActivation] = {} + + def schedule( + self, + scene_id: str, + target_kind: TargetKind = TargetKind.ALL, + target_id: str | None = None, + now_show_ns: int | None = None, + lead_ns: int = DEFAULT_LEAD_NS, + ) -> PlannedActivation: + """Plant eine Aktivierung; Ziel-Nodes werden sofort aufgelöst.""" + now = now_show_ns if now_show_ns is not None else now_monotonic_ns() + planned = PlannedActivation.plan( + scene_id=scene_id, + target_kind=target_kind, + target_id=target_id, + now_show_ns=now, + lead_ns=lead_ns, + ) + self._planned[planned.command_id] = planned + return planned + + def expected_nodes(self, planned: PlannedActivation) -> frozenset[str]: + """Ziel-Nodes dieser Aktivierung (§17.5: vor Commit sichtbar).""" + return self._router.preview_targets(planned.target_kind, planned.target_id) + + def acknowledge_arm(self, command_id: str, node_id: str) -> ArmState: + """Node hat geprüft und armed (§6.5).""" + planned = self._planned.get(command_id) + if planned is None: + raise KeyError(f"unbekannter Command {command_id}") + planned.armed_nodes = frozenset(set(planned.armed_nodes) | {node_id}) + return planned.state + + def acknowledge_execute(self, command_id: str, node_id: str) -> ArmState: + """Node hat ausgeführt; Ist-Zeit wird vom Node gemeldet (§6.5).""" + planned = self._planned.get(command_id) + if planned is None: + raise KeyError(f"unbekannter Command {command_id}") + planned.executed_nodes = frozenset(set(planned.executed_nodes) | {node_id}) + if planned.executed_nodes >= self.expected_nodes(planned): + planned.state = ArmState.EXECUTED + return planned.state + + def due(self, now_show_ns: int | None = None) -> list[PlannedActivation]: + """Aktivierungen, deren Showzeit gekommen ist – nur für armed. + + Nicht voll armbare Aktivierungen werden FAILED, nicht still + ausgeführt (§6.5: fehlende Bestätigung blockiert standardmäßig). + """ + now = now_show_ns if now_show_ns is not None else now_monotonic_ns() + due_list: list[PlannedActivation] = [] + for planned in list(self._planned.values()): + if planned.state is not ArmState.CREATED: + continue + if now >= planned.execute_at_show_ns: + expected = self.expected_nodes(planned) + if expected and planned.armed_nodes >= expected: + planned.state = ArmState.ARMED + due_list.append(planned) + else: + planned.state = ArmState.FAILED # §6.5: Blockade sichtbar + return due_list + + def get(self, command_id: str) -> PlannedActivation | None: + return self._planned.get(command_id) diff --git a/packages/cluster/hms_cluster/clock.py b/packages/cluster/hms_cluster/clock.py new file mode 100644 index 0000000..7ef1fa4 --- /dev/null +++ b/packages/cluster/hms_cluster/clock.py @@ -0,0 +1,117 @@ +"""Clock-Offset- und Drift-Messung (PLAN.md §6.4 Clock Sync). + +V1-Verfahren: PTP, wenn verfügbar; sonst gemessene Offset-/Drift- +Schätzung gegen den Coordinator (§6.4). Das hier ist die softwareseitige +Messung: Round-Trip-basierte Offset-Schätzung mit Min-Filterung (NTP-artig) +und linearer Drift-Schätzung über Probenpaare. + +Grenzen (§6.4): messbare Softwarezeit, kein Genlock; p95 ≤ 10 ms für +vorgepufferte Preset-/Command-Starts im verkabelten Referenz-LAN ist ein +Abnahmeziel von Gate 2, keine Zusage für beliebige Netze. +""" + +from __future__ import annotations + +import time +from dataclasses import dataclass, field + + +@dataclass(frozen=True) +class ClockSample: + """Eine RTT-Messprobe zwischen Coordinator und Node. + + t0/t1 in lokaler monotoner Zeit des Coordinators; node_time ist die + vom Node zurückgemeldete eigene monotone Zeit (normiert). + """ + + t0_ns: int + t1_ns: int + node_time_ns: int + + @property + def rtt_ns(self) -> int: + return self.t1_ns - self.t0_ns + + @property + def offset_ns(self) -> int: + """Min-RTT-Näherung: Offset = node_time - (t0 + rtt/2).""" + return self.node_time_ns - (self.t0_ns + self.rtt_ns // 2) + + +@dataclass +class ClockEstimator: + """Schätzt Offset und Drift aus RTT-Proben (§6.4). + + - feed(): neue Probe; behält die Proben mit kleinster RTT (Min-Filter, + weil geringe RTT ≈ geringe Warteschlangen-Verzögerung) + - offset: geglätteter Offset gegen den Coordinator + - drift_ppm: Änderung des Offsets über die Zeit (μs/s) + - Window begrenzt (kein unbeschränkter Zustand, §33) + """ + + max_samples: int = 64 + _samples: list[ClockSample] = field(default_factory=list) + _best: list[ClockSample] = field(default_factory=list) # kleinste RTTs + + def feed(self, sample: ClockSample) -> None: + if sample.rtt_ns < 0: + raise ValueError("negative RTT unmöglich") + self._samples.append(sample) + if len(self._samples) > self.max_samples: + self._samples.pop(0) + # Min-RTT-Filter: nur Proben mit RTT ≤ 2× Minimum sind belastbar; + # hohe RTT bedeutet Warteschlangen-Jitter, der den Offset verfälscht + min_rtt = min(s.rtt_ns for s in self._samples) + ranked = sorted( + (s for s in self._samples if s.rtt_ns <= 2 * min_rtt), + key=lambda s: s.rtt_ns, + )[:8] + self._best = ranked + + @property + def offset_ns(self) -> int | None: + """Aktueller Offset-Schätzer (Mittel über Best-Proben).""" + if not self._best: + return None + return sum(s.offset_ns for s in self._best) // len(self._best) + + @property + def rtt_ns(self) -> int | None: + """Beste (kleinste) gemessene RTT.""" + if not self._best: + return None + return min(s.rtt_ns for s in self._best) + + @property + def drift_ppm(self) -> float | None: + """Lineare Drift-Schätzung über die Best-Proben (μs/s). + + Offset-Änderung geteilt durch verstrichene RTT-Mittezeit; None bei + weniger als zwei Best-Proben oder zu kurzem Fenster (< 1 s). + """ + if len(self._best) < 2: + return None + ordered = sorted(self._best, key=lambda s: s.t0_ns) + first, last = ordered[0], ordered[-1] + dt_ns = last.t0_ns - first.t0_ns + if dt_ns < 1_000_000_000: # < 1 s: Drift nicht belastbar + return None + d_offset = last.offset_ns - first.offset_ns + return (d_offset / dt_ns) * 1_000_000.0 + + def map_show_time(self, show_time_ns: int) -> int | None: + """Bildet Coordinator-Showzeit auf lokale Node-Zeit ab (§6.4: + Showzeit → lokale Monotonic). + + Voraussetzung: dieser Estimator läuft Node-seitig mit Proben, + deren node_time die eigene Uhr ist. + """ + offset = self.offset_ns + if offset is None: + return None + return show_time_ns + offset + + +def now_monotonic_ns() -> int: + """Gemeinsame monotone Zeitbasis (§12.2: Audio/Video gemeinsame Basis).""" + return time.monotonic_ns() diff --git a/packages/cluster/hms_cluster/groups.py b/packages/cluster/hms_cluster/groups.py new file mode 100644 index 0000000..872dffe --- /dev/null +++ b/packages/cluster/hms_cluster/groups.py @@ -0,0 +1,116 @@ +"""Servergruppen und Zielrouting (PLAN.md §6.3 Bedienmodelle, §10.1 ServerGroup). + +- Coordinator routet Commands an `All`, eine `ServerGroup`, einen einzelnen + `Node` oder einen `Output` (§6.3 Control-Center-Modell) +- Zielregel je Gruppe: all | selected | tag_query | feste Node-Liste +- Zielauflösung ist deterministisch und testbar; Gruppenänderungen zeigen + vor dem Commit, welche Nodes sie erhalten (§17.5) +""" + +from __future__ import annotations + +from dataclasses import dataclass, field +from enum import StrEnum + + +class TargetKind(StrEnum): + """Command-Ziele (§6.3): All, ServerGroup, Node oder Output.""" + + ALL = "all" + SERVER_GROUP = "server_group" + NODE = "node" + OUTPUT = "output" + + +class GroupRule(StrEnum): + """Zielregel je ServerGroup (§10.1).""" + + ALL = "all" + SELECTED = "selected" + TAG_QUERY = "tag_query" + NODE_LIST = "node_list" + + +@dataclass +class ServerGroup: + """Servergruppe (§10.1): Node-Auswahl mit Zielregel. + + - node_ids: feste Mitglieder bei SELECTED/NODE_LIST + - tags: Node-Tags für TAG_QUERY (Node-Einträge tragen passende Tags) + - output_map: optionaler Layer-/Output-Zuordnungshinweis (§10.1) + """ + + id: str + name: str + rule: GroupRule = GroupRule.SELECTED + node_ids: frozenset[str] = field(default_factory=frozenset) + tags: frozenset[str] = field(default_factory=frozenset) + output_map: dict[str, str] = field(default_factory=dict) # layer_id → output_id + + +class GroupRouter: + """Löst Command-Ziele auf Node-Mengen auf (§6.3). + + Der Coordinator nutzt resolve() vor jedem Senden; die UI kann + resolve() für die Commit-Vorschau verwenden (§17.5: „welche Nodes + erhalten dies?"). + """ + + def __init__(self) -> None: + self._groups: dict[str, ServerGroup] = {} + self._node_tags: dict[str, frozenset[str]] = {} # node_id → tags + self._node_outputs: dict[str, list[str]] = {} # node_id → output_ids + + # ---------- Verwaltung ---------- + + def upsert_group(self, group: ServerGroup) -> None: + self._groups[group.id] = group + + def remove_group(self, group_id: str) -> None: + self._groups.pop(group_id, None) + + def get_group(self, group_id: str) -> ServerGroup | None: + return self._groups.get(group_id) + + def set_node_tags(self, node_id: str, tags: frozenset[str]) -> None: + self._node_tags[node_id] = frozenset(tags) + + def set_node_outputs(self, node_id: str, output_ids: list[str]) -> None: + self._node_outputs[node_id] = list(output_ids) + + def known_nodes(self) -> frozenset[str]: + return frozenset(self._node_tags) + + # ---------- Zielauflösung (§6.3) ---------- + + def resolve(self, kind: TargetKind, target_id: str | None = None) -> frozenset[str]: + """Löst ein Ziel auf eine Node-Menge auf; leer bei unbekanntem Ziel.""" + if kind is TargetKind.ALL: + return self.known_nodes() + if kind is TargetKind.NODE: + return frozenset({target_id}) if target_id in self._node_tags else frozenset() + if kind is TargetKind.OUTPUT: + return frozenset( + node_id + for node_id, outputs in self._node_outputs.items() + if target_id in outputs + ) + if kind is TargetKind.SERVER_GROUP: + group = self._groups.get(target_id or "") + if group is None: + return frozenset() + if group.rule is GroupRule.ALL: + return self.known_nodes() + if group.rule is GroupRule.SELECTED or group.rule is GroupRule.NODE_LIST: + return group.node_ids & self.known_nodes() + if group.rule is GroupRule.TAG_QUERY: + return frozenset( + node_id + for node_id, tags in self._node_tags.items() + if group.tags & tags + ) + return frozenset() + + def preview_targets(self, kind: TargetKind, target_id: str | None = None) -> frozenset[str]: + """Commit-Vorschau: identisch zu resolve (§17.5).""" + return self.resolve(kind, target_id) diff --git a/packages/persistence/hms_persistence/state_store.py b/packages/persistence/hms_persistence/state_store.py new file mode 100644 index 0000000..31caa0d --- /dev/null +++ b/packages/persistence/hms_persistence/state_store.py @@ -0,0 +1,165 @@ +"""Autoritativer Projekt- und Showzustand (PLAN.md §6.4 State Sync, §24.2). + +- Vollständiger Snapshot nach Verbindung; danach inkrementelle Deltas +- monotone Revisionen: kein halber Zustand (§11.4, §6.4) +- Livezustand und dauerhafter Projektzustand sind getrennt (§24.2) +- Szenenaktivierung erzeugt eine neue State-Revision +""" + +from __future__ import annotations + +import time +from dataclasses import dataclass, field + +from hms_domain.model import PresetScene, Project + + +@dataclass(frozen=True) +class StateDelta: + """Inkrementelle Änderung mit monotoner Revision (§6.4). + + - project_revision: Projekt-Inhaltsrevision (Manifest-Ebene) + - state_revision: monotone Showzustands-Revision + - changes: Parameterpfad → Wert; gelöschte Pfade als None markiert + """ + + state_revision: int + project_revision: int + changes: dict[str, float | None] + monotonic_ns: int + + def to_dict(self) -> dict: + return { + "state_revision": self.state_revision, + "project_revision": self.project_revision, + "changes": self.changes, + "monotonic_ns": self.monotonic_ns, + } + + +@dataclass +class ProjectStateStore: + """Autoritative Instanz im Control Core (§6.1B, §6.4). + + Trennung (§24.2): + - project: dauerhafter Projektzustand (Domänenmodell, persistiert) + - live: Show-Livezustand (Parameterwerte je Pfad, nicht persistent) + + Revisions: + - state_revision steigt bei jeder Livezustandsänderung monoton + - project_revision steigt bei Projektinhalts-Änderungen (z. B. neue + Szenen, Medien-Revision) – Grundlage für Preflight (§6.5) + """ + + _state_revision: int = 0 + _project_revision: int = 0 + _live: dict[str, float] = field(default_factory=dict) + _project: Project | None = None + _project_name: str = "" + + # ---------- Projekt ---------- + + def activate_project(self, project: Project) -> int: + """Aktiviert ein Projekt als autoritative Basis; erhöht die + Projektrevision. Livezustand wird zurückgesetzt (kein Mischzustand).""" + self._project = project + self._project_name = project.name + self._project_revision += 1 + self._live.clear() + self._state_revision += 1 # neuer Zustand nach Projektwechsel + return self._state_revision + + @property + def project(self) -> Project | None: + return self._project + + @property + def project_revision(self) -> int: + return self._project_revision + + @property + def state_revision(self) -> int: + return self._state_revision + + # ---------- Livezustand ---------- + + def set_value(self, path: str, value: float) -> int: + """Setzt einen Liveparameter; gibt die neue State-Revision zurück.""" + self._live[path] = float(value) + self._state_revision += 1 + return self._state_revision + + def clear_value(self, path: str) -> int: + """Entfernt einen Liveparameter (z. B. Release); neue Revision.""" + self._live.pop(path, None) + self._state_revision += 1 + return self._state_revision + + def get_value(self, path: str) -> float | None: + return self._live.get(path) + + # ---------- Snapshot / Delta (§6.4) ---------- + + def snapshot(self) -> dict: + """Vollständiger Zustand nach Verbindungsaufbau (§6.2, §6.4).""" + return { + "state_revision": self._state_revision, + "project_revision": self._project_revision, + "project_name": self._project_name, + "values": dict(self._live), + "monotonic_ns": time.monotonic_ns(), + } + + def delta_since(self, last_seen_revision: int, pending: dict[str, float]) -> StateDelta | None: + """Delta seit einer gesehenen Revision; None, wenn nichts Neues. + + pending: letzter BEKANNTER Zustand des Empfängers (Pfad → Wert, + wie er bei last_seen_revision beim Client stand). Das Delta enthält: + - neue Pfade (in live, nicht in pending) mit ihrem Wert + - geänderte Pfade mit dem neuen Wert + - gelöschte Pfade als None + Ein vollständiger Re-Sync (neuer Snapshot) ist Aufgabe des + Transports, wenn last_seen_revision zu alt ist (§6.2). + """ + if last_seen_revision > self._state_revision: + raise ValueError( + f"gesehene Revision {last_seen_revision} liegt in der Zukunft" + ) + if last_seen_revision == self._state_revision: + return None # nichts Neues + changes: dict[str, float | None] = {} + for path, value in self._live.items(): + if path not in pending: + changes[path] = value # neu seit last_seen + elif pending[path] != value: + changes[path] = value # geändert + for path in pending: + if path not in self._live: + changes[path] = None # gelöscht + return StateDelta( + state_revision=self._state_revision, + project_revision=self._project_revision, + changes=changes, + monotonic_ns=time.monotonic_ns(), + ) + + # ---------- Szenen (§18) ---------- + + def apply_scene(self, scene: PresetScene) -> int: + """Aktiviert eine Szene direkt (§18.1): ÜBERNAHME der Snapshot-Werte + in den Livezustand als neue State-Revision. Übergänge (Crossfade + etc.) berechnet der Renderer aus vorher/nachher – hier entsteht nur + der Zielzustand (§18.1: Diff zur Laufzeit). + """ + snapshot_values = scene.composition_snapshot.get("values", {}) + if not isinstance(snapshot_values, dict): + raise ValueError("Szene enthält keine Werte") + for path, value in snapshot_values.items(): + self._live[str(path)] = float(value) + self._state_revision += 1 + return self._state_revision + + def bump_project_revision(self) -> int: + """Projektinhalt geändert (Medien/Plugins/Szenen) → neue Revision.""" + self._project_revision += 1 + return self._project_revision diff --git a/tests/unit/test_cluster_routing.py b/tests/unit/test_cluster_routing.py new file mode 100644 index 0000000..23a7f99 --- /dev/null +++ b/tests/unit/test_cluster_routing.py @@ -0,0 +1,258 @@ +"""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: 100–300 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 diff --git a/tests/unit/test_state_store.py b/tests/unit/test_state_store.py new file mode 100644 index 0000000..022731e --- /dev/null +++ b/tests/unit/test_state_store.py @@ -0,0 +1,112 @@ +"""Unit-Tests Project-State-Store (PLAN.md §6.4 State Sync, §24.2, §18.1).""" + +from __future__ import annotations + +import pytest +from hms_domain import PresetScene +from hms_persistence.state_store import ProjectStateStore + + +def test_revisions_start_at_zero() -> None: + store = ProjectStateStore() + assert store.state_revision == 0 + assert store.project_revision == 0 + assert store.snapshot()["values"] == {} + + +def test_activate_project_resets_live_state() -> None: + """Kein Mischzustand: Projektwechsel leert den Livezustand (§24.2).""" + from hms_domain import Project + + store = ProjectStateStore() + store.set_value("master/intensity", 0.5) + store.activate_project(Project(name="Show A")) + assert store.project_revision == 1 + assert store.project is not None and store.project.name == "Show A" + assert store.get_value("master/intensity") is None # Live geleert + assert store.state_revision == 2 # set_value (1) + Projektwechsel (2) + + +def test_set_and_clear_value_increase_state_revision() -> None: + store = ProjectStateStore() + r1 = store.set_value("master/intensity", 0.7) + assert r1 == 1 + r2 = store.set_value("master/intensity", 0.9) + assert r2 == 2 + r3 = store.clear_value("master/intensity") + assert r3 == 3 + assert store.get_value("master/intensity") is None + # Löschen eines nicht existierenden Pfades erhöht trotzdem deterministisch + r4 = store.clear_value("master/intensity") + assert r4 == 4 + + +def test_snapshot_contains_revisions_and_values() -> None: + store = ProjectStateStore() + store.set_value("master/intensity", 1.0) + snap = store.snapshot() + assert snap["state_revision"] == 1 + assert snap["project_revision"] == 0 + assert snap["values"] == {"master/intensity": 1.0} + assert snap["monotonic_ns"] > 0 + + +def test_delta_since_reports_changes_and_deletions() -> None: + """Delta überträgt Änderungen; gelöschte Pfade als None (§6.4).""" + store = ProjectStateStore() + store.set_value("a", 1.0) + seen = store.state_revision + # nach „Disconnect": neue Änderungen + store.set_value("a", 2.0) + store.set_value("b", 3.0) + store.clear_value("c") # c war nie gesetzt – bleibt neutral + delta = store.delta_since(seen, pending={"a": 1.0, "c": 0.0}) + assert delta is not None + assert delta.state_revision == store.state_revision + assert delta.changes["a"] == 2.0 # geändert + assert "b" in delta.changes and delta.changes["b"] == 3.0 # neu + assert delta.changes.get("c") is None # gelöscht + + +def test_delta_none_when_up_to_date() -> None: + store = ProjectStateStore() + store.set_value("a", 1.0) + assert store.delta_since(store.state_revision, pending={}) is None + + +def test_delta_rejects_future_revision() -> None: + store = ProjectStateStore() + with pytest.raises(ValueError, match="Zukunft"): + store.delta_since(99, pending={}) + + +def test_apply_scene_overwrites_live_as_target_state() -> None: + """Szenenabruf erzeugt Zielzustand; Übergänge macht der Renderer (§18.1).""" + store = ProjectStateStore() + store.set_value("layer/x/opacity", 0.2) + scene = PresetScene( + name="Look", + composition_snapshot={ + "values": {"layer/x/opacity": 0.8, "master/intensity": 1.0} + }, + ) + rev = store.apply_scene(scene) + assert rev == 2 # set + scene + assert store.get_value("layer/x/opacity") == pytest.approx(0.8) + assert store.get_value("master/intensity") == pytest.approx(1.0) + + +def test_apply_scene_rejects_scene_without_values() -> None: + scene = PresetScene(name="Kaputt", composition_snapshot={"values": "kein dict"}) + store = ProjectStateStore() + with pytest.raises(ValueError, match="Werte"): + store.apply_scene(scene) + + +def test_bump_project_revision_independent_of_state() -> None: + """Projektinhalt (Medien/Szenen) und Livezustand getrennt (§24.2).""" + store = ProjectStateStore() + state_before = store.state_revision + store.bump_project_revision() + assert store.project_revision == 1 + assert store.state_revision == state_before # Live unverändert