Skip to content

@entrypoint / @task : API fonctionnelle

源码版本1.2.9

Responsabilités

@entrypoint et @task sont les deux pierres angulaires de l'API fonctionnelle de LangGraph, définies dans libs/langgraph/langgraph/func/__init__.py. @entrypoint transforme une fonction Python ordinaire (synchrone ou asynchrone) en une instance Pregel — la valeur de retour est un graphe compilé, utilisable directement via .invoke / .astream(entrypoint.__call__:516). @task transforme une fonction en _TaskFunction(_TaskFunction:59) ; à l'appel, il renvoie un SyncAsyncFuture, et ne peut être appelé qu'à l'intérieur d'un @entrypoint ou d'un nœud StateGraph.

Sa position est au-dessus de StateGraph, comme autre entrée de haut niveau — pas besoin de construire un StateGraph, d'appeler add_node / add_edge / compile, il suffit d'écrire une fonction. Le moteur reste Pregel : entrypoint.__call__ finit par return Pregel(...), en emballant la fonction entrypoint comme un graphe à nœud unique (func.__name__ est le seul nom de nœud, son trigger est START) ; le nœud exécute le corps de la fonction entrypoint, et à son retour la valeur de retour est écrite dans le canal END(Pregel construction:572).

Motivation de conception

  • Moins de boilerplate : StateGraph exige de déclarer explicitement le schéma d'état, les fonctions de nœud, les arêtes et les arêtes conditionnelles ; @entrypoint permet de n'écrire qu'une fonction, dont la séquence d'appels task dans le corps constitue la structure du graphe. Pour les workflows simples (quelques appels, enchaîner quelques outils), l'API fonctionnelle est plus légère.
  • Toutes les capacités de Pregel préservées : l'API fonctionnelle n'est pas une version amputée ; le moteur est Pregel, donc checkpointer, store, cache, retry_policy, context_schema, interrupt(), Command(resume=...) sont tous disponibles — appeler interrupt() dans une fonction entrypoint emprunte exactement le même mécanisme qu'un interrupt() dans un nœud StateGraph.
  • entrypoint.final découple valeur de retour et valeur persistée : beaucoup de workflows conversationnels veulent « renvoyer X à l'appelant cette fois, mais stocker Y dans le checkpoint pour la prochaine fois » (le paramètre previous lit la dernière save). entrypoint.final(value=X, save=Y)(entrypoint.final:477) sépare explicitement ces deux choses, plutôt que de se reposer sur une convention via dict.
  • Modèle de paramètre previous : la signature de la fonction entrypoint peut ajouter *, previous: Any = None ; au prochain appel avec le même thread_id, la valeur récupérée est celle de la dernière save, équivalente à une « variable de mémoire entre appels », sans avoir à déclarer soi-même un schéma d'état et un reducer.
  • Appel de task = arête du graphe : dans le corps de la fonction entrypoint, task_a(...) renvoie un future, et .result() / await récupère le résultat — cette relation d'appel est naturellement une arête ; Pregel la traite comme une tâche PUSH insérée dans le canal TASKS, et prepare_next_tasks l'ordonnance naturellement au pas suivant. Pas besoin de add_edge.

Fichiers clés

  • _TaskFunction:59 — produit du décorateur @task, détient func / retry_policy / cache_policy / timeout.
  • _TaskFunction.__call__:86 — à l'appel d'une task, passe par _call_with_options, récupère la callback CONFIG_KEY_CALL depuis la config, et remet la tâche au runner courant.
  • task decorator:110 — entrée du décorateur @task, supporte @task / @task(retry_policy=...).
  • entrypoint class:262 — classe du décorateur @entrypoint, final est une dataclass imbriquée accrochée à la classe.
  • entrypoint.__init__:437 — reçoit checkpointer / store / cache / retry_policy / timeout / context_schema, stocké dans self.
  • entrypoint.final:477 — dataclass final(value=R, save=S), renvoie value à l'appelant, stocke save dans le checkpoint.
  • entrypoint.__call__:516 — transforme la fonction en instance Pregel, la logique de conversion centrale est ici.
  • Pregel construction:572 — construit Pregel(nodes={func.__name__: PregelNode(bound=bound, triggers=[START], ...)}), les canaux sont START / END / PREVIOUS, trois LastValue.
  • get_runnable_for_entrypoint:175 — emballe la fonction entrypoint en RunnableCallable, la version synchrone passe par run_in_executor.
  • get_runnable_for_task:200 — emballe la fonction task en RunnableSeq(run, ChannelWrite([RETURN])), la valeur de retour de la task est écrite dans le canal RETURN.

