Skip to content

Broadcast de chat y entrega de eventos

源码版本v2026.6.11

Responsabilidad

server-chat.ts entrega de vuelta al cliente los eventos (agent event) que emite el bucle principal del agent. No produce eventos, solo hace despacho por tipo: recibe AgentEventPayload (delta, final, tool, lifecycle, etc.) y, según el tipo de evento, la pertenencia a sesión (session) y la suscripción del cliente, decide si hace broadcast a todas las conexiones visibles, broadcast dirigido (broadcastToConnIds) a un conjunto de connIds, o entrega entre nodos (nodeSendToSession) a un proceso de otro nodo que haya suscrito a esa sesión.

Esta capa mantiene tres tipos de estado de suscripción: toolEventRecipients (suscriptores de eventos de herramienta a nivel de runId, registrados por el cliente antes de iniciar el run), sessionEventSubscribers (nivel sesión, añadidos por la UI del operator cuando el run ya está en curso) y sessionMessageSubscribers (suscripción de mensajes de chat a nivel sesión, entrega dirigida cuando la UI de control está oculta).

Motivación de diseño

¿Por qué no hacer directamente socket.send desde el bucle principal del agent? Una misma sesión puede estar suscrita por múltiples clientes; los clientes pueden estar repartidos entre varios nodos del gateway; los eventos del agent deben clasificarse por tipo de stream (agent / chat / session.tool / sessions.changed); algunos eventos necesitan throttle y otros deben descartarse si el consumidor es lento (drop if slow). Meter todas esas políticas en el bucle principal lo convertiría en un galimatías — el bucle solo debería ocuparse de la máquina de estados del agent; las políticas de broadcast se abstraen en una capa aparte.

Las tres primitivas las inyecta desde fuera createGatewayNodeSessionRuntime (createGatewayNodeSessionRuntime:971) y createAgentEventHandler solo las consume, sin preocuparse por la implementación — el módulo de chat se puede probar sin un WebSocket real.

Archivos clave

Flujo de datos

createAgentEventHandler (createAgentEventHandler:299) recibe las tres primitivas de broadcast y los tres registros de suscripción, y devuelve el dispatcher (evt: AgentEventPayload) => void. Internamente mantiene throttle, lifecycle retry, cache spawnedBy, etc.

broadcast y broadcastToConnIds provienen de createGatewayBroadcaster (createGatewayBroadcaster:100), cuya implementación recorre clients: Set<GatewayWsClient> con el bucle núcleo:

typescript
for (const c of params.clients) {
  if (targetConnIds && !targetConnIds.has(c.connId)) continue;
  if (!hasEventScope(c, event)) continue;        // filtrado por scope
  const slow = c.socket.bufferedAmount > MAX_BUFFERED_BYTES;
  if (slow && opts?.dropIfSlow) {
    if (!isTargeted) clientSeq.set(c, (clientSeq.get(c) ?? 0) + 1);  // dirigido no consume seq
    continue;                                                       // si se puede drop, se dropa
  }
  if (slow) { c.socket.close(1008, "slow consumer"); continue; }    // si no, se corta
  c.socket.send(frame);  // {"type":"event","event":...,"payload":...,"seq":N}
}

El backpressure tiene dos niveles: dropIfSlow: true (deltas de chat de alta frecuencia, stream descartable) dropa pero actualiza seq para evitar desalineación; dropIfSlow: false (lifecycle crítico, resultados de tool) hace directamente socket.close(1008, "slow consumer"); se prefiere cortar la conexión antes que dropar eventos críticos. seq solo se incrementa en broadcasts no dirigidos (isTargeted ? undefined : nextSeq); los broadcasts dirigidos son subconjuntos ya filtrados, no participan en la reordenación del cliente.

nodeSendToSession proviene de createGatewayNodeSessionRuntime (nodeSendToSession:31); su implementación es una sola línea: nodeSubscriptions.sendToSession(sessionKey, event, payload, nodeSendEvent). El comentario del código «session fanout goes through the subscription manager so node reconnects and explicit unsubscribes keep both node->session indexes in sync» explica por qué la entrega entre nodos debe ir por el subscription manager y no por nodeRegistry.send directo — al reconectar un nodo o al hacer un unsubscribe explícito, los índices bidireccionales deben mantenerse sincronizados; de lo contrario, se entregaría a un nodo ya desconectado o se perdería un nodo recién conectado.

