85.1.2. pepsi-dispatch

drive messages through the stage pipeline

Section du manuel:

1

85.1.2.1.1. Nom

pepsi-dispatch - revendiquer les messages en attente et exécuter leurs programmes d’étape.

85.1.2.1.2. Synopsis

pepsi-dispatch [GLOBAL-OPTIONS] serve

85.1.2.1.3. Description

pepsi-dispatch est le dispatcher de longue durée qui fait avancer les messages à travers le pipeline d’étapes (voir pepsi.conf(5)). Il revendique les lignes pending de pepsi.workqueue, met chacune à running, recherche l’étape du message dans sa section [stage-<stage>], et transmet l’identifiant de message à un processus worker de cette étape. Chaque étape s’exécute comme un pool de workers persistants démarrés par PROGRAM worker (PROGRAM -c FILE worker lorsque CONFIG_FILE est défini — le drapeau global précède la sous-commande, comme partout dans Pepsi ; le programme est trouvé sur le PATH sauf si PROGRAM est un chemin absolu) ; un worker se connecte une fois à la base de données puis lit les identifiants de message sur son entrée standard et écrit une ligne de statut par message sur sa sortie standard (0 en cas de succès), dans l’ordre d’entrée. Un statut peut être suivi d’un espace et d’un rapport JSON des étapes qui ont été fusionnées dans cette passe, [["<stage>",<microseconds>], …], ce qui permet aux statistiques par étape de rester entières sans que les workers les écrivent (voir Statistiques ci-dessous). Seul un statut illisible est fatal — il désynchronise le protocole positionnel, si bien que le dispatcher démantèle le worker ; un rapport illisible est abandonné avec un avertissement, de sorte qu’un worker et un dispatcher issus de builds différents perdent des statistiques le temps d’un redémarrage plutôt que des messages. N’exécutez qu”exactement un dispatcher par système ; c’est le seul processus qui démarre les programmes d’étape (bien que de nombreux workers s’exécutent ensuite en concurrence contre la base de données partagée). Le nouveau travail est récupéré depuis un LISTEN sur le canal workqueue, avec un battement de cœur périodique de filet de sécurité (POLL_INTERVAL).

Les notifications sur ce canal sont émises explicitement, et seulement par des écrivains extérieurs à la boucle propre du dispatcher : pepsi-ingress(1) (une fois par rafale de messages acceptés, non une fois par message), les fonctions SQL workqueue_inject et workqueue_resume, pepsi-keydisc(1) libérant du courrier garé, pepsi-failure-bouncer(1), les commandes de réparation de pepsi-queue(1) et les commandes de modération et de condensé de pepsi-list(1). Une étape qui fait avancer un message ne notifie délibérément pas : le worker rapporte le message terminé sur sa sortie standard, et le dispatcher y répond par la même revendication complète qu’une notification aurait déclenchée ; la notification ne porterait donc aucune information. Un déclencheur de table qui notifierait sur chaque ligne devenant pending serait pire qu’inutile : PostgreSQL détient un verrou à l’échelle de la base de données depuis le moment où une transaction met une notification en file jusqu’à sa validation, de sorte que chaque saut d’étape sérialiserait les validations de toute la base de données — voir le chapitre Performance du manuel.

Revendication par lots. À chaque réveil (une notification, un worker libéré ou mort, le battement de cœur de sondage, une (re)connexion du listener), le dispatcher draine d’abord chaque événement interne en attente afin que ses capacités par étape soient à jour, puis exécute un UPDATE … RETURNING qui revendique, pour chaque étape disposant de capacité libre à la fois, jusqu’à la propre capacité restante de cette étape (pending``→``running ; quelles lignes, c’est l’objet du paragraphe suivant — ordonnancement équitable désactivé, il s’agit d’une sous-requête LATERAL par étape dont l”ORDER BY workqueue_id LIMIT cap arrête le parcours de l’arriéré de chaque étape à son plafond, piloté par l’index partiel sur (stage, workqueue_id) WHERE status='pending'). La notification porte une charge utile vide — une notification n’est qu’un réveil, et une rafale de notifications se réduit à une seule revendication par lots (PostgreSQL fusionne aussi les notifications identiques mises en file).

