Saltar a contenido

Tiempo real (SSE)

Esta página explica cómo funciona la capa de tiempo real de Camarero IA: un event bus en memoria que distribuye eventos a dos audiencias distintas (el cliente en la mesa y el dashboard de sala del manager) mediante Server-Sent Events (SSE).

Diátaxis: explicación

Esta página es conceptual. Para el catálogo exhaustivo de eventos, sus triggers y la forma exacta de cada payload, consulta la Referencia de eventos SSE.

Qué es y por qué SSE

SSE (Server-Sent Events) es un canal HTTP unidireccional servidor → cliente: el cliente abre una conexión GET que queda viva y el servidor va emitiendo "frames" de texto a medida que ocurren eventos. A diferencia de WebSocket, es half-duplex (solo baja datos), viaja sobre HTTP normal y el navegador lo consume con la API EventSource.

Camarero IA lo usa para dos flujos que necesitan empuje del servidor:

  • Cliente en la mesa: ver llegar el streaming del chat de la IA, cambios del carrito, cambios de estado de la comanda, llamadas al manager, etc.
  • Dashboard de sala (manager): ver el estado de todas las mesas activas en tiempo real (entregas, alertas de mesa silenciosa, transiciones de ciclo de vida).

Tres flujos SSE, un solo formato

Además de los dos canales push del event bus, el endpoint de chat (POST .../messages) responde con un StreamingResponse SSE propio (tokens content, done, etc.). Comparte el mismo formato de cable que el bus, pero no pasa por el event bus: lo genera el motor de chat inline. Esta página cubre los dos canales del bus; el stream de chat se detalla en la Referencia de eventos SSE.

El event bus en memoria

El corazón del sistema es EventService, definido en app/shared/infrastructure/event_bus.py. Es un pub/sub in-memory basado en colas asyncio.Queue, expuesto como singleton mediante get_event_service().

Cada módulo accede al bus a través de un adaptador façade que satisface su puerto de dominio y resuelve el singleton de forma lazy (en cada llamada, no en la construcción) para que los monkeypatches de los tests sigan siendo efectivos:

  • Cliente: LegacyEventServiceBus (app/modules/session/infrastructure/adapters/realtime_bus.py), satisface ISessionRealtimeBus.
  • Manager: EventServiceManagerBus (app/modules/manager/infrastructure/adapters/realtime_bus.py), satisface IManagerRealtimeBus.

In-memory: no escala horizontalmente

Al vivir en memoria del proceso, el bus no funciona con varias instancias del backend tras un balanceador: un evento publicado en el proceso A no llega a un subscriber conectado al proceso B. El propio código lo documenta como deuda: para escalado horizontal habría que migrar a Redis pub/sub (ver comentario en EventService). Hoy asume un único proceso.

Conceptos

  • Topic (session_token en la firma del método): la clave de enrutado. Cada publicación va dirigida a un topic y solo lo reciben los subscribers de ese topic. Pese a que el parámetro se llama session_token, el valor es un string arbitrario — y eso es lo que permite tener dos canales (ver abajo).
  • Subscriber: un consumidor con su propia asyncio.Queue(maxsize=100).
  • Buffer de pendientes: por cada topic se guardan los últimos 50 eventos (_pending_events). Al suscribirse, el bus vacía primero ese buffer en la cola del nuevo subscriber (drenaje de reconexión) bajo el mismo lock y antes de registrar al subscriber, de modo que ningún publish_event se cuele entre el drenaje y el alta.
  • Estado de typing: el bus recuerda si la IA estaba "escribiendo" (typing_started) para que un cliente que se conecta tarde reciba un replay del estado (get_typing_state). El estado se limpia al recibir typing_stopped o chat_message.

API principal

Método Propósito
publish_event(topic, event_type, data) Crea un SessionEvent, lo añade al buffer (máx. 50) y lo encola en cada subscriber del topic. Marca el topic como activo y actualiza el estado de typing.
subscribe(topic, timeout=30.0) Async generator: drena pendientes, luego espera eventos vivos hasta timeouts; si la cola desborda emite un frame overflow y termina.
get_typing_state(topic) Devuelve el data del último typing_started o None. Sync (sin lock).
format_sse_message(event_type, data) Serializa un evento al formato de cable SSE.

Además, EventService expone un guard de stream concurrente desacoplado del pub/sub (mark_stream_active / mark_stream_inactive / is_stream_active sobre conversation_id), usado por el endpoint de chat para devolver 429 si ya hay un stream activo para la misma conversación.

Backpressure: el frame overflow

Si un subscriber no consume lo bastante rápido y su cola (100 ítems) se llena, publish_event marca overflow=True; en el siguiente ciclo subscribe emite un frame overflow y cierra el stream. El cliente debe interpretarlo como "te has quedado atrás, refresca el estado" y reconectar.

Formato de cable (wire format)

Todos los eventos del bus viajan con el mismo sobre. format_sse_message produce:

data: {"type": <str>, "data": <obj>, "timestamp": <ISO-8601>}\n\n

Claves del formato:

  • Solo se usa el campo SSE data:. No se usan event:, id: ni retry:.
  • El "tipo" del evento viaja dentro del JSON como campo type (no como la cabecera SSE event:).
  • El sobre siempre es {type, data, timestamp}, donde timestamp es ISO-8601 UTC autogenerado por el SessionEvent en el momento de su creación.
  • Cada frame termina en doble salto de línea (\n\n), como exige la especificación SSE.

Serialización de datetime (canal cliente)

El adaptador del cliente convierte recursivamente cualquier datetime del payload a ISO-8601 (_jsonify) antes de json.dumps, porque los eventos de dominio tipados (dataclasses) pueden contener datetime crudos que romperían la serialización JSON del bus.

