Diffusion du chat et livraison d'événements
Responsabilités
server-chat.ts livre les événements émis par la boucle principale de l'agent (agent event) aux clients. Il ne produit pas d'événements, il fait du dispatch par distribution: il reçoit AgentEventPayload (delta, final, tool, lifecycle, etc.), et selon le type d'événement, l'appartenance à une session (session), et les abonnements du client, il décide de broadcaster (broadcast) à toutes les connexions visibles, d'envoyer ciblé (broadcastToConnIds) à un ensemble de connIds, ou de livrer inter-nœud (nodeSendToSession) à un processus d'une autre machine qui a souscrit à cette session.
Cette couche maintient trois types d'abonnements: toolEventRecipients (abonnés aux événements d'outils au niveau runId, enregistrés par le client avant de lancer un run), sessionEventSubscribers (niveau session, attachés par l'UI operator après qu'un run a démarré), sessionMessageSubscribers (abonnés aux messages de chat niveau session, livraison ciblée quand l'UI de contrôle est masquée).
Motivation de conception
Pourquoi ne pas appeler socket.send directement dans la boucle principale de l'agent? Parce qu'une même session peut être abonnée par plusieurs clients, que les clients peuvent être sur plusieurs nœuds de la passerelle, que les événements agent doivent être routés par type de stream (agent / chat / session.tool / sessions.changed), que certains événements doivent être throttle et d'autres droppés si le client est lent. Mettre toutes ces stratégies dans la boucle principale la transformerait en bouillie — la boucle ne doit s'occuper que de la machine à états de l'agent; la stratégie de diffusion est extraite dans une couche dédiée.
Les trois primitives sont injectées par createGatewayNodeSessionRuntime (createGatewayNodeSessionRuntime:971) à l'extérieur, et createAgentEventHandler ne fait que les consommer sans se soucier de l'implémentation — le module chat peut être testé sans WebSocket réel.
Fichiers clés
AgentEventHandlerOptions:255-293— Trois primitives de diffusion + types des trois registries d'abonnements.signature createAgentEventHandler:299-317— Injecte les primitives de diffusion, retourne l'event dispatcher.resolveNodeSessionDeliveryKeys:495— Une session global est dépliée enagent:DEFAULT:global+globaldeux delivery keys.sendChatPayload:890— Livraison d'événement chat; si l'UI de contrôle est visible on broadcast, sinon on cible.sendAgentPayload:967— Livraison d'événement agent; même logique de branchement.closure dispatcher:1136—(evt: AgentEventPayload) => {...}entrée d'événement.createGatewayBroadcaster:100-209— Implémentation réelle debroadcast/broadcastToConnIds, incluant filtrage de scope et backpressure.createGatewayNodeSessionRuntime:14-53—nodeSendToSessionrouté vianodeSubscriptions.
Flux de données
createAgentEventHandler (createAgentEventHandler:299) reçoit les trois primitives de diffusion et les trois registries d'abonnements, et retourne un event dispatcher (evt: AgentEventPayload) => void. Il maintient en interne throttle, lifecycle retry, spawnedBy cache, etc.
broadcast et broadcastToConnIds viennent de createGatewayBroadcaster (createGatewayBroadcaster:100) qui itère sur clients: Set<GatewayWsClient>; la boucle centrale:
for (const c of params.clients) {
if (targetConnIds && !targetConnIds.has(c.connId)) continue;
if (!hasEventScope(c, event)) continue; // Filtrage scope
const slow = c.socket.bufferedAmount > MAX_BUFFERED_BYTES;
if (slow && opts?.dropIfSlow) {
if (!isTargeted) clientSeq.set(c, (clientSeq.get(c) ?? 0) + 1); // Le ciblé ne consomme pas de seq
continue; // Si on peut dropper, on droppe
}
if (slow) { c.socket.close(1008, "slow consumer"); continue; } // Sinon on ferme la connexion
c.socket.send(frame); // {"type":"event","event":...,"payload":...,"seq":N}
}La backpressure a deux modes: dropIfSlow: true (chat delta haute fréquence, peut être droppé) on droppe mais on met à jour le seq pour éviter un décalage; dropIfSlow: false (lifecycle critiques, résultats d'outils) on appelle socket.close(1008, "slow consumer") pour couper la connexion — on préfère couper que dropper un événement critique. Le seq n'est incrémenté qu'en diffusion non ciblée (isTargeted ? undefined : nextSeq); la diffusion ciblée est un sous-ensemble filtré, qui ne participe pas au réordonnancement côté client.
nodeSendToSession vient de createGatewayNodeSessionRuntime (nodeSendToSession:31) et se résume à nodeSubscriptions.sendToSession(sessionKey, event, payload, nodeSendEvent). Le commentaire du code « session fanout goes through the subscription manager so node reconnects and explicit unsubscribes keep both node->session indexes in sync » explique pourquoi la livraison inter-nœud passe par le subscription manager plutôt que par nodeRegistry.send direct — quand un nœud se reconnecte ou unsubscribe explicite, les index bidirectionnels doivent rester synchronisés, sinon on livrerait à un nœud déconnecté, ou on raterait un nœud nouvellement connecté.
L'entrée d'événement est la closure retournée (dispatcher:1136): elle résout d'abord sessionKey, agentId, clientRunId (beaucoup d'événements ne portent pas ces champs, il faut remonter via chatRunState.registry), puis dispatche par type de stream. sendChatPayload (sendChatPayload:890) dispatche selon la visibilité de l'UI de contrôle: si visible, broadcast("chat", payload) + sendNodeSessionPayloadForAgent pour la livraison inter-nœud; si masquée, on consulte sessionMessageSubscribers.get(deliverySessionKey), et s'il y a des abonnés, on fait broadcastToConnIds ciblé — typiquement l'outil operator en arrière-plan, le panneau principal ne voit pas le run. sendAgentPayload (sendAgentPayload:967) a une logique symétrique, seul le nom d'événement devient "agent".
La livraison inter-nœud a un détail (resolveNodeSessionDeliveryKeys:495): si sessionKey est "global", on déplie en deux delivery keys agent:DEFAULT:global et global. Les événements globaux doivent atteindre à la fois les clients abonnés à « un certain agent global » et ceux abonnés à « global »; mais seulement quand agentId est l'agent par défaut on push les deux, sinon on push seulement l'agent-scoped — pour éviter une double livraison.
Limites et modes d'échec
- Gestion des consommateurs lents:
socket.bufferedAmount > MAX_BUFFERED_BYTESdéclenche la backpressure.dropIfSlow: true(delta streamé,chat.send_timing) skip l'envoi mais met à jour le seq;dropIfSlow: false(lifecycle critique, résultats d'outils) appellesocket.close(1008, "slow consumer")— sans cela le buffer gonflerait et tuerait le processus (backpressure:157). - Retry des erreurs lifecycle:
pendingTerminalLifecycleErrorsMap trace l'état de retry des événements de phase end/error dans le grace period.lifecycleErrorRetryGraceMscontrôle la fenêtre; un nouveau lifecycle generation annule le timer précédent. - Cache spawnedBy:
resolveSpawnedByne fait de requête de stockage que pour les sessionKey subagent/ACP; les non-lineage retournent null sans polluer le cache (spawnedByCache:383). Les chemins delta/flush/final haute fréquence évitent ainsi un hit au stockage de session à chaque fois. - Suppression du restart recovery:
resolveRestartRecoveryLifecycleStatevérifie si l'événement lifecycle est un produit de restart recovery; si oui,suppress: true— pour éviter qu'après redémarrage les anciens événements de run restaurés ne ressurgissent aux clients. - Repli de résolution sessionKey: d'abord
chatLink.sessionKey, puisevt.sessionKey, enfinresolveSessionKeyForRun(evt.runId). Trois niveaux de repli pour garantir qu'un événement trouve toujours sa session. - La livraison inter-nœud ne consomme pas de seq:
isTargeted ? undefined : nextSeq; si la diffusion ciblée incrémentait, le client verrait des sauts de seq et penserait à tort qu'il a perdu des événements. - Synchronisation des abonnements de nœuds:
nodeSendToSessionpasse par le subscription manager plutôt que parnodeRegistry.senddirect, pour que les index bidirectionnels restent synchronisés lors des reconnexions/déconnexions/unsubscribe. Court-circuiter le manager ferait perdre des événements après reconnexion.
Résumé
La couche de diffusion détache la stratégie de livraison des événements agent hors de la boucle principale, l'abstrait en trois primitives injectées par createGatewayNodeSessionRuntime à l'extérieur, et createAgentEventHandler ne fait que consommer. L'implémentation réelle (createGatewayBroadcaster) gère le filtrage de scope, la backpressure, la numérotation seq; la livraison inter-nœud est routée via le subscription manager pour garantir la synchro des index après reconnexion. L'entrée d'événement est la closure dispatcher qui route par type de stream vers chat / agent / session.tool / sessions.changed, chaque stream pouvant à son tour broadcaster complet ou cibler selon la visibilité de l'UI de contrôle.
Comment les requêtes entrent et sont réparties, voir Table de méthodes RPC et répartition des requêtes; comment les primitives de diffusion sont injectées à toute la passerelle, voir Cœur de passerelle; comment l'état d'abonnement de session est lié au cycle de vie d'un run agent, voir Session Manager.
Pour comparer avec la documentation officielle: Gateway docs · README.