Ordonnancement équitable entre expéditeurs. Le traitement du plus ancien d’abord permettrait à la rafale d’un seul expéditeur de retarder chaque message mis en file derrière elle ; aussi, par défaut (FAIR_SCHEDULING, voir pepsi.conf(5)), la revendication choisit quelles lignes remplissent la capacité libre d’une étape comme l’ordonnanceur de Linux (CFS/EEVDF) choisit quel processus obtient un CPU : une étape est le CPU et un expéditeur est le processus. L’expéditeur d’une ligne est la colonne générée workqueue.sender_key :

  • user:<login> — le compte authentifié, pour une submission (SASL ou les identifiants du pair de la socket locale), quel que soit l’expéditeur d’enveloppe qu’il a choisi ;

  • ip:<address> — sinon le pair SMTP qui se connecte, un pair IPv6 par son /64 (un seul hôte reçoit couramment un /64 entier). Délibérément pas le domaine d’enveloppe, qu’un expéditeur distant choisit librement pour chaque message ;

  • from:<domain> — pour un message que Pepsi a lui-même créé (sans origine SMTP : un avis, un rapport), le domaine de l’expéditeur d’enveloppe ;

  • null:<domain> — pour un message à expéditeur nul (une DSN, une réponse automatique), le domaine de son premier destinataire.

Comme la clé est générée à partir de la ligne, chaque écrivain — ingress, un message injecté, une scission, un clone — l’obtient sans savoir qu’elle existe. Le dispatcher conserve, en mémoire et par étape, l”avance de chaque expéditeur : le temps de worker que ses messages y ont consommé au-delà de l’expéditeur en attente le moins servi. Un message revendiqué est immédiatement débité du temps moyen par message de l’étape, corrigé au temps mesuré lorsque son worker rend compte, et remboursé s’il est remis en file sans avoir été exécuté ; un message en timeout est débité de la totalité de MAX_RUNTIME. La revendication classe ensuite la k-ième plus ancienne ligne en attente d’un expéditeur d’avance L selon L + k × average, la plus petite d’abord, de sorte qu’un expéditeur ayant mille lignes en attente alterne avec un expéditeur qui n’en a qu’une au lieu de passer en premier — et, comme le débit est du temps plutôt qu’un nombre, un expéditeur dont les messages sont coûteux à une étape obtient d’autant moins de créneaux. Un expéditeur nouveau pour l’étape, ou dont l’avance s’est dissipée, part à égalité avec le moins servi, de sorte que rester silencieux n’achète aucune priorité ultérieure. Les avances sont divisées par deux toutes les FAIR_HALF_LIFE (60 s par défaut) et une avance valant moins d’un message est oubliée ; un redémarrage les oublie toutes, ce qui est sans conséquence. Pour garantir qu’aucune ligne n’attende indéfiniment — quel que soit le nombre de nouveaux expéditeurs à égalité qui continuent d’arriver — FAIR_FIFO_PERCENT (10 par défaut) des créneaux revendiqués de chaque étape vont à ses plus anciennes lignes en attente quel que soit l’expéditeur, de sorte qu’un message attend au plus que les lignes mises en file avant lui à cette étape se soient écoulées à ce débit réservé.

La revendication reste une seule instruction pour toutes les étapes à la fois, l’écriture pending``→``running de l’unique revendicateur, bornée par étape par sa capacité : la fonction SQL workqueue_claim_fair reçoit les avances sous forme de tableaux, énumère les expéditeurs en attente à chaque étape par un parcours d’index lâche (loose index scan) sur l’index partiel (stage, sender_key, workqueue_id) WHERE status='pending' (une sonde par expéditeur, au plus 64 expéditeurs par étape et par revendication — au-delà, la fenêtre tourne d’une revendication à la suivante), et lit au plus la capacité de l’étape parmi les plus anciennes lignes de chaque expéditeur, de sorte que son coût ne croît pas avec la profondeur de l’arriéré d’un expéditeur donné. FAIR_SCHEDULING = no rétablit la simple revendication du plus ancien d’abord.

