diff --git a/backend/fluksio/flow/pipeline.py b/backend/fluksio/flow/pipeline.py index 2de189d..e08863c 100644 --- a/backend/fluksio/flow/pipeline.py +++ b/backend/fluksio/flow/pipeline.py @@ -1352,12 +1352,30 @@ class Pipeline: continue ready_inputs: dict[str, Any] | None = None if check_ready: + # The values are kept whether or not the answer is yes: + # a node that may run reads them below instead of + # asking state for the same keys again. ready, ready_inputs = self._is_node_ready(n, state) if not ready: - if n.synchronous: - # Not ready now; a later trigger may make it - # ready. - skipped.add(n) + # Completed, not merely passed over: a later trigger + # may still make it ready, but not running is not + # the same as blocking. Whatever it last published + # is in state, and a consumer reading that alongside + # something this wave *did* publish has every right + # to run on the pair. + # + # Passing over it held that consumer for the whole + # wave, and a node that can never be ready — an + # input nothing has ever published — took its + # consumers out for good. On this house + # `plugs.pump_run` waits on a watering pulse that + # only exists at 02:00, and it sits upstream of + # `dmx.switches`: every boiler, plug and appliance + # command reached its own message and stopped there, + # reaching the DMX universe only in the occasional + # wave the pump node happened not to be part of. + complete(n) + progressed = True continue if replay and entry_id and self._already_done(entry_id, n): # Its side effect happened on an earlier delivery; its @@ -1379,17 +1397,23 @@ class Pipeline: for n in completed: result = node_futures.pop(n).result() - # A node returning nothing (rate limiting, an error) stops - # propagation along its branch. - if result is not None: - # Before the decrement: a consumer freed by this node - # is judged on what it just published. - if fresh is not None: - fresh.update(result) - for consumer in edges[n]: - if consumer in target_nodes: - in_degree[consumer] -= 1 - submit_ready(executor) + # Before the decrement: a consumer freed by this node is + # judged on what it just published. + if result is not None and fresh is not None: + fresh.update(result) + # Decremented whatever it returned. A node that published + # nothing — rate limited, unchanged, or failed — has not + # *blocked* its consumers; it has merely given them nothing + # to read, and `untouched` above is the accurate test of + # that. Holding the count instead made a quiet producer + # veto a noisy sibling: `dmx.switches` reads both boilers + # through report-by-exception nodes, so the kitchen one + # staying silent — which is almost always — kept the main + # boiler's command from ever reaching the DMX universe. + for consumer in edges[n]: + if consumer in target_nodes: + in_degree[consumer] -= 1 + submit_ready(executor) if self._node_pool is not None: # The execution service owns a long-lived pool; building one per diff --git a/backend/tests/conftest.py b/backend/tests/conftest.py index 3c24418..f606e90 100644 --- a/backend/tests/conftest.py +++ b/backend/tests/conftest.py @@ -7,7 +7,7 @@ from typing import Any import pytest from cryptography.hazmat.primitives.asymmetric import rsa from fastapi.testclient import TestClient -from sqlalchemy import delete +from sqlalchemy import delete, inspect from sqlalchemy.engine import make_url from sqlmodel import Session, select @@ -89,6 +89,11 @@ def _clean_engine_rows() -> Generator[None, None, None]: between modules bounds it to where the writes actually come from. """ yield + # `tests/flow` overrides `db` with a no-op — the engine holds no database + # state, so those tests build no schema — and there is nothing to clear + # there. Running that directory on its own is how that shows. + if not inspect(engine).has_table("run"): + return with Session(engine) as session: for model in _ENGINE_TABLES: session.execute(delete(model)) diff --git a/backend/tests/flow/test_pipeline.py b/backend/tests/flow/test_pipeline.py index 65cb81b..af55cf6 100644 --- a/backend/tests/flow/test_pipeline.py +++ b/backend/tests/flow/test_pipeline.py @@ -206,3 +206,150 @@ def test_only_advisory_codes_are_flagged_advisory(): assert hook.advisory is True assert cycle.advisory is False + + +def test_a_producer_that_cannot_run_does_not_strand_its_consumer(): + """One wave, two producers of one consumer, and one of them never ready. + + The house's shape: a battery reading wakes the boiler arbiter and the + watering-pump trigger in the same wave, and the DMX encoder reads both. The + pump waits on a pulse that only exists at 02:00, so it is never ready — and + passing over it used to leave the encoder unreachable, which put every + boiler and plug command into its message and no further. + """ + ran = [] + + power = make_node( + "cerbo", "power", lambda params: {"batt_v": 50.3}, provides=[spec("batt_v")] + ) + + def arbitrate(batt_v, params): + ran.append("boiler") + return {"boiler_on": 1.0} + + boiler = make_node( + "arbiter", + "boiler", + arbitrate, + requires=[spec("power.batt_v")], + provides=[spec("boiler_on")], + ) + + # Woken by the same wave; `pump_for` has never been published. + pump = make_node( + "pump_run", + "plugs", + lambda batt_v, pump_for, params: {"pump": 1.0}, + requires=[spec("power.batt_v"), spec("plugs.pump_for")], + provides=[spec("pump")], + ) + + def encode(boiler_on, pump, params): + ran.append("encoder") + return {"dmx_boiler": boiler_on} + + encoder = make_node( + "switches", + "dmx", + encode, + requires=[spec("boiler.boiler_on"), spec("plugs.pump")], + provides=[spec("dmx_boiler")], + ) + + pipeline = Pipeline(nodes=[power, boiler, pump, encoder]) + pipeline.state["plugs.pump"] = 0.0 + + power.inject() + + assert ran == ["boiler", "encoder"] + assert pipeline.state["dmx.dmx_boiler"] == 1.0 + + +def test_a_quiet_producer_does_not_veto_a_noisy_sibling(): + """Two report-by-exception producers, one consumer, and only one with news. + + The house reads both boilers into one DMX encoder through `rbe` nodes. The + kitchen boiler is unchanged almost always, so it returns nothing — and + while that held the encoder's dependency count up, the main boiler's + command reached its message and got no further. + """ + ran = [] + + power = make_node( + "cerbo", "power", lambda params: {"batt_v": 50.3}, provides=[spec("batt_v")] + ) + water = make_node( + "water_changed", + "boiler", + lambda batt_v, params: {"water_boiler": 1.0}, + requires=[spec("power.batt_v")], + provides=[spec("water_boiler")], + ) + # Unchanged, so its rbe publishes nothing this wave. + kitchen = make_node( + "kitchen_changed", + "boiler", + lambda batt_v, params: None, + requires=[spec("power.batt_v")], + provides=[spec("kitchen_boiler")], + ) + + def encode(water_boiler, kitchen_boiler, params): + ran.append("switches") + return {"dmx_water_boiler": 255.0 if water_boiler else 0.0} + + encoder = make_node( + "switches", + "dmx", + encode, + requires=[spec("boiler.water_boiler"), spec("boiler.kitchen_boiler")], + provides=[spec("dmx_water_boiler")], + ) + + pipeline = Pipeline(nodes=[power, water, kitchen, encoder]) + pipeline.state["boiler.water_boiler"] = 0.0 + pipeline.state["boiler.kitchen_boiler"] = 0.0 + + power.inject() + + assert ran == ["switches"] + assert pipeline.state["dmx.dmx_water_boiler"] == 255.0 + + +def test_a_consumer_of_only_quiet_producers_still_does_not_run(): + """The other half of it: freeing a consumer is not the same as running it. + + Nothing it reads was refreshed, so it must stay put — otherwise every wave + would re-run the whole graph on values it has already seen. + """ + ran = [] + + power = make_node( + "cerbo", "power", lambda params: {"batt_v": 50.3}, provides=[spec("batt_v")] + ) + quiet = make_node( + "changed", + "boiler", + lambda batt_v, params: None, + requires=[spec("power.batt_v")], + provides=[spec("water_boiler")], + ) + + def encode(water_boiler, params): + ran.append("switches") + return {"dmx_water_boiler": 255.0} + + encoder = make_node( + "switches", + "dmx", + encode, + requires=[spec("boiler.water_boiler")], + provides=[spec("dmx_water_boiler")], + ) + + pipeline = Pipeline(nodes=[power, quiet, encoder]) + pipeline.state["boiler.water_boiler"] = 0.0 + + power.inject() + + assert ran == []