Sesión y entrega
Responsabilidad
Dos tareas: ① Sesión (Session) — agrupar «de dónde viene el mensaje» (plataforma, chat, remitente) bajo una clave estable y persistir la conversación a disco, para que el agent continúe el contexto (context) entre reinicios; ② Libro de entrega (Delivery Ledger) — anota para la respuesta final del agent una «obligación de entrega», de modo que tras un crash pueda identificar respuestas no entregadas y reenviarlas, garantizando que el usuario las reciba. Ambas viven en la capa del gateway, desacopladas del agent.
Motivo de diseño
El agent en sí es stateless — solo corre un turno (turn) del bucle; de dónde viene el contexto y a dónde va la respuesta final no son asunto suyo. La capa de sesión resuelve «los mensajes de la misma persona en el mismo grupo sobre el mismo tema pertenecen a la misma conversación»; esta capa debe persistir con una clave estable, si no, tras un reinicio el agent pierde el hilo. El libro de entrega resuelve «el agent terminó de generar la respuesta, pero el proceso del gateway se cayó a mitad del send()»: esa respuesta a medio morir no la ve completa el usuario y además provoca regeneración duplicada. El ledger anota en sqlite cada obligación «pendiente de entregar»; tras restart, un barrido la reenvía o la marca como fallida, en lugar de tirarla y empezar de cero.
Archivos clave
class SessionSource:149-299— dataclass del origen del mensaje (platform/chat/sender, con hash)class SessionContext:299-700— contexto de sesión (key, origen, estrategia de continuación)hash de id:40-74—_hash_id/_hash_sender_id/_hash_chat_id(privacidad + estabilidad)validación de seguridad de ruta:100-149—_is_path_unsafe/_is_session_key_unsafe(la session key fluye a una ruta del sistema de archivos)ventana de auto_continue:26-40—auto_continue_freshness_windowdoc del libro de entrega:1-58— objetivo de diseño: best-effort, los fallos nunca bloquean el flujo principalcompute_obligation_id:146-155— id único de obligación (session key + ref de mensaje + contenido)máquina de estados:155-203—record_obligation/mark_attempting/mark_delivered/mark_failedsweep_recoverable:203-276— tras un crash, explora obligaciones recuperables y las reenvíaconexión SQLite:77-102—_db_path/_connect, persistencia en sqlite local
Flujo de datos
- El adaptador (adapter) recibe un evento de plataforma, del que extrae
SessionSourcevíaMessageEvent:1759. SessionSource:149hashea un id estable de chat/sender (gateway/session.py:40).- Se construye un
SessionContext:299, obteniendo la session key (trasvalidación de seguridad de ruta:109, se usa como ruta en disco). - La session key determina desde qué archivo se carga el historial, que se entrega al agent para continuar el contexto.
- Tras producir la respuesta final, el gateway anota en
record_obligation:155→ envía →mark_delivered. - Si el envío cae a mitad, en el siguiente arranque
sweep_recoverable(gateway/delivery_ledger.py:203) explora obligaciones en estadoattemptingcuyo proceso propietario ya no existe, y las reenvía o marca como fallidas.
SessionSource hashea campos sensibles — plataforma + chat + sender — con un sha256 truncado a 12 caracteres:
def _hash_id(value: str) -> str:
"""Deterministic 12-char hex hash of an identifier."""
return hashlib.sha256(value.encode("utf-8")).hexdigest()[:12]
def _hash_sender_id(value: str) -> str:
"""Hash a sender ID to ``user_<12hex>``."""
return f"user_{_hash_id(value)}"El hash persigue dos cosas: privacidad (el user_id de la plataforma no se vuelca a disco tal cual) y estabilidad — la misma persona que cambia el nombre del dispositivo sigue en la misma sesión. La session key luego fluye a una ruta del sistema de archivos, así que la validación es la siguiente línea de defensa:
def _is_path_unsafe(value: object) -> bool:
"""Return True if ``value`` could traverse outside the sessions dir."""
if not value:
return False
s = str(value)
if ".." in s or "/" in s or "\\" in s:
return True
# Leading Windows drive path, e.g. "C:\\..." or "d:/...".
return len(s) >= 2 and s[0].isalpha() and s[1] == ":"
def _is_session_key_unsafe(value: object) -> bool:
"""True if ``value`` could be a real traversal vector in a session_key.
``session_key`` is a *logical* routing key (e.g.
``agent:main:google_chat:group:spaces/<id>``) — it never touches the
filesystem, so the strict separator-rejecting guard from
``_is_path_unsafe`` is over-broad"""Que las dos funciones de validación sean distintas es deliberado — session_key es una clave lógica de enrutado (Google Chat con spaces/<id>/threads/<id> lleva / embebido); rechazar / mataría casos legítimos. Pero los campos que fluyen a una ruta de archivo, como session_id, sí deben rechazarlos. Mantenerlos separados evita el agujero de compatibilidad de «una validación para todo».
La máquina de estados del libro de entrega se reduce a cuatro métodos + una tabla sqlite:
def compute_obligation_id(session_key: str, message_ref: str, content: str) -> str:
"""Stable id: same turn + same content re-records idempotently, while
distinct threads/topics on the same chat can never collide."""
payload = f"{session_key}|{message_ref}|{content}"
return hashlib.sha256(payload.encode("utf-8", "replace")).hexdigest()[:24]
def record_obligation(*, obligation_id, session_key, platform, chat_id, thread_id, content) -> None:
"""Record a final response as owed to the platform (state='pending')."""obligation_id hashea session_key + message_ref + content en 24 chars hex — registrar el mismo turno con el mismo contenido es idempotente (evita que regeneraciones repetidas hinchen el ledger), y threads distintos en el mismo chat nunca colisionan.
El sweep_recoverable para recuperación tras crash es el núcleo del diseño:
def sweep_recoverable(now=None, *, deliverable_platforms=None) -> List[Dict[str, Any]]:
"""Claim undelivered rows owned by dead processes; return them for
redelivery.
Claiming atomically re-stamps the owner to THIS process and increments
``attempts``, so a second gateway racing the same sweep cannot
double-claim (the UPDATE is guarded on the previous owner stamp).
Rows over the attempts cap or older than the stale cutoff transition to
'abandoned' instead of being returned.
"""La clave es el «claim» — cuando dos instancias de gateway barren la misma fila, la que primero cambia owner_pid al suyo gana el derecho de reenvío; la cláusula WHERE del UPDATE garantiza que solo una lo consigue. deliverable_platforms evita «quemar attempts en vano»: una plataforma que no está conectada en este boot no se reclama, se deja para un próximo boot que sí pueda enviar, en lugar de empujarla hacia MAX_ATTEMPTS en cada restart.
Límites y fallos
- Mezclar ruta y clave lógica:
_is_path_unsaferechaza/;_is_session_key_unsafesolo rechaza..y/inicial. Elspaces/<id>de Google Chat es unsession_keyválido pero no unsession_id— confundirlos provoca o bien un falso rechazo o bien path traversal. - Entrega duplicada: si dos gateways barren la misma fila a la vez, el claim debe ser atómico (UPDATE WHERE owner_pid = valor_anterior), si no, la respuesta se envía dos veces. Las filas reenviadas además llevan
needs_marker(porque el send original pudo quedar semi-exitoso). - Plataforma no conectada:
deliverable_platformsfiltra las plataformas offline en este boot, para no quemar attempts en cada restart hasta que lleguen aMAX_ATTEMPTSy quedenabandoned. Se aguarda al boot que sí las tenga conectadas. - Sesión caducada:
auto_continue_freshness_windowes la ventana para «continuar la conversación anterior»; fuera de la ventana no se auto-continue, para no tratar una conversación de hace horas como continuación actual.
Resumen
La capa de sesión se encarga de «a qué conversación continua pertenece este mensaje»: con hash + validación de ruta proyecta identidades externas a archivos en disco de forma segura. El libro de entrega se encarga de «si la respuesta final que el agent quería enviar realmente llegó al usuario»: una máquina de estados en sqlite hace la recuperación tras crash. Ambas son best-effort pero imprescindibles como capa de fiabilidad.