Skip to content

BaseChannel : le canal comme primitive d'état

源码版本1.2.9

Responsabilités

Dans le runtime de LangGraph, le « canal (channel) » est l'unité minimale d'état (state). L'instance Pregel produite par la compilation d'un StateGraph ne détient pas directement le dict utilisateur : elle détient un ensemble d'instances de BaseChannel — un canal par clé — et chaque canal stocke sa valeur, fusionne les écritures, expose une lecture, et sait aussi sérialiser son état courant sous forme de point de contrôle (checkpoint). BaseChannel est la classe de base abstraite partagée par tous les types de canaux (BaseChannel:19).

Elle ne stocke aucune donnée elle-même (__slots__ ne contient que key et typ, voir __slots__:22-26). Elle fixe seulement quatre membres abstraits que les sous-classes doivent implémenter : les deux propriétés de type ValueType / UpdateType, et les trois méthodes update / get / from_checkpoint (abstract methods:60-99). S'y ajoutent checkpoint / is_available / consume / finish / copy, cinq méthodes avec une implémentation par défaut qui soit se rabat sur get(), soit renvoie False en no-op.

Autrement dit, BaseChannel fixe en une fois l'interface de « comment l'état est stocké, modifié et tracé ». Le moteur Pregel n'a qu'à appeler update / get / checkpoint, sans avoir à savoir que la fusion est faite par un LastValue, un Topic ou un BinaryOperatorAggregate. C'est ce qui découple totalement le moteur de la forme de l'état.

Motivation de conception

Pourquoi modéliser l'état comme des « canaux » plutôt que de passer un dict entre les nœuds ?

  • Sémantique unifiée de fusion d'écritures : au sein d'un même superpas (superstep), plusieurs nœuds peuvent écrire concurrently sur la même clé. Un dict ne sait pas exprimer les différentes règles de fusion — « liste à étendre, scalaire à écraser, compteur à additionner » ; le canal réunit ces stratégies derrière un update(values: Sequence[Update]) -> bool, et la sous-classe décide du fold.
  • Vue immutable intra-superpas : un nœud lit le snapshot pris en début de superpas ; les valeurs écrites par d'autres nœuds dans le même superpas ne deviennent visibles qu'au superpas suivant. get est une lecture de la valeur courante, update n'est appelé qu'aux frontières de superpas par apply_writes (apply_writes:317-323), ce qui élimine toute race.
  • Sérialisable : chaque canal sait produire une valeur sérialisable via checkpoint() (checkpoint:49-58) et se reconstruire via from_checkpoint (from_checkpoint:60-65). C'est sur ces deux interfaces que la famille BaseCheckpointSaver s'appuie pour écrire l'état du thread en SQLite/Postgres.
  • Hooks de cycle de vie : consume et finish (consume / finish:101-121) donnent aux canaux spéciaux (par ex. Topic(accumulate=False) ou LastValueAfterFinish) une chance de modifier leur état à la frontière de superpas ou en fin d'exécution ; le no-op par défaut indique que la plupart des canaux n'ont pas besoin de cette sémantique.
  • Auto-description des types : les deux propriétés ValueType / UpdateType (ValueType / UpdateType:28-36) permettent à la compilation comme à l'exécution de retrouver « ce que ce canal stocke et ce qu'il accepte en écriture ». StateGraph les utilise dans _is_field_channel:1862-1887 pour reconnaître les déclarations Annotated[..., SomeChannel].

Fichiers clés

  • BaseChannel class:19-26 — définition de la classe de base abstraite, paramètres génériques Value / Update / Checkpoint, __slots__ ne déclarant que key et typ.
  • ValueType / UpdateType:28-36 — deux propriétés abstraites où la sous-classe déclare le type stocké et le type d'écriture.
  • checkpoint:49-58 — implémentation par défaut : renvoie self.get(), ou MISSING pour un canal vide.
  • from_checkpoint:60-65 — méthode abstraite : construire un canal équivalent depuis une valeur de checkpoint, à implémenter par la sous-classe.
  • get:69-73 — lecture abstraite, lève EmptyChannelError sur un canal vide.
  • is_available:75-85 — implémentation par défaut via get() + capture de EmptyChannelError ; généralement surchargée par un test plus rapide.
  • update:89-99 — écriture abstraite, appelée par Pregel une fois à la fin de chaque superpas, dans n'importe quel ordre ; renvoie True si le canal a été modifié.
  • consume / finish:101-121 — hooks de cycle de vie, no-op par défaut renvoyant False.
  • apply_writes:315-345 — point d'entrée unique où Pregel réécrit les writes du superpas dans les canaux ; update / consume / finish y sont tous ordonnancés.
  • _is_field_channel:1862-1887 — à la compilation, StateGraph traduit Annotated[..., channel] en instance de canal via isinstance(item, BaseChannel).