La entrada de eventos es el closure devuelto (dispatcher:1136), que primero resuelve sessionKey, agentId y clientRunId (muchos eventos no traen estos campos; hay que mirarlos en chatRunState.registry), y luego despide por tipo de stream. sendChatPayload (sendChatPayload:890) clasifica según la visibilidad de la UI de control: si es visible, broadcast("chat", payload) y sendNodeSessionPayloadForAgent entrega entre nodos; si está oculta, consulta sessionMessageSubscribers.get(deliverySessionKey) y, si hay suscriptores, broadcastToConnIds entrega dirigida — normalmente se trata de una herramienta del operator que el panel principal no ve en este run. sendAgentPayload (sendAgentPayload:967) tiene una lógica simétrica, solo cambia el nombre del evento a "agent".

La entrega entre nodos tiene además un detalle (resolveNodeSessionDeliveryKeys:495): si sessionKey es "global", se expande en dos delivery keys: agent:DEFAULT:global y global. Los eventos globales deben llegar tanto a los clientes suscritos a «el global de un agent» como a los suscritos a «lo global»; pero solo se empujan ambos cuando el agentId es el agent por defecto — en caso contrario, solo se empuja el agent-scoped, para evitar entregas duplicadas.

Límites y fallos

  • Tratamiento de consumidores lentos: cuando socket.bufferedAmount > MAX_BUFFERED_BYTES se dispara el backpressure. dropIfSlow: true (deltas de stream, chat.send_timing) salta el envío pero actualiza seq; dropIfSlow: false (lifecycle crítico, resultados de tool) hace directamente socket.close(1008, "slow consumer"); si no se corta, el buffer crecería sin límite y colapsaría el proceso (backpressure:157).
  • Reintento de errores lifecycle: el Map pendingTerminalLifecycleErrors registra el estado de reintento de eventos end/error phase dentro del grace period. lifecycleErrorRetryGraceMs controla la ventana; cuando llega una nueva generación de lifecycle, se cancela el timer anterior.
  • Cache spawnedBy: resolveSpawnedBy solo hace consulta al store para sessionKey de subagent/ACP; los no lineage devuelven null directamente para no contaminar la cache (spawnedByCache:383). Las rutas de alta frecuencia como delta/flush/final evitan golpear el store de sesión en cada llamada.
  • Supresión de restart recovery: resolveRestartRecoveryLifecycleState comprueba si un evento lifecycle es producto de restart recovery; si lo es, suppress: true y no se hace broadcast — evita que, tras un reinicio, los eventos de un run viejo recuperado reaparezcan al cliente.
  • Fallback de resolución de sessionKey: primero chatLink.sessionKey, luego evt.sessionKey, por último resolveSessionKeyForRun(evt.runId). Esta cascada de tres niveles garantiza que un evento siempre encuentre la sesión a la que pertenece.
  • La entrega entre nodos no consume seq: isTargeted ? undefined : nextSeq; si los broadcasts dirigidos también incrementaran seq, el cliente vería saltos y pensaría que se perdieron eventos.
  • Sincronización de suscripciones de nodo: nodeSendToSession pasa por el subscription manager y no por nodeRegistry.send directo, para que los índices bidireccionales se mantengan sincronizados al reconectar/desconectar/desuscribir un nodo. Saltarse el manager provocaría pérdidas tras una reconexión.

Resumen

La capa de broadcast separa la política de entrega de eventos del agent del bucle principal, abstrayéndola en tres primitivas inyectadas por createGatewayNodeSessionRuntime, que createAgentEventHandler se limita a consumir. La implementación real de broadcast (createGatewayBroadcaster) gestiona filtrado de scope, backpressure y numeración seq; la entrega entre nodos se enruta por el subscription manager para garantizar índices sincronizados tras reconexiones. La entrada de eventos es el dispatcher closure, que despide por tipo de stream a chat / agent / session.tool / sessions.changed, y dentro de cada stream elige entre broadcast completo o dirigido según la visibilidad de la UI de control.

Cómo llegan las peticiones a esta capa y cómo se despachan se ve en Tabla de métodos RPC y despacho de peticiones; cómo se inyectan las primitivas de broadcast en todo el gateway se ve en Núcleo del gateway; cómo se vincula el estado de suscripción de sesión al ciclo de vida de un run del agent se ve en Session Manager.

Referencias oficiales: Documentación del gateway · README.