Pools de workers pipelinés et élastiques. Le pool de chaque étape croît à la demande jusqu’à son PARALLELISM (par défaut 4) et rétrécit à l’inactivité : aucun worker ne s’exécute tant que l’étape ne voit pas de trafic ; lorsqu’une étape a du travail et de la capacité libre, le dispatcher remplit le créneau libre d’un worker existant ou démarre un nouveau worker ; le travail est tassé sur le moins de workers possible, de sorte qu’un surplus se refroidit et est arrêté après WORKER_IDLE_TIMEOUT (par défaut 5 s). Sous une charge légère et stable, une étape oscille donc autour d’un seul worker. Un worker reçoit jusqu’à QUEUE_LIMIT (par défaut 4) identifiants de message à la fois — écrits sur son entrée standard en un seul tampon — plutôt qu’un par un ; il les traite tout de même strictement en série, mais l’identifiant suivant est déjà en file, de sorte que terminer un message le libère sans aller-retour préalable avec le coordinateur. La capacité totale en vol d’une étape est ainsi QUEUE_LIMIT × PARALLELISM. Comme le nombre de workers (et donc le nombre de connexions de base de données) est inchangé, relever QUEUE_LIMIT échange un peu de latence de tête de file contre du débit sans dépenser plus de connexions. Après MAX_MESSAGES (par défaut 1000), un worker recycle son enfant — une fois son pipeline drainé — avec un nouveau processus pour borner la croissance de la mémoire.

Étapes sans corps par lots. Un worker d’une étape qui ne charge ni le bloc d’en-têtes ni le corps (les étapes à métadonnées seules — srs, if, auto-whitelist, discard, aliases …) traite ses identifiants pipelinés en un lot plutôt qu’un à la fois : il prend avidement chaque identifiant déjà mis en tampon sur son entrée standard (jusqu’à QUEUE_LIMIT, s’arrêtant à l’instant où une lecture bloquerait, de sorte qu’un identifiant solitaire est tout de même traité immédiatement), les charge tous en un seul SELECT … WHERE workqueue_id = ANY(...), exécute chaque corps d’étape, puis écrit les avancées, les échecs et les remises en file en un UPDATE … FROM jsonb_to_recordset(...) sur tableau par forme d’issue (typiquement une seule). Il écrit tout de même une ligne de statut par identifiant, dans l’ordre, de sorte que le protocole worker est inchangé. Lorsque la file est pleine, cela réduit les allers-retours de base de données d’une étape sans corps de deux par message à environ deux par QUEUE_LIMIT messages. Les étapes qui chargent les en-têtes ou le corps conservent le chemin un identifiant à la fois (leurs grandes colonnes ne se mettent pas en tableau à bas coût).

Avancement. Une étape fait avancer un message en mettant sa colonne stage à l’étape suivante et la ligne de nouveau à pending en une seule mise à jour ; le dispatcher ne réécrit jamais une ligne en cas de succès. La revendication (pending``→``running) est la seule écriture du dispatcher dans le status d’une ligne pour la progression, et il est le seul à l’effectuer. Lorsqu’un worker rapporte un succès (statut 0), le dispatcher rebalaie à la recherche de travail — c’est précisément pourquoi un avancement n’émet aucune notification propre — de sorte que le pool de l’étape suivante revendique la ligne désormais pending au tour d’ordonnancement suivant. Le message est terminé lorsque le worker supprime la ligne, en pause lorsqu’il la laisse paused, ou terminal lorsqu’elle est failed/timeout. Une étape écrit son issue en un seul aller-retour — un UPDATE/DELETE, ou une fonction stockée qui déploie le travail côté serveur (un découpage de destinataires, ou un finish/pause qui engendre aussi des DSN de rebond, de délai ou de succès à partir de tableaux de clones) — jamais une transaction client multi-instructions. La connexion d’un worker s’exécute en outre avec le synchronous_commit de PostgreSQL désactivé ([pepsi-postgres] WORKER_SYNCHRONOUS_COMMIT, par défaut no), ce qui est ce qui permet aux workers concurrents de valider en groupe. Cela n’abandonne rien d’autre que la durabilité des dernières centaines de millisecondes, et perdre un avancement d’étape de cette manière est indiscernable d’un worker tué juste avant lui : la ligne est trouvée running à l’étape précédente, remise à pending au démarrage suivant, et l’étape s’exécute de nouveau. pepsi-ingress(1) conserve délibérément la valeur par défaut, car perdre une admission à laquelle il a déjà répondu 250 serait un message perdu plutôt qu’une étape répétée.

