Files
app/backend/tests/flow/test_queue.py
T
stroblmeandClaude Opus 5 2554488a73 Backend: real alert test results, state cleanup on delete/rename, node trigger errors, queue and collector fixes
`AlertManager.send` swallowed every delivery failure, so the alerts screen's
Test button answered 200 whatever happened — the one thing it exists for. It
takes `raise_on_error` now, which only the test route passes; the per-channel
loop keeps the swallow, because one dead channel must not stop the others
hearing about the same fault. A refused delivery answers 502 with whatever the
sender said.

Renaming a flow left its values under the old name for good: the delete path
already swept them, the rename path never did. It calls the same `forget_flow`,
which covers the messages and the `__ts__`/`__version__`/`__history__`
bookkeeping keyed by message name. Cleanup, not migration — they repopulate
under the new name on the next run.

Triggering a node by hand ran `Node.__call__` with nothing catching it, so a
node that raised produced a 500 and a stack trace in the server log, and
nothing at all on the canvas. `Pipeline.publish_error` is the reporting half of
`_execute_node` lifted out; both paths go through it, so a manual failure now
reads the same on the canvas and in the metrics as a queued one. The route
answers 400 with the node's error.

`MemoryWorkQueue.stats()` counts claimed-but-unacknowledged work rather than
reporting zero, so the health tile means something without Redis. The metrics
collector's held tracebacks are capped at `DETAIL_CAP` and swept on the same
`RUN_STALE_S` cutoff the open runs use, instead of one untruncated traceback
per node kept for the life of the process — a traceback still survives the
flush between the log and the failure it belongs to.

`GET /observability/events` takes `since`/`until`, the window `/runs` already
took, so a failures list can cover the span the charts beside it are drawn from.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01XC2jX6Hdj7pxGGKzBTrbqB
2026-08-17 11:33:25 +02:00

265 lines
7.8 KiB
Python

