Session et livraison
Responsabilité
Deux choses : (1) Session (session) — associer « d'où vient le message » (plateforme, chat, expéditeur) à une clé stable et persister la conversation sur disque, pour que l'agent puisse reprendre le contexte (context) entre redémarrages ; (2) Registre de livraison (Delivery Ledger) — enregistrer pour la réponse finale de l'agent une « obligation de livraison » ; après un crash, identifier les réponses non livrées et les réexpédier, pour garantir que l'utilisateur finit par les recevoir. Les deux vivent dans la couche passerelle (gateway), découplés de l'agent.
Mot de conception
L'agent lui-même est sans état — il ne fait que tourner une boucle ; d'où vient le contexte et où va la réponse finale ne devrait pas être son problème. La couche session résout « les messages d'une même personne dans un même groupe sur un même sujet appartiennent à une même conversation » ; elle doit utiliser une clé stable persistée sur disque, sinon l'agent perd le fil après un redémarrage. Le registre de livraison résout « l'agent a fini de générer sa réponse, mais le processus passerelle a planté au milieu de send() » — une mort en plein vol que l'utilisateur ne voit ni en entier, ni sans risquer une régénération en double. Le registre enregistre en sqlite chaque obligation « à livrer » ; au redémarrage un simple balaye permet de réexpédier ou marquer échec, plutôt que de tout recommencer.
Fichiers clés
class SessionSource:149-299— dataclass de la source du message (platform/chat/sender, avec hash)class SessionContext:299-700— contexte de session (clé, source, stratégie de reprise)hash d'id:40-74—_hash_id/_hash_sender_id/_hash_chat_id(confidentialité + stabilité)validation de sécurité des chemins:100-149—_is_path_unsafe/_is_session_key_unsafe(la clé de session se retrouve dans un chemin du système de fichiers)fenêtre de fraîcheur auto_continue:26-40—auto_continue_freshness_windowdoc du registre de livraison:1-58— objectif de conception : best-effort, un échec ne bloque jamais le flux principalcompute_obligation_id:146-155— id unique d'obligation (clé de session + ref message + contenu)machine à états:155-203—record_obligation/mark_attempting/mark_delivered/mark_failedsweep_recoverable:203-276— après un crash, balaye les obligations récupérables et les réexpédieconnexion SQLite:77-102—_db_path/_connect, persistance en sqlite locale
Flux de données
- L'adaptateur (adapter) reçoit l'événement de la plateforme, remonte un
SessionSourceviaMessageEvent:1759. SessionSource:149produit par hash des id stables de chat/sender (gateway/session.py:40).- Construction d'un
SessionContext:299pour obtenir la clé de session (utilisée comme chemin disque aprèsvalidation de sécurité des chemins:109). - La clé de session décide depuis quel fichier l'historique est chargé, puis alimente l'agent pour reprendre le contexte.
- Une fois la réponse finale produite par l'agent, la passerelle enregistre l'obligation dans
record_obligation:155→ envoi →mark_delivered. - En cas de crash pendant l'envoi, au prochain démarrage
sweep_recoverable(gateway/delivery_ledger.py:203) balaie les obligations à l'étatattemptingdont le processus propriétaire est mort, les réexpédie ou les marque comme échouées.
SessionSource hash des champs sensibles comme plateforme + chat + sender, par sha256 tronqué à 12 caractères :
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)}"Le hash sert d'abord à la confidentialité (le user_id de la plateforme n'est pas écrit tel quel sur disque), ensuite à la stabilité — une même personne qui change de nom d'appareil reprend dans la même session. La clé de session finit dans un chemin du système de fichiers, donc une seconde validation :
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"""Les deux validateurs sont volontairement différents — session_key est une clé de routage logique (Google Chat utilise spaces/<id>/threads/<id> qui contient légitimement des /), refuser / tuerait des clés valides ; mais un champ comme session_id qui se retrouve dans un chemin doit refuser /. Distinguer les deux évite le piège de compatibilité d'« une seule validation pour tout ».
La machine à états du registre de livraison, ce sont quatre méthodes + une table 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 hashe session_key + message_ref + content en 24 hex — un même tour (turn) + un même contenu se réenregistre de façon idempotente (évite de gonfler le registre par régénération), et des threads distincts dans le même chat ne collent jamais.
Le cœur de la conception est sweep_recoverable pour la récupération post-crash :
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.
"""Le point clé est le « claim » — quand deux instances de passerelle balaient la même ligne, celle qui réussit à remplacer owner_pid par elle-même gagne le droit de réexpédier ; la clause WHERE du UPDATE garantit qu'un seul des deux réussit. deliverable_platforms évite de « brûler des attempts pour rien » : une plateforme non connectée à ce boot n'est pas réclamée, elle attend un boot capable de l'envoyer, plutôt que d'être poussée à chaque redémarrage vers la limite MAX_ATTEMPTS.
Limites et échecs
- Mélange chemin et clé logique :
_is_path_unsaferefuse strictement/,_is_session_key_unsafene refuse que..et les/de tête.spaces/<id>de Google Chat est unsession_keylégitime mais pas unsession_id; les confondre mène soit à un faux rejet, soit à une traversée de chemin. - Livraison en double : quand deux passerelles balaient la même ligne, le claim doit être atomique (UPDATE WHERE owner_pid = ancienne valeur), sinon la même réponse part deux fois. Les lignes réexpédiées portent aussi
needs_marker(l'envoi initial a pu partiellement réussir). - Plateforme non connectée :
deliverable_platformsfiltre les plateformes hors ligne à ce boot, pour éviter qu'à chaque redémarrage on brûle un attempt jusqu'àMAX_ATTEMPTSpuisabandoned. La ligne attend un boot capable de livrer. - Expiration de session :
auto_continue_freshness_windowest la fenêtre pendant laquelle « continuer la conversation précédente » est acceptable ; au-delà, plus d'auto-continue, pour éviter de traiter une conversation de plusieurs heures avant comme la suite courante.
Résumé
La couche de session gère « à quelle conversation continue appartient ce message », en mappant l'identité externe vers un fichier disque via hash + validation de chemin. Le registre de livraison gère « la réponse finale que l'agent voulait envoyer est-elle vraiment arrivée à l'utilisateur », via une machine à états sqlite pour la récupération après crash. Les deux sont best-effort mais indispensables à la fiabilité.