Flux de données

Les écritures sur les canaux ont lieu dans apply_writes : Pregel regroupe les writes de toutes les tâches du superpas par canal, puis appelle update (apply_writes:317-323) :

python
# Apply writes to channels
updated_channels: set[str] = set()
for chan, vals in pending_writes_by_channel.items():
    if chan in channels:
        if channels[chan].update(vals) and next_version is not None:
            checkpoint["channel_versions"][chan] = next_version
            # unavailable channels can't trigger tasks, so don't add them
            if channels[chan].is_available():
                updated_channels.add(chan)

# Channels that weren't updated in this step are notified of a new step
if bump_step:
    for chan in channels:
        if channels[chan].is_available() and chan not in updated_channels:
            if channels[chan].update(EMPTY_SEQ) and next_version is not None:
                checkpoint["channel_versions"][chan] = next_version
                if channels[chan].is_available():
                    updated_channels.add(chan)

Notez le second bloc : les canaux non écrits par aucune tâche dans le superpas reçoivent tout de même un appel update(EMPTY_SEQ). C'est un contrat implicite de l'interface BaseChannel — la méthode update d'une sous-classe doit accepter une séquence vide, généralement en renvoyant False pour « pas de modification ». La sémantique « effacer à chaque superpas » de Topic(accumulate=False) s'appuie justement sur cet appel à vide : l'update vide lui fait perdre la valeur du superpas précédent, qui devient invisible au superpas suivant.

Position du cycle de vie d'un canal dans un tick de Pregel :

Limites et échecs

  • EmptyChannelError : levé par get() quand le canal n'a jamais été écrit (get:70-73). L'implémentation par défaut de is_available s'appuie sur la capture de cette exception ; une sous-classe qui peut décider à moindre coût devrait la surcharger.
  • Sentinelle MISSING : checkpoint() renvoie langgraph._internal._typing.MISSING — pas None — quand get() lève EmptyChannelError (checkpoint fallback:55-58), de sorte que None puisse être une valeur légale sans être confondu avec « vide ».
  • update renvoie False : un retour False signifie que le canal n'a pas bougé, Pregel ne met donc pas à jour channel_versions et prepare_next_tasks ne déclenchera aucun nœud depuis ce canal (update call:319). C'est la base du critère de convergence de Pregel.
  • Appel à séquence vide : apply_writes appelle aussi update(EMPTY_SEQ) sur les canaux inchangés (EMPTY_SEQ update:329). Le contrat de BaseChannel.update est qu'une entrée vide renvoie False, mais Topic(accumulate=False) en profite pour s'effacer et peut renvoyer True : une sous-classe ne peut donc pas supposer « entrée vide = ne rien faire ».
  • consume et finish en no-op par défaut (consume / finish:101-121) : le retour False signifie que la plupart des canaux ne participent pas à la gestion du cycle de vie. apply_writes appelle finish à la fin de bump_step (finish call:338) ; seuls les canaux spéciaux comme LastValueAfterFinish l'utilisent pour la sémantique « visible seulement après la fin de l'exécution ».
  • copy se rabat sur checkpoint par défaut : BaseChannel.copy fait par défaut self.from_checkpoint(self.checkpoint()) (copy:40-47). Si le checkpoint est coûteux pour une sous-classe (par ex. une deep-copy d'une grosse liste), elle devrait surcharger copy avec une implémentation moins chère — c'est ce que fait Topic.copy.

Résumé

BaseChannel fixe en une fois « comment l'état est stocké, fusionné, sérialisé », et le moteur Pregel n'a qu'à suivre la séquence update → is_available → get → checkpoint → finish pour fonctionner sur n'importe quelle implémentation de canal. Pour les implémentations concrètes, voir /channels/last-value et /channels/topic-binop ; le déroulé des lectures/écritures par Pregel est dans /pregel/algo ; la persistance de la valeur renvoyée par checkpoint() est détaillée dans /checkpoint/base-saver.

Voir la documentation officielle : documentation LangGraph · README