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
This commit is contained in:
HMS MediaEngine Agent
2026-09-11 01:31:26 +02:00
parent d24fe3a963
commit cc119d3751
8 changed files with 958 additions and 6 deletions
+19 -2
View File
@@ -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",
]
+148
View File
@@ -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
100300 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 100300 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)
+117
View File
@@ -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()
+116
View File
@@ -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)
@@ -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