Flux de données

@entrypoint transforme une fonction en Pregel en construisant un graphe minimaliste à trois canaux :

python
graph: Pregel[Any, ContextT, Any, Any] = Pregel(
    nodes={
        func.__name__: PregelNode(
            bound=bound,
            triggers=[START],
            channels=START,
            timeout=self.timeout,
            writers=[
                ChannelWrite(
                    [
                        ChannelWriteEntry(END, mapper=_pluck_return_value),
                        ChannelWriteEntry(PREVIOUS, mapper=_pluck_save_value),
                    ]
                )
            ],
        )
    },
    channels={
        START: EphemeralValue(input_type),
        END: LastValue(output_type, END),
        PREVIOUS: LastValue(save_type, PREVIOUS),
    },
    input_channels=START,
    output_channels=END,
    stream_channels=END,
    stream_mode=stream_mode,
    stream_eager=True,
    checkpointer=self.checkpointer,
    store=self.store,
    cache=self.cache,
    # ...
)

START est un EphemeralValue (réinitialisé à chaque pas), END et PREVIOUS sont des LastValue (ne gardent que la dernière valeur). La valeur de retour de la fonction entrypoint est décomposée par deux mappers : _pluck_return_value écrit entrypoint.final.value (ou la valeur nue) dans END pour l'appelant ; _pluck_save_value écrit entrypoint.final.save (ou la valeur nue) dans le canal PREVIOUS, utilisé comme previous au prochain appel(pluck mappers:547). stream_mode="updates" est la valeur par défaut de l'API fonctionnelle (pas values), car l'entrypoint peut contenir plusieurs tasks, et updates reflète mieux la progression intermédiaire.

Chaîne d'appel d'une task : task_fn(arg)_TaskFunction.__call___call_with_options(_call_with_options:276) → récupère la callback impl depuis config[CONF][CONFIG_KEY_CALL]impl(func, args, kwargs, retry_policy=..., callbacks=..., timeout=...) renvoie un SyncAsyncFuture. Cet impl est injecté via partial(...) par PregelRunner.tick / atick au moment de l'ordonnancement(CONFIG_KEY_CALL inject:211) — donc la task ne s'exécute pas elle-même, elle remet la fonction et les arguments au runner courant, qui décide quand l'exécuter en concurrence, puis renvoie le résultat via le future à l'appelant.

Limites et échecs

  • Générateurs non supportés : la fonction entrypoint ne peut pas être def gen / async def gen (avec yield) ; __call__ lève explicitement raise NotImplementedError au début(no generators:525). Pour du streaming, utiliser StreamWriter (stream_mode="custom") plutôt que yield.
  • Task non appelable depuis l'extérieur : appeler task_fn(...) en dehors d'un nœud @entrypoint / StateGraph échoue, car get_config()[CONF][CONFIG_KEY_CALL] n'existe pas. L'existence d'une task dépend du contexte du runner à l'exécution.
  • Task synchrone ne supporte pas timeout : task(timeout=...) n'a d'effet que pour les fonctions asynchrones ; pour les fonctions synchrones, cela lève sync_timeout_unsupported(sync timeout unsupported:237) — une task synchrone s'exécute dans un thread executor, et Python ne peut pas l'annuler proprement intra-processus.
  • previous est None par défaut : au premier appel, previous est None ; la fonction entrypoint doit gérer elle-même la branche « pas de valeur précédente » ; si entrypoint.final n'est pas utilisé, la valeur de retour est écrite à la fois dans END et PREVIOUS, et le previous du prochain appel récupère la dernière valeur de retour elle-même.
  • final doit être correctement paramétré : l'annotation de retour -> entrypoint.final[int, str] doit fournir les deux paramètres de type ; n'en donner qu'un seul ou aucun lève TypeError dans __call__(final param check:562).
  • Dépendance au checkpointer : si entrypoint(checkpointer=...) n'est pas configuré avec un checkpointer, interrupt() et la mémoire previous entre appels sont inutilisables — la première lève RuntimeError, la seconde renvoie toujours None.

Résumé

L'API fonctionnelle est du sucre syntaxique au-dessus de StateGraph ; le moteur reste Pregel avec un graphe à nœud unique + canal TASKS + CONFIG_KEY_CALL injecté par le runner. La séquence d'appels task(...) dans le corps de la fonction entrypoint est la structure implicite du graphe. Pour les détails de concurrence voir /func/concurrency, pour la reprise d'interruption voir /interrupt/interrupt. Voir la documentation officielle : LangGraph 文档 · README