Ambos endpoints SSE responden con media_type="text/event-stream" y estas cabeceras:

Content-Type: text/event-stream
Cache-Control: no-cache
Connection: keep-alive
X-Accel-Buffering: no

(X-Accel-Buffering: no desactiva el buffering de Nginx para que los frames lleguen sin retraso.)

Los dos canales

El mismo EventService sirve dos audiencias usando dos espacios de topic distintos:

Canal Topic Audiencia Endpoint SSE
Cliente {session_token} Cliente en la mesa GET /api/v1/sessions/{session_token}/events
Sala manager:{restaurant_id} Dashboard de sala del manager GET /api/v1/admin/manager/restaurants/{restaurant_id}/events

El topic del cliente es el propio session_token (string opaco de la sesión); el del manager se construye con el helper manager_topic(restaurant_id) (app/modules/manager/domain/ports.py), que antepone el prefijo manager:. Al estar en espacios de nombres separados, una publicación al cliente nunca llega al manager y viceversa — salvo cuando se publica a propósito a ambos (p. ej. una acción del manager que dispara un table_state_changed en manager:{id} y, en paralelo, un lifecycle_state_changed en el topic del cliente para que el chat del cliente reaccione).

Canal del cliente

GET /api/v1/sessions/{session_token}/events (app/modules/session/interface/sessions_router.py, función session_events).

Antes de abrir el stream, el endpoint verifica que la sesión existe y está activa con GetSessionUseCase (404 si no). El generador entonces:

  1. Emite un frame connected (handshake) directo (connected_frame(), payload {"status": "ok"}).
  2. Si la IA estaba escribiendo, hace replay de typing_started (get_typing_state).
  3. Entra en bucle: comprueba request.is_disconnected(), se suscribe al bus con timeout de 30s y reenvía cada evento publicado al topic del cliente con format(event["type"], event["data"]).
  4. Al expirar el timeout sin eventos, emite un heartbeat y vuelve a empezar (esto mantiene viva la conexión).
  5. Si detecta que el cliente se ha desconectado (request.is_disconnected()), corta (lo comprueba tanto al inicio de cada ciclo como tras cada evento).

Cierre tras session_closed (SSE-CLOSE01)

Cuando el stream reenvía un evento session_closed, termina la conexión. La sesión ya no es válida; mantenerla abierta filtraría heartbeats, dejaría subscribers zombi y arriesgaría reproducir eventos viejos a una sesión nueva que reutilizara el token.

Canal del manager (sala)

GET /api/v1/admin/manager/restaurants/{restaurant_id}/events (app/modules/manager/interface/admin_manager_router.py, función manager_events).

Particularidades respecto al canal del cliente:

  • Autenticación por query param: la API EventSource del navegador no puede enviar cabeceras, así que el token JWT de admin viaja como ?token=... (parámetro token requerido). El endpoint lo valida con el autenticador de admin (get_admin_authenticator().authenticate, 401 si es inválido o ha expirado) y luego exige require_restaurant_access(restaurant_id, db, admin) (multitenancy: el admin debe ser dueño del restaurante).
  • Snapshot inicial: el primer frame es sala_snapshot con el estado de todas las mesas activas ({"sessions": [...]}), calculado en una sesión de BD fresca abierta con async_session_maker() (la sesión de la request no se mantiene abierta durante toda la vida del stream).
  • Sin handshake ni heartbeat: a diferencia del canal del cliente, el manager no emite connected ni heartbeat, y no hace replay de typing. Tras el snapshot entra al bucle de subscribe(topic, timeout=30s) y, al expirar el timeout sin eventos, simplemente vuelve a suscribirse en silencio.
  • Reenvía table_state_changed / silent_table_alert a medida que ocurren.
  • Cierre limpio ante cancelación: el bucle captura asyncio.CancelledError (además de comprobar request.is_disconnected()) para terminar el generador sin propagar la excepción cuando el cliente se va.

Reconexión, heartbeat y timeout

  • Timeout de 30s: subscribe espera como mucho 30s por evento. Si no llega nada, el ciclo termina y el llamante vuelve a entrar al bucle. El canal del cliente emite entonces un heartbeat como keep-alive (para que proxies y el navegador no corten una conexión "ociosa"); el canal del manager se re-suscribe sin emitir nada.
  • Reconexión y replay: si el cliente reconecta (nueva conexión GET), el drenaje del buffer de pendientes (últimos 50 eventos del topic) le pone al día con lo que se perdió mientras estuvo desconectado. El handshake connected y el replay de typing_started completan la puesta al día. El manager se pone al día con el sala_snapshot inicial de cada reconexión.
  • Sin retry: de cliente: como el wire format no incluye la directiva SSE retry:, la cadencia de reconexión la decide el cliente (el EventSource del navegador reconecta solo por defecto).
  • Desconexión: ambos generadores comprueban request.is_disconnected() para salir limpiamente cuando el cliente se va; al terminar el generador, el bus elimina el subscriber (remove_subscriber, en el finally de subscribe) y limpia el topic si queda vacío.

Keep-alive durante el razonamiento del LLM

El streaming del chat (endpoint POST .../messages) emite además comentarios SSE : thinking mientras el modelo razona. Son líneas que empiezan por : que el cliente ignora por especificación, pero evitan que la conexión se considere ociosa durante respuestas largas. Detalle en la Referencia de eventos SSE.

Páginas relacionadas

  • Referencia de eventos SSE — catálogo completo de tipos de evento, triggers y forma de cada payload.
  • Control de sala — el dashboard del manager que consume el canal de sala.
  • Chat híbrido — el motor de chat que alimenta el stream de content / done y publica chat_message por el bus.
  • Arquitectura hexagonal — cómo se organizan los módulos session y manager que exponen estos endpoints.