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), satisfaceISessionRealtimeBus. - Manager:
EventServiceManagerBus(app/modules/manager/infrastructure/adapters/realtime_bus.py), satisfaceIManagerRealtimeBus.
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_tokenen 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 llamasession_token, el valor es un string arbitrario — y eso es lo que permite tener dos canales (ver abajo). Subscriber: un consumidor con su propiaasyncio.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únpublish_eventse 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 recibirtyping_stoppedochat_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:
Claves del formato:
- Solo se usa el campo SSE
data:. No se usanevent:,id:niretry:. - El "tipo" del evento viaja dentro del JSON como campo
type(no como la cabecera SSEevent:). - El sobre siempre es
{type, data, timestamp}, dondetimestampes ISO-8601 UTC autogenerado por elSessionEventen 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:
- Emite un frame
connected(handshake) directo (connected_frame(), payload{"status": "ok"}). - Si la IA estaba escribiendo, hace replay de
typing_started(get_typing_state). - 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 conformat(event["type"], event["data"]). - Al expirar el timeout sin eventos, emite un
heartbeaty vuelve a empezar (esto mantiene viva la conexión). - 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
EventSourcedel navegador no puede enviar cabeceras, así que el token JWT de admin viaja como?token=...(parámetrotokenrequerido). El endpoint lo valida con el autenticador de admin (get_admin_authenticator().authenticate, 401 si es inválido o ha expirado) y luego exigerequire_restaurant_access(restaurant_id, db, admin)(multitenancy: el admin debe ser dueño del restaurante). - Snapshot inicial: el primer frame es
sala_snapshotcon el estado de todas las mesas activas ({"sessions": [...]}), calculado en una sesión de BD fresca abierta conasync_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
connectedniheartbeat, y no hace replay de typing. Tras el snapshot entra al bucle desubscribe(topic, timeout=30s)y, al expirar el timeout sin eventos, simplemente vuelve a suscribirse en silencio. - Reenvía
table_state_changed/silent_table_alerta medida que ocurren. - Cierre limpio ante cancelación: el bucle captura
asyncio.CancelledError(además de comprobarrequest.is_disconnected()) para terminar el generador sin propagar la excepción cuando el cliente se va.
Reconexión, heartbeat y timeout¶
- Timeout de 30s:
subscribeespera 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 unheartbeatcomo 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 handshakeconnectedy el replay detyping_startedcompletan la puesta al día. El manager se pone al día con elsala_snapshotinicial de cada reconexión. - Sin
retry:de cliente: como el wire format no incluye la directiva SSEretry:, la cadencia de reconexión la decide el cliente (elEventSourcedel 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 elfinallydesubscribe) 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/doney publicachat_messagepor el bus. - Arquitectura hexagonal — cómo se organizan los módulos
sessionymanagerque exponen estos endpoints.