Échecs et réessais. Un worker d’étape s’annonce par une ligne ready une fois qu’il a ouvert son pool de base de données et vérifié sa configuration, puis répond une ligne de statut par message. Un worker qui rapporte un statut non nul laisse ce message failed avec la raison dans state et reste en vie pour l’identifiant suivant ; mais un worker ne le rapporte que pour une erreur que son étape a marquée permanente (voir Erreurs d’étape ci-dessous) ou après avoir renoncé aux réessais, et jamais pour les deux statuts de base de données. Un worker dont l’enfant ferme sa sortie (un plantage) après s’être déclaré prêt est démantelé, et son message en tête de file reçoit une pénalité ; il en va de même du message en tête de file d’un enfant qui ne répond pas dans le délai MAX_RUNTIME (temps réel), lequel est tué. Un message pénalisé est mis en pause et réessayé au bout d’une minute (deux après la deuxième pénalité), avec la raison dans state.last_error et le compte dans state.strikes ; seule sa troisième pénalité est imputée au message lui-même, et le met à failed (après un plantage) ou à timeout (après un blocage), avec state.failure_class crashed ou timed-out. Un worker meurt aussi pour des raisons qui lui sont propres – le tueur OOM, un helper en cours de mise à niveau, un résolveur qui se bloque – et une telle mort ne doit pas coûter le message qu’il se trouvait tenir, tandis qu’un message qui tue chaque worker auquel il est confié ne doit pas être réessayé indéfiniment. Dans ces deux cas, les autres messages déjà pipelinés vers ce worker n’ont jamais été exécutés : ils sont donc réinitialisés à pending et de nouveau dispatchés (un unique message bloqué ou empoisonné n’abandonne donc pas ses frères en vol). Un enfant qui meurt avant d’être prêt n’impute rien : il a échoué à son propre démarrage – un PostgreSQL injoignable ou à court de connexions, un secret illisible – de sorte que chaque message qui lui a été confié retourne à pending et que l’étape est mise en attente pendant un court refroidissement avant qu’un autre worker ne soit démarré, au lieu d’être relancée aussitôt. Un enfant qui se termine avant d’avoir répondu avec le statut 78 (EX_CONFIG) a refusé de s’exécuter – parce qu’il a rencontré un schéma de base de données d’une autre version, dans la fenêtre d’une mise à niveau de paquet entre l’arrivée des nouveaux binaires et la mise à niveau du schéma, ou parce que sa configuration ne s’analyse pas – et est traité de la même manière, avec une erreur nommant la cause ; l’étape est réessayée après le refroidissement jusqu’à ce que le schéma corresponde ou que la configuration soit corrigée. L’horloge MAX_RUNTIME est par message et démarre lorsqu’un message atteint la tête de la file du worker, de sorte qu’un message attendant derrière un frère lent n’est pas facturé pour l’attente. Lorsqu’une étape laisse un message paused, elle note, dans timeout, quand le message devra être réessayé ; le dispatcher dort jusqu’à la première de ces échéances, rebascule les lignes dues à pending et rebalaie à la recherche de travail. Au démarrage, il réinitialise à pending toute ligne running résiduelle (orpheline d’un dispatcher précédent). Comme il y a exactement un dispatcher et que son coordinateur revendique en série — étant seul à effectuer le passage pending``→``running — la revendication par lots par étape n’a besoin d’aucun FOR UPDATE SKIP LOCKED. « Exactement un » est imposé plutôt que supposé : au démarrage, le dispatcher prend sur sa base de données un verrou consultatif PostgreSQL de portée session et, si un autre dispatcher le détient encore au bout d’environ dix secondes, refuse de s’exécuter. Sans cela, le balayage de démarrage d’un second dispatcher annulerait aussitôt la revendication des lignes en vol du premier et tous deux les dispatcheraient — remise en double, rebonds en double, règlement de paiement en double. Le verrou est libéré par PostgreSQL à la fin de la session, si bien qu’un dispatcher tué ne laisse rien à nettoyer. Comme au plus un écrivain touche jamais une ligne donnée, la connexion de base de données partagée fonctionne en READ COMMITTED plutôt qu’en SERIALIZABLE (les opérations de file sont tout de même acheminées via un helper de réessai comme assurance bon marché).

