diff --git a/backend/fluksio/flow/controller.py b/backend/fluksio/flow/controller.py index 980d7e7..29cc1f1 100644 --- a/backend/fluksio/flow/controller.py +++ b/backend/fluksio/flow/controller.py @@ -85,6 +85,13 @@ NODE_STOP_TIMEOUT = 5.0 # that never ends. REBUILD_WAIT = 15.0 +# How long a caller waits for a *flow* rebuild that is already running. A +# per-flow rebuild reconnects one flow's nodes rather than the installation's, +# so this is a queueing budget — several of them back to back, which is what +# seeding does — not the room a single one needs. Fifteen seconds was sized for +# the whole-pipeline rebuild and would let a wedge sit unreported. +FLOW_REBUILD_WAIT = 5.0 + class RebuildBusy(RuntimeError): """A rebuild could not start because the one before it has not finished. @@ -343,6 +350,9 @@ class FlowController: self.disabled: set[str] = set() #: Flows that only run when a run asks them to. self.batch: set[str] = set() + #: Each flow's declared inputs, kept so rebuilding one flow can + #: validate the whole graph without re-reading every flow off the disk. + self._flow_inputs: dict[str, dict[str, bool]] = {} self.supervisor = Supervisor(events) self.history_limits: dict[str, int] = {} self._lock = asyncio.Lock() @@ -420,22 +430,27 @@ class FlowController: await run_in_threadpool(self.store.write_enabled, flow, enabled) await self.reload() - async def reload(self, wait: float | None = REBUILD_WAIT) -> None: - """Rebuild the whole pipeline from what is currently stored. + async def _acquire(self, wait: float | None) -> None: + """Take the rebuild lock, or say the engine is busy. - Only one rebuild runs at a time. A caller waits *wait* seconds for the - one in front of it and then gives up with ``RebuildBusy`` — hanging on - a rebuild that is stuck is worse than saying so. ``None`` waits. + Only one rebuild runs at a time whatever its scope. A caller waits + *wait* seconds for the one in front of it and then gives up with + ``RebuildBusy`` — hanging on a rebuild that is stuck is worse than + saying so. ``None`` waits. """ if wait is None: await self._lock.acquire() - else: - try: - await asyncio.wait_for(self._lock.acquire(), wait) - except asyncio.TimeoutError: - raise RebuildBusy( - f"A pipeline rebuild is still running after {wait:.0f}s" - ) from None + return + try: + await asyncio.wait_for(self._lock.acquire(), wait) + except asyncio.TimeoutError: + raise RebuildBusy( + f"A pipeline rebuild is still running after {wait:.0f}s" + ) from None + + async def reload(self, wait: float | None = REBUILD_WAIT) -> None: + """Rebuild the whole pipeline from what is currently stored.""" + await self._acquire(wait) try: # Work already claimed belongs to the pipeline it was claimed # against; let it finish there before swapping the graph out. @@ -463,6 +478,9 @@ class FlowController: nodes, loaded, initial_values, flow_inputs = await run_in_threadpool( self._build_flows, [(flow, False) for flow in published] ) + self._flow_inputs = { + flow.name: _declared_inputs(flow)[0] for flow in published + } # A rebuild is a fresh set of nodes, but not a fresh history: every # publish rebuilds every flow, so dropping the failures here would @@ -510,6 +528,159 @@ class FlowController: } ) + async def reload_flow( + self, name: str, wait: float | None = FLOW_REBUILD_WAIT + ) -> None: + """Rebuild one flow, leaving every other flow's nodes connected. + + A flow is not a subgraph — its nodes can read and write messages + another flow owns — but the wiring is derived from message names, so + swapping one flow's nodes into the graph and deriving the edges again + is enough. What that saves is the reconnecting: the cost of a rebuild + on a populated installation is every node opening its socket again, + and only one flow's have changed. + + A flow the store no longer has is taken out instead of replaced. + + The failure shape is the whole-pipeline one narrowed to a flow, and + deliberately no better: a build that raises changes nothing, and + anything failing after the teardown leaves that flow's nodes stopped + in the graph — which is exactly what a failed ``reload`` leaves behind + for all of them. + """ + pipeline = self.pipeline + if pipeline is None: + # Nothing built yet, so there is nothing to splice one flow into. + await self.reload(wait=wait) + return + + await self._acquire(wait) + try: + # Built before anything is stopped, so the window in which the + # flow is not running is its own teardown and nothing more. Off + # the loop for the same reason the full build is: compiling a + # python node waits for a worker slot. + built = await run_in_threadpool(self._build_one, name) + await self._teardown(name) + + if built is None: + pipeline.remove_flow(name) + self.loaded = { + node_id: entry + for node_id, entry in self.loaded.items() + if entry.flow != name + } + # Or the flow would go on declaring inputs nothing provides + # any more, and the consumers it left behind would look fine. + self._flow_inputs.pop(name, None) + self.disabled = self.disabled - {name} + self.batch = self.batch - {name} + else: + flow, nodes, loaded, initial_values, enabled = built + pipeline.replace_flow(name, nodes, initial_values) + + # A rebuild is a fresh set of nodes, but not a fresh history: + # a failure nobody has dismissed is the operator's to keep. + for node_id, entry in loaded.items(): + previous = self.loaded.get(node_id) + if previous is not None and previous.last_error: + entry.last_error = previous.last_error + entry.last_error_ts = previous.last_error_ts + + # Rebound rather than mutated: request threads and the failure + # watcher read these without holding anything. + self.loaded = { + **{ + node_id: entry + for node_id, entry in self.loaded.items() + if entry.flow != name + }, + **loaded, + } + self._flow_inputs[name] = _declared_inputs(flow)[0] + self.disabled = ( + self.disabled - {name} if enabled else self.disabled | {name} + ) + self.batch = ( + self.batch | {name} if flow.mode == "batch" else self.batch - {name} + ) + + pipeline.set_disabled(self.disabled) + # Over the whole graph, because it has to be: taking a producer out + # of one flow is what leaves another flow's input unconnected. + self.issues = _collect_issues( + self.loaded, pipeline, self._all_flow_inputs() + ) + + await self._activate(name) + # This flow's pause is cleared by the replace, so no resume will + # ever come for what it parked. Release it here or it is lost. + self._release_parked(name) + finally: + self._lock.release() + + self._publish( + { + "type": "pipeline_rebuilt", + "issues": [issue.model_dump() for issue in self.issues], + "nodes": [status.model_dump() for status in self.node_statuses()], + # Every flow's, not this one's: the canvas replaces its paused + # set from this, so a narrowed list would clear the markers of + # flows the rebuild never touched. + "paused": self.paused_flows(), + } + ) + + async def reload_failed_flows(self) -> list[str]: + """Rebuild only the flows holding a node that would not load. + + What installing a module wants: a node red because its import was + missing is built again against the packages just installed, and every + other flow keeps the connections it already has. + """ + broken = sorted( + { + entry.flow + for entry in self.loaded.values() + if entry.status is NodeStatus.ERROR + } + ) + for flow in broken: + await self.reload_flow(flow) + return broken + + async def reload_lib_users(self, ref: str) -> list[str]: + """Rebuild the flows running a shared node whose source has changed.""" + flows = sorted({usage.split(".", 1)[0] for usage in self.store.usages(ref)}) + for flow in flows: + await self.reload_flow(flow) + return flows + + def _build_one( + self, name: str + ) -> tuple[FlowDef, list[Node], dict[str, LoadedNode], dict[str, Any], bool] | None: + """Read one published flow and build its nodes. None when it is gone. + + Blocking: this is the part of a rebuild that talks to the worker pool. + """ + try: + flow = self.store.read_flow(name) + except FlowNotFound: + return None + nodes, loaded, initial_values, _ = self._build_flows([(flow, False)]) + return flow, nodes, loaded, initial_values, self.store.read_enabled(name) + + def _all_flow_inputs(self) -> dict[str, bool]: + """Every flow's declared inputs, merged the way one build would see them. + + In name order, which is the order ``read_all`` returns flows in, so a + name two flows both declare resolves to the same one either way. + """ + merged: dict[str, bool] = {} + for name in sorted(self._flow_inputs): + merged.update(self._flow_inputs[name]) + return merged + async def _teardown(self, flow: str | None = None) -> None: """Stop everything the previous pipeline started, or one flow's share.""" for entry in self.loaded.values(): @@ -583,13 +754,9 @@ class FlowController: loaded[entry.id] = entry if entry.node is not None: nodes.append(entry.node) - for flow_input in flow.inputs: - name = qualify(flow.name, flow_input.spec.name) - if not name: - continue - flow_inputs[name] = flow_input.initial is not None - if flow_input.initial is not None: - initial_values[name] = flow_input.initial + declared, initial = _declared_inputs(flow) + flow_inputs.update(declared) + initial_values.update(initial) return nodes, loaded, initial_values, flow_inputs @@ -1225,6 +1392,20 @@ def with_settings( return call +def _declared_inputs(flow: FlowDef) -> tuple[dict[str, bool], dict[str, Any]]: + """A flow's declared inputs, and the ones that start with a value.""" + declared: dict[str, bool] = {} + initial: dict[str, Any] = {} + for flow_input in flow.inputs: + name = qualify(flow.name, flow_input.spec.name) + if not name: + continue + declared[name] = flow_input.initial is not None + if flow_input.initial is not None: + initial[name] = flow_input.initial + return declared, initial + + def _collect_issues( loaded: dict[str, LoadedNode], pipeline: Pipeline, diff --git a/backend/tests/flow/test_runtime_control.py b/backend/tests/flow/test_runtime_control.py index 2ce3738..19177b1 100644 --- a/backend/tests/flow/test_runtime_control.py +++ b/backend/tests/flow/test_runtime_control.py @@ -10,11 +10,16 @@ from fastapi.testclient import TestClient from fluksio.api.deps import get_current_user from fluksio.api.routes.flows import router -from fluksio.flow.controller import FlowController, LoadedNode +from fluksio.flow.controller import ( + NODE_TYPES, + FlowController, + LoadedNode, + NodeType, +) 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 +from fluksio.flow.schemas import FlowDef, FlowInput, NodeDef from fluksio.flow.store import FlowStore @@ -198,3 +203,97 @@ def test_a_node_that_will_not_stop_does_not_wedge_the_rebuild( await engine.reload() asyncio.run(asyncio.wait_for(scenario(), timeout=5)) + + +class _Lifecycle(Node): + """A node type that writes down every lifecycle call a rebuild makes.""" + + calls: list[str] = [] + + def __init__(self, **kwargs: Any) -> None: + super().__init__(f=lambda params: None, **kwargs) + + async def start(self, app: FastAPI | None = None) -> None: + _Lifecycle.calls.append(f"start {self.id}") + + async def stop(self, app: FastAPI | None = None) -> None: + _Lifecycle.calls.append(f"stop {self.id}") + + +@pytest.fixture +def lifecycle(monkeypatch: pytest.MonkeyPatch) -> list[str]: + """A node type on the registry, and the log of what it was asked to do.""" + monkeypatch.setitem( + NODE_TYPES, + "lifecycle", + NodeType( + title="Lifecycle", description="Counts start and stop.", cls=_Lifecycle + ), + ) + _Lifecycle.calls.clear() + return _Lifecycle.calls + + +def test_rebuilding_one_flow_leaves_another_flows_node_running( + tmp_path: Path, lifecycle: list[str] +): + """The whole point: reconnecting one flow's nodes, not the installation's.""" + store = FlowStore(tmp_path / "flows") + for name in ("a", "b"): + store.write_flow(FlowDef(name=name, nodes=[NodeDef(id="io", type="lifecycle")])) + engine = FlowController(store) + + async def scenario() -> None: + await engine.reload() + untouched = engine.loaded["b.io"].node + lifecycle.clear() + + await engine.reload_flow("a") + + # Not rebuilt, not restarted, and the very same object still bound. + assert engine.loaded["b.io"].node is untouched + assert lifecycle == ["stop a.io", "start a.io"] + assert engine.loaded["a.io"].node is not None + + asyncio.run(asyncio.wait_for(scenario(), timeout=10)) + + +def test_a_deleted_flow_leaves_its_consumers_reporting_a_missing_input( + tmp_path: Path, lifecycle: list[str] +): + """A flow is not a subgraph, so deleting one is another flow's problem.""" + store = FlowStore(tmp_path / "flows") + store.write_flow( + FlowDef( + name="a", + nodes=[NodeDef(id="meter", type="lifecycle", provides=[spec("temp")])], + inputs=[FlowInput(spec=spec("spare"), initial=1.0)], + ) + ) + store.write_flow( + FlowDef( + name="b", + nodes=[ + NodeDef(id="load", type="lifecycle", requires=[spec("a.temp")]), + NodeDef(id="watch", type="lifecycle", requires=[spec("a.spare")]), + ], + ) + ) + engine = FlowController(store) + + async def scenario() -> None: + await engine.reload() + assert engine.issues == [] + + store.delete_flow("a") + await engine.reload_flow("a") + + # Both the producer and the declared input went with the flow. + assert sorted((issue.code, issue.node) for issue in engine.issues) == [ + ("unconnected_input", "b.load"), + ("unconnected_input", "b.watch"), + ] + assert engine.pipeline is not None + assert engine.pipeline.get_node_by_id("a.meter") is None + + asyncio.run(asyncio.wait_for(scenario(), timeout=10))