Files
hms-mediaengine/packages/cluster/hms_cluster/message.py
T
HMS MediaEngine Agent c4574927fd Phase 1: Cluster-Fundament – Nachrichten, Registry, Paarung, Discovery
- hms_cluster: ClusterMessage (§6.5: Pflichtfelder, Sequenzen, Revisionen,
  execute_at, Trace-ID), CommandTracker (idempotent, vorwärts-only:
  accepted->armed->executed/failed)
- NodeRegistry: online/degraded/stale/offline über Schwellen, UI-Kategorien
  discovered/paired/unknown/incompatible/offline, Doppel-Node-ID blockiert,
  IP-Wechsel erhaelt node_id (§6.3, §6.5)
- Pairing: PIN (TTL 120s, Versuchslimit), Fingerprint (SHA-256 gruppiert),
  Token nur als Hash, Scopes read/control/content_sync/admin, Ablauf +
  sofortiger Widerruf (§6.3, §27.1, ADR-0010)
- Discovery-Modell: _hmsmedia._tcp.local. TXT ohne Secrets, Capability-
  Digest, persistente manuelle Fallback-Liste (ADR-0009)
- ADR-0009 (mDNS + Fallback) und ADR-0010 (Paarung) dokumentiert
- 25 Unit-Tests; Gesamtsuite 183 Tests gruen, Ruff gruen
2026-09-11 01:06:26 +02:00

165 lines
5.5 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.
"""Cluster-Nachrichtenhülle (PLAN.md §6.5).
Jede Cluster-Nachricht enthält mindestens: cluster_id, node_id, command_id,
Sequenz, Projekt-Revision, Absenderzeit, optionale execute_at-Showzeit und
Trace-ID. Zustandsändernde Commands sind idempotent und werden mit
accepted/armed/executed/failed bestätigt.
"""
from __future__ import annotations
import time
import uuid
from enum import StrEnum
class CommandStatus(StrEnum):
"""Bestätigungsstufen zustandsändernder Commands (§6.5)."""
ACCEPTED = "accepted"
ARMED = "armed"
EXECUTED = "executed"
FAILED = "failed"
class ClusterMessage:
"""Versionierte Cluster-Nachricht mit Pflichtfeldern (§6.5).
- sequence: je Absender monoton; Lücken signalisieren Paketverlust
- project_revision: Zustandsrevision, auf die sich der Command bezieht
- sender_time_ns: monotone Absenderzeit (nicht Wanduhr)
- execute_at_show_time_ns: optional; 100300 ms Vorlauf für Sync-Starts
- trace_id: Korrelations-ID über Log-Grenzen (§28.1)
"""
__slots__ = (
"cluster_id",
"node_id",
"command_id",
"sequence",
"project_revision",
"sender_time_ns",
"execute_at_show_time_ns",
"trace_id",
"status",
"payload",
)
def __init__(
self,
cluster_id: str,
node_id: str,
command_id: str,
sequence: int,
project_revision: int,
sender_time_ns: int | None = None,
execute_at_show_time_ns: int | None = None,
trace_id: str | None = None,
status: CommandStatus | None = None,
payload: dict | None = None,
) -> None:
for field_name, value in (
("cluster_id", cluster_id),
("node_id", node_id),
("command_id", command_id),
):
try:
uuid.UUID(value)
except (ValueError, AttributeError) as exc:
raise ValueError(f"{field_name} must be a UUID") from exc
if sequence < 0:
raise ValueError("sequence must be >= 0")
if project_revision < 0:
raise ValueError("project_revision must be >= 0")
self.cluster_id = cluster_id
self.node_id = node_id
self.command_id = command_id
self.sequence = sequence
self.project_revision = project_revision
self.sender_time_ns = (
sender_time_ns if sender_time_ns is not None else time.monotonic_ns()
)
self.execute_at_show_time_ns = execute_at_show_time_ns
self.trace_id = trace_id or str(uuid.uuid4())
self.status = status
self.payload = payload or {}
def to_dict(self) -> dict:
return {
"cluster_id": self.cluster_id,
"node_id": self.node_id,
"command_id": self.command_id,
"sequence": self.sequence,
"project_revision": self.project_revision,
"sender_time_ns": self.sender_time_ns,
"execute_at_show_time_ns": self.execute_at_show_time_ns,
"trace_id": self.trace_id,
"status": self.status.value if self.status else None,
"payload": self.payload,
}
@classmethod
def from_dict(cls, data: dict) -> ClusterMessage:
status_raw = data.get("status")
status = CommandStatus(status_raw) if status_raw else None
return cls(
cluster_id=data["cluster_id"],
node_id=data["node_id"],
command_id=data["command_id"],
sequence=int(data["sequence"]),
project_revision=int(data["project_revision"]),
sender_time_ns=data.get("sender_time_ns"),
execute_at_show_time_ns=data.get("execute_at_show_time_ns"),
trace_id=data.get("trace_id"),
status=status,
payload=data.get("payload", {}),
)
class CommandTracker:
"""Idempotenz je (node_id, command_id) mit Statusübergängen (§6.5).
- register: neuer Command → True; Duplikat → False (gleiches Ack)
- advance: nur vorwärts accepted → armed → executed/failed
- veraltete Einträge werden nach Kapazität begrenzt (kein unbeschränkter
Cache, §33)
"""
_ORDER = {
CommandStatus.ACCEPTED: 1,
CommandStatus.ARMED: 2,
CommandStatus.EXECUTED: 3,
CommandStatus.FAILED: 3,
}
def __init__(self, capacity: int = 4096) -> None:
if capacity <= 0:
raise ValueError("capacity must be positive")
self._capacity = capacity
self._states: dict[str, CommandStatus] = {}
def register(self, node_id: str, command_id: str) -> bool:
"""True, wenn der Command neu ist; False bei Duplikat."""
key = f"{node_id}:{command_id}"
if key in self._states:
return False
self._states[key] = CommandStatus.ACCEPTED
if len(self._states) > self._capacity:
oldest = next(iter(self._states))
del self._states[oldest]
return True
def advance(self, node_id: str, command_id: str, new_status: CommandStatus) -> bool:
"""Nur vorwärts; False bei unbekanntem Command oder Rückschritt."""
key = f"{node_id}:{command_id}"
current = self._states.get(key)
if current is None:
return False
if self._ORDER[new_status] <= self._ORDER[current]:
return False
self._states[key] = new_status
return True
def status(self, node_id: str, command_id: str) -> CommandStatus | None:
return self._states.get(f"{node_id}:{command_id}")