Interrupteur de télémétrie. Le dispatcher écoute aussi telemetry_changed, que quiconque réécrit [pepsi] SHARE_TELEMETRY dans le fichier de configuration envoie ensuite, et retire chaque worker vivant : un worker installe son puits (sink) de télémétrie des fonctionnalités une seule fois, à partir du fichier, de sorte que son remplaçant commence (ou cesse) de rendre compte en conséquence. La table des étapes n’est pas modifiée. Voir pepsi-telemetry-client(1).

Rechargement de la configuration. Chaque section [stage-*] est un réglage à chaud : à la notification config_changed, le dispatcher reconstruit sa propre table d’étapes à partir de la surcouche de configuration et retire tout worker vivant, de sorte qu’une étape peut être ajoutée, modifiée ou supprimée pendant que le pipeline tourne et que le message suivant est à la fois acheminé et exécuté avec les nouvelles valeurs. Un graphe d’étapes qui ne s’analyse pas met fin au processus avec le code 78 après un arrêt ordonné (workers arrêtés, leurs lignes de nouveau à pending), et l’erreur d’analyse est journalisée. Conserver la table précédente serait pire : les workers sont retirés par le même événement et leurs remplaçants lisent bien la nouvelle surcouche, si bien que le dispatcher acheminerait selon un pipeline pendant que les étapes en exécuteraient un autre. Le pepsi-dispatch.service livré utilise Restart=always sans limite de démarrages : sans le dispatcher, le courrier accepté ne fait que s’accumuler dans la file, de sorte que l’unité continue de réessayer, en espaçant les tentatives de 2 s jusqu’à une par minute (RestartSteps=, RestartMaxDelaySec=), et prend en compte d’elle-même une configuration corrigée. Une surcouche qui ne peut pas être lue (une erreur de base de données) n’est pas une nouvelle configuration : le dispatcher conserve son graphe d’étapes actuel et retente le rechargement, plutôt que d’abandonner chaque étape définie uniquement dans la base de données. La section [pepsi-dispatch] elle-même n’est pas rechargée — le pool de connexions doit exister avant que la surcouche puisse être lue — la modifier exige donc un redémarrage.

Contre-pression en cas d’engorgement de la base de données. Un worker qui ne peut atteindre la base de données parce que sa limite de connexions est épuisée (too_many_connections de PostgreSQL, ou son pool expire en acquérant une connexion) rapporte un statut distinct (EX_TEMPFAIL, 75) plutôt que le code d’échec générique — c’est une pression d’infrastructure, pas un défaut du message. Le dispatcher ne fait alors pas échouer le message : il le remet en file à pending pour un réessai ultérieur, et, pour délester des connexions, il réduit temporairement le parallélisme de l’étape qui exécute actuellement le plus de processus workers (le plus gros contributeur à la pression, qui n’est pas forcément l’étape ayant rapporté l’erreur) — divisant par deux son plafond, plancher 1, pendant cinq minutes. Des rapports répétés abaissent encore cette étape et rafraîchissent la fenêtre ; les workers existants se vident par la récupération à l’inactivité et l’étape remonte une fois la fenêtre passée. Le remède à une pression soutenue est de relever le max_connections de PostgreSQL ou d’abaisser le PARALLELISM des étapes afin que la somme du pool de chaque composant tienne sur le serveur (voir la note sur le budget de connexions [pepsi-postgres] dans pepsi.conf(5)).

