Journal work before running it, so a crash stops losing messages

Execution was fire-and-forget: an MQTT message or webhook ran a cascade
on a ThreadPoolExecutor built for that one wave, and an engine that died
halfway through simply lost whatever was in flight. Concurrent triggers
each built their own pool, so load meant unbounded threads.

Every external trigger is now journaled to a Redis Streams queue before
anything runs, and acknowledged only once its cascade finishes. A
consumer thread drives cascades on one long-lived pool while node bodies
run on another, so a cascade cannot starve the nodes it is waiting for.
A reaper reclaims what a dead consumer never acknowledged — verified end
to end: work journaled while the engine was stopped runs on restart, and
work abandoned mid-cascade comes back as a second delivery.

At-least-once needs a guard, so nodes that reach outside are marked
non-idempotent and skipped on a redelivery they already completed.
Without Redis the queue degrades to an in-memory one that does not
pretend to be durable, and interactive callers still run inline.

Also fixes two things this turned up: a delay node was sleeping on a
worker thread, where a handful of them could occupy the whole pool, and
webhooks 404'd whenever MCP was enabled because the app mounted at /
answered first for every path.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_011LF61rxW1FG5YCD2J9YqjY
This commit is contained in:
root
2026-08-16 07:57:07 +02:00
co-authored by Claude Fable 5
parent 2cc25970e3
commit 04329149b3
17 changed files with 1271 additions and 66 deletions
+2
View File
@@ -311,6 +311,8 @@ async def delete_flow(name: str, controller: FlowControllerDep) -> Any:
await run_in_threadpool(controller.store.delete_flow, name)
except FlowNotFound:
raise HTTPException(status_code=404, detail=f"No flow named '{name}'")
# Its files are gone; its values and queued work would otherwise linger.
await run_in_threadpool(controller.forget_flow, name)
await controller.reload()
return Message(message=f"Deleted flow '{name}'")
+56 -5
View File
@@ -23,6 +23,7 @@ from fastapi import FastAPI
from fastapi.concurrency import run_in_threadpool
from app.flow.events import EventBus
from app.flow.executor import ExecutionService
from app.flow.messages import MessageSpec, qualify
from app.flow.nodes import (
DelayNode,
@@ -197,12 +198,15 @@ class FlowController:
events: EventBus | None = None,
max_workers: int | None = None,
fastapi_app: FastAPI | None = None,
execution: ExecutionService | None = None,
) -> None:
self.store = store
self.state = state if state is not None else MemoryState()
self.events = events
self.max_workers = max_workers
self.app = fastapi_app
# Without one, every trigger runs inline where it was raised.
self.execution = execution
self.pipeline: Pipeline | None = None
self.loaded: dict[str, LoadedNode] = {}
@@ -216,10 +220,17 @@ class FlowController:
# -------------------------------------------------------------------------
async def start(self) -> None:
# Build first: the consumer must have a pipeline to execute against
# before it claims anything, or work waiting from the last run would be
# taken and dropped — which is the very case the queue exists for.
await self.reload()
if self.execution is not None:
self.execution.start()
async def stop(self) -> None:
await self._teardown()
if self.execution is not None:
await run_in_threadpool(self.execution.stop)
async def set_enabled(self, flow: str, enabled: bool) -> None:
"""Stop or start one flow. Rebuilding is what applies it."""
@@ -229,6 +240,10 @@ class FlowController:
async def reload(self) -> None:
"""Rebuild the whole pipeline from what is currently stored."""
async with self._lock:
# Work already claimed belongs to the pipeline it was claimed
# against; let it finish there before swapping the graph out.
if self.execution is not None:
await run_in_threadpool(self.execution.pause_intake)
await self._teardown()
# A fresh supervisor per build, so a flow quarantined by the last
# one gets another chance once its author has changed something.
@@ -252,11 +267,18 @@ class FlowController:
max_workers=self.max_workers,
initial_values=initial_values,
disabled_flows=self.disabled,
work_queue=self.execution.queue if self.execution else None,
node_pool=self.execution.node_pool if self.execution else None,
)
if self.execution is not None:
self.execution.bind(self.pipeline)
self.issues = _collect_issues(loaded, self.pipeline, flow_inputs)
await self._activate()
if self.execution is not None:
self.execution.resume_intake()
self._publish(
{
"type": "pipeline_rebuilt",
@@ -498,8 +520,33 @@ class FlowController:
def resume_flow(self, flow: str) -> None:
"""Blocking — call from a worker thread: held-back nodes run on resume."""
if self.pipeline is not None:
self.pipeline.resume(flow)
if self.pipeline is None:
return
self.pipeline.resume(flow)
if self.execution is not None:
# Whatever arrived while the flow was held is queued again, oldest
# first, so a pause loses nothing.
for item in self.execution.queue.unpark(flow):
self.execution.queue.add(item)
def queue_stats(self) -> dict[str, Any]:
return self.execution.stats() if self.execution is not None else {}
def forget_flow(self, flow: str) -> None:
"""Drop what a deleted flow left behind. Blocking."""
if self.execution is not None:
self.execution.queue.clear_flow(flow)
prefix = f"{flow}."
with self.state.lock():
stale = [
key
for key in self.state.keys()
# Both the messages themselves and the engine's own bookkeeping
# about them, which is keyed by message name too.
if key.startswith(prefix) or f":{prefix}" in key
]
for key in stale:
self.state.delete(key)
def run_flow(self, flow: str, inputs: dict[str, Any] | None = None) -> None:
"""Run every node of one flow. Blocking — call from a worker thread."""
@@ -527,14 +574,18 @@ class FlowController:
pipeline.run(inputs or {})
def trigger_node(self, node_id: str, values: dict[str, Any] | None = None) -> None:
"""Feed values into one node. Blocking — call from a worker thread."""
"""Feed values into one node. Blocking — call from a worker thread.
Runs here rather than through the queue: the caller is a person waiting
on the response, and wants the state it produced.
"""
node = self.get_node(node_id)
if node is None:
raise KeyError(node_id)
if node.requires and values:
node.trigger(values)
node.trigger(values, durable=False)
else:
node.inject(values or {})
node.inject(values or {}, durable=False)
def _publish(self, event: dict[str, Any]) -> None:
if self.events is not None:
+272
View File
@@ -0,0 +1,272 @@
"""The execution service: what turns journaled work into node runs.
A consumer thread claims items from the work queue and hands each one to a
dispatch pool, which drives the wave it starts. Node bodies run on a second,
separate pool: if cascade drivers and node bodies shared one, a wave waiting
for its own nodes could occupy every thread and deadlock.
A reaper takes back items claimed by an engine that died before acknowledging
them, which is the mechanism that makes a crash mid-cascade recoverable rather
than lossy.
"""
from __future__ import annotations
import logging
import threading
import time
from concurrent.futures import ThreadPoolExecutor
from typing import TYPE_CHECKING, Any
from app.flow.queue import MAX_DELIVERIES, WorkItem, WorkQueue
if TYPE_CHECKING:
from app.flow.events import EventBus
from app.flow.pipeline import Pipeline
logger = logging.getLogger(__name__)
CLAIM_COUNT = 4
CLAIM_BLOCK_MS = 1000
# Long enough that a busy cascade is not mistaken for a dead one.
RECLAIM_IDLE_MS = 60_000
RECLAIM_INTERVAL_S = 30.0
DELAYED_INTERVAL_S = 1.0
MAX_CASCADES = 4
# How long a reload waits for claimed work to finish before rebuilding anyway.
DRAIN_TIMEOUT_S = 10.0
class ExecutionService:
"""Owns the engine's worker threads and the queue they read from."""
def __init__(
self,
queue: WorkQueue,
max_workers: int | None = None,
events: EventBus | None = None,
) -> None:
self.queue = queue
self._events = events
self._pipeline: Pipeline | None = None
self._stop = threading.Event()
self._intake = threading.Event()
self._intake.set()
self._inflight = 0
self._inflight_lock = threading.Condition()
self.node_pool = ThreadPoolExecutor(
max_workers=max_workers or 4, thread_name_prefix="node"
)
self._cascade_pool = ThreadPoolExecutor(
max_workers=MAX_CASCADES, thread_name_prefix="cascade"
)
self._consumer: threading.Thread | None = None
self._timers: threading.Thread | None = None
# -------------------------------------------------------------------------
# Lifecycle
# -------------------------------------------------------------------------
def start(self) -> None:
if self._consumer is not None:
return
self._consumer = threading.Thread(
target=self._consume, name="queue-consumer", daemon=True
)
self._consumer.start()
self._timers = threading.Thread(
target=self._tick, name="queue-timers", daemon=True
)
self._timers.start()
def stop(self) -> None:
self._stop.set()
for thread in (self._consumer, self._timers):
if thread is not None:
thread.join(timeout=5)
self._consumer = None
self._timers = None
self._cascade_pool.shutdown(wait=False)
self.node_pool.shutdown(wait=False)
self.queue.close()
def bind(self, pipeline: Pipeline) -> None:
"""Point the service at the pipeline it should execute against."""
self._pipeline = pipeline
def pause_intake(self) -> None:
"""Stop claiming, and wait for what is already claimed to finish.
Called around a rebuild: items claimed against the old pipeline should
finish there rather than half-run against the new one.
"""
self._intake.clear()
deadline = time.monotonic() + DRAIN_TIMEOUT_S
with self._inflight_lock:
while self._inflight:
remaining = deadline - time.monotonic()
if remaining <= 0:
logger.warning(
"Rebuild did not wait out %d cascades", self._inflight
)
return
self._inflight_lock.wait(remaining)
def resume_intake(self) -> None:
self._intake.set()
def alive(self) -> bool:
return self._consumer is not None and self._consumer.is_alive()
# -------------------------------------------------------------------------
# Threads
# -------------------------------------------------------------------------
def _consume(self) -> None:
failures = 0
while not self._stop.is_set():
if not self._intake.is_set():
self._intake.wait(timeout=0.5)
continue
try:
items = self.queue.claim(CLAIM_COUNT, CLAIM_BLOCK_MS)
failures = 0
except Exception as exc:
failures += 1
logger.error("Could not claim work: %s", exc)
self._publish_unavailable(exc)
# Backing off hard: a queue that is down stays down for a while.
self._stop.wait(min(30.0, 2.0**failures))
continue
for item in items:
self._dispatch(item)
def _tick(self) -> None:
"""Promote delayed items, and take back what a dead engine dropped."""
last_reclaim = 0.0
while not self._stop.is_set():
self._stop.wait(DELAYED_INTERVAL_S)
if self._stop.is_set():
break
try:
self.queue.move_due(time.time())
except Exception as exc:
logger.error("Could not promote delayed work: %s", exc)
now = time.monotonic()
if now - last_reclaim < RECLAIM_INTERVAL_S:
continue
last_reclaim = now
try:
for item in self.queue.reclaim_stale(RECLAIM_IDLE_MS):
logger.info(
"Reclaimed work for '%s' (delivery %d)",
item.node,
item.deliveries,
)
self._dispatch(item)
except Exception as exc:
logger.error("Could not reclaim stale work: %s", exc)
def _dispatch(self, item: WorkItem) -> None:
with self._inflight_lock:
self._inflight += 1
try:
self._cascade_pool.submit(self._handle, item)
except RuntimeError:
# Pool already shutting down.
self._done()
def _done(self) -> None:
with self._inflight_lock:
self._inflight -= 1
self._inflight_lock.notify_all()
# -------------------------------------------------------------------------
# Handling one item
# -------------------------------------------------------------------------
def _handle(self, item: WorkItem) -> None:
handled = True
try:
handled = self._run_item(item)
except Exception:
logger.exception("Work item for '%s' failed", item.node)
finally:
# Leaving it unacknowledged is how it comes back: the reaper hands
# it to whoever can actually run it.
if handled:
try:
self.queue.ack(item)
except Exception as exc:
logger.error(
"Could not acknowledge work for '%s': %s", item.node, exc
)
self._done()
def _run_item(self, item: WorkItem) -> bool:
"""Run one item. False means it was not handled and must come back."""
pipeline = self._pipeline
if pipeline is None:
logger.warning("No pipeline bound; leaving work for '%s'", item.node)
return False
if item.deliveries > MAX_DELIVERIES:
self.queue.dead_letter(item, f"{item.deliveries} deliveries")
self._publish(
{
"type": "cascade_dropped",
"flow": item.flow,
"node": item.node,
"deliveries": item.deliveries,
"ts": time.time(),
}
)
return True
node = pipeline.get_node_by_id(item.node)
if node is None:
# The flow was edited while this was queued; its values are already
# in state, so there is nothing to salvage.
logger.debug("Work item for unknown node '%s', dropped", item.node)
return True
if pipeline.is_disabled(node.flow):
return True
if pipeline.is_paused(node.flow):
self.queue.park(node.flow, item)
return True
pipeline.apply_outputs(node, item.outputs or None)
pipeline.run_downstream(
node, entry_id=item.entry_id, replay=item.deliveries > 1
)
return True
# -------------------------------------------------------------------------
# Reporting
# -------------------------------------------------------------------------
def stats(self) -> dict[str, Any]:
try:
stats = self.queue.stats()
except Exception as exc:
return {"error": str(exc), "consumer_alive": self.alive()}
stats["consumer_alive"] = self.alive()
stats["cascades_busy"] = self._inflight
return stats
def _publish_unavailable(self, exc: Exception) -> None:
self._publish(
{
"type": "queue_unavailable",
"error": f"{type(exc).__name__}: {exc}",
"ts": time.time(),
}
)
def _publish(self, event: dict[str, Any]) -> None:
if self._events is not None:
self._events.publish(event)
+19 -6
View File
@@ -157,6 +157,11 @@ class Node:
# particular node type is.
# -------------------------------------------------------------------------
# Whether running this node twice is the same as running it once. False
# for anything that reaches outside — a second publish is a second command
# to the device, which at-least-once delivery must not cause.
idempotent: bool = True
async def start(self, app: FastAPI | None = None) -> None:
"""Begin whatever this node listens to. Called when its flow starts."""
@@ -237,7 +242,9 @@ class Node:
kwargs = self._to_kwargs(inputs or {})
return self._to_messages(self.f(**kwargs, params=self.params))
def trigger(self, inputs: dict[str, Any] | None = None) -> NodeResult:
def trigger(
self, inputs: dict[str, Any] | None = None, durable: bool | None = None
) -> NodeResult:
"""
Trigger this node externally, executing downstream nodes if dependencies are met.
@@ -252,9 +259,11 @@ class Node:
"""
if self._pipeline is None:
raise RuntimeError("Node must be bound to a pipeline to trigger")
return self(inputs)
return self(inputs, durable=durable)
def inject(self, outputs: dict[str, Any] | None = None) -> NodeResult:
def inject(
self, outputs: dict[str, Any] | None = None, durable: bool | None = None
) -> NodeResult:
"""
Inject data into the pipeline as if this node produced it.
@@ -285,9 +294,11 @@ class Node:
outputs = self.f(params=self.params) or {}
self._pipeline.publish_log(self, collected, "")
return self._pipeline.trigger(self, self._to_messages(outputs))
return self._pipeline.trigger(self, self._to_messages(outputs), durable=durable)
def __call__(self, inputs: dict[str, Any] | None = None) -> NodeResult:
def __call__(
self, inputs: dict[str, Any] | None = None, durable: bool | None = None
) -> NodeResult:
"""
Execute the node and trigger downstream nodes if bound to a pipeline.
@@ -302,4 +313,6 @@ class Node:
:rtype: dict | None
"""
outputs = self.execute(inputs)
return self._pipeline.trigger(self, outputs) if self._pipeline else outputs
if self._pipeline is None:
return outputs
return self._pipeline.trigger(self, outputs, durable=durable)
+11 -4
View File
@@ -145,10 +145,6 @@ class DelayNode(Node):
return None
self.ts = ts
# Apply fixed delay
if self.delay > 0:
time.sleep(self.delay)
if not kwargs:
return None
@@ -158,6 +154,17 @@ class DelayNode(Node):
if in_port in kwargs
}
if self.delay > 0:
# Hand the wait to the queue rather than sit on a worker thread:
# a handful of delay nodes would otherwise occupy the whole pool
# and the engine would stop dead until they woke up.
if self._pipeline is not None and self._pipeline.defer(
self, self._to_messages(output) or {}, self.delay
):
logger.info("[%s] Sending %s in %ss", self.name, output, self.delay)
return None
time.sleep(self.delay)
logger.info("[%s] Sending %s", self.name, output)
return output
+22 -1
View File
@@ -44,6 +44,20 @@ def close_shared_client() -> None:
_client = None
def _first_catch_all(app: FastAPI) -> int:
"""Where a webhook has to go in to be reachable.
A mount at the root answers for every path, so anything registered after
it is dead. Returns that mount's index, or the end of the table.
"""
from starlette.routing import Mount
for index, route in enumerate(app.routes):
if isinstance(route, Mount) and route.path in ("", "/"):
return index
return len(app.routes)
class HttpNode(Node):
"""
HTTP node that can act as a trigger (receiver) or sender based on configuration.
@@ -105,6 +119,10 @@ class HttpNode(Node):
TRIGGER = "trigger" # Receives HTTP requests
SENDER = "sender" # Sends HTTP requests
# A repeated request is a repeated request, whatever the endpoint does
# with it.
idempotent = False
__slots__ = (
"url",
"method",
@@ -363,7 +381,10 @@ class HttpNode(Node):
methods=[self.method],
name=self.id,
)
app.routes.append(route)
# Ahead of any catch-all mount. Starlette takes the first route that
# matches, and the MCP app is mounted at "/", so appending would put
# every webhook behind something that answers for every path.
app.routes.insert(_first_catch_all(app), route)
self._route_registered = True
logger.info(
+3
View File
@@ -128,6 +128,9 @@ class InfluxDbNode(Node):
... )
"""
# Writing the same point twice doubles it in the series.
idempotent = False
__slots__ = (
"url",
"token",
+3
View File
@@ -112,6 +112,9 @@ class MqttNode(Node):
SUBSCRIBER = "subscriber" # Receives MQTT messages (trigger)
PUBLISHER = "publisher" # Sends MQTT messages (sender)
# Publishing again is a second command to whatever is listening.
idempotent = False
__slots__ = (
"topics",
"mode",
+193 -42
View File
@@ -21,6 +21,7 @@ from app.flow import logs
from app.flow.events import EventBus
from app.flow.messages import flow_of
from app.flow.nodes import Node
from app.flow.queue import WorkQueue
from app.flow.state import MemoryState, StateBackend
logger = logging.getLogger(__name__)
@@ -61,6 +62,8 @@ class Pipeline:
"_disabled",
"_paused",
"_gate_lock",
"_queue",
"_node_pool",
)
def __init__(
@@ -71,6 +74,8 @@ class Pipeline:
max_workers: int | None = None,
initial_values: dict[str, Any] | None = None,
disabled_flows: set[str] | None = None,
work_queue: WorkQueue | None = None,
node_pool: ThreadPoolExecutor | None = None,
) -> None:
self._nodes = nodes or []
# Stopped flows are stored and survive a restart; paused ones are a
@@ -84,6 +89,11 @@ class Pipeline:
self._state: StateBackend = state if state is not None else MemoryState()
self._events = events
self._max_workers = max_workers
# Without a queue the pipeline runs everything inline, which is what
# tests, previews and manual runs want.
self._queue = work_queue
# A pool owned by the execution service, so a wave does not build one.
self._node_pool = node_pool
# A message may have several producers; every one of them is upstream
# of the nodes consuming it.
@@ -435,7 +445,20 @@ class Pipeline:
}
)
def _execute_node(self, node: Node, state: StateBackend) -> dict[str, Any] | None:
def _already_done(self, entry_id: str, node: Node) -> bool:
"""Did this node's side effect already happen for this work item?"""
if node.idempotent or self._queue is None:
return False
try:
return self._queue.was_done(entry_id, node.id)
except Exception:
# Not knowing means running it again, which is the safer default
# for a value that may never have been delivered at all.
return False
def _execute_node(
self, node: Node, state: StateBackend, entry_id: str = ""
) -> dict[str, Any] | None:
"""Run one node and record its outputs. Never raises."""
started = time.perf_counter()
collected = logs.Collector()
@@ -447,6 +470,14 @@ class Pipeline:
result = node.execute(inputs)
self.publish_log(node, collected, "")
if entry_id and not node.idempotent and self._queue is not None:
# Written after the fact: a crash between the side effect and
# this marker is the one window at-least-once cannot close.
try:
self._queue.mark_done(entry_id, node.id)
except Exception as exc:
logger.warning("Could not mark '%s' done: %s", node.id, exc)
if result:
result = self._throttled(node, result)
@@ -532,11 +563,17 @@ class Pipeline:
with self._gate_lock:
return node.flow in self._paused
def is_paused(self, flow: str) -> bool:
with self._gate_lock:
return flow in self._paused
def _execute_parallel(
self,
nodes_subset: set[Node] | None,
state: StateBackend,
check_ready: bool = False,
entry_id: str = "",
replay: bool = False,
) -> StateBackend:
"""Execute nodes concurrently, scheduling each as its inputs arrive."""
target_nodes = nodes_subset if nodes_subset is not None else set(self._nodes)
@@ -569,12 +606,23 @@ class Pipeline:
continue
if is_ready(n):
submitted.add(n)
node_futures[n] = executor.submit(self._execute_node, n, state)
if replay and entry_id and self._already_done(entry_id, n):
# Its side effect happened on an earlier delivery; its
# outputs are still in state, so downstream carries on.
skipped.add(n)
submitted.discard(n)
for consumer in self.edges[n]:
if consumer in target_nodes:
in_degree[consumer] -= 1
continue
node_futures[n] = executor.submit(
self._execute_node, n, state, entry_id
)
elif n.synchronous and in_degree[n] == 0:
# Not ready now; a later trigger may make it ready.
skipped.add(n)
with ThreadPoolExecutor(max_workers=self._max_workers) as executor:
def drain(executor: ThreadPoolExecutor) -> None:
submit_ready(executor)
while node_futures:
@@ -591,6 +639,14 @@ class Pipeline:
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
# wave is what used to spawn threads without bound under load.
drain(self._node_pool)
else:
with ThreadPoolExecutor(max_workers=self._max_workers) as executor:
drain(executor)
return state
def run(
@@ -604,9 +660,82 @@ class Pipeline:
return self._execute_parallel(nodes, self._state, check_ready=False)
def trigger(self, node: Node, outputs: dict[str, Any] | None) -> StateBackend:
def apply_outputs(self, node: Node, outputs: dict[str, Any] | None) -> None:
"""Record what a node emitted: state, history, versions and events.
Shared by the direct path and by the execution service replaying a
journaled item, so a value looks the same on the canvas either way.
"""
if outputs:
# This is where a chatty subscriber gets thinned out, so a port set
# to publish every 60s does so whatever the broker sends.
outputs = self._throttled(node, outputs)
if not outputs:
return
state = self._state
ts = time.time()
with state.lock():
state.update(outputs)
state.update({self._timestamp_key(name): ts for name in outputs})
state.append_history(outputs, ts)
self._increment_message_versions(outputs)
for name, value in outputs.items():
self._publish(
{
"type": "message_value",
"flow": flow_of(name),
"name": name,
"value": value,
"ts": ts,
}
)
# An injecting node — an MQTT subscriber, a webhook — publishes
# without going through the executor, but it did emit.
self._publish(
{
"type": "node_executed",
"flow": node.flow,
"node": node.id,
"outputs": len(outputs),
"duration_ms": 0,
"ts": ts,
}
)
def run_downstream(
self, node: Node, entry_id: str = "", replay: bool = False
) -> StateBackend:
"""Run everything downstream of a node that has just published.
``entry_id`` identifies the journaled item this run belongs to, so a
node with outside side effects can record that it ran. On a ``replay``
— the same item handed back after a crash — that record is checked
first: at-least-once delivery must not mean two of the same request.
"""
downstream = set(self._get_downstream(node))
if not downstream:
return self._state
return self._execute_parallel(
downstream,
self._state,
check_ready=True,
entry_id=entry_id,
replay=replay,
)
def trigger(
self, node: Node, outputs: dict[str, Any] | None, durable: bool | None = None
) -> StateBackend:
"""Publish a node's outputs and run everything downstream of it.
With a work queue attached the event is journaled and the caller
returns immediately — that is the path every external trigger takes, so
a crash mid-cascade loses nothing. Interactive callers (a manual run, a
draft preview) pass ``durable=False`` and get the old synchronous
behaviour, because they are waiting for the result.
A stopped flow drops the event: its subscriptions and schedules are torn
down anyway, and anything still arriving from another thread would be
work the flow was explicitly told not to do. A *paused* flow still
@@ -618,46 +747,68 @@ class Pipeline:
if node.flow in self._disabled:
return state
if outputs:
# This is where a chatty subscriber gets thinned out, so a port set
# to publish every 60s does so whatever the broker sends.
outputs = self._throttled(node, outputs)
if outputs:
ts = time.time()
with state.lock():
state.update(outputs)
state.update({self._timestamp_key(name): ts for name in outputs})
state.append_history(outputs, ts)
self._increment_message_versions(outputs)
for name, value in outputs.items():
self._publish(
{
"type": "message_value",
"flow": flow_of(name),
"name": name,
"value": value,
"ts": ts,
}
)
# An injecting node — an MQTT subscriber, a webhook — publishes
# without going through the executor, but it did emit.
self._publish(
{
"type": "node_executed",
"flow": node.flow,
"node": node.id,
"outputs": len(outputs),
"duration_ms": 0,
"ts": ts,
}
)
downstream = set(self._get_downstream(node))
if not downstream:
if durable is None:
durable = self._queue is not None
if durable and self._queue is not None:
self._enqueue_cascade(node, outputs)
return state
return self._execute_parallel(downstream, state, check_ready=True)
self.apply_outputs(node, outputs)
return self.run_downstream(node)
def defer(self, node: Node, outputs: dict[str, Any], seconds: float) -> bool:
"""Publish a node's outputs later, without holding a worker thread.
Returns False when there is no queue to hold the item, in which case
the caller has to wait however it waited before.
"""
from app.flow.queue import WorkItem
if self._queue is None or seconds <= 0:
return False
item = WorkItem(
kind="cascade",
node=node.id,
flow=node.flow,
outputs=outputs,
cause="delay",
)
try:
self._queue.add_delayed(item, time.time() + seconds)
return True
except Exception as exc:
logger.error("Could not defer work for '%s': %s", node.id, exc)
return False
def _enqueue_cascade(self, node: Node, outputs: dict[str, Any] | None) -> None:
"""Journal a trigger, or fall back to running it here if that fails."""
from app.flow.queue import WorkItem
item = WorkItem(
kind="cascade",
node=node.id,
flow=node.flow,
outputs=outputs or {},
cause="external",
)
assert self._queue is not None
try:
self._queue.add(item)
return
except Exception as exc:
logger.error("Could not journal work for '%s': %s", node.id, exc)
self._publish(
{
"type": "queue_unavailable",
"flow": node.flow,
"node": node.id,
"error": f"{type(exc).__name__}: {exc}",
"ts": time.time(),
}
)
# Losing the value outright would be worse than running it here.
self.apply_outputs(node, outputs)
self.run_downstream(node)
def values(self, flow: str | None = None) -> dict[str, dict[str, Any]]:
"""Last value and timestamp of every message, optionally one flow's."""
+414
View File
@@ -0,0 +1,414 @@
"""The engine's work queue: what survives a crash.
Every external event — an MQTT message, a webhook, a cron tick, a connector
poll — becomes a work item before anything runs. The item is journaled first
and acknowledged only once the wave it started has quiesced, so an engine that
dies mid-cascade picks the work up again on the way back rather than losing it.
Two implementations: Redis Streams, which is what makes the above true, and an
in-memory one for tests and for running without Redis, where "durable" degrades
honestly to "not".
"""
from __future__ import annotations
import heapq
import json
import logging
import threading
import time
from abc import ABC, abstractmethod
from collections import deque
from dataclasses import dataclass, field
from typing import Any, cast
import redis
logger = logging.getLogger(__name__)
# The stream is a buffer, not an archive: anything this far behind is long
# superseded by a newer value of the same message.
STREAM_MAXLEN = 10_000
# An item redelivered this often is not going to succeed. Park it where someone
# can look rather than letting it loop forever.
MAX_DELIVERIES = 3
GROUP = "engine"
@dataclass
class WorkItem:
"""One unit of journaled work.
:param kind: ``cascade`` replays a node's outputs and runs what is
downstream; ``node`` executes exactly one node.
:param node: The node the item is about — the source for a cascade, the
target for a node item.
:param flow: The flow that node belongs to, so gating needs no lookup.
:param outputs: What the source node emitted (cascade only).
:param cause: Where the work came from, for logs and debugging.
:param not_before: Epoch seconds before which the item must not run.
:param entry_id: Set by the queue on claim; stable across redeliveries,
which is what makes it usable as an idempotency key.
:param deliveries: How many times this item has been handed out.
"""
kind: str
node: str
flow: str
outputs: dict[str, Any] = field(default_factory=dict)
cause: str = "system"
not_before: float = 0.0
entry_id: str = ""
deliveries: int = 1
def to_fields(self) -> dict[str, str]:
return {
"kind": self.kind,
"node": self.node,
"flow": self.flow,
"outputs": json.dumps(self.outputs),
"cause": self.cause,
"not_before": str(self.not_before),
}
@classmethod
def from_fields(
cls, fields: dict[str, str], entry_id: str, deliveries: int = 1
) -> WorkItem:
return cls(
kind=fields.get("kind", "cascade"),
node=fields.get("node", ""),
flow=fields.get("flow", ""),
outputs=json.loads(fields.get("outputs") or "{}"),
cause=fields.get("cause", "system"),
not_before=float(fields.get("not_before") or 0.0),
entry_id=entry_id,
deliveries=deliveries,
)
class WorkQueue(ABC):
"""What the execution service pulls from."""
@abstractmethod
def add(self, item: WorkItem) -> None:
"""Journal an item for immediate execution."""
@abstractmethod
def add_delayed(self, item: WorkItem, not_before: float) -> None:
"""Journal an item that must not run before ``not_before``."""
@abstractmethod
def claim(self, count: int, block_ms: int) -> list[WorkItem]:
"""Take up to ``count`` items, waiting up to ``block_ms`` for one."""
@abstractmethod
def ack(self, item: WorkItem) -> None:
"""Mark an item done, so it is never redelivered."""
@abstractmethod
def reclaim_stale(self, min_idle_ms: int) -> list[WorkItem]:
"""Take back items claimed by a consumer that never acknowledged them."""
@abstractmethod
def move_due(self, now: float) -> int:
"""Promote delayed items whose time has come. Returns how many."""
@abstractmethod
def dead_letter(self, item: WorkItem, reason: str) -> None:
"""Set an item aside after it has failed too often."""
@abstractmethod
def park(self, flow: str, item: WorkItem) -> None:
"""Hold an item while its flow is paused."""
@abstractmethod
def unpark(self, flow: str) -> list[WorkItem]:
"""Return a paused flow's held items to the queue, oldest first."""
@abstractmethod
def clear_flow(self, flow: str) -> None:
"""Forget anything held for a flow that no longer exists."""
@abstractmethod
def stats(self) -> dict[str, Any]:
"""Queue depth and age, for the health endpoint."""
@abstractmethod
def mark_done(self, entry_id: str, node: str) -> None:
"""Record that a side effect already happened for this delivery."""
@abstractmethod
def was_done(self, entry_id: str, node: str) -> bool:
"""Did this exact item already run that side-effecting node?"""
@abstractmethod
def close(self) -> None:
"""Release whatever the queue holds."""
class MemoryWorkQueue(WorkQueue):
"""In-process queue. No durability: a crash loses whatever is in flight."""
def __init__(self) -> None:
self._items: deque[WorkItem] = deque()
self._delayed: list[tuple[float, int, WorkItem]] = []
self._parked: dict[str, list[WorkItem]] = {}
self._done: set[tuple[str, str]] = set()
self._counter = 0
self._seq = 0
self._lock = threading.Lock()
self._wake = threading.Condition(self._lock)
def _next_id(self) -> str:
self._seq += 1
return f"mem-{self._seq}"
def add(self, item: WorkItem) -> None:
with self._wake:
if not item.entry_id:
item.entry_id = self._next_id()
self._items.append(item)
self._wake.notify()
def add_delayed(self, item: WorkItem, not_before: float) -> None:
with self._wake:
if not item.entry_id:
item.entry_id = self._next_id()
item.not_before = not_before
self._counter += 1
heapq.heappush(self._delayed, (not_before, self._counter, item))
self._wake.notify()
def claim(self, count: int, block_ms: int) -> list[WorkItem]:
deadline = time.monotonic() + block_ms / 1000.0
with self._wake:
while not self._items:
remaining = deadline - time.monotonic()
if remaining <= 0:
return []
self._wake.wait(remaining)
return [self._items.popleft() for _ in range(min(count, len(self._items)))]
def ack(self, item: WorkItem) -> None:
"""Nothing to acknowledge: claiming already removed it."""
def reclaim_stale(self, min_idle_ms: int) -> list[WorkItem]:
"""Nothing to reclaim: an item lost here died with the process."""
return []
def move_due(self, now: float) -> int:
moved = 0
with self._wake:
while self._delayed and self._delayed[0][0] <= now:
self._items.append(heapq.heappop(self._delayed)[2])
moved += 1
if moved:
self._wake.notify()
return moved
def dead_letter(self, item: WorkItem, reason: str) -> None:
logger.error("Dropping work item for '%s': %s", item.node, reason)
def park(self, flow: str, item: WorkItem) -> None:
with self._lock:
self._parked.setdefault(flow, []).append(item)
def unpark(self, flow: str) -> list[WorkItem]:
with self._lock:
return self._parked.pop(flow, [])
def clear_flow(self, flow: str) -> None:
with self._lock:
self._parked.pop(flow, None)
def stats(self) -> dict[str, Any]:
with self._lock:
return {
"depth": len(self._items),
"pending": 0,
"delayed": len(self._delayed),
"parked": sum(len(v) for v in self._parked.values()),
"oldest_pending_s": 0.0,
"durable": False,
}
def mark_done(self, entry_id: str, node: str) -> None:
with self._lock:
self._done.add((entry_id, node))
def was_done(self, entry_id: str, node: str) -> bool:
with self._lock:
return (entry_id, node) in self._done
def close(self) -> None:
"""Nothing is held outside the process."""
class RedisWorkQueue(WorkQueue):
"""Redis Streams queue: the one that survives the process.
Items live in a stream read through a consumer group, so an item handed to
a consumer that dies before acknowledging it stays pending and is reclaimed
by whoever comes next.
"""
def __init__(
self,
host: str,
port: int = 6379,
namespace: str = "pipeline",
consumer: str | None = None,
) -> None:
self._redis = redis.Redis(host=host, port=port, decode_responses=True)
self._ns = namespace
self._consumer = consumer or f"engine-{int(time.time() * 1000) % 1_000_000}"
self._stream = f"{namespace}:__queue__"
self._delayed_key = f"{namespace}:__delayed__"
self._dead_key = f"{namespace}:__dead__"
self._ensure_group()
def _ensure_group(self) -> None:
try:
self._redis.xgroup_create(self._stream, GROUP, id="0", mkstream=True)
except redis.ResponseError as exc:
if "BUSYGROUP" not in str(exc):
raise
def _parked_key(self, flow: str) -> str:
return f"{self._ns}:__parked__:{flow}"
def _done_key(self, entry_id: str, node: str) -> str:
return f"{self._ns}:__done__:{entry_id}:{node}"
def add(self, item: WorkItem) -> None:
self._redis.xadd(
self._stream,
cast(Any, item.to_fields()),
maxlen=STREAM_MAXLEN,
approximate=True,
)
def add_delayed(self, item: WorkItem, not_before: float) -> None:
item.not_before = not_before
self._redis.zadd(self._delayed_key, {json.dumps(item.to_fields()): not_before})
def claim(self, count: int, block_ms: int) -> list[WorkItem]:
response = cast(
list[Any],
self._redis.xreadgroup(
GROUP,
self._consumer,
{self._stream: ">"},
count=count,
block=block_ms,
),
)
items: list[WorkItem] = []
for _stream, entries in response or []:
for entry_id, fields in entries:
items.append(WorkItem.from_fields(fields, entry_id))
return items
def ack(self, item: WorkItem) -> None:
if item.entry_id:
self._redis.xack(self._stream, GROUP, item.entry_id)
def reclaim_stale(self, min_idle_ms: int) -> list[WorkItem]:
"""Take over entries a dead consumer never acknowledged."""
_cursor, entries, _deleted = cast(
tuple[Any, list[Any], Any],
self._redis.xautoclaim(
self._stream,
GROUP,
self._consumer,
min_idle_time=min_idle_ms,
count=32,
),
)
if not entries:
return []
# xautoclaim does not report delivery counts, so ask xpending for them.
counts: dict[str, int] = {}
for record in cast(
list[dict[str, Any]],
self._redis.xpending_range(self._stream, GROUP, min="-", max="+", count=64),
):
counts[record["message_id"]] = record["times_delivered"]
return [
WorkItem.from_fields(fields, entry_id, counts.get(entry_id, 1))
for entry_id, fields in entries
if fields
]
def move_due(self, now: float) -> int:
due = cast(
list[str],
self._redis.zrangebyscore(self._delayed_key, "-inf", now, start=0, num=100),
)
moved = 0
for raw in due:
# Whoever removes it owns it: a second engine would get 0 here.
if self._redis.zrem(self._delayed_key, raw):
self._redis.xadd(
self._stream,
cast(Any, json.loads(raw)),
maxlen=STREAM_MAXLEN,
approximate=True,
)
moved += 1
return moved
def dead_letter(self, item: WorkItem, reason: str) -> None:
fields = item.to_fields()
fields["reason"] = reason
self._redis.xadd(
self._dead_key, cast(Any, fields), maxlen=1000, approximate=True
)
logger.error("Dead-lettered work item for '%s': %s", item.node, reason)
def park(self, flow: str, item: WorkItem) -> None:
self._redis.rpush(self._parked_key(flow), json.dumps(item.to_fields()))
def unpark(self, flow: str) -> list[WorkItem]:
key = self._parked_key(flow)
raw = cast(list[str], self._redis.lrange(key, 0, -1))
self._redis.delete(key)
return [WorkItem.from_fields(json.loads(r), "") for r in raw]
def clear_flow(self, flow: str) -> None:
self._redis.delete(self._parked_key(flow))
def stats(self) -> dict[str, Any]:
pending = cast(dict[str, Any], self._redis.xpending(self._stream, GROUP))
oldest = 0.0
count = pending.get("pending", 0) if isinstance(pending, dict) else 0
if count:
records = cast(
list[dict[str, Any]],
self._redis.xpending_range(
self._stream, GROUP, min="-", max="+", count=1
),
)
if records:
oldest = records[0]["time_since_delivered"] / 1000.0
return {
"depth": cast(int, self._redis.xlen(self._stream)),
"pending": count,
"delayed": cast(int, self._redis.zcard(self._delayed_key)),
"parked": 0,
"oldest_pending_s": round(oldest, 1),
"durable": True,
}
def mark_done(self, entry_id: str, node: str) -> None:
# An hour outlives any redelivery; after that the marker is noise.
self._redis.set(self._done_key(entry_id, node), "1", ex=3600)
def was_done(self, entry_id: str, node: str) -> bool:
return bool(self._redis.exists(self._done_key(entry_id, node)))
def close(self) -> None:
self._redis.close()
+13
View File
@@ -77,6 +77,11 @@ class StateBackend(ABC):
"""
...
@abstractmethod
def delete(self, key: str) -> None:
"""Forget one key and any history kept for it."""
...
@abstractmethod
def clear(self) -> None:
"""Clear all keys in the state."""
@@ -243,6 +248,11 @@ class MemoryState(StateBackend):
with self._lock:
return key in self._data
def delete(self, key: str) -> None:
with self._lock:
self._data.pop(key, None)
self._history.pop(key, None)
def clear(self) -> None:
with self._lock:
self._data.clear()
@@ -390,6 +400,9 @@ class RedisState(StateBackend):
def exists(self, key: str) -> bool:
return bool(self._client.exists(self._key(key)))
def delete(self, key: str) -> None:
self._client.delete(self._key(key), self._history_key(key))
def clear(self) -> None:
"""Clear all keys in the namespace."""
pattern = f"{self._namespace}:*"
+4 -2
View File
@@ -208,8 +208,10 @@ class FlowStore:
self._lib_file(lib_name).parent.mkdir(parents=True, exist_ok=True)
self._lib_file(lib_name).write_text(code)
for path in (self._node_file(flow, node_id),
self._draft_node_file(flow, node_id)):
for path in (
self._node_file(flow, node_id),
self._draft_node_file(flow, node_id),
):
if path.exists():
path.unlink()
+15
View File
@@ -15,8 +15,10 @@ from app.core.config import settings
from app.flow import logs
from app.flow.controller import FlowController
from app.flow.events import event_bus
from app.flow.executor import ExecutionService
from app.flow.nodes.http import close_shared_client
from app.flow.plugins import load_plugins
from app.flow.queue import MemoryWorkQueue, RedisWorkQueue, WorkQueue
from app.flow.secrets import init_secrets
from app.flow.state import MemoryState, RedisState, StateBackend
from app.flow.store import FlowStore
@@ -37,6 +39,13 @@ def _state_backend() -> StateBackend:
return MemoryState()
def _work_queue() -> WorkQueue:
"""Redis makes queued work survive the process; memory does not pretend to."""
if settings.REDIS_HOST:
return RedisWorkQueue(host=settings.REDIS_HOST, port=settings.REDIS_PORT)
return MemoryWorkQueue()
def _mcp_sessions() -> AbstractAsyncContextManager[None]:
"""The MCP session manager's run scope, or nothing when MCP is off."""
if not settings.MCP_ENABLED:
@@ -56,12 +65,18 @@ async def lifespan(app: FastAPI) -> AsyncIterator[None]:
# Connectors register their node types before any flow is built with them.
load_plugins()
execution = ExecutionService(
queue=_work_queue(),
max_workers=settings.FLOW_MAX_WORKERS,
events=event_bus,
)
controller = FlowController(
store=FlowStore(settings.FLOWS_DIR),
state=_state_backend(),
events=event_bus,
max_workers=settings.FLOW_MAX_WORKERS,
fastapi_app=app,
execution=execution,
)
app.state.flow_controller = controller
watchdog = LoopWatchdog(event_bus)
@@ -61,6 +61,21 @@ def test_a_hook_without_a_secret_still_answers_on_its_plain_url():
assert pipeline.values()["house.temp"]["value"] == 21.5
def test_a_hook_is_reachable_past_an_app_mounted_at_the_root():
"""The MCP app is mounted at "/" and answers for every path, so a hook
registered after it would never be reached."""
node = hook_node()
pipeline = Pipeline(nodes=[node])
app = FastAPI()
app.mount("/", FastAPI())
node.register_route(app)
response = TestClient(app).post(HOOK, json={"temp": 21.5})
assert response.status_code == 200
assert pipeline.values()["house.temp"]["value"] == 21.5
@pytest.mark.parametrize(
"path",
[
+220
View File
@@ -0,0 +1,220 @@
"""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 _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_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]