29. pepsi-dispatch¶
The single long-lived coordinator that drives messages through the stages.
29.1. Role¶
pepsi-dispatch is the only process that starts stage programs. Run exactly
one per system, as pepsi-dispatch serve (the one subcommand; the global
-c/-L/-v flags come before it). It claims newly queued rows and feeds
each to an idle worker process of the row’s stage, requeues paused messages
when their retry time is due, and recovers from crashes and timeouts. It processes no
message content itself. Started as root it drops to the pepsi service
account, which the workers inherit. Reference: pepsi-dispatch(1).
29.2. Features¶
Event-driven, batched claiming.
LISTENs on theworkqueueandconfig_changedchannels with a periodicPOLL_INTERVALheartbeat. On any wake it drains its pending internal events (so per-stage capacities are current) and then claims work for every stage with spare capacity in oneUPDATE … RETURNING— each stage up to its own remaining capacity (a per-stageLATERAL … ORDER BY workqueue_id LIMIT capthat stops at the cap via the partialpendingindex). With exactly one dispatcher claiming serially as the solepending``→``runningwriter, the claim needs noFOR UPDATE SKIP LOCKED.Persistent, pipelined workers. Each stage runs
PROGRAM worker(found onPATHunless absolute) as long-lived processes that connect to the database once and read message ids from standard input, answering with a one-line status (0on success) per message, in order. A worker is handed up toQUEUE_LIMITids at once (written to its stdin in a single buffer); it processes them strictly serially, but the next id is already queued, removing a coordinator round-trip per message. A stage’s in-flight capacity is thereforeQUEUE_LIMIT × PARALLELISM, and because the worker count is unchanged it costs no extra database connections.Elastic per-stage pools. Workers are started on demand up to a stage’s
PARALLELISM; work is packed onto the fewest workers, and surplus workers are stopped afterWORKER_IDLE_TIMEOUT— so a quiet stage collapses toward zero workers and a busy one grows to its cap. AfterMAX_MESSAGESa worker recycles its child (once its pipeline drains); a stage whosePROGRAMcannot start is held off with a short cooldown instead of spinning.Advancing is the stage’s write. A stage that advances a message sets the row
pendingat the next stage itself, in the same update that movesstage; the dispatcher never rewrites a row on success. The worker’s status report brings the coordinator round to its next claim, which picks the row up for the next stage’s pool (so an advance needs no notification). A worker may also fuse a fast successor into the same pass (see pepsi-dispatch(1), Stage fusion).Retry scheduling. Sleeps until the earliest
pausedrow’stimeout(or a due worker reap / spawn cooldown / the heartbeat), then flips due rows back topendingand re-scans.Crash & watchdog recovery. A worker that reports non-zero fails that one message and survives; a crashed child fails its head-of-line message, and a head-of-line message unanswered within
MAX_RUNTIME(timed from when it reached the head) →timeout. In the crash and timeout cases the worker’s other pipelined messages never ran and are requeued topendingrather than lost. Reasons are written tostate.dispatch_errorunder astatus = 'running'guard so a stage that already transitioned the row is never clobbered.Start-up / shutdown safety. Resets orphaned
runningrows (from a previous dispatcher) topendingon start; onSIGINT/SIGTERMstops workers and resets their (and any claimed-but-unassigned) rows.Isolation. At most one writer ever touches a given queue row, so the shared pool runs
READ COMMITTED; queue operations are still wrapped in a retry helper that re-runs a transient40001/40P01failure as cheap insurance.Config propagation. Passes
CONFIG_FILEto workers via-c(before theworkersubcommand) so they load the same configuration.Hot reload of the pipeline. On a
config_changednotification it rebuilds its own stage table from the database overlay and retires every worker — idle ones at once, busy ones when their current message reports back — so a[stage-*]change takes effect without a restart. A reloaded graph that does not parse ends the process rather than routing by a pipeline the workers no longer share. Its own[pepsi-dispatch]options are not reloaded.Statistics. Accumulates per-stage throughput, kills/timeouts, crashes and processing time plus the global stage/message totals in memory, and flushes the deltas to the
stage_stats/dispatch_statstables in a single transaction everySTATS_INTERVAL, whenever the pipeline goes idle (at most once a second, so a finished burst’s numbers do not wait for the next tick), and on shutdown; nothing is written when there are no deltas. It is the only writer of those tables — stages fused into another stage’s worker pass are reported back on that pass’s status line instead of being written by the worker, which keeps every worker of a stage from contending for that stage’s single row. These are exported by pepsi-httpd’s/metrics.
29.3. Configuration¶
[pepsi-dispatch]: MAX_RUNTIME (default 300 s), WORKER_IDLE_TIMEOUT
(default 5 s), POLL_INTERVAL (default 30 s), STATS_INTERVAL (default 60 s),
DB_POOL_SIZE (default 8) and CONFIG_FILE (no default; without it the
workers fall back to the default configuration search). Per-stage worker pools are sized by PARALLELISM (default 4),
QUEUE_LIMIT (default 4, the messages pipelined to one worker at once) and
MAX_MESSAGES (default 1000) in each [stage-<name>] section. See
pepsi.conf(5) and Configuration.
29.4. See also¶
Architecture, pepsi-queue, pepsi-ingress, pepsi-dispatch(1).