The poll loop remembered what it read rather than what it published, so a value the node could not publish counted as said: the next poll skipped it, succeeded, and health went back to ok with the port still dark. Remember it only after inject returns, and report ok last. A node reporting itself down is now derived into its flow's issues on read and counted on the health summary, so the canvas marks it and Home says so. Being down does not stop the flow, and the issue clears by itself when the node reports well again. The repeating poll warning is logged once per outage rather than once per tick. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01K1moruzue2kTJd3uVisgNk
261 lines
7.9 KiB
Python
261 lines
7.9 KiB
Python
"""The connector contract: polling, deduplication, health and discovery."""
|
|
|
|
import asyncio
|
|
from types import SimpleNamespace
|
|
from typing import Any
|
|
|
|
import pytest
|
|
|
|
from fluksio.flow.artifacts import ArtifactStore
|
|
from fluksio.flow.connector import CONTRACT_VERSION, ConnectorNode
|
|
from fluksio.flow.controller import NODE_TYPES
|
|
from fluksio.flow.messages import DType, MessageSpec
|
|
from fluksio.flow.nodes import Node
|
|
from fluksio.flow.pipeline import Pipeline
|
|
from fluksio.flow.plugins import load_plugins
|
|
|
|
|
|
class Sensor(ConnectorNode):
|
|
contract = CONTRACT_VERSION
|
|
title = "Test sensor"
|
|
description = "Reads whatever it is told to."
|
|
|
|
class Params(ConnectorNode.Params):
|
|
secret_token: str | None = None
|
|
|
|
def __init__(self, readings: list[Any] | None = None, **kwargs: Any) -> None:
|
|
super().__init__(**kwargs)
|
|
self._readings = list(readings or [])
|
|
self.polls = 0
|
|
|
|
async def poll(self) -> dict[str, Any] | None:
|
|
self.polls += 1
|
|
if not self._readings:
|
|
return None
|
|
value = self._readings.pop(0)
|
|
if isinstance(value, Exception):
|
|
raise value
|
|
return {"reading": value}
|
|
|
|
|
|
def a_sensor(readings: list[Any], **params: Any) -> Sensor:
|
|
node = Sensor(
|
|
readings=readings,
|
|
provides=[MessageSpec(name="reading", dtype=DType.FLOAT)],
|
|
params={"poll_interval": 0.01, **params},
|
|
)
|
|
node.assign_flow("demo", "sensor")
|
|
return node
|
|
|
|
|
|
def run_briefly(node: ConnectorNode, seconds: float = 0.12, app: Any = None) -> None:
|
|
"""Start the poll loop, let it tick a few times, stop it."""
|
|
|
|
async def cycle() -> None:
|
|
await node.start(app)
|
|
await asyncio.sleep(seconds)
|
|
await node.stop()
|
|
|
|
asyncio.run(cycle())
|
|
|
|
|
|
def test_polling_publishes_what_it_reads():
|
|
node = a_sensor([21.5])
|
|
pipeline = Pipeline(nodes=[node])
|
|
|
|
run_briefly(node)
|
|
|
|
assert pipeline.state["demo.reading"] == 21.5
|
|
|
|
|
|
def test_an_unchanged_reading_is_not_republished():
|
|
node = a_sensor([21.5, 21.5, 21.5])
|
|
consumer_ran: list[float] = []
|
|
|
|
def consume(reading, params):
|
|
consumer_ran.append(reading)
|
|
return None
|
|
|
|
consumer = Node(
|
|
f=consume,
|
|
requires=[MessageSpec(name="reading", dtype=DType.FLOAT)],
|
|
name="consumer",
|
|
)
|
|
consumer.assign_flow("demo", "consumer")
|
|
Pipeline(nodes=[node, consumer])
|
|
|
|
run_briefly(node)
|
|
|
|
# Polled repeatedly, but the value never changed, so downstream ran once.
|
|
assert node.polls > 1
|
|
assert consumer_ran == [21.5]
|
|
|
|
|
|
def test_a_failing_poll_reports_down_and_keeps_going():
|
|
health: list[tuple[str, str | None]] = []
|
|
node = a_sensor([RuntimeError("device unplugged"), 21.5])
|
|
node._on_health = lambda _node, status, detail: health.append((status, detail))
|
|
Pipeline(nodes=[node])
|
|
|
|
run_briefly(node)
|
|
|
|
assert ("down", "RuntimeError: device unplugged") in health
|
|
# It recovered rather than giving up.
|
|
assert health[-1][0] == "ok"
|
|
|
|
|
|
def test_an_undeclared_port_keeps_failing_until_the_node_declares_it():
|
|
"""A publication that raised is retried, not remembered as published.
|
|
|
|
The loop remembers what it published. If it remembered what it read, a
|
|
value the node cannot publish would be skipped on the next poll, the poll
|
|
would succeed, and the node would go back to reporting itself healthy with
|
|
its port still dark.
|
|
"""
|
|
|
|
class Chatty(Sensor):
|
|
"""Reads a port it never declared."""
|
|
|
|
async def poll(self) -> dict[str, Any]:
|
|
self.polls += 1
|
|
return {"reading": 21.5, "lat": 48.1}
|
|
|
|
health: list[tuple[str, str | None]] = []
|
|
node = Chatty(
|
|
provides=[MessageSpec(name="reading", dtype=DType.FLOAT)],
|
|
params={"poll_interval": 0.01},
|
|
)
|
|
node.assign_flow("demo", "sensor")
|
|
node._on_health = lambda _node, status, detail: health.append((status, detail))
|
|
pipeline = Pipeline(nodes=[node])
|
|
|
|
run_briefly(node)
|
|
|
|
assert pipeline.state.get("demo.reading") is None
|
|
assert health[-1][0] == "down"
|
|
assert "NodeOutputError" in (health[-1][1] or "")
|
|
# Still failing on the last poll, not just the first.
|
|
assert len([entry for entry in health if entry[0] == "down"]) > 1
|
|
|
|
|
|
class Actuator(ConnectorNode):
|
|
"""A connector that commands something instead of reading it."""
|
|
|
|
contract = CONTRACT_VERSION
|
|
title = "Test actuator"
|
|
|
|
def __init__(self, **kwargs: Any) -> None:
|
|
super().__init__(**kwargs)
|
|
self.commands: list[dict[str, Any]] = []
|
|
|
|
def write(self, **ports: Any) -> None:
|
|
self.commands.append(ports)
|
|
return None
|
|
|
|
|
|
def test_an_incoming_message_reaches_a_connector_that_writes():
|
|
node = Actuator(requires=[MessageSpec(name="level", dtype=DType.INT)])
|
|
node.assign_flow("demo", "actuator")
|
|
|
|
assert node.execute({"demo.level": 255}) is None
|
|
assert node.commands == [{"level": 255}]
|
|
|
|
|
|
def test_a_read_only_connector_ignores_what_reaches_it():
|
|
node = a_sensor([])
|
|
assert node.execute({}) is None
|
|
|
|
|
|
def test_a_credential_param_is_marked_for_the_editor():
|
|
schema = Sensor.Params.model_json_schema()
|
|
assert schema["properties"]["poll_interval"]["default"] == 0
|
|
assert "secret_token" in schema["properties"]
|
|
|
|
|
|
def test_a_connector_is_discovered_from_its_entry_point(monkeypatch):
|
|
class FakeDist:
|
|
name = "fluksio-connector-test"
|
|
version = "0.1.0"
|
|
|
|
class FakeEntry:
|
|
name = "test_sensor"
|
|
dist = FakeDist()
|
|
|
|
def load(self):
|
|
return Sensor
|
|
|
|
monkeypatch.setattr(
|
|
"fluksio.flow.plugins.entry_points", lambda group: [FakeEntry()]
|
|
)
|
|
try:
|
|
assert load_plugins() == ["test_sensor"]
|
|
assert NODE_TYPES["test_sensor"].plugin == "fluksio-connector-test 0.1.0"
|
|
assert NODE_TYPES["test_sensor"].title == "Test sensor"
|
|
finally:
|
|
NODE_TYPES.pop("test_sensor", None)
|
|
|
|
|
|
def test_a_connector_written_for_another_contract_is_refused(monkeypatch):
|
|
class Outdated(ConnectorNode):
|
|
contract = CONTRACT_VERSION + 1
|
|
|
|
class FakeEntry:
|
|
name = "outdated"
|
|
dist = None
|
|
|
|
def load(self):
|
|
return Outdated
|
|
|
|
monkeypatch.setattr(
|
|
"fluksio.flow.plugins.entry_points", lambda group: [FakeEntry()]
|
|
)
|
|
assert load_plugins() == []
|
|
assert "outdated" not in NODE_TYPES
|
|
|
|
|
|
def test_a_connector_may_not_take_over_a_built_in_type(monkeypatch):
|
|
class FakeEntry:
|
|
name = "mqtt"
|
|
dist = None
|
|
|
|
def load(self): # pragma: no cover - never reached
|
|
raise AssertionError("should not be loaded")
|
|
|
|
monkeypatch.setattr(
|
|
"fluksio.flow.plugins.entry_points", lambda group: [FakeEntry()]
|
|
)
|
|
assert load_plugins() == []
|
|
assert NODE_TYPES["mqtt"].plugin is None
|
|
|
|
|
|
def test_a_connector_publishes_bytes_as_a_media_reference(tmp_path):
|
|
"""A camera's reading is bytes, and bytes never travel as a message."""
|
|
|
|
class Camera(ConnectorNode):
|
|
contract = CONTRACT_VERSION
|
|
|
|
async def poll(self) -> dict[str, Any] | None:
|
|
return {"frame": self.save_artifact(b"\x89PNG...", "f.png", "image/png")}
|
|
|
|
node = Camera(
|
|
provides=[MessageSpec(name="frame", dtype=DType.IMAGE)],
|
|
params={"poll_interval": 0.01},
|
|
)
|
|
node.assign_flow("demo", "camera")
|
|
pipeline = Pipeline(nodes=[node])
|
|
|
|
# Before it starts there is no store to write to, and saying so beats an
|
|
# AttributeError from inside somebody's connector.
|
|
with pytest.raises(RuntimeError, match="artifact store"):
|
|
node.save_artifact(b"x")
|
|
|
|
store = ArtifactStore(tmp_path / "artifacts")
|
|
app = SimpleNamespace(state=SimpleNamespace(artifact_store=store))
|
|
run_briefly(node, app=app)
|
|
|
|
reference = pipeline.state["demo.frame"]
|
|
# It typechecks against the port it was published on, which is the whole
|
|
# point of a media dtype.
|
|
MessageSpec(name="demo.frame", dtype=DType.IMAGE).check(reference)
|
|
assert store.path(reference["digest"]).read_bytes() == b"\x89PNG..."
|