Connexion à la base de données perdue. Un worker dont la connexion à la base de données se rompt ou ne peut être établie pendant qu’il traite un message (PostgreSQL qui redémarre ou s’arrête, une interruption réseau : une erreur d’E/S ou TLS, la classe SQLSTATE 08, 57P01–57P03) rapporte le statut 74 (EX_IOERR). Ce n’est pas non plus un défaut du message, si bien que le dispatcher ne le fait pas échouer : il le met à paused pendant 30 secondes, après quoi le balayage ordinaire des réessais en pause le rend à la même étape. Cela se répète tant que la base de données reste injoignable. Rien de ce qu’a fait l’étape n’a été validé, donc le réessai exécute l’étape depuis le début ; un effet de bord extérieur à la base de données dont la validation a été perdue (un relais qui a remis le message juste avant la rupture de la connexion) peut se produire une seconde fois. Chaque report est journalisé comme avertissement sous pepsi-dispatch et n’est compté nulle part ailleurs : un compteur devrait être écrit dans la base de données qui vient justement de disparaître. Si le dispatcher ne peut pas non plus écrire la pause, il réessaie pendant environ une demi-minute et laisse sinon le message running, que son prochain démarrage remet à pending.

Erreurs d’étape. Une erreur qu’une étape renvoie est considérée comme une défaillance de l’hôte, à moins que l’étape ne la marque permanente – une entrée qu’elle ne sait pas analyser, une construction qu’elle refuse. Une défaillance de l’hôte, c’est un helper qui ne peut pas être démarré, un modèle, une table, une clé ou un magasin de confiance qui ne peut pas être lu, une instruction de base de données refusée pour une autre raison que la connexion, un service qui ne répond pas. Le worker ne la signale pas comme un échec. Il met le message en pause à l’étape à laquelle il a été revendiqué, enregistre l’erreur dans state.last_error et le nombre de ces réessais à cette étape dans state.temporary_failures, et rapporte un succès ; le balayage des réessais en pause rend le message au bout d’une minute, puis de deux, quatre et ainsi de suite jusqu’à une heure. Chaque réessai est journalisé comme un avertissement sous le nom de l’étape. Une fois que le message a été réessayé plus longtemps que le MAX_LIFETIME de l’étape (par défaut 120 heures, voir pepsi.conf(5) ; compté depuis l’arrivée ou, pour un DSN, depuis sa construction par pepsi-stage-bounce(1)), le worker abandonne : il réachemine le message vers le BOUNCE_STAGE de l’étape, qui le signale à l’expéditeur, ou – sans BOUNCE_STAGE, ou pour une copie de liste de diffusion – le fait échouer, et pepsi-failure-bouncer(1) prend le relais. Une erreur permanente fait échouer le message immédiatement. Dans les deux cas, state.last_error, state.failed_stage et state.failure_class (permanent ou retries-exhausted) disent ce qui s’est passé ; le DSN lui-même ne porte qu’une raison générique, car l’erreur enregistrée peut nommer des fichiers, des objets de base de données et des hôtes.

La valeur par défaut est dans ce sens parce que les deux erreurs ne se valent pas : une défaillance de l’hôte traitée comme un défaut du message annonce à un expéditeur que son courrier n’a pu être remis parce que le serveur a été mal configuré pendant dix minutes, et cela ne peut être repris, tandis qu’un défaut traité comme une défaillance de l’hôte coûte des réessais et un rapport tardif. La configuration propre d’une étape est vérifiée au démarrage de son worker, non par message : une section ou un secret qui ne s’analyse pas fait refuser au worker de démarrer (statut 78, ci-dessus), de sorte que la file attend la correction au lieu que chaque message échoue de son côté. Une surcharge par adresse qui rend la section d’une étape inanalysable est la seule erreur de configuration permanente, car elle relève du correspondant et non de l’hôte.

Lorsqu’une étape fusionnée dans la passe d’une autre étape échoue, la chaîne jusqu’à elle est d’abord validée – la ligne déplacée vers l’étape en échec, toujours running – et l’erreur y est traitée, avec le MAX_LIFETIME et le BOUNCE_STAGE de cette étape ; un réessai reprend alors à cette étape au lieu de répéter les étapes qui la précèdent.

