"""Message specifications: the typed contract between nodes. A node declares ports; each port binds to a message name. The message name is the wiring: a node consuming ``heating.temperature`` receives whatever any node provides under that name. Names are namespaced per flow — a bare name is qualified with the owning flow (``temperature`` in flow ``heating`` becomes ``heating.temperature``), a dotted name is used as written, which is how flows consume each other's messages. """ from __future__ import annotations import json from enum import Enum from typing import Any from pydantic import BaseModel, ConfigDict, Field, model_validator class DType(str, Enum): """Serializable payload types. The scalars carry what a single reading can say. The three structured ones are declared shapes rather than "some JSON": a widget or a downstream node knows what it is getting before anything runs, which is what lets the dashboard picker offer a message and refuse a wrong binding. Binary payloads — tensors, checkpoints, images — travel as ``artifact``: the bytes go to the artifact store and the message carries a reference to them. That keeps everything on the wire JSON, which is what the state backend, the queue and the worker protocol all rely on, and it means a thirty-megabyte checkpoint never sits in Redis. Inline codecs would only be needed for payloads too small to be worth a round trip, and nothing asks for that yet. """ FLOAT = "float" INT = "int" STR = "str" BOOL = "bool" JSON = "json" #: ``{"lines": [{"label": str, "points": [[ts, value], ...]}], ...}``. #: Keys beside ``lines`` are carried through untouched — a chart's query #: puts the range and interval it asked for there and reads them back. SERIES = "series" #: Flat named scalars: ``{"title": "Boiler", "severity": "error"}``. RECORD = "record" #: Ordered items of one declared shape; see :attr:`MessageSpec.item`. LIST = "list" #: A reference to stored bytes: #: ``{"digest": "sha256:…", "size": int, "media_type": str, "name": str}``. ARTIFACT = "artifact" _JSON_TYPES = (dict, list, str, int, float, bool, type(None)) #: What a record may hold. Nesting is deliberately out: a record that can #: contain a record is a schema language, and the shape stops being readable #: from the declaration alone. _SCALARS = (str, int, float, bool) #: Item types a list may declare. Recursion is refused for the same reason. _ITEM_TYPES = frozenset( {DType.FLOAT, DType.INT, DType.STR, DType.BOOL, DType.JSON, DType.RECORD} ) def _is_number(value: Any) -> bool: """A measurement. bool is an int subclass; a flag is not a number here.""" return isinstance(value, (int, float)) and not isinstance(value, bool) def _is_record(value: Any) -> bool: return isinstance(value, dict) and all( isinstance(key, str) and (item is None or isinstance(item, _SCALARS)) for key, item in value.items() ) def _is_series(value: Any) -> bool: """Labelled lines of ``(timestamp, value)`` pairs. Checked all the way down. That is O(n) in the number of points, but so is the JSON encoding every message already pays for. """ if not isinstance(value, dict) or not isinstance(value.get("lines"), list): return False return all( isinstance(line, dict) and isinstance(line.get("label"), str) and isinstance(line.get("points"), list) and all( isinstance(point, (list, tuple)) and len(point) == 2 and _is_number(point[0]) and _is_number(point[1]) for point in line["points"] ) for line in value["lines"] ) def _is_artifact(value: Any) -> bool: """A reference to stored bytes, not the bytes themselves. The digest is what makes it one: it names content rather than a location, so the same file produced twice is stored once and a reference stays valid wherever the store is reachable from. """ return ( isinstance(value, dict) and isinstance(value.get("digest"), str) and value["digest"].startswith("sha256:") and isinstance(value.get("size"), int) ) def _matches(dtype: DType, value: Any) -> bool: """Whether one value satisfies a scalar or record type.""" if dtype is DType.BOOL: return isinstance(value, bool) if dtype is DType.INT: return isinstance(value, int) and not isinstance(value, bool) if dtype is DType.FLOAT: return _is_number(value) if dtype is DType.STR: return isinstance(value, str) if dtype is DType.RECORD: return _is_record(value) if dtype is DType.ARTIFACT: return _is_artifact(value) return isinstance(value, _JSON_TYPES) class MessageSpec(BaseModel): """A single port of a node, and the message it is bound to. :param name: The message this port binds to. Bare names are qualified with the flow name at load time; empty means the port is unbound. :param port: The identifier the node function sees. Defaults to the last segment of ``name``, so unqualified flows read naturally. :param dtype: Payload type, validated on every message that passes through. :param item: The type of each item of a ``list`` port, ignored otherwise. Unset means ``record``, which is what the agenda and forecast widgets read; ``float`` is the numeric list a pipeline passes around. A list of lists, or of series, is refused — one declared level is the point. :param interval: Deliver at most every this many seconds; 0 is every time. On an output it holds back publishing, on an input it holds back waking the node. The value is never lost — state keeps the latest — only the delivery is skipped. :param trigger: Whether arriving values wake the node. An input with this off is read when the node runs for some other reason, but never causes a run and never makes the node wait — which is how a node reads a message it also produces without depending on itself. :param stream: On an output, that this port produces repeatedly *during* one execution rather than once at the end — a training loss, a progress fraction. A node emits on it by being a generator and yielding, or by calling ``fluksio.emit``. What it means downstream is nothing special: a value published mid-execution is a value like any other. What it means to a run is that the whole series is kept, which is how a run's metrics are simply its streaming outputs rather than something logged beside them. """ model_config = ConfigDict(frozen=True) name: str = "" port: str = "" dtype: DType = DType.FLOAT item: DType | None = None interval: float = Field(default=0, ge=0) trigger: bool = True stream: bool = False @model_validator(mode="after") def _default_port(self) -> MessageSpec: if not self.port and self.name: object.__setattr__(self, "port", self.name.rsplit(".", 1)[-1]) if self.item is not None and self.item not in _ITEM_TYPES: raise ValueError(f"a list cannot hold '{self.item.value}' items") return self @property def item_dtype(self) -> DType: """What each item of a ``list`` port is, declared or defaulted.""" return self.item or DType.RECORD def check(self, value: Any) -> None: """Raise if ``value`` does not match this port's declared type.""" where = self.name or self.port if self.dtype is DType.SERIES: ok = _is_series(value) elif self.dtype is DType.LIST: if not isinstance(value, list): raise TypeError(f"{where}: expected list, got {type(value).__name__}") item = self.item_dtype for index, element in enumerate(value): if not _matches(item, element): raise TypeError( f"{where}: expected list[{item.value}], got " f"{type(element).__name__} at index {index}" ) return else: ok = _matches(self.dtype, value) if not ok: raise TypeError( f"{where}: expected {self.dtype.value}, got {type(value).__name__}" ) def coerce(self, value: Any) -> Any: """Best-effort conversion of an external value into this port's type. Used where payloads arrive as text (HTTP query strings, MQTT), never on the path between nodes — there a wrong type is an error, not a hint. """ if self.dtype is DType.FLOAT: return float(value) if self.dtype is DType.INT: return int(value) if self.dtype is DType.BOOL: if isinstance(value, bool): return value return str(value).lower() in ("true", "1", "yes", "on") if self.dtype is DType.STR: return value if isinstance(value, str) else json.dumps(value) if self.dtype in (DType.SERIES, DType.RECORD, DType.LIST, DType.ARTIFACT): # A structured payload arriving as text is the same hint a numeric # one is; the shape itself is still checked afterwards. return json.loads(value) if isinstance(value, str) else value return value def __repr__(self) -> str: return f"MessageSpec({self.name or self.port})" def __hash__(self) -> int: return hash((self.name, self.port)) def qualify(flow: str, name: str) -> str: """Resolve a message name against its flow namespace.""" if not name: return "" return name if "." in name else f"{flow}.{name}" def flow_of(qualified: str) -> str: """The flow a qualified message name belongs to.""" return qualified.split(".", 1)[0]