PregelLoop:スーパーステップループ駆動器
役割
PregelLoop は LangGraph 実行時の「メトロノーム」です。Pregel 自身はリソースの組立と invoke / astream の入口提供のみを担い、グラフを一圈々々回すのは一つの PregelLoop インスタンスです。現在のチェックポイント (checkpoint)、チャネル (channel) の集合、pending_writes、step カウンタ、status ステートマシン、そして tasks 辞書を保持します——これらが単一のスーパーステップ (superstep) 内のすべてのコンテキストです。外層の while loop.tick(): ... loop.after_tick() という形はこの上に構築されます。
二つの主要メソッドは tick(tick:599)と after_tick(after_tick:683)です。前者は「走る前」を担います: ステップ数上限に達していないか確認、prepare_next_tasks でこのステップが実行するノードを self.tasks に詰め込み、中断をチェック、drain_requested を処理、前回の再開で残った pending_writes を成功ノードに貼り戻す。後者は「走った後」を担います: すべての task.writes を集め、apply_writes でチャネルへマージし、pending をクリアし、チェックポイントを落とし、interrupt_after が必要か判定する。両者は前後で「入力を読む → 调度 → 実行 → 書き戻し → 永続化」という BSP ループを挟み込みます。
言い換えると、PregelLoop は自らノードを走らせません——それは PregelRunner の仕事です。PregelLoop は「このステップで誰を見るか、走り終わったらどう収尾するか、次のステップを続けるか」だけを決めます。
設計動機
なぜループロジックを tick と after_tick の二段に分け、一つの step() メソッドで済ませないのか?
- 実行権を呼び出し元に返す:
tickがTrueを返した後、呼び出し側(Pregel.astreamの実装体、sync main loop:2964-2984参照)が自らrunner.tickを駆動してこのバッチの task を走らせ、その後loop.after_tickを呼びます。これによりPregelLoopはPregelRunnerへの参照を持たず、両オブジェクトは疎結合な協力関係にあって入れ子になりません。 - ステートマシンの明示化:
statusフィールドは"input" / "pending" / "done" / "draining" / "interrupt_before" / "interrupt_after" / "out_of_steps"の七値を取ります(status Literal:256-264)。各状態は一つの離脱理由に対応し、外部コードは戻り値のセマンティクスを推測せずに resume すべきか raise すべきか break すべきか判断できます。 - 中断再開のセマンティクスが明確:
is_replayingは最初の tick でのみ真——このステップが前回中断前に走り終えた task の「再再生」であることを示し、after_tickが終わると閉じます(is_replaying = False:716)。_reapply_writes_to_succeeded_nodesと組み合わせてチェックポイントに残った成功 writes をメモリの task に貼り戻すことで、成功ノードを再実行せず、失敗ノードを error handler に進めます。 - Drain は外部介入の入口:
RunControl.drain_requestedは stream の消費側が設定し(例えばクライアントが途中停止を要求)、tickはこれを見るとdraining状態に入りFalseを返し、外層のwhileを自然に終了させ、ノード途中で強制 kill しません。 - Exit-mode delta 書き込みとチェックポイント時序:
_delta_write_futs(_delta_write_futs:207)はすべての delta チャネルのput_writesfuture を集め、_checkpointer_put_after_previousは次の checkpoint を書く前にこれらを必ず drain します——「書き込みが checkpoint を産み、checkpoint が書き込みを上書きしない」という因果順序を保証します。
主要ファイル
class PregelLoop:158— クラス定義。フィールドはほぼ素の attribute で property ラップがなく、サブクラスSyncPregelLoop/AsyncPregelLoopが直接上書きしやすい。status Literal:256-264— 七状態の列挙。"out_of_steps"はステップ数上限到達、"draining"は外部停止要求、"done"は正常終了。tick method:599-681— 一つの tick の全ロジック: ステップ数チェック、prepare_next_tasks、空タスク時の done、drain、_reapply_writes_to_succeeded_nodes、should_interrupt。out_of_steps:607-609—self.step > self.stopのときout_of_steps状態に入りFalseを返し、外層のwhileが終了。done branch:653-655—self.tasksが空のときdone状態に入る。グラフの「自然死」。draining branch:657-659—control.drain_requestedが真のときdraining状態に入る。reapply + resume handlers:662-664— 再開時に pending writes を成功ノードへ貼り戻し、_resume_error_handlers_if_applicableで失敗ノードに handler task を準備。interrupt_before:667-671—should_interruptでノードがinterrupt_beforeリストに入るかチェックし、命中すればGraphInterruptを raise。after_tick method:683-726— writes 収集 →apply_writes→ values emit → pending クリア →_put_checkpoint→interrupt_afterチェック。_reapply_writes_to_succeeded_nodes:736-749— pending writes をメモリ task に貼り戻すが、ERROR / ERROR_SOURCE_NODE / INTERRUPT / RESUMEの四種の制御信号はスキップ。_resume_error_handlers_if_applicable:751-816—ERROR_SOURCE_NODEマーカーを走査し、失敗ノードに error-handler task を構築して runner が元タスクをスキップし handler を直接走らせる。put_writes:415-508— task が writes を産出する入口。制御チャネルの重複除去、null task の蓄積、checkpointer のput_writesfuture を起動。SyncPregelLoop:1469/AsyncPregelLoop:1722— 二つのサブクラス。__enter__/__exit__、accept_push、schedule_error_handlerなどの同期/非同期差分を上書き。
データフロー
tick の入口はまずステップ数上限をチェックし、続いて prepare_next_tasks で「次のステップでどの task を走らせるか」を算出します。このステップがスーパーステップ全体のシード——以降のすべての動作はこのバッチの task の加工です。
def tick(self) -> bool:
"""Execute a single iteration of the Pregel loop.
Returns:
True if more iterations are needed.
"""
# check if iteration limit is reached
if self.step > self.stop:
self.status = "out_of_steps"
return False
# prepare next tasks
self.tasks = prepare_next_tasks(
self.checkpoint,
self.checkpoint_pending_writes,
self.nodes,
self.channels,
self.managed,
self.config,
self.step,
self.stop,
for_execution=True,
manager=self.manager,
store=self.store,
checkpointer=self.checkpointer,
trigger_to_nodes=self.trigger_to_nodes,
updated_channels=self.updated_channels,
retry_policy=self.retry_policy,
cache_policy=self.cache_policy,
)これは tick 頭部:599-629 からの引用です。self.step > self.stop は「ステップ数上限」の離脱条件——stop は PregelLoop.__init__ 時に recursion_limit から設定され、再帰深度が超えると out_of_steps で止まります。prepare_next_tasks は現在の checkpoint、pending_writes、ノード表、チャネル表から本ステップで実行する task を算出し self.tasks に詰めます。for_execution=True はこれらの task を本当に走らせる(dry-run でストリーミングプレビューではない)ことを示します。
続いて三つの並列離脱条件——空 task / drain / 再開時の writes 貼り戻し——を経て should_interrupt に入ります:
# if no more tasks, we're done
if not self.tasks:
self.status = "done"
return False
if self.control is not None and self.control.drain_requested:
self.status = "draining"
return False
# if there are pending writes from a previous loop, apply them
if not self.is_replaying and self.checkpoint_pending_writes:
self._reapply_writes_to_succeeded_nodes(self.tasks)
self._resume_error_handlers_if_applicable()
# before execution, check if we should interrupt
if self.interrupt_before and should_interrupt(
self.checkpoint, self.interrupt_before, self.tasks.values()
):
self.status = "interrupt_before"
raise GraphInterrupt()これは tick 中段:652-671 からの引用です。interrupt_before のチェックは pending_writes の再開後に行われます——再開時はチャネルのバージョン番号が前回の書き込みを反映しているため、should_interrupt が「前回の interrupt 以降に新たな更新があるか」を正しく判定できます。
after_tick は runner が task をすべて走らせ終わった後に呼ばれ、収尾を担います:
def after_tick(self) -> None:
# finish superstep
writes = [w for t in self.tasks.values() for w in t.writes]
self._delta_channels_with_overwrite.update(
ch
for ch, v in writes
if isinstance(self.specs.get(ch), DeltaChannel) and _get_overwrite(v)[0]
)
# all tasks have finished
self.updated_channels = apply_writes(
self.checkpoint,
self.channels,
self.tasks.values(),
self.checkpointer_get_next_version,
self.trigger_to_nodes,
)これは after_tick 頭部:683-698 からの引用です。まずすべての task の writes を平坦化し、同時にどの delta チャネルで overwrite が発生したかを集計し(sparse replay の開始値に影響)、apply_writes を呼んでこれらの writes を本当にチャネルへマージします。戻り値 updated_channels は次ステップの prepare_next_tasks で「次に触発されるノードを見つける」の高速化に使われます。after_tick の尾部では checkpoint_pending_writes を空にし、is_replaying を False にし、_put_checkpoint({"source": "loop"}) でチェックポイントを落とし、interrupt_after 判定を行います(after_tick 尾部:714-724)。
スーパーステップ全体のライフサイクルはこのように描けます:
境界と失敗
out_of_stepsはエラーではない:step > stopで直接Falseを返し、status をout_of_stepsに設定し、外層のwhileが自然終了します。Pregelの上層は status に基づきRecursionErrorを throw するか黙って収尾するかを判断します——これは「ソフト上限」セマンティクスで、例外パスではありません(out_of_steps:607-609)。drainingでは lifecycle events を emit できない:_emit_graph_lifecycle_eventはdraining状態での呼び出しを明示的に拒否します(draining guard:384-385)。drain はユーザ能動の停止で、interrupt / resume セマンティクスではないからです。- 再開時の pending writes に制御信号が含まれる場合はスキップ:
_reapply_writes_to_succeeded_nodesはERROR / ERROR_SOURCE_NODE / INTERRUPT / RESUMEを必ずスキップします(skip control signals:746-747)。さもないと、失敗した task が成功扱いされ、走った handler が再び触発されます。 - Resume 時の error handler 再スケジュール:
_resume_error_handlers_if_applicableはerror_handler_nodeを持つノードにのみ作用します(handler_node check:791-793)。handler を持たない失敗ノードは原状のままで、runner の_should_stop_othersが panic パスへ進みます。 put_writesの null task 蓄積セマンティクス:NULL_TASK_IDの writes は上書きせず蓄積します(null task accumulate:422-431)。input writes のような「どのノードにも属さないが checkpoint に入るべき」書き込みに使われます。- Delta channel の future は次の checkpoint より先に落盤されなければならない:
_delta_write_futsはすべての delta チャネルのput_writesfuture を集め、次の checkpoint 落盤前にこれを drain します(_delta_write_futs:201-207)——順序を逆にすると sparse replay が checkpoint を産んだ書き込みを見逃します。
まとめ
PregelLoop はステートマシン駆動のスーパーステップメトロノームで、「task 準備 → 走る → 収尾」の三段で PregelRunner を挟み込み、自身は実行に触れません。tick / after_tick の二段、status の七状態、is_replaying + _reapply_writes_to_succeeded_nodes の再開メカニズムを理解すれば、LangGraph 実行モデルの主軸を掴めます。Runner がどう task を本当に走らせるかは /pregel/runner を、アルゴリズム層 prepare_next_tasks / apply_writes / should_interrupt の内部詳細は /pregel/algo を、全体の組立は /pregel/pregel を参照してください。チャネルの書き戻しセマンティクスは /channel/base-channel にあります。
公式資料:LangGraph ドキュメント · README