**SQLite is the database, and now says so.** `metric_minute` and every run
table are written with `sqlalchemy.dialects.sqlite.insert(...)
.on_conflict_do_update` and with `max(a, b)`, neither of which another
dialect has — so pointing `DATABASE_URL` at Postgres migrated cleanly,
served, logged in, and then lost every observability flush into the
collector's hold buffer and failed every run. It refuses at startup
instead. (The Postgres in the compose stack is Umami's; the engine's own
database has been a file beside the flows since 2026-08-21.)
**Every integer query parameter is bounded.** The caps were written as
`min(limit, 500)`, which a negative walks straight through — `?limit=-1`
compiles to `LIMIT -1` and SQLite returns the whole table. Ten signatures,
now `Query(ge=…, le=…)`. `hours=0` still means an hour, which
`_window_hours` was already deliberate about.
**Exports are capped at 10 000 runs** and say so with `X-Truncated`. The
filters bounded a sensible request and nothing bounded an unfiltered one,
which read every row into memory before a byte was streamed. `_series`
resolves cached curves in two queries rather than a `Run` lookup and a
`RunMetric` query per restored node — a comparison of twenty runs was
calling that twenty times over.
**`PUT /artifacts` has a size limit** (`MAX_ARTIFACT_BYTES`, 2 GiB, 0 to
disable), checked against `Content-Length` and again against the stream for
a chunked body, and its writes moved off the event loop.
**`/observability/timeseries` takes `since`/`until`**, the same window
`/runs` and `/events` take, capped at 2000 points — `hours=720&bucket_s=60`
was 43 200 of them in one array. It is also what a dragged chart needs to
re-fetch at its own resolution rather than magnifying buckets it has.
**Composite indexes** for the three list screens: `run(flow, created_at)`
and `(status, created_at)`, `flow_run(flow, started_at)`,
`engine_event(type, ts)`. Every index was single-column, so SQLite picked
one and sorted the rest by hand. Verified against a copy of a live database
(250k `flow_run` rows): the planner takes all four.
**Redis clients have socket timeouts.** A Redis that stops answering
without closing the connection hung the caller until the kernel gave up —
including `/utils/health/`, whose job is to notice.
**The panels file is written under one lock.** `save_panels` and
`unpair_panel` are both read-modify-write, and a save that read before an
unpair wrote put the old nonce back — silently un-revoking a screen that
had just been unpaired. The nonce carry-forward was written to make that
impossible; the gap between its read and its write is where it happened.
**Startup releases what it acquired.** Everything past `event_bus.bind`
registers how to close itself and the `finally` walks that list backwards;
a failure part-way through used to reach none of the shutdown steps and
leave the worker pool's subprocesses and every background task behind —
under `--reload`, once per bad edit. `modules.reconcile` moved into the
background: `uv` gets five minutes twice over, the healthcheck allows
eighty seconds, and the autoheal restarted the container before it could
finish installing.
`delete_run` takes SQLite's write lock up front (`core.db.writing`) rather
than upgrading a deferred transaction and losing to whichever flush
committed in between. `modules.sync` is serialised — two applies mutated
one venv at once. The proxied-call and stream dicts are bounded, and a
reused id cancels its predecessor instead of dropping the reference.
Security, found in passing and small enough to fix here:
- **`/secrets/` required only a signed-in user.** The names alone say what
this installation talks to, and `PUT /{name}` takes any name, so any
account could overwrite the credential a flow authenticates with.
Superuser now — which `/search` already assumed and said so.
- **`POST /login/access-token` had no rate limit.** Argon2 is deliberately
expensive and the route is unauthenticated and runs in the shared
threadpool. Ten *failed* attempts per address per five minutes; a
successful sign-in spends nothing.
- **a password reset link worked repeatedly for 48 hours.** The token now
carries a digest of the password hash it was minted against, so it stops
verifying once it has set one. No table of spent tokens needed.
- **enrolment accepted `http://`**, sending the claim code and then this
installation's credential in clear. https, or a local address.
- the rate limiter read `request.client.host`, which behind Traefik is the
proxy — so every per-address limit was one global bucket and one caller
could lock out everyone. It reads the forwarded address, and its
bucket table is capped rather than growing one key per address forever.
- SMTP has a timeout and sends after the response, so an unreachable mail
host cannot pin a threadpool worker, and a reply's timing no longer says
whether the address exists.
Test suite: engine-written rows are cleared between modules. A `FlowRun`
left `running` by one module turned up in another's query. Per-test
rollback is not available here — the module-scoped `client` runs the real
lifespan and its collector and run service write through sessions of their
own — so this bounds it where the writes come from. Three consecutive
green runs, orders randomised.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01M6hPWS6YEbT1P8LxhhFb2T
778 lines
30 KiB
Python
778 lines
30 KiB
Python
"""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 logging
|
|
import threading
|
|
import time
|
|
from abc import ABC, abstractmethod
|
|
from collections import OrderedDict, deque
|
|
from collections.abc import Callable
|
|
from dataclasses import dataclass, field
|
|
from typing import Any, cast
|
|
|
|
import orjson
|
|
import redis
|
|
|
|
from fluksio.flow.state import CONNECT_TIMEOUT_S, SOCKET_TIMEOUT_S
|
|
|
|
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"
|
|
# Entry ids from the due lane wear this, so one id still says which stream it
|
|
# came from. Deterministic per stream, which is what keeps it usable as an
|
|
# idempotency key across redeliveries.
|
|
DUE_PREFIX = "due:"
|
|
# A consumer this far past its last read belongs to an engine that is gone. A
|
|
# live one interacts every claim, so nothing in service comes close.
|
|
STALE_CONSUMER_IDLE_MS = 3_600_000
|
|
# How many "this side effect already happened" markers the memory queue keeps.
|
|
# The Redis one expires each after an hour; this is the same idea sized by
|
|
# count, since redelivery happens seconds after the claim and never later.
|
|
DONE_MARKERS = 10_000
|
|
|
|
|
|
@dataclass
|
|
class WorkItem:
|
|
"""One unit of journaled work.
|
|
|
|
:param kind: ``cascade`` replays a node's outputs and runs what is
|
|
downstream; ``emission`` is a value a node published while still
|
|
running, already in state, carried so its readers get *that* value
|
|
rather than whichever is newest when they run; ``flush`` lets out what
|
|
a node's rate limits held back; ``run`` is a whole batch run, and
|
|
carries only its id.
|
|
:param node: The node the item is about — the source for a cascade, the
|
|
one whose rate limits are released for a flush.
|
|
:param flow: The flow that node belongs to, so gating needs no lookup.
|
|
:param outputs: What the source node emitted (cascade and emission).
|
|
:param cause: Where the work came from, for logs and debugging.
|
|
:param not_before: Epoch seconds before which the item must not run.
|
|
:param guard: ``(key, value)`` the target node must still remember for
|
|
this item to be worth running — how a rescheduled wait cancels the
|
|
one it replaced.
|
|
: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.
|
|
:param enqueued_at: When the item was made. Carried rather than read off
|
|
the entry id, because only the Redis queue has timestamps in its ids.
|
|
"""
|
|
|
|
kind: str
|
|
node: str
|
|
flow: str
|
|
outputs: dict[str, Any] = field(default_factory=dict)
|
|
cause: str = "system"
|
|
not_before: float = 0.0
|
|
guard_key: str = ""
|
|
guard_value: str = ""
|
|
entry_id: str = ""
|
|
deliveries: int = 1
|
|
enqueued_at: float = field(default_factory=time.time)
|
|
#: The run this item is, on a ``run`` item. Everything else about a run is
|
|
#: in its database row, so the journal only has to name it.
|
|
run_id: str = ""
|
|
|
|
def to_fields(self) -> dict[str, str]:
|
|
return {
|
|
"kind": self.kind,
|
|
"node": self.node,
|
|
"flow": self.flow,
|
|
"outputs": orjson.dumps(self.outputs).decode(),
|
|
"cause": self.cause,
|
|
"not_before": str(self.not_before),
|
|
"guard_key": self.guard_key,
|
|
"guard_value": self.guard_value,
|
|
"enqueued_at": str(self.enqueued_at),
|
|
"run_id": self.run_id,
|
|
}
|
|
|
|
@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=orjson.loads(fields.get("outputs") or "{}"),
|
|
cause=fields.get("cause", "system"),
|
|
not_before=float(fields.get("not_before") or 0.0),
|
|
guard_key=fields.get("guard_key", ""),
|
|
guard_value=fields.get("guard_value", ""),
|
|
entry_id=entry_id,
|
|
deliveries=deliveries,
|
|
enqueued_at=float(fields.get("enqueued_at") or time.time()),
|
|
run_id=fields.get("run_id", ""),
|
|
)
|
|
|
|
|
|
class WorkQueue(ABC):
|
|
"""What the execution service pulls from."""
|
|
|
|
#: Called after an item is journaled for later, so whoever is waiting on
|
|
#: the next deadline can find out it just moved closer. Set by the
|
|
#: execution service; a queue nobody is timing leaves it alone.
|
|
on_delayed: Callable[[], None] | None = None
|
|
|
|
def _woken(self) -> None:
|
|
"""Say a deadline changed. Never lets a listener break the enqueue."""
|
|
wake = self.on_delayed
|
|
if wake is None:
|
|
return
|
|
try:
|
|
wake()
|
|
except Exception:
|
|
logger.exception("Delayed-work listener failed")
|
|
|
|
@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 next_due(self) -> float | None:
|
|
"""When the soonest delayed item comes due, or None if there is none.
|
|
|
|
What lets the engine sleep until a timer rather than poll for it: a
|
|
fixed poll makes every delayed item late by up to its interval, which
|
|
a motor driven for a measured number of seconds notices.
|
|
"""
|
|
|
|
@abstractmethod
|
|
def claim(
|
|
self, count: int, block_ms: int, due_only: bool = False
|
|
) -> list[WorkItem]:
|
|
"""Take up to ``count`` items, waiting up to ``block_ms`` for one.
|
|
|
|
``due_only`` takes nothing but promoted timers. A saturated engine
|
|
uses it to let a motor's stop past the cascade limit while every slot
|
|
is held by a long node.
|
|
"""
|
|
|
|
@abstractmethod
|
|
def ack(self, item: WorkItem) -> None:
|
|
"""Mark an item done, so it is never redelivered."""
|
|
|
|
def touch(self, entry_ids: list[str]) -> None: # noqa: B027
|
|
"""Say these items are still being worked on, not abandoned.
|
|
|
|
A node with no timeout may run far longer than the reclaim window, and
|
|
nothing else distinguishes that from an engine that died holding the
|
|
item. Concrete rather than abstract: a queue with no redelivery has
|
|
nothing to answer here.
|
|
"""
|
|
|
|
@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 unpark_one(self, flow: str) -> WorkItem | None:
|
|
"""Take the oldest held item, leaving the rest parked. That is a step."""
|
|
|
|
@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]:
|
|
"""Backlog, in-flight, delayed and parked counts plus the oldest
|
|
pending age, for the health endpoint."""
|
|
|
|
def backlog(self) -> int:
|
|
"""How much work is waiting to be claimed.
|
|
|
|
Distinct from ``pending``, which is what has been handed to a consumer
|
|
and not yet acknowledged — work in progress. An engine hours behind has
|
|
a small ``pending`` and a large backlog, which is why one is not the
|
|
other. Concrete rather than abstract so the degradation watcher can ask
|
|
any queue.
|
|
"""
|
|
return 0
|
|
|
|
def backlog_flows(self, sample: int = 100) -> dict[str, int]:
|
|
"""Which flows the waiting work belongs to, as far as can be sampled."""
|
|
return {}
|
|
|
|
@abstractmethod
|
|
def dead_letters(self, count: int = 50) -> list[dict[str, Any]]:
|
|
"""What was given up on, newest first."""
|
|
|
|
@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()
|
|
# Promoted timers, kept apart from the backlog the way Redis keeps
|
|
# two streams — so they can be claimed on their own.
|
|
self._due: deque[WorkItem] = deque()
|
|
self._delayed: list[tuple[float, int, WorkItem]] = []
|
|
self._parked: dict[str, list[WorkItem]] = {}
|
|
# Insertion-ordered and capped, because this grows one entry per
|
|
# non-idempotent node per item and nothing ever removed one — the
|
|
# Redis side expires its markers after an hour, this one leaked for
|
|
# the life of the process. Redelivery is what the marker guards, and
|
|
# that happens within seconds of the claim.
|
|
self._done: OrderedDict[tuple[str, str], None] = OrderedDict()
|
|
# Claimed and not yet acknowledged, which is what Redis's `pending` is.
|
|
self._in_flight = 0
|
|
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()
|
|
# Outside the lock: the listener is whoever is asleep until the old
|
|
# deadline, and it comes straight back here to read the new one.
|
|
self._woken()
|
|
|
|
def next_due(self) -> float | None:
|
|
with self._lock:
|
|
return self._delayed[0][0] if self._delayed else None
|
|
|
|
def claim(
|
|
self, count: int, block_ms: int, due_only: bool = False
|
|
) -> list[WorkItem]:
|
|
deadline = time.monotonic() + block_ms / 1000.0
|
|
with self._wake:
|
|
while not self._due and not (self._items and not due_only):
|
|
remaining = deadline - time.monotonic()
|
|
if remaining <= 0:
|
|
return []
|
|
self._wake.wait(remaining)
|
|
# Due first: an item that has waited out a deadline is late by
|
|
# however long it queues here, while work merely enqueued is not
|
|
# waiting on a clock.
|
|
claimed = [self._due.popleft() for _ in range(min(count, len(self._due)))]
|
|
if not due_only:
|
|
claimed += [
|
|
self._items.popleft()
|
|
for _ in range(min(count - len(claimed), len(self._items)))
|
|
]
|
|
self._in_flight += len(claimed)
|
|
return claimed
|
|
|
|
def ack(self, item: WorkItem) -> None:
|
|
"""Nothing to take off the queue — claiming did that. Only counted."""
|
|
with self._lock:
|
|
self._in_flight = max(0, self._in_flight - 1)
|
|
|
|
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:
|
|
due: list[WorkItem] = []
|
|
with self._wake:
|
|
while self._delayed and self._delayed[0][0] <= now:
|
|
due.append(heapq.heappop(self._delayed)[2])
|
|
if due:
|
|
self._due.extend(due)
|
|
self._wake.notify()
|
|
return len(due)
|
|
|
|
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 unpark_one(self, flow: str) -> WorkItem | None:
|
|
with self._lock:
|
|
held = self._parked.get(flow)
|
|
return held.pop(0) if held else None
|
|
|
|
def clear_flow(self, flow: str) -> None:
|
|
with self._lock:
|
|
self._parked.pop(flow, None)
|
|
|
|
def backlog(self) -> int:
|
|
with self._lock:
|
|
return len(self._items) + len(self._due)
|
|
|
|
def stats(self) -> dict[str, Any]:
|
|
with self._lock:
|
|
return {
|
|
"backlog": len(self._items) + len(self._due),
|
|
"pending": self._in_flight,
|
|
"delayed": len(self._delayed),
|
|
"parked": sum(len(v) for v in self._parked.values()),
|
|
"oldest_pending_s": 0.0,
|
|
"durable": False,
|
|
}
|
|
|
|
def dead_letters(self, count: int = 50) -> list[dict[str, Any]]:
|
|
"""Nothing is kept: a dropped item here died with the process."""
|
|
return []
|
|
|
|
def mark_done(self, entry_id: str, node: str) -> None:
|
|
with self._lock:
|
|
self._done[(entry_id, node)] = None
|
|
while len(self._done) > DONE_MARKERS:
|
|
self._done.popitem(last=False)
|
|
|
|
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,
|
|
# A root of its own: the state backend owns `pipeline:`, and sharing the
|
|
# prefix meant `RedisState.clear()` could DEL the queue stream.
|
|
namespace: str = "queue",
|
|
consumer: str | None = None,
|
|
) -> None:
|
|
self._redis = redis.Redis(
|
|
host=host,
|
|
port=port,
|
|
decode_responses=True,
|
|
# Longer than any block this asks for, so a claim still waits its
|
|
# second — but bounded, so a Redis that goes silent without
|
|
# closing the socket fails rather than hanging the consumer.
|
|
socket_timeout=SOCKET_TIMEOUT_S,
|
|
socket_connect_timeout=CONNECT_TIMEOUT_S,
|
|
health_check_interval=30,
|
|
)
|
|
self._ns = namespace
|
|
self._consumer = consumer or f"engine-{int(time.time() * 1000) % 1_000_000}"
|
|
self._stream = f"{namespace}:__queue__"
|
|
# Promoted timers, read before the backlog. A stream of its own rather
|
|
# than a flag on the entry, because a stream is ordered by id and there
|
|
# is no way to put a late item in front of one already in it.
|
|
self._due_stream = f"{namespace}:__due__"
|
|
self._delayed_key = f"{namespace}:__delayed__"
|
|
self._dead_key = f"{namespace}:__dead__"
|
|
# Which flows have something parked. Maintained rather than discovered:
|
|
# finding them with a keyspace scan walked every state and idempotency
|
|
# key in the database, twice per health poll.
|
|
self._parked_flows_key = f"{namespace}:__parked_flows__"
|
|
self._ensure_group()
|
|
self._drop_stale_consumers()
|
|
|
|
def _ensure_group(self) -> None:
|
|
for stream in (self._stream, self._due_stream):
|
|
try:
|
|
self._redis.xgroup_create(stream, GROUP, id="0", mkstream=True)
|
|
except redis.ResponseError as exc:
|
|
if "BUSYGROUP" not in str(exc):
|
|
raise
|
|
|
|
def _drop_stale_consumers(self) -> None:
|
|
"""Forget the consumers of engines that are gone.
|
|
|
|
One is registered per engine start and none was ever removed, so the
|
|
list grew without bound — and `XPENDING`, `XAUTOCLAIM` and `XINFO` all
|
|
walk it. Only consumers holding nothing and untouched for an hour are
|
|
dropped, so a live engine blocked on a read is never one of them; a
|
|
consumer is recreated by its next read anyway.
|
|
"""
|
|
for stream in (self._stream, self._due_stream):
|
|
try:
|
|
consumers = cast(
|
|
list[dict[str, Any]], self._redis.xinfo_consumers(stream, GROUP)
|
|
)
|
|
except redis.ResponseError:
|
|
continue
|
|
for row in consumers:
|
|
name = str(row.get("name", ""))
|
|
if name == self._consumer or row.get("pending"):
|
|
continue
|
|
if int(row.get("idle") or 0) < STALE_CONSUMER_IDLE_MS:
|
|
continue
|
|
try:
|
|
self._redis.xgroup_delconsumer(stream, GROUP, name)
|
|
except redis.ResponseError:
|
|
pass
|
|
|
|
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, {orjson.dumps(item.to_fields()): not_before}
|
|
)
|
|
self._woken()
|
|
|
|
def next_due(self) -> float | None:
|
|
soonest = cast(
|
|
list[tuple[str, float]],
|
|
self._redis.zrange(self._delayed_key, 0, 0, withscores=True),
|
|
)
|
|
return soonest[0][1] if soonest else None
|
|
|
|
def _route(self, entry_id: str) -> tuple[str, str]:
|
|
"""The stream an entry id belongs to, and its plain id there."""
|
|
if entry_id.startswith(DUE_PREFIX):
|
|
return self._due_stream, entry_id[len(DUE_PREFIX) :]
|
|
return self._stream, entry_id
|
|
|
|
def claim(
|
|
self, count: int, block_ms: int, due_only: bool = False
|
|
) -> list[WorkItem]:
|
|
"""Take up to ``count`` items, due timers before anything queued.
|
|
|
|
One read over both streams rather than a read each: the block has to
|
|
cover the due lane, or an item promoted while this waits sits there
|
|
until the block expires — which is the second of lateness the sleeping
|
|
timer thread exists to remove. ``COUNT`` is per stream, so a read may
|
|
return up to twice what was asked for; the excess is bounded by the
|
|
cascade pool's own free count, and starting a due timer is worth more
|
|
than holding that number exactly.
|
|
"""
|
|
response = cast(
|
|
list[Any],
|
|
self._redis.xreadgroup(
|
|
GROUP,
|
|
self._consumer,
|
|
(
|
|
{self._due_stream: ">"}
|
|
if due_only
|
|
else {self._due_stream: ">", self._stream: ">"}
|
|
),
|
|
count=count,
|
|
block=block_ms,
|
|
),
|
|
)
|
|
due: list[WorkItem] = []
|
|
backlog: list[WorkItem] = []
|
|
for stream, entries in response or []:
|
|
is_due = stream == self._due_stream
|
|
into = due if is_due else backlog
|
|
prefix = DUE_PREFIX if is_due else ""
|
|
for entry_id, fields in entries:
|
|
into.append(WorkItem.from_fields(fields, prefix + entry_id))
|
|
return due + backlog
|
|
|
|
def ack(self, item: WorkItem) -> None:
|
|
if item.entry_id:
|
|
stream, entry_id = self._route(item.entry_id)
|
|
self._redis.xack(stream, GROUP, entry_id)
|
|
|
|
def touch(self, entry_ids: list[str]) -> None:
|
|
if not entry_ids:
|
|
return
|
|
# Claiming an entry we already hold resets how long it has been idle,
|
|
# which is the only thing `reclaim_stale` reads. `justid` keeps the
|
|
# delivery count where it is, so a long node neither redelivers nor
|
|
# spends its way towards the dead-letter cap. An entry already
|
|
# acknowledged is simply not there, and this is a no-op for it.
|
|
by_stream: dict[str, list[str]] = {}
|
|
for entry_id in entry_ids:
|
|
stream, plain = self._route(entry_id)
|
|
by_stream.setdefault(stream, []).append(plain)
|
|
for stream, ids in by_stream.items():
|
|
self._redis.xclaim(
|
|
stream,
|
|
GROUP,
|
|
self._consumer,
|
|
min_idle_time=0,
|
|
# redis-py types `message_ids` as an invariant list of every id
|
|
# shape it accepts, which a list[str] is not assignable to.
|
|
message_ids=cast("list[int | bytes | str | memoryview]", ids),
|
|
justid=True,
|
|
)
|
|
|
|
def reclaim_stale(self, min_idle_ms: int) -> list[WorkItem]:
|
|
"""Take over entries a dead consumer never acknowledged."""
|
|
return self._reclaim(self._due_stream, DUE_PREFIX, min_idle_ms) + self._reclaim(
|
|
self._stream, "", min_idle_ms
|
|
)
|
|
|
|
def _reclaim(self, stream: str, prefix: str, min_idle_ms: int) -> list[WorkItem]:
|
|
_cursor, entries, _deleted = cast(
|
|
tuple[Any, list[Any], Any],
|
|
self._redis.xautoclaim(
|
|
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(stream, GROUP, min="-", max="+", count=64),
|
|
):
|
|
counts[record["message_id"]] = record["times_delivered"]
|
|
|
|
return [
|
|
WorkItem.from_fields(fields, prefix + 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),
|
|
)
|
|
if not due:
|
|
return 0
|
|
# Whoever removes it owns it: a second engine gets 0 for that member and
|
|
# leaves it alone. One round trip for the batch rather than one each,
|
|
# which matters because this runs every second.
|
|
pipe = self._redis.pipeline()
|
|
for raw in due:
|
|
pipe.zrem(self._delayed_key, raw)
|
|
claimed = cast(list[int], pipe.execute())
|
|
|
|
pipe = self._redis.pipeline()
|
|
moved = 0
|
|
for raw, owned in zip(due, claimed, strict=True):
|
|
if owned:
|
|
# Into the due lane, which `claim` reads first: what has waited
|
|
# out a deadline should not then queue behind a backlog.
|
|
pipe.xadd(
|
|
self._due_stream,
|
|
cast(Any, orjson.loads(raw)),
|
|
maxlen=STREAM_MAXLEN,
|
|
approximate=True,
|
|
)
|
|
moved += 1
|
|
if moved:
|
|
pipe.execute()
|
|
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:
|
|
pipe = self._redis.pipeline()
|
|
pipe.rpush(self._parked_key(flow), orjson.dumps(item.to_fields()))
|
|
pipe.sadd(self._parked_flows_key, flow)
|
|
pipe.execute()
|
|
|
|
def unpark(self, flow: str) -> list[WorkItem]:
|
|
key = self._parked_key(flow)
|
|
raw = cast(list[str], self._redis.lrange(key, 0, -1))
|
|
pipe = self._redis.pipeline()
|
|
pipe.delete(key)
|
|
pipe.srem(self._parked_flows_key, flow)
|
|
pipe.execute()
|
|
return [WorkItem.from_fields(orjson.loads(r), "") for r in raw]
|
|
|
|
def unpark_one(self, flow: str) -> WorkItem | None:
|
|
key = self._parked_key(flow)
|
|
pipe = self._redis.pipeline()
|
|
pipe.lpop(key)
|
|
pipe.llen(key)
|
|
raw, remaining = cast(tuple["str | None", int], pipe.execute())
|
|
if not remaining:
|
|
self._redis.srem(self._parked_flows_key, flow)
|
|
return WorkItem.from_fields(orjson.loads(raw), "") if raw else None
|
|
|
|
def clear_flow(self, flow: str) -> None:
|
|
pipe = self._redis.pipeline()
|
|
pipe.delete(self._parked_key(flow))
|
|
pipe.srem(self._parked_flows_key, flow)
|
|
pipe.execute()
|
|
|
|
def _group_info(self, stream: str = "") -> dict[str, Any]:
|
|
"""This consumer group's row of ``XINFO GROUPS``, or an empty one."""
|
|
try:
|
|
groups = cast(
|
|
list[dict[str, Any]], self._redis.xinfo_groups(stream or self._stream)
|
|
)
|
|
except redis.ResponseError:
|
|
# No stream yet: nothing has ever been enqueued.
|
|
return {}
|
|
return next((g for g in groups if g.get("name") == GROUP), {})
|
|
|
|
def backlog(self) -> int:
|
|
"""Entries in either stream this group has never been handed.
|
|
|
|
Redis calls it the group's ``lag``. It is nil rather than zero when the
|
|
stream has been trimmed under the group — entries that were dropped
|
|
before anyone read them — and unknown is reported as none waiting,
|
|
since the alternative is a health screen crying wolf after a trim.
|
|
"""
|
|
return sum(
|
|
int(self._group_info(stream).get("lag") or 0)
|
|
for stream in (self._stream, self._due_stream)
|
|
)
|
|
|
|
def backlog_flows(self, sample: int = 100) -> dict[str, int]:
|
|
"""Which flows the waiting work belongs to, from a sample of the tail.
|
|
|
|
The group's lag is one number for the whole stream, and the actionable
|
|
half of "the engine is behind" is always *what* is producing the work.
|
|
Reads forward from the last entry the group was handed, which is where
|
|
the undelivered entries start.
|
|
"""
|
|
after = self._group_info().get("last-delivered-id")
|
|
if not after:
|
|
return {}
|
|
entries = cast(
|
|
list[tuple[str, dict[str, str]]],
|
|
self._redis.xrange(self._stream, min=f"({after}", max="+", count=sample),
|
|
)
|
|
counts: dict[str, int] = {}
|
|
for _entry_id, fields in entries:
|
|
flow = fields.get("flow", "")
|
|
if flow:
|
|
counts[flow] = counts.get(flow, 0) + 1
|
|
return counts
|
|
|
|
def stats(self) -> dict[str, Any]:
|
|
count = 0
|
|
oldest = 0.0
|
|
for stream in (self._stream, self._due_stream):
|
|
pending = cast(dict[str, Any], self._redis.xpending(stream, GROUP))
|
|
waiting = pending.get("pending", 0) if isinstance(pending, dict) else 0
|
|
count += waiting
|
|
if not waiting:
|
|
continue
|
|
records = cast(
|
|
list[dict[str, Any]],
|
|
self._redis.xpending_range(stream, GROUP, min="-", max="+", count=1),
|
|
)
|
|
if records:
|
|
oldest = max(oldest, records[0]["time_since_delivered"] / 1000.0)
|
|
flows = cast(set[str], self._redis.smembers(self._parked_flows_key))
|
|
parked = sum(cast(int, self._redis.llen(self._parked_key(f))) for f in flows)
|
|
return {
|
|
"backlog": self.backlog(),
|
|
"pending": count,
|
|
"delayed": cast(int, self._redis.zcard(self._delayed_key)),
|
|
"parked": parked,
|
|
"oldest_pending_s": round(oldest, 1),
|
|
"durable": True,
|
|
}
|
|
|
|
def dead_letters(self, count: int = 50) -> list[dict[str, Any]]:
|
|
"""What was set aside, newest first. The entry id carries the time."""
|
|
entries = cast(
|
|
list[tuple[str, dict[str, str]]],
|
|
self._redis.xrevrange(self._dead_key, count=count),
|
|
)
|
|
return [
|
|
{
|
|
"id": entry_id,
|
|
"ts": int(entry_id.split("-")[0]) / 1000.0,
|
|
"flow": fields.get("flow", ""),
|
|
"node": fields.get("node", ""),
|
|
"cause": fields.get("cause", ""),
|
|
"reason": fields.get("reason", ""),
|
|
}
|
|
for entry_id, fields in entries
|
|
]
|
|
|
|
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()
|