Three things a pause and an interval were quietly losing:
- A rebuild builds a fresh pipeline, so nothing is paused any more and no
resume ever comes for what the old one parked. Release it on rebuild.
- POST /flows/{name}/step takes the oldest parked item and runs that one wave
while the flow stays paused, so a held-back cascade can be walked through.
Nothing parked answers plainly rather than failing.
- A per-port interval was leading-edge only: a producer going quiet inside the
window left the consumer on the value before it. The held value is kept and
a flush item scheduled on the queue's existing timer, so the window ends with
a delivery. One timer in flight per node, and none without a queue to run it.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01H7LwYgJfpkbLCTeiAf8U4A
167 lines
5.2 KiB
Python
167 lines
5.2 KiB
Python
"""Per-port intervals: deliver at most every x seconds."""
|
|
|
|
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
|
|
|
|
# Short enough to wait out in a test, long enough not to race the engine.
|
|
WINDOW = 0.05
|
|
|
|
|
|
def spec(name: str, interval: float = 0) -> MessageSpec:
|
|
return MessageSpec(name=name, dtype=DType.FLOAT, interval=interval)
|
|
|
|
|
|
def make_node(node_id: str, f, requires=(), provides=()) -> Node:
|
|
node = Node(f=f, requires=list(requires), provides=list(provides), name=node_id)
|
|
node.assign_flow("demo", node_id)
|
|
return node
|
|
|
|
|
|
def test_a_limited_output_publishes_once_inside_its_window():
|
|
readings = iter([1.0, 2.0, 3.0])
|
|
source = make_node(
|
|
"source",
|
|
lambda params: {"temp": next(readings)},
|
|
provides=[spec("temp", interval=60)],
|
|
)
|
|
pipeline = Pipeline(nodes=[source])
|
|
|
|
pipeline.run({})
|
|
assert pipeline.state["demo.temp"] == 1.0
|
|
|
|
# Same window: the reading is taken but not published.
|
|
pipeline.run({})
|
|
assert pipeline.state["demo.temp"] == 1.0
|
|
|
|
|
|
def test_an_unlimited_output_publishes_every_time():
|
|
readings = iter([1.0, 2.0])
|
|
source = make_node(
|
|
"source",
|
|
lambda params: {"temp": next(readings)},
|
|
provides=[spec("temp")],
|
|
)
|
|
pipeline = Pipeline(nodes=[source])
|
|
|
|
pipeline.run({})
|
|
pipeline.run({})
|
|
assert pipeline.state["demo.temp"] == 2.0
|
|
|
|
|
|
def test_a_limited_input_wakes_its_node_once_inside_the_window():
|
|
seen: list[float] = []
|
|
source = make_node("source", lambda params: {"temp": 20.0}, provides=[spec("temp")])
|
|
consumer = make_node(
|
|
"consumer",
|
|
lambda temp, params: seen.append(temp),
|
|
requires=[spec("temp", interval=60)],
|
|
)
|
|
# Binding the nodes is what the pipeline is for here.
|
|
Pipeline(nodes=[source, consumer])
|
|
|
|
source.inject({"temp": 20.0})
|
|
source.inject({"temp": 21.0})
|
|
|
|
assert seen == [20.0]
|
|
|
|
|
|
def test_an_unthrottled_input_still_wakes_a_node_beside_a_throttled_one():
|
|
seen: list[tuple[float, float]] = []
|
|
fast = make_node("fast", lambda params: None, provides=[spec("quick")])
|
|
slow = make_node("slow", lambda params: None, provides=[spec("rare")])
|
|
consumer = make_node(
|
|
"consumer",
|
|
lambda quick, rare, params: seen.append((quick, rare)),
|
|
requires=[spec("quick"), spec("rare", interval=60)],
|
|
)
|
|
pipeline = Pipeline(nodes=[fast, slow, consumer])
|
|
pipeline.state["demo.rare"] = 1.0
|
|
|
|
fast.inject({"quick": 1.0})
|
|
fast.inject({"quick": 2.0})
|
|
|
|
# The throttled port holds back only itself.
|
|
assert [quick for quick, _ in seen] == [1.0, 2.0]
|
|
|
|
|
|
def _running(nodes: list[Node]) -> tuple[Pipeline, MemoryWorkQueue, ExecutionService]:
|
|
"""A pipeline with the timer the engine uses for held-back values."""
|
|
queue = MemoryWorkQueue()
|
|
pipeline = Pipeline(nodes=nodes, work_queue=queue)
|
|
service = ExecutionService(queue)
|
|
service.bind(pipeline)
|
|
return pipeline, queue, service
|
|
|
|
|
|
def _run_due(queue: MemoryWorkQueue, service: ExecutionService) -> int:
|
|
"""What the engine's timer thread does once the window has passed."""
|
|
moved = queue.move_due(time.time())
|
|
for item in queue.claim(10, 10):
|
|
service._run_item(item)
|
|
return moved
|
|
|
|
|
|
def test_a_limited_output_publishes_its_last_value_when_the_window_ends():
|
|
"""A producer going quiet must not strand the reading it held back."""
|
|
readings = iter([1.0, 2.0])
|
|
source = make_node(
|
|
"source",
|
|
lambda params: {"temp": next(readings)},
|
|
provides=[spec("temp", interval=WINDOW)],
|
|
)
|
|
pipeline, queue, service = _running([source])
|
|
|
|
pipeline.run({})
|
|
pipeline.run({})
|
|
# Inside the window, so the second reading is held rather than published.
|
|
assert pipeline.state["demo.temp"] == 1.0
|
|
|
|
time.sleep(WINDOW * 2)
|
|
assert _run_due(queue, service) == 1
|
|
|
|
assert pipeline.state["demo.temp"] == 2.0
|
|
|
|
|
|
def test_a_limited_input_wakes_its_node_when_the_window_ends():
|
|
seen: list[float] = []
|
|
source = make_node("source", lambda params: None, provides=[spec("temp")])
|
|
consumer = make_node(
|
|
"consumer",
|
|
lambda temp, params: seen.append(temp),
|
|
requires=[spec("temp", interval=WINDOW)],
|
|
)
|
|
_pipeline, queue, service = _running([source, consumer])
|
|
|
|
source.inject({"temp": 20.0}, durable=False)
|
|
source.inject({"temp": 21.0}, durable=False)
|
|
assert seen == [20.0]
|
|
|
|
time.sleep(WINDOW * 2)
|
|
assert _run_due(queue, service) == 1
|
|
|
|
# The value that arrived inside the window is delivered at the end of it.
|
|
assert seen == [20.0, 21.0]
|
|
|
|
|
|
def test_a_manual_run_is_never_throttled_on_its_inputs():
|
|
seen: list[float] = []
|
|
source = make_node("source", lambda params: {"temp": 20.0}, provides=[spec("temp")])
|
|
consumer = make_node(
|
|
"consumer",
|
|
lambda temp, params: seen.append(temp),
|
|
requires=[spec("temp", interval=3600)],
|
|
)
|
|
pipeline = Pipeline(nodes=[source, consumer])
|
|
|
|
# Pressing Run is an explicit ask; the interval governs the flow's own
|
|
# traffic, not what the person in front of it asked for.
|
|
pipeline.run({})
|
|
pipeline.run({})
|
|
|
|
assert len(seen) == 2
|