Docs / docs (push) Successful in 20s
Playwright Tests / test-playwright (1, 2) (push) Failing after 1m59s
Playwright Tests / test-playwright (2, 2) (push) Failing after 1m40s
pre-commit / pre-commit (push) Failing after 2m49s
Test Backend / test-backend (push) Successful in 2m20s
Compose Smoke Test / test-compose (push) Successful in 30s
Playwright Tests / merge-reports (push) Failing after 1m6s
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01UytviPMJbXzD8P84nLvXcq
200 lines
6.2 KiB
Python
200 lines
6.2 KiB
Python
"""The wiring fundamentals: name binding, fan-in, namespaces, validation."""
|
|
|
|
from fluksio.flow.controller import _declared_inputs
|
|
from fluksio.flow.events import EventBus
|
|
from fluksio.flow.messages import DType, MessageSpec
|
|
from fluksio.flow.nodes import Node
|
|
from fluksio.flow.pipeline import Pipeline
|
|
from fluksio.flow.schemas import FlowDef, FlowInput
|
|
|
|
|
|
def spec(name: str, dtype: DType = DType.FLOAT, port: str = "") -> MessageSpec:
|
|
return MessageSpec(name=name, dtype=dtype, port=port)
|
|
|
|
|
|
def make_node(node_id: str, flow: str, f, requires=(), provides=()) -> Node:
|
|
node = Node(f=f, requires=list(requires), provides=list(provides), name=node_id)
|
|
node.assign_flow(flow, node_id)
|
|
return node
|
|
|
|
|
|
def test_bare_names_are_scoped_to_their_flow():
|
|
source = make_node(
|
|
"source", "heating", lambda params: {"temp": 20.0}, provides=[spec("temp")]
|
|
)
|
|
assert "heating.temp" in source.provides
|
|
|
|
|
|
def test_consumer_receives_from_any_producer():
|
|
# Two producers of one message: each publication reaches the consumer.
|
|
seen = []
|
|
|
|
def emit_a(params):
|
|
return {"temp": 1.0}
|
|
|
|
def emit_b(params):
|
|
return {"temp": 2.0}
|
|
|
|
def consume(temp, params):
|
|
seen.append(temp)
|
|
return None
|
|
|
|
a = make_node("a", "heating", emit_a, provides=[spec("temp")])
|
|
b = make_node("b", "heating", emit_b, provides=[spec("temp")])
|
|
c = make_node("c", "heating", consume, requires=[spec("temp")])
|
|
|
|
pipeline = Pipeline(nodes=[a, b, c])
|
|
|
|
assert pipeline.produces["heating.temp"] == [a, b]
|
|
assert pipeline.dependencies[c] == frozenset({a, b})
|
|
|
|
a.inject()
|
|
b.inject()
|
|
|
|
assert seen == [1.0, 2.0]
|
|
# Latest value wins.
|
|
assert pipeline.state["heating.temp"] == 2.0
|
|
|
|
|
|
def test_flows_connect_through_qualified_names():
|
|
def emit(params):
|
|
return {"power": 500.0}
|
|
|
|
def consume(power, params):
|
|
return {"used": power}
|
|
|
|
producer = make_node("meter", "solar", emit, provides=[spec("power")])
|
|
# A bare name would be heating.power; the dotted one crosses the flow.
|
|
consumer = make_node(
|
|
"load",
|
|
"heating",
|
|
consume,
|
|
requires=[spec("solar.power")],
|
|
provides=[spec("used")],
|
|
)
|
|
|
|
pipeline = Pipeline(nodes=[producer, consumer])
|
|
assert pipeline.dependencies[consumer] == frozenset({producer})
|
|
|
|
producer.inject()
|
|
assert pipeline.state["heating.used"] == 500.0
|
|
|
|
|
|
def test_ports_keep_function_arguments_local():
|
|
def convert(celsius, params):
|
|
return {"fahrenheit": celsius * 9 / 5 + 32}
|
|
|
|
node = make_node(
|
|
"convert",
|
|
"heating",
|
|
convert,
|
|
requires=[spec("solar.celsius", port="celsius")],
|
|
provides=[spec("fahrenheit")],
|
|
)
|
|
Pipeline(nodes=[node])
|
|
|
|
assert node.execute({"solar.celsius": 100.0}) == {"heating.fahrenheit": 212.0}
|
|
|
|
|
|
def test_validate_reports_cycles_and_dangling_inputs():
|
|
a = make_node(
|
|
"a", "f", lambda b, params: {"a": b}, requires=[spec("b")], provides=[spec("a")]
|
|
)
|
|
b = make_node(
|
|
"b", "f", lambda a, params: {"b": a}, requires=[spec("a")], provides=[spec("b")]
|
|
)
|
|
lonely = make_node(
|
|
"lonely", "f", lambda missing, params: None, requires=[spec("missing")]
|
|
)
|
|
|
|
issues = Pipeline(nodes=[a, b, lonely]).validate()
|
|
codes = {issue.code for issue in issues}
|
|
|
|
assert "cycle" in codes
|
|
assert "unconnected_input" in codes
|
|
dangling = next(i for i in issues if i.code == "unconnected_input")
|
|
assert dangling.message_name == "f.missing"
|
|
assert dangling.node == "f.lonely"
|
|
|
|
|
|
def test_declared_flow_inputs_are_not_dangling():
|
|
node = make_node(
|
|
"n", "f", lambda setpoint, params: None, requires=[spec("setpoint")]
|
|
)
|
|
issues = Pipeline(nodes=[node]).validate({"f.setpoint": True})
|
|
assert issues == []
|
|
|
|
|
|
def test_a_flow_input_without_a_starting_value_is_reported():
|
|
node = make_node(
|
|
"n", "f", lambda setpoint, params: None, requires=[spec("setpoint")]
|
|
)
|
|
issues = Pipeline(nodes=[node]).validate({"f.setpoint": False})
|
|
assert [issue.code for issue in issues] == ["missing_initial_value"]
|
|
|
|
|
|
def test_a_batch_flows_input_is_a_run_parameter_not_a_missing_value():
|
|
"""A flow that only runs when asked gets its inputs from the run.
|
|
|
|
Calling that a message nothing ever sets had the health summary count a
|
|
training flow as one that cannot run while the Runs screen showed it
|
|
running. A live flow, which nothing is going to start on its own, still
|
|
reports it.
|
|
"""
|
|
node = make_node(
|
|
"n", "f", lambda setpoint, params: None, requires=[spec("setpoint")]
|
|
)
|
|
batch = FlowDef(name="f", mode="batch", inputs=[FlowInput(spec=spec("setpoint"))])
|
|
live = batch.model_copy(update={"mode": "live"})
|
|
|
|
assert Pipeline(nodes=[node]).validate(_declared_inputs(batch)[0]) == []
|
|
issues = Pipeline(nodes=[node]).validate(_declared_inputs(live)[0])
|
|
assert [issue.code for issue in issues] == ["missing_initial_value"]
|
|
|
|
|
|
def test_a_failing_node_does_not_stop_its_siblings():
|
|
ran = []
|
|
|
|
def boom(params):
|
|
raise RuntimeError("nope")
|
|
|
|
def fine(params):
|
|
ran.append("fine")
|
|
return {"ok": 1.0}
|
|
|
|
bad = make_node("bad", "f", boom, provides=[spec("bad_out")])
|
|
good = make_node("good", "f", fine, provides=[spec("ok")])
|
|
|
|
pipeline = Pipeline(nodes=[bad, good])
|
|
pipeline.run()
|
|
|
|
assert ran == ["fine"]
|
|
assert pipeline.state["f.ok"] == 1.0
|
|
|
|
|
|
def test_a_node_returning_something_other_than_a_dict_says_what_is_wrong():
|
|
"""Outputs are keyed by port, so a bare value cannot be one of them."""
|
|
events = []
|
|
|
|
def wrong(params):
|
|
return 42.0
|
|
|
|
bus = EventBus()
|
|
bus.publish = events.append # type: ignore[method-assign]
|
|
node = make_node("n", "f", wrong, provides=[spec("out")])
|
|
Pipeline(nodes=[node], events=bus).run()
|
|
|
|
(error,) = [e for e in events if e["type"] == "node_error"]
|
|
assert "NodeOutputError" in error["error"]
|
|
assert "returned float" in error["error"]
|
|
|
|
|
|
def test_values_carry_timestamps():
|
|
node = make_node("n", "f", lambda params: {"out": 1.0}, provides=[spec("out")])
|
|
pipeline = Pipeline(nodes=[node])
|
|
node.inject()
|
|
|
|
values = pipeline.values("f")
|
|
assert values["f.out"]["value"] == 1.0
|
|
assert values["f.out"]["ts"] > 0
|