Paramètres par adresse. Avant qu’un worker n’exécute la logique d’une étape, il consulte la table pepsi.settings (voir pepsi-settings(1)) pour le correspondant du message — l’expéditeur d’enveloppe d’un message state.local_origin, sinon ses destinataires — et superpose toute redéfinition par adresse aux options [stage-<name>] de l’étape. Les redéfinitions sont récupérées avec la ligne du message dans la même requête (une sous-requête corrélée sur les adresses pertinentes), de sorte que charger un message et ses paramètres est un seul aller-retour. Lorsque les destinataires d’un message entrant se résolvent en redéfinitions différentes pour l’étape sur le point de s’exécuter, le worker scinde le message paresseusement : en un aller-retour (la fonction workqueue_split), il groupe les destinataires par leur redéfinition effective, garde le premier groupe sur la ligne courante, et déploie chaque groupe restant comme une nouvelle ligne pending à la même étape (aucune notification n’est émise pour elles : l’achèvement propre du worker ramène le coordinateur à la même revendication, comme ci-dessus). Chaque ligne exécute ensuite l’étape sous ses propres paramètres. Un message dont les destinataires ne divergent jamais n’est jamais scindé.

Fusion d’étapes. De nombreuses étapes sont très rapides et n’ont besoin que des métadonnées du message (l’enveloppe, les From:/Subject: extraits et le JSON state), pas du bloc d’en-têtes ni du corps. Lorsqu’une telle étape avance vers un successeur qui n’a besoin d’aucune donnée de plus que ce qui est déjà chargé, faire repasser la ligne par la base de données et le dispatcher est un pur surcoût. La fusion d’étapes le supprime : au lieu d’écrire la ligne pending à l’étape suivante et d’attendre qu’elle soit revendiquée, le worker exécute le corps du successeur dans le même processus, en réutilisant la ligne qu’il a déjà chargée. Toute une chaîne d’étapes fusionnées est un SELECT au début, le travail propre des étapes, et une seule écriture terminale, rapportée au dispatcher comme un achèvement.

Un successeur n’est fusionné que lorsque toutes ces conditions tiennent : la fusion d’étapes est activée globalement ([pepsi] ALLOW_FUSION, activée par défaut) ; la section du successeur est marquée FUSION = yes (la valeur par défaut pour les étapes rapides et sans corps — 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)) ; le PROGRAM du successeur est plié dans le même binaire pepsi unifié (de sorte qu’il peut s’exécuter dans le processus — la fusion est donc inerte dans un build multibin par programme) ; et le successeur n’a besoin d’aucune colonne de message que le prédécesseur n’ait chargée. Tout manquement fait de l’avancement un saut dispatché normal : la fusion ne change donc jamais l’issue d’un message — seulement si le passage de main touche la base de données. Une étape peut en outre appliquer une barrière dépendante des données qui refuse la fusion pour certains messages : pepsi-stage-edit-settings(1) charge le corps complet, pourtant elle ne fait rien à un message qui n’est pas un message de contrôle des paramètres ; elle fusionne donc à travers chaque message non-contrôle sur les seules métadonnées et ne décline la fusion que pour un véritable message de contrôle — lequel est alors écrit et dispatché normalement afin que le worker charge son corps. Les paramètres par adresse sont tout de même appliqués à chaque étape fusionnée (leurs redéfinitions ont déjà été récupérées avec la ligne), et un découpage de destinataires divergents a tout de même lieu là où il le faut — mais un saut où les deux pourraient s’appliquer n’est pas fusionné : lorsqu’une étape de la chaîne a réécrit le message (expéditeur d’enveloppe, From:, Subject:, bloc d’en-têtes ou corps) et que la ligne porte encore plus d’un destinataire, l’avancement est écrit normalement, de sorte qu’un découpage par adresse ultérieur clone les frères depuis une ligne qui porte déjà la réécriture. Un cycle d’étapes mal configuré est borné par un plafond de profondeur de fusion, après quoi l’avancement est écrit normalement. Les étapes fusionnées sont tout de même comptées individuellement dans pepsi.stage_stats, de sorte que les statistiques par étape restent exactes — le worker les rapporte sur la ligne de statut de la passe et le dispatcher les y intègre (voir Statistiques) ; seul le benchmark de coût par étape désactive la fusion (avec ALLOW_FUSION = no) afin de pouvoir chronométrer chaque étape comme son propre worker dispatché.