"""The work queue, and what the execution service does with it."""
import time
from app.flow.executor import ExecutionService
from app.flow.messages import DType, MessageSpec
from app.flow.nodes import Node
from app.flow.pipeline import Pipeline
from app.flow.queue import MemoryWorkQueue, WorkItem
from app.flow.state import MemoryState
def test_items_come_back_in_the_order_they_went_in():
queue = MemoryWorkQueue()
for i in range(3):
queue.add(WorkItem(kind="cascade", node=f"f.n{i}", flow="f"))
claimed = queue.claim(10, 10)
assert [item.node for item in claimed] == ["f.n0", "f.n1", "f.n2"]
# Every item gets an id, which is what idempotency markers hang off.
assert all(item.entry_id for item in claimed)
def test_claiming_an_empty_queue_waits_and_gives_up():
queue = MemoryWorkQueue()
started = time.monotonic()
assert queue.claim(1, 50) == []
assert time.monotonic() - started >= 0.04
def test_a_delayed_item_stays_put_until_it_is_due():
queue = MemoryWorkQueue()
queue.add_delayed(WorkItem(kind="cascade", node="f.n", flow="f"), time.time() + 60)
assert queue.claim(1, 10) == []
assert queue.move_due(time.time()) == 0
assert queue.move_due(time.time() + 61) == 1
assert [i.node for i in queue.claim(1, 10)] == ["f.n"]
def test_parked_work_comes_back_oldest_first():
queue = MemoryWorkQueue()
for i in range(3):
queue.park("heating", WorkItem(kind="cascade", node=f"f.n{i}", flow="heating"))
assert [i.node for i in queue.unpark("heating")] == ["f.n0", "f.n1", "f.n2"]
# Unparking empties it, so a second resume does not replay the same work.
assert queue.unpark("heating") == []
def test_a_deleted_flow_leaves_nothing_parked():
queue = MemoryWorkQueue()
queue.park("gone", WorkItem(kind="cascade", node="gone.n", flow="gone"))
queue.clear_flow("gone")
assert queue.unpark("gone") == []
def test_claimed_work_counts_as_in_flight_until_it_is_acknowledged():
"""The health tile's "in flight" reads zero without this."""
queue = MemoryWorkQueue()
queue.add(WorkItem(kind="cascade", node="f.n", flow="f"))
assert queue.stats()["pending"] == 0
(item,) = queue.claim(1, 10)
assert queue.stats()["pending"] == 1
queue.ack(item)
assert queue.stats()["pending"] == 0
def _pipeline_with_a_consumer() -> tuple[Pipeline, Node, MemoryState, list]:
"""A source whose message a consumer records."""
seen: list[float] = []
def consume(reading, params):
seen.append(reading)
return {"doubled": reading * 2}
source = Node(
f=lambda params: None,
provides=[MessageSpec(name="reading", port="reading", dtype=DType.FLOAT)],
name="source",
)
consumer = Node(
f=consume,
requires=[MessageSpec(name="reading", port="reading", dtype=DType.FLOAT)],
provides=[MessageSpec(name="doubled", port="doubled", dtype=DType.FLOAT)],
name="consumer",
)
source.assign_flow("f", "source")
consumer.assign_flow("f", "consumer")
state = MemoryState()
queue = MemoryWorkQueue()
pipeline = Pipeline(nodes=[source, consumer], state=state, work_queue=queue)
return pipeline, source, state, seen
def test_a_trigger_is_journaled_rather_than_run_on_the_spot():
pipeline, source, state, seen = _pipeline_with_a_consumer()
source.inject({"reading": 3.0})
# Nothing ran yet: the value is in the queue, not in state.
assert seen == []
assert "f.reading" not in state
service = ExecutionService(pipeline._queue)
service.bind(pipeline)
for item in pipeline._queue.claim(10, 10):
service._run_item(item)
assert seen == [3.0]
assert state["f.doubled"] == 6.0
def test_work_for_a_paused_flow_is_held_and_released_on_resume():
pipeline, source, state, seen = _pipeline_with_a_consumer()
service = ExecutionService(pipeline._queue)
service.bind(pipeline)
pipeline.pause("f")
source.inject({"reading": 1.0})
for item in pipeline._queue.claim(10, 10):
service._run_item(item)
assert seen == []
pipeline.resume("f")
for item in pipeline._queue.unpark("f"):
pipeline._queue.add(item)
for item in pipeline._queue.claim(10, 10):
service._run_item(item)
assert seen == [1.0]
def test_a_step_runs_one_held_item_and_leaves_the_flow_paused():
pipeline, source, _state, seen = _pipeline_with_a_consumer()
service = ExecutionService(pipeline._queue)
service.bind(pipeline)
pipeline.pause("f")
source.inject({"reading": 1.0})
source.inject({"reading": 2.0})
for item in pipeline._queue.claim(10, 10):
service._run_item(item)
assert seen == []
assert service.step("f") == "f.source"
assert seen == [1.0]
# Still paused, and the second value is still waiting for the next step.
assert pipeline.is_paused("f")
assert [i.outputs for i in pipeline._queue.unpark("f")] == [{"f.reading": 2.0}]
def test_stepping_a_flow_with_nothing_held_says_so_rather_than_failing():
pipeline, _source, _state, _seen = _pipeline_with_a_consumer()
service = ExecutionService(pipeline._queue)
service.bind(pipeline)
pipeline.pause("f")
assert service.step("f") is None
def test_work_for_a_stopped_flow_is_dropped():
pipeline, source, state, seen = _pipeline_with_a_consumer()
stopped = Pipeline(
nodes=pipeline.nodes,
state=pipeline.state,
work_queue=pipeline._queue,
disabled_flows={"f"},
)
service = ExecutionService(pipeline._queue)
service.bind(stopped)
# Reaching the queue at all takes a direct add: trigger drops it earlier.
stopped._queue.add(
WorkItem(kind="cascade", node="f.source", flow="f", outputs={"f.reading": 1.0})
)
for item in stopped._queue.claim(10, 10):
service._run_item(item)
assert seen == []
def test_an_item_that_keeps_coming_back_is_dead_lettered():
pipeline, _source, _state, seen = _pipeline_with_a_consumer()
service = ExecutionService(pipeline._queue)
service.bind(pipeline)
item = WorkItem(
kind="cascade",
node="f.source",
flow="f",
outputs={"f.reading": 1.0},
deliveries=4,
)
service._run_item(item)
# Given up on rather than run again, so a poison item cannot loop forever.
assert seen == []
def test_an_item_for_a_node_that_no_longer_exists_is_dropped():
pipeline, _source, _state, seen = _pipeline_with_a_consumer()
service = ExecutionService(pipeline._queue)
service.bind(pipeline)
service._run_item(WorkItem(kind="cascade", node="f.removed", flow="f"))
assert seen == []
def test_a_replayed_item_does_not_repeat_a_side_effect():
"""At-least-once delivery must not mean two of the same outgoing request."""
calls: list[float] = []
def send(reading, params):
calls.append(reading)
return None
source = Node(
f=lambda params: None,
provides=[MessageSpec(name="reading", port="reading", dtype=DType.FLOAT)],
name="source",
)
class SendingNode(Node):
"""Stands in for the built-ins that reach outside."""
idempotent = False
sender = SendingNode(
f=send,
requires=[MessageSpec(name="reading", port="reading", dtype=DType.FLOAT)],
name="sender",
)
source.assign_flow("f", "source")
sender.assign_flow("f", "sender")
queue = MemoryWorkQueue()
pipeline = Pipeline(nodes=[source, sender], state=MemoryState(), work_queue=queue)
service = ExecutionService(queue)
service.bind(pipeline)
source.inject({"reading": 5.0})
(item,) = queue.claim(10, 10)
service._run_item(item)
assert calls == [5.0]
# The same item again, as a reaper would hand it back after a crash.
item.deliveries = 2
service._run_item(item)
assert calls == [5.0]