Docs / docs (push) Successful in 25s
Playwright Tests / test-playwright (1, 2) (push) Successful in 2m23s
Playwright Tests / test-playwright (2, 2) (push) Successful in 2m0s
pre-commit / pre-commit (push) Failing after 4m31s
Test Backend / test-backend (push) Successful in 2m55s
Compose Smoke Test / test-compose (push) Successful in 35s
Playwright Tests / merge-reports (push) Successful in 1m11s
The timer thread promoted due work on a fixed one-second tick, so every delayed item was 0-1000ms late whatever the load — measured on the house at 705ms mean on a rollershutter stop, which is 2-4% of a 26-second travel and accumulates in the position the motor node believes it is at. It now sleeps to the soonest deadline and is woken when a nearer one is scheduled, which measures 0.9ms end to end through Redis. A promoted timer also went to the back of the queue. It goes into a due lane of its own that `claim` reads first, so work that has waited out a deadline is not held up by work that is merely queued. Beside it, in the same code: seeding a message now bumps its version, so a re-put flow's synchronous nodes no longer wait forever on a value that is sitting in state; the consumer group drops the consumers of engines that are gone (138 had accumulated on this installation); and the cast that closes the long-standing `xclaim` mypy error. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
497 lines
19 KiB
Python
497 lines
19 KiB
Python
"""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. It claims only what that pool
|
|
can start, so a backlog waits in the queue rather than inside the process.
|
|
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 fluksio.flow.queue import MAX_DELIVERIES, WorkItem, WorkQueue
|
|
|
|
if TYPE_CHECKING:
|
|
from fluksio.flow.events import EventBus
|
|
from fluksio.flow.pipeline import Pipeline
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
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
|
|
# How often to tell the queue that what we hold is still being worked on. A
|
|
# node may run for as long as it likes, so what marks an item abandoned is this
|
|
# stopping — which is what an engine that died does.
|
|
TOUCH_INTERVAL_S = 20.0
|
|
# The longest the timer thread sleeps with nothing due. It is a housekeeping
|
|
# cadence and a backstop for a deadline written by another process, not the
|
|
# resolution of a delay: a delayed item is waited for exactly, so what a timer
|
|
# fires late by is a wake-up and a promotion rather than up to a whole second.
|
|
DELAYED_INTERVAL_S = 1.0
|
|
#: How many cascades may be in flight, unless the service is given a number.
|
|
#: Sustained throughput is this over the mean cascade time, so an installation
|
|
#: whose nodes wait on a network rather than a CPU may want more of them —
|
|
#: `FLOW_MAX_CASCADES` is where that is said.
|
|
MAX_CASCADES = 4
|
|
# How long a reload waits for claimed work to finish before rebuilding anyway.
|
|
DRAIN_TIMEOUT_S = 10.0
|
|
# Work waiting in the stream, undelivered. A burst is normal — the pool claims
|
|
# only what it can start — so what marks an engine as falling behind is the
|
|
# backlog staying up across several checks rather than any one reading.
|
|
BACKLOG_INTERVAL_S = 5.0
|
|
BACKLOG_DEGRADED = 50
|
|
BACKLOG_STRIKES = 3
|
|
|
|
|
|
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,
|
|
max_cascades: int | None = None,
|
|
) -> None:
|
|
self.queue = queue
|
|
self._events = events
|
|
self.max_cascades = max_cascades or MAX_CASCADES
|
|
self._pipeline: Pipeline | None = None
|
|
self._stop = threading.Event()
|
|
# Set when a deadline moves closer, so the timer thread stops waiting
|
|
# on the one it read and goes back for the new one.
|
|
self._timer_wake = threading.Event()
|
|
queue.on_delayed = self._timer_wake.set
|
|
self._intake = threading.Event()
|
|
self._intake.set()
|
|
self._inflight = 0
|
|
self._inflight_lock = threading.Condition()
|
|
# Entry ids claimed and still running, under _inflight_lock.
|
|
self._active: set[str] = set()
|
|
self.node_pool = ThreadPoolExecutor(
|
|
max_workers=max_workers or 4, thread_name_prefix="node"
|
|
)
|
|
self._cascade_pool = ThreadPoolExecutor(
|
|
max_workers=self.max_cascades, thread_name_prefix="cascade"
|
|
)
|
|
self._consumer: threading.Thread | None = None
|
|
self._timers: threading.Thread | None = None
|
|
# Consecutive backlog readings over the threshold, and whether the last
|
|
# of them said so out loud.
|
|
self._backlog_strikes = 0
|
|
self.behind = False
|
|
|
|
# -------------------------------------------------------------------------
|
|
# 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()
|
|
# The timer thread sleeps on this, not on _stop.
|
|
self._timer_wake.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()
|
|
|
|
@property
|
|
def inflight(self) -> int:
|
|
"""Cascades claimed and still running."""
|
|
return self._inflight
|
|
|
|
# -------------------------------------------------------------------------
|
|
# 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
|
|
free = self._await_capacity()
|
|
if not free:
|
|
continue
|
|
try:
|
|
items = self.queue.claim(free, 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 _sleep_until_due(self) -> None:
|
|
"""Wait for the soonest deadline, the housekeeping cap, or a new one.
|
|
|
|
A fixed poll here made every delayed item late by 0-1000ms whatever the
|
|
load — on a rollershutter driven for a measured 26 seconds, 2-4% of its
|
|
travel every time, accumulating in the position its node believes it is
|
|
at. Sleeping to the deadline instead leaves a wake-up and a promotion,
|
|
which is milliseconds.
|
|
"""
|
|
self._timer_wake.clear()
|
|
try:
|
|
# Read after the clear: a deadline arriving in between sets the
|
|
# event again, so the wait below returns immediately rather than
|
|
# sleeping through work that landed in the gap.
|
|
due = self.queue.next_due()
|
|
except Exception as exc:
|
|
logger.error("Could not read the next deadline: %s", exc)
|
|
due = None
|
|
wait = DELAYED_INTERVAL_S if due is None else due - time.time()
|
|
self._timer_wake.wait(min(max(wait, 0.0), DELAYED_INTERVAL_S))
|
|
|
|
def _tick(self) -> None:
|
|
"""Promote delayed items, and take back what a dead engine dropped."""
|
|
last_reclaim = 0.0
|
|
last_touch = 0.0
|
|
last_backlog = 0.0
|
|
while not self._stop.is_set():
|
|
self._sleep_until_due()
|
|
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)
|
|
# The item is still due, so the wait above would be zero and
|
|
# this would spin on a queue that is down. Back off to what a
|
|
# fixed poll used to cost.
|
|
self._stop.wait(DELAYED_INTERVAL_S)
|
|
|
|
now = time.monotonic()
|
|
if now - last_backlog >= BACKLOG_INTERVAL_S:
|
|
last_backlog = now
|
|
try:
|
|
self._check_backlog()
|
|
except Exception as exc:
|
|
logger.error("Could not read the queue backlog: %s", exc)
|
|
|
|
if now - last_touch >= TOUCH_INTERVAL_S:
|
|
last_touch = now
|
|
with self._inflight_lock:
|
|
running = list(self._active)
|
|
try:
|
|
self.queue.touch(running)
|
|
except Exception as exc:
|
|
logger.error("Could not touch claimed work: %s", exc)
|
|
|
|
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 _check_backlog(self) -> None:
|
|
"""Say so when work has been waiting in the stream for a while.
|
|
|
|
A flow enqueuing faster than the pool drains produces no event of its
|
|
own: the backlog simply grows, every timer and connector poll drifts
|
|
behind it, and nothing on the health screen moves. This is that event.
|
|
The flow named is the one most of the waiting work belongs to, which is
|
|
the half somebody can act on.
|
|
"""
|
|
backlog = self.queue.backlog()
|
|
if backlog < BACKLOG_DEGRADED:
|
|
self._backlog_strikes = 0
|
|
self.behind = False
|
|
return
|
|
|
|
self._backlog_strikes += 1
|
|
if self._backlog_strikes < BACKLOG_STRIKES or self.behind:
|
|
return
|
|
|
|
self.behind = True
|
|
flows = self.queue.backlog_flows()
|
|
worst = max(flows, key=lambda f: flows[f], default="")
|
|
logger.warning("engine behind: %d items waiting (%s)", backlog, worst or "?")
|
|
self._publish(
|
|
{
|
|
"type": "engine_degraded",
|
|
"reason": f"{backlog} items waiting in the queue",
|
|
"flow": worst,
|
|
"ts": time.time(),
|
|
}
|
|
)
|
|
|
|
def _await_capacity(self) -> int:
|
|
"""How many cascades may be claimed now. Zero means the service stops.
|
|
|
|
Claiming past what the pool can run makes nothing faster: the extra
|
|
items queue up inside the pool, count as in flight and hold their
|
|
journal entries open the whole time, which is how four cascade threads
|
|
came to report hundreds busy on a healthy engine. Work left in the
|
|
stream is work that is still anyone's to take; work that is claimed is
|
|
work that is actually being run.
|
|
"""
|
|
with self._inflight_lock:
|
|
while self._inflight >= self.max_cascades and not self._stop.is_set():
|
|
self._inflight_lock.wait(0.5)
|
|
return 0 if self._stop.is_set() else self.max_cascades - self._inflight
|
|
|
|
def _dispatch(self, item: WorkItem) -> None:
|
|
with self._inflight_lock:
|
|
self._inflight += 1
|
|
if item.entry_id:
|
|
self._active.add(item.entry_id)
|
|
try:
|
|
self._cascade_pool.submit(self._handle, item)
|
|
except RuntimeError:
|
|
# Pool already shutting down.
|
|
self._done(item)
|
|
|
|
def _done(self, item: WorkItem) -> None:
|
|
with self._inflight_lock:
|
|
self._inflight -= 1
|
|
self._active.discard(item.entry_id)
|
|
self._inflight_lock.notify_all()
|
|
|
|
# -------------------------------------------------------------------------
|
|
# Handling one item
|
|
# -------------------------------------------------------------------------
|
|
|
|
def step(self, flow: str) -> str | None:
|
|
"""Run one item a pause is holding, and hold everything else still.
|
|
|
|
Blocking, so the caller sees the wave finish. Returns the node the item
|
|
came from, or None when nothing is parked for this flow.
|
|
"""
|
|
pipeline = self._pipeline
|
|
if pipeline is None:
|
|
return None
|
|
item = self.queue.unpark_one(flow)
|
|
if item is None:
|
|
return None
|
|
with pipeline.stepping(flow):
|
|
try:
|
|
self._run_item(item)
|
|
except Exception:
|
|
logger.exception("Step of '%s' failed", item.node)
|
|
return item.node
|
|
|
|
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(item)
|
|
|
|
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) and not pipeline.is_stepping(node.flow):
|
|
self.queue.park(node.flow, item)
|
|
return True
|
|
|
|
if item.kind == "flush":
|
|
# A rate-limit window ended; nothing to replay, only to let out.
|
|
# ponytail: no run record for a flush — it is the tail of the run
|
|
# that scheduled it, not a run of its own.
|
|
pipeline.flush(node)
|
|
return True
|
|
|
|
if item.guard_key and str(node.recall(item.guard_key, "")) != item.guard_value:
|
|
# The node moved on while this waited — a restarted timer, say.
|
|
logger.debug("Guard no longer holds for '%s', dropped", item.node)
|
|
return True
|
|
|
|
# Only here is the item certain to run, which is what a run record is.
|
|
now = time.time()
|
|
self._publish(
|
|
{
|
|
"type": "cascade_started",
|
|
"run": item.entry_id,
|
|
"flow": item.flow,
|
|
"node": item.node,
|
|
"cause": item.cause,
|
|
"deliveries": item.deliveries,
|
|
"ts": now,
|
|
}
|
|
)
|
|
if item.deliveries == 1:
|
|
# A redelivery waited for the reaper, not for the engine.
|
|
self._publish(
|
|
{
|
|
"type": "work_latency",
|
|
"flow": item.flow,
|
|
"node": item.node,
|
|
"lag_ms": max(
|
|
0.0,
|
|
(now - max(item.enqueued_at, item.not_before)) * 1000,
|
|
),
|
|
"ts": now,
|
|
}
|
|
)
|
|
|
|
replay = item.deliveries > 1
|
|
emission = item.kind == "emission"
|
|
try:
|
|
# An emission's values went into state when the node produced them;
|
|
# this item carries them so its readers get the chunk that caused
|
|
# the wave rather than whichever is newest by the time they run.
|
|
# Applying them again would let a mid-node emission overwrite what
|
|
# the node returned at the end.
|
|
published = (
|
|
set(item.outputs)
|
|
if emission
|
|
else pipeline.apply_outputs(node, item.outputs or None)
|
|
)
|
|
pipeline.run_downstream(
|
|
node,
|
|
entry_id=item.entry_id,
|
|
replay=replay,
|
|
# A redelivery has to finish a walk that may be half done, and
|
|
# an item with no payload is the value already being in state.
|
|
# Neither can say what changed, so neither filters on it.
|
|
changed=None if replay or not item.outputs else published,
|
|
overrides=item.outputs if emission else None,
|
|
)
|
|
finally:
|
|
# Paired, or a cascade that raised — state backend gone, say — is a
|
|
# run left open until the abandoned sweep ten minutes later.
|
|
self._publish(
|
|
{
|
|
"type": "cascade_finished",
|
|
"run": item.entry_id,
|
|
"flow": item.flow,
|
|
"ts": time.time(),
|
|
}
|
|
)
|
|
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
|
|
stats["behind"] = self.behind
|
|
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)
|