Docs / docs (push) Successful in 30s
Playwright Tests / test-playwright (1, 2) (push) Successful in 3m7s
Playwright Tests / test-playwright (2, 2) (push) Successful in 1m54s
pre-commit / pre-commit (push) Failing after 4m24s
Test Backend / test-backend (push) Successful in 3m8s
Compose Smoke Test / test-compose (push) Successful in 40s
Playwright Tests / merge-reports (push) Successful in 1m33s
A port may now declare `image`, `audio` or `video`. Each is the artifact
reference the engine already had, narrowed by the `media_type` on it, so a
speech recogniser declares what it eats rather than taking any bytes at all and
finding out. Bytes still never travel as a message and nothing on the wire
stops being JSON: a camera publishes one reference per frame, a microphone one
per chunk, and a reference may carry a `meta` dict nothing here interprets.
Streaming media is therefore an ordinary streaming port — with one change to
what that means. An emission used to journal an item with no payload, so
downstream read whatever was current when the item was claimed; a consumer
slower than its producer saw only the newest chunk and the ones between were
lost. That is right for a training curve and wrong for a second of speech, so
an emission now journals a `kind="emission"` item carrying its values, and the
executor hands them to the nodes reading that message instead of writing them
to state again. The value in state stays the latest, which is what everything
else reads, and the wave is filtered by what actually changed rather than
walking everything reachable. No queue serialization change — the existing
`outputs` field carries it.
Continuous media makes the store's missing GC a real problem, so this closes
it: `sweep_artifacts` runs hourly, keeps every digest a `run_artifact` row
records or a live message holds, spares anything written in the last hour, and
stands aside entirely while a run is in flight, since a node may store a
checkpoint long before it returns the reference to it. That also collects the
orphans a deleted flow has always left behind. `ARTIFACT_GC_INTERVAL_S=0` turns
it off.
Around the edges: `GET /artifacts/{digest}` serves the media type the caller
passes and answers ranged requests, so a browser plays a clip rather than
downloading it; `PUT` spools to disk instead of holding the whole body in
memory, as does `save_artifact` given a path; a Media widget draws whatever its
message points at, and a wall panel may fetch the bytes its own tiles are
showing and nothing else; and a connector gets `save_artifact`, for a device
whose readings are bytes.
What this cannot do is live video: a frame every second or two is a glance, and
the honest answer above that is the camera's own stream, which the widget takes
as a URL and the browser plays from source.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
462 lines
17 KiB
Python
462 lines
17 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
|
|
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()
|
|
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()
|
|
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 _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._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_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)
|