Statistiques. Les compteurs cumulés de pepsi.stage_stats et pepsi.dispatch_stats sont accumulés en mémoire par le dispatcher et écrits en une seule transaction : à chaque STATS_INTERVAL, chaque fois que le pipeline devient inactif (de sorte que les chiffres d’une rafale sont visibles dès qu’elle se termine plutôt qu’un intervalle plus tard — avec au plus un tel vidage par seconde), et à l’arrêt. Un dispatcher inactif n’effectue aucun travail de base de données, car un vidage sans delta accumulé n’écrit rien.

Corps de message. Le même tic STATS_INTERVAL balaie pepsi.workqueue_body à la recherche de corps qu’aucune ligne en file ne référence plus (pepsi.workqueue_body_gc(), au plus 10 000 par balayage). Un corps est partagé par toutes les lignes scindées d’un même message et est normalement supprimé par un trigger dès que la dernière d’entre elles disparaît ; le balayage recueille le rare corps que deux suppressions concurrentes ont chacune laissé à l’autre. Jusque-là, les orphelins ne coûtent que de l’espace disque, jamais un message.

Le dispatcher est délibérément le seul écrivain de ces tables, ce qui explique qu’un saut fusionné lui soit rapporté sur la ligne de statut du worker plutôt qu’écrit par le worker : stage_stats contient une ligne par étape, si bien qu’un compteur écrit par message ferait sérialiser tous les workers d’une étape sur cette unique ligne. Porter les chiffres sur une ligne que le worker écrit déjà supprime entièrement l’écriture. Le compromis est l’ordinaire en matière de statistiques : les deltas non encore vidés sont perdus si le dispatcher meurt, ce pourquoi ces compteurs sont documentés comme fournis au mieux. stages_executed compte les passes de worker, de sorte qu’une chaîne fusionnée compte pour une, quel que soit le nombre d’étapes qu’elle couvre.

pepsi-dispatch ne traite pas lui-même les messages — il ne fait qu’exécuter les workers d’étape, qui notent leur propre issue sur la ligne.

85.1.2.1.4. Commandes

serve

Exécuter le dispatcher jusqu’à interruption. Nécessite que le schéma ait été installé avec pepsi-setup(1).

85.1.2.1.5. Options globales

Ces options globales précèdent la sous-commande (un indicateur placé à la fin est rejeté).

-c FILE, –config FILE

Lit la configuration depuis FILE au lieu de parcourir les emplacements par défaut. Réglez CONFIG_FILE dans [pepsi-dispatch] sur le même chemin afin que les programmes d’étape lancés en héritent (voir pepsi.conf(5)).

-L LOGLEVEL, –log LOGLEVEL

Règle la verbosité de journalisation (par défaut info).

-v, –verbose

Affiche les messages de journal de toutes les sources.

-h, –help ; -V, –version

Affiche un résumé d’utilisation / la version et quitte.

85.1.2.1.6. Signaux

SIGINT, SIGTERM

Amorcer l’arrêt : cesser de revendiquer du nouveau travail, arrêter chaque worker (en tuant son enfant), réinitialiser à pending toute ligne en vol ou revendiquée mais non assignée, et quitter.

85.1.2.1.7. Code de sortie

0

Arrêt propre.

1

Une erreur s’est produite (par exemple un fichier de configuration malformé ou une connexion de base de données échouée). La raison est écrite dans le journal.

78

EX_CONFIG : un rechargement de configuration a produit un graphe d’étapes qui ne s’analyse pas, le dispatcher s’est donc arrêté plutôt que d’acheminer selon un pipeline que ses workers d’étape ne partagent plus. Corrigez les sections [stage-*] (dans le fichier ou dans la surcouche pepsi.config_override) et redémarrez l’unité. Également renvoyé au démarrage lorsque le schéma de la base de données n’est pas celui avec lequel cette version a été compilée (plus ancien, plus récent ou construit à partir d’autres fichiers) ; le message précise lequel, et pepsi-setup schema met à niveau un schéma plus ancien (voir pepsi-setup(1)).

85.1.2.1.8. Exemples

Exécuter le dispatcher

pepsi-dispatch -c /etc/pepsi/pepsi.conf serve

85.1.2.1.9. Voir aussi

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. Bogues

Signalez les bogues au gestionnaire de tickets de Pepsi.