85.1.2. pepsi-dispatch¶
drive messages through the stage pipeline
- Manual section:
1
85.1.2.1.1. Name¶
pepsi-dispatch - claim pending messages and run their stage programs.
85.1.2.1.2. Synopsis¶
pepsi-dispatch [GLOBAL-OPTIONS] serve
85.1.2.1.3. Description¶
pepsi-dispatch is the long-lived dispatcher that advances messages through
the stage pipeline (see pepsi.conf(5)). It claims pending rows
from pepsi.workqueue, sets each to running, looks up the message’s stage in
its [stage-<stage>] section, and hands the message id to a worker process
of that stage. Each stage runs as a pool of persistent workers started as
PROGRAM worker (PROGRAM -c FILE worker when CONFIG_FILE is set — the
global flag precedes the subcommand, as everywhere in Pepsi; the program is found
on the PATH unless PROGRAM is an
absolute path); a worker connects to the database once and then reads message ids
from its standard input and writes one status line per message to its
standard output (0 on success), in input order. A status may be followed
by a space and a JSON
report of the stages that were fused into that pass,
[["<stage>",<microseconds>], …], which is how per-stage statistics stay whole
without the workers writing them (see Statistics below). Only an unreadable
status is fatal — it desynchronises the positional protocol, so the dispatcher
tears the worker down; an unreadable report is dropped with a warning, so a worker
and a dispatcher from different builds lose statistics for the length of a restart
window rather than messages. Run exactly one dispatcher per system; it
is the only process that starts stage programs (though many workers then run
concurrently against the shared database). New work is picked up from a
LISTEN on the workqueue channel, with a periodic safety-net heartbeat
(POLL_INTERVAL).
Notifications on that channel are issued explicitly, and only by writers
outside the dispatcher’s own loop: pepsi-ingress(1) (once per burst of
accepted messages, not once per message), the workqueue_inject and
workqueue_resume SQL functions, pepsi-keydisc(1) releasing parked mail,
pepsi-failure-bouncer(1), the pepsi-queue(1) repair commands and
the pepsi-list(1) moderation and digest commands.
A stage advancing a message deliberately does not notify: the worker reports
the finished message on its standard output, and the dispatcher answers that with
the same full claim a notification would have triggered, so the notification would
carry no information. A table trigger notifying on every row that becomes
pending would be worse than useless: PostgreSQL holds a database-wide lock
from the moment a transaction queues a notification until it commits, so every
stage hop would serialise the whole database’s commits — see the Performance
chapter of the manual.
Batched claiming. On any wake (a notification, a freed or dead worker, the
poll heartbeat, a listener (re)connect) the dispatcher first drains every pending
internal event so its per-stage capacities are current, then runs one
UPDATE … RETURNING that claims, for every stage with spare capacity at once,
up to that stage’s own remaining capacity (pending``→``running; which rows
is the subject of the next paragraph — with fair scheduling turned off it is a
per-stage LATERAL subquery whose ORDER BY workqueue_id LIMIT cap stops
scanning each stage’s backlog at its cap, driven by the partial index on
(stage, workqueue_id) WHERE status='pending'). The notification carries an empty
payload — a notification is simply a wake, and a burst of them collapses into one
batched claim (PostgreSQL also coalesces identical queued notifications).
Fair scheduling across senders. Oldest-first would let one sender’s burst
delay every message queued behind it, so by default (FAIR_SCHEDULING, see
pepsi.conf(5)) the claim chooses which rows fill a stage’s spare capacity
the way the Linux scheduler (CFS/EEVDF) chooses which process gets a CPU: a stage
is the CPU and a sender is the process. The sender of a row is the generated
column workqueue.sender_key:
user:<login>— the authenticated account, for a submission (SASL or the local socket’s peer credentials), whatever envelope sender it chose;ip:<address>— otherwise the connecting SMTP peer, an IPv6 one by its /64 (a single host is routinely given a whole /64). Deliberately not the envelope domain, which a remote sender picks freely per message;from:<domain>— for a message Pepsi created itself (no SMTP origin: a notice, a report), the envelope sender’s domain;null:<domain>— for a null-sender one (a DSN, an auto-reply), its first recipient’s domain.
Because the key is generated from the row, every writer — ingress, an injected
message, a split, a clone — gets it without knowing it exists. The dispatcher
keeps, in memory and per stage, each sender’s lead: the worker time its
messages have used there beyond the least-served waiting sender. A claimed
message is charged the stage’s average per-message time at once, corrected to
the measured time when its worker reports, and refunded if it is requeued without
having run; a timed-out message is charged the whole MAX_RUNTIME. The claim
then ranks the k-th oldest waiting row of a sender with lead L by
L + k × average, smallest first, so a sender with a thousand rows waiting
takes turns with one that has a single row instead of going first — and, since
the charge is time rather than a count, a sender whose messages are expensive at a
stage gets correspondingly fewer of its slots. A sender new to the stage, or one
whose lead has decayed away, starts level with the least-served one, so being
quiet buys no later priority. Leads halve every FAIR_HALF_LIFE (default 60 s)
and a lead worth less than one message is forgotten; a restart forgets all of
them, which is harmless. To guarantee that no row waits for ever — however many
new, level senders keep arriving — FAIR_FIFO_PERCENT (default 10) of each
stage’s claimed slots go to its oldest waiting rows regardless of sender, so a
message waits at most until the rows queued at that stage before it have drained
at that reserved rate.
The claim is still one statement for every stage at once, the single claimer’s
pending``→``running write, bounded per stage by its capacity: the
workqueue_claim_fair SQL function receives the leads as arrays, enumerates the
senders waiting at each stage with a loose index scan on the partial index
(stage, sender_key, workqueue_id) WHERE status='pending' (one probe per sender,
at most 64 senders per stage per claim — beyond that the window rotates from one
claim to the next), and reads at most the stage’s capacity of each sender’s
oldest rows, so its cost does not grow with the depth of any one sender’s
backlog. FAIR_SCHEDULING = no restores the plain oldest-first claim.
Pipelined, elastic worker pools. Each stage’s pool grows on demand up to its
PARALLELISM (default 4) and shrinks when idle: no worker runs until the stage
sees traffic; when a stage has work and spare capacity the dispatcher fills an
existing worker’s free slot or starts a new worker; work is packed onto the
fewest workers so a surplus goes cold and is stopped after WORKER_IDLE_TIMEOUT
(default 5 s). Under steady light load a stage therefore hovers around a single
worker. A worker is handed up to QUEUE_LIMIT (default 4) message ids at once —
written to its standard input in a single buffer — rather than one-at-a-time; it
still processes them strictly serially, but the next id is already queued, so
completing a message frees it without a coordinator round-trip first. A stage’s
total in-flight capacity is thus QUEUE_LIMIT × PARALLELISM. Because the worker
count (and so the database-connection count) is unchanged, raising QUEUE_LIMIT
trades a little head-of-line latency for throughput without spending more
connections. After MAX_MESSAGES (default 1000) a worker recycles its child —
once its pipeline has drained — with a fresh process to bound memory growth.
Batched body-free stages. A worker for a stage that loads neither the header
block nor the body (the metadata-only stages — srs, if,
auto-whitelist, discard, aliases …) processes its pipelined
ids in one batch rather than one at a time: it greedily takes every id already
buffered on its standard input (up to QUEUE_LIMIT, stopping the instant a read
would block, so a lone id is still handled immediately), loads them all in a
single SELECT … WHERE workqueue_id = ANY(...), runs each stage body, then
commits the advances/fails/requeues in one arrayed UPDATE … FROM
jsonb_to_recordset(...) per outcome shape (typically one). It still writes one
status line per id, in order, so the worker protocol is unchanged. When the queue
is full this cuts a body-free stage’s database round-trips from two per message to
roughly two per QUEUE_LIMIT messages. Stages that load the headers/body keep
the one-id-at-a-time path (their large columns do not array cheaply).
Advancing. A stage advances a message by setting its stage column to
the next stage and the row back to pending in one update; the dispatcher
never rewrites a row on success. Claiming (pending``→``running) is the
dispatcher’s only write to a row’s status for forward progress, and it is the
only writer that performs it. When a worker reports success (status 0) the
dispatcher re-scans for work — which is precisely why an advance issues no
notification of its own — so the next stage’s pool claims the now
pending row on the following scheduling round. The message is
done when the worker deletes the row, paused when it leaves it paused, or
terminal when failed/timeout. A stage commits its outcome in a single
round-trip — one UPDATE/DELETE, or one stored function that fans work out
server-side (a recipient split, or a finish/pause that also spawns bounce, delay
or success DSNs from arrays of clones) — never a multi-statement client
transaction. A worker’s connection also runs with PostgreSQL’s
synchronous_commit off ([pepsi-postgres] WORKER_SYNCHRONOUS_COMMIT,
default no), which is what lets concurrent workers group-commit. It gives up
nothing but the durability of the last few hundred milliseconds, and losing a
stage advance that way is indistinguishable from a worker killed just before
it: the row is found running at the previous stage, reset to pending on
the next start-up, and the stage runs again. pepsi-ingress(1) deliberately
keeps the default instead, because losing an admission it has already answered
250 to would be a lost message rather than a repeated stage.
Failures and retries. A stage worker announces itself with a ready
line once it has opened its database pool and checked its configuration, and
then answers one status line per message. A worker that reports a non-zero
status leaves that message failed with the reason in state and stays
alive for the next id; but a worker reports that only for an error its stage
marked permanent (see Stage errors below) or after giving up on retries,
and never for the two database statuses. A worker whose child closes its output
(a crash) after it was ready is torn down, and its head-of-line message takes
a strike; so does the head-of-line message of a child that does not answer
within MAX_RUNTIME (real time), which is killed. A message with a strike is
paused and retried after a minute (two after the second strike), with the reason
in state.last_error and the count in state.strikes; only its third strike
is taken to be the message’s own doing, and sets it failed (after a crash)
or timeout (after a hang), with state.failure_class crashed or
timed-out. A worker dies for reasons of its own too – the OOM killer, a
helper being upgraded, a resolver that hangs – and one such death must not cost
the message it happened to hold, while a message that kills every worker it is
given must not be retried for ever. In both cases the other messages already
pipelined to that worker never ran, so they are reset to pending and
re-dispatched (a single stuck or poison message therefore does not strand its
in-flight siblings). A child that dies before it was ready blames nothing: it
failed at its own start-up – an unreachable or connection-exhausted PostgreSQL,
an unreadable secret – so every message it was handed goes back to pending
and the stage is held off for a short cooldown before another worker is started,
instead of being respawned at once. A child that exits with status 78
(EX_CONFIG) before answering has refused to run – because it met a
database schema from another release, in the window of a package upgrade
between the new binaries landing and the schema being upgraded, or because its
configuration does not parse – and is handled the same way, with an error
naming the cause; the stage is retried after the cooldown until the schema
matches or the configuration is fixed. The MAX_RUNTIME clock is
per message and starts when a message reaches the head of the worker’s queue, so a
message waiting behind a slow sibling is not charged for the wait. When a stage
leaves a message paused it records, in timeout, when the message should be
retried; the dispatcher sleeps until the earliest such time, flips the due rows
back to pending, and re-scans for work. At start-up it resets any leftover
running rows (orphaned by a previous dispatcher) back to pending. Because
there is exactly one dispatcher and its coordinator claims serially — the sole
writer that sets pending``→``running — the batched per-stage claim needs no
FOR UPDATE SKIP LOCKED. “Exactly one” is enforced rather than assumed: at
start-up the dispatcher takes a session-scoped PostgreSQL advisory lock on its
database and, if another dispatcher still holds it after about ten seconds,
refuses to run. Without
that, a second dispatcher’s own start-up sweep would immediately un-claim the
first’s in-flight rows and both would dispatch them — duplicate delivery,
duplicate bounces, duplicate payment settlement. The lock is released by
PostgreSQL when the session ends, so a killed dispatcher leaves nothing to clean
up. Because at most one writer ever touches a given row,
the shared database connection runs READ COMMITTED rather than
SERIALIZABLE (queue operations are still routed through a retry helper as
cheap insurance).
Telemetry switch. The dispatcher also listens on telemetry_changed,
which whoever rewrites [pepsi] SHARE_TELEMETRY in the configuration file
sends afterwards, and retires every live worker: a worker installs its feature
telemetry sink once, from the file, so its replacement starts (or stops)
reporting accordingly. The stage table is untouched. See
pepsi-telemetry-client(1).
Configuration reload. Every [stage-*] section is a hot setting: on the
config_changed notification the dispatcher rebuilds its own stage table from
the configuration overlay and retires every live worker, so a stage may be added,
edited or removed while the pipeline runs and the next message is both routed and
executed with the new values. A stage graph that does not parse ends the
process with exit 78 after an orderly shutdown (workers stopped, their rows
back to pending), and the parse error is logged. Keeping the previous table
would be worse: the workers are retired by the same event and
their replacements do read the new overlay, so the dispatcher would be routing
by one pipeline while the stages ran another. The shipped
pepsi-dispatch.service uses Restart=always with no start limit: without
the dispatcher, accepted mail only accumulates in the queue, so the unit keeps
retrying, backing off from 2 s to once a minute (RestartSteps=,
RestartMaxDelaySec=), and picks up a corrected configuration by itself. An
overlay that cannot be read (a database error) is not a new configuration:
the dispatcher keeps its current stage graph and retries the reload, rather
than dropping every stage defined only in the database. The
[pepsi-dispatch] section itself is not reloaded — the connection pool has to exist before the overlay can
be read — so changing it needs a restart.
Database-overload backpressure. A worker that cannot reach the database
because its connection limit is exhausted (PostgreSQL too_many_connections,
or its pool times out acquiring a connection) reports a distinct status
(EX_TEMPFAIL, 75) rather than the generic failure code — this is
infrastructure pressure, not a defect in the message. The dispatcher then does
not fail the message: it requeues it to pending for a later retry, and to
shed connections it temporarily reduces the parallelism of the stage currently
running the most worker processes (the largest contributor to the pressure,
which need not be the stage that reported the error) — halving its cap, floor 1,
for five minutes. Repeated reports ratchet that stage down further and refresh
the window; existing workers drain via the idle reaper and the stage ramps back
up once the window passes. The remedy for sustained pressure is to raise
PostgreSQL’s max_connections or lower the stages’ PARALLELISM so that the
sum of every component’s pool fits the server (see the [pepsi-postgres]
connection-budget note in pepsi.conf(5)).
Lost database connection. A worker whose database connection breaks or
cannot be made while it handles a message (PostgreSQL restarting or shutting
down, a network interruption: an I/O or TLS error, SQLSTATE class 08,
57P01–57P03) reports status 74 (EX_IOERR). That is not a defect
in the message either, so the dispatcher does not fail it: it sets it
paused for 30 seconds, after which the ordinary paused-retry sweep hands it
back to the same stage. This repeats for as long as the database stays
unreachable. Nothing the stage did was committed, so the retry runs the stage
from the start; a side effect outside the database whose commit was lost (a
relay that delivered the message just before the connection broke) can happen
a second time. Each deferral is logged as a warning under
pepsi-dispatch and counted nowhere else: a counter would have to be written
to the database that has just gone away. If the dispatcher cannot write the
pause either, it retries for about half a minute and otherwise leaves the
message running, which its next start resets to pending.
Stage errors. An error a stage returns is taken to be a fault of the host
unless the stage marks it permanent – input it cannot parse, a construct it
refuses. A host fault is a helper that cannot be started, a template, map, key or
trust store that cannot be read, a database statement refused for a reason other
than the connection, a service that does not answer. The worker does not report
it as a failure. It pauses the message at the stage it was claimed at, records
the error as state.last_error and the number of such retries at this stage
as state.temporary_failures, and reports success; the paused-retry sweep
hands the message back after one minute, then two, four and so on up to an hour.
Each retry is logged as a warning under the stage’s name. Once the message has
been retried longer than the stage’s MAX_LIFETIME (default 120 hours, see
pepsi.conf(5); counted from arrival, or for a DSN from when
pepsi-stage-bounce(1) built it) the worker gives up: it reroutes the
message to the stage’s BOUNCE_STAGE, which reports it to the sender, or –
with no BOUNCE_STAGE, or for a mailing-list copy – fails it, and
pepsi-failure-bouncer(1) takes over. A permanent error fails the message at
once. Either way state.last_error, state.failed_stage and
state.failure_class (permanent or retries-exhausted) say what
happened; the DSN itself carries only a generic reason, because the recorded
error can name files, database objects and hosts.
The default is this way round because the two mistakes are not alike: a host fault treated as a defect of the message tells a sender their mail was undeliverable because the server was misconfigured for ten minutes, and cannot be taken back, while a defect treated as a host fault costs retries and a late report. A stage’s own configuration is checked when its worker starts, not per message: a section or secret that does not parse makes the worker refuse to start (status 78, above), so the queue waits for the fix instead of each message failing on its own. A per-address override that makes a stage’s section unparseable is the one configuration error that is permanent, because it is the correspondent’s, not the host’s.
When a stage fused into another stage’s pass fails, the chain up to it is
committed first – the row moved to the failing stage, still running – and
the error is handled there, with that stage’s MAX_LIFETIME and
BOUNCE_STAGE; a retry then resumes at that stage instead of repeating the
stages before it.
Per-address settings. Before a worker runs a stage’s logic it consults the
pepsi.settings table (see pepsi-settings(1)) for the message’s
correspondent — the envelope sender of a state.local_origin message,
otherwise its recipients — and layers any per-address overrides onto the stage’s
[stage-<name>] options. The overrides are fetched with the message row in
the same query (a correlated subquery on the relevant addresses), so loading a
message and its settings is a single round-trip. When the recipients of one inbound message resolve to
different overrides for the stage about to run, the worker splits the message
lazily: in one round-trip (the workqueue_split function) it groups the
recipients by their effective override, keeps the first group on the current row,
and fans each remaining group out as a new pending row at the same stage (no
notification is issued for them: the worker’s own completion brings the
coordinator round to the same claim, as above). Each row then runs the stage under
its own settings. A message whose recipients never diverge is never split.
Stage fusion. Many stages are very fast and need only the message metadata
(the envelope, the extracted From:/Subject: and the state JSON), not
the header block or body. When such a stage advances to a successor that needs no
more data than is already loaded, routing the row back through the database and
the dispatcher is pure overhead. Stage fusion removes it: instead of writing the
row pending at the next stage and waiting for it to be claimed, the worker
runs the successor’s body in the same process, reusing the row it already
loaded. A whole chain of fused stages is one SELECT at the front, the stages’
own work, and a single terminal write, reported to the dispatcher as one
completion.
A successor is fused only when all hold: stage fusion is enabled globally
([pepsi] ALLOW_FUSION, on by default); the successor’s section is marked
FUSION = yes (the default for the fast, body-free stages —
pepsi-stage-if(1), pepsi-stage-discard(1), pepsi-stage-srs(1),
pepsi-stage-check-whitelist(1), pepsi-stage-auto-whitelist(1),
pepsi-stage-block-language(1), pepsi-stage-list(1),
pepsi-stage-edit-settings(1)); the
successor’s PROGRAM is folded into
the same unified pepsi binary (so it can run in-process — fusion is therefore
inert in a per-program multibin build); and the successor needs no message
column the predecessor did not load. Any miss makes the advance a normal
dispatched hop, so fusion never changes a message’s outcome — only whether the
hand-off touches the database. A stage may additionally apply a data-dependent
gate that refuses fusion for some messages: pepsi-stage-edit-settings(1)
loads the full body, yet it does nothing to any message that is not a settings
control message, so it fuses through every non-control message on metadata alone
and declines fusion only for a genuine control message — which is then committed
and dispatched normally so the worker loads its body. Per-address settings are
still applied to each
fused stage (their overrides were already fetched with the row), and a divergent
recipient split still happens where needed — but a hop where both could
apply is not fused: when a stage in the chain has rewritten the message
(envelope sender, From:, Subject:, header block or body) and the row
still carries more than one recipient, the advance is committed normally, so a
subsequent per-address split clones siblings from a row that already holds the
rewrite. A misconfigured stage cycle is bounded
by a fusion-depth cap, after which the advance is committed normally. Fused stages
are still counted individually in pepsi.stage_stats, so per-stage statistics
stay accurate — the worker reports them on the pass’s status line and the
dispatcher folds them in (see Statistics); only the per-stage cost benchmark
turns fusion off (with ALLOW_FUSION = no) so it can time each stage as its own
dispatched worker.
Statistics. The cumulative counters in pepsi.stage_stats and
pepsi.dispatch_stats are accumulated in memory by the dispatcher and
written in a single transaction: every STATS_INTERVAL, whenever the pipeline
goes idle (so a burst’s numbers are visible as soon as it ends rather than up to
an interval later — rate-limited to at most one such flush a second), and on
shutdown. An idle dispatcher performs no database work at all, because a flush
with no accumulated deltas writes nothing.
Message bodies. The same STATS_INTERVAL tick sweeps
pepsi.workqueue_body for bodies no queued row references any more
(pepsi.workqueue_body_gc(), at most 10 000 a sweep). A body is shared by every
row split off one message and is normally deleted by a trigger the moment the
last of them goes; the sweep collects the rare body two concurrent deletes each
left for the other. Orphans cost only disk until then, never a message.
The dispatcher is deliberately the only writer of these tables, which is why
a fused hop is reported to it on the worker’s status line rather than written by
the worker: stage_stats holds one row per stage, so a counter written per
message would make every worker of a stage serialise on that one row. Carrying
the figures on a line the worker already writes removes the write entirely. The
trade is the usual one for statistics: deltas not yet flushed are lost if the
dispatcher dies, which is why these counters are documented as best-effort.
stages_executed counts worker passes, so a fused chain counts once however
many stages it covers.
pepsi-dispatch does not itself process messages — it only runs the stage workers, which record their own outcome on the row.
85.1.2.1.4. Commands¶
- serve
Run the dispatcher until interrupted. Requires that the schema has been installed with pepsi-setup(1).
85.1.2.1.5. Global Options¶
These global options precede the subcommand (a trailing flag is rejected).
- -c FILE, –config FILE
Read the configuration from FILE instead of searching the default locations. Set CONFIG_FILE in
[pepsi-dispatch]to the same path so spawned stage programs inherit it (see pepsi.conf(5)).- -L LOGLEVEL, –log LOGLEVEL
Set the logging verbosity (default
info).- -v, –verbose
Show log messages from all sources.
- -h, –help; -V, –version
Print a usage summary / the version and exit.
85.1.2.1.6. Signals¶
- SIGINT, SIGTERM
Initiate shutdown: stop claiming new work, stop every worker (killing its child), reset any in-flight or claimed-but-unassigned row back to
pending, and exit.
85.1.2.1.7. Exit Status¶
- 0
Clean shutdown.
- 1
An error occurred (for example a malformed configuration file or a failed database connection). The reason is written to the log.
- 78
EX_CONFIG: a configuration reload produced a stage graph that does not parse, so the dispatcher shut down rather than route by a pipeline its stage workers no longer share. Fix the[stage-*]sections (in the file or in thepepsi.config_overrideoverlay) and start the unit again. Also returned at start-up when the database schema is not the one this release was built with (older, newer or built from other files); the message says which, andpepsi-setup schemaupgrades an older one (see pepsi-setup(1)).
85.1.2.1.8. Examples¶
Run the dispatcher:
pepsi-dispatch -c /etc/pepsi/pepsi.conf serve
85.1.2.1.9. See Also¶
pepsi-config(1), pepsi.conf(5), pepsi-ingress(1), pepsi-stage-srs(1), pepsi-stage-bounce(1), pepsi-stage-dkim-sign(1), pepsi-stage-relay-to-smarthost(1), pepsi-stage-relay-to-internet(1), pepsi-settings(1), pepsi-failure-bouncer(1), pepsi-queue(1), pepsi-setup(1)
85.1.2.1.10. Bugs¶
Report bugs to the Pepsi issue tracker.