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 the pepsi.config_override overlay) 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, and pepsi-setup schema upgrades 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.