ちゃんと動いているのに止まって見える
目標
同じグラフを5つの方式で流しながら、何がいつ出てくるかを自分で数えます。更新だけを集めて最後の状態を再構築してみて、流した最後の状態がinvokeの答えと同じかを突き合わせます。
なぜ重要なのか
エージェントで流すべきものは、トークンだけではありません。モデルを最後に1回呼び出すグラフなら、トークンのストリーミングを付けても、前の沈黙はそのままです。ユーザーに必要なのは、「今どの段階か」と「何が新しく決まったか」で、それはモデルではなく、グラフが知っています。
モードごとに単位が違います。updatesはノード単位なので、並んで動くノード2つが同じスーパーステップにあっても、イベントが別々に出て、valuesはスーパーステップ単位なので、その2つの結果が統合されたあとに1回出ます。debugは、スーパーステップ番号まで教えてくれます。この違いを知らないと、「なぜ2回出るのか」をバグと誤解します。
更新だけを集めて状態を再構築するには、リデューサーを知っている必要があります。グラフの中ではリデューサーがしてくれていたことを、外では自分で行う必要があり、連結するキーを上書きすると、足跡が1つしか残りません。
そして、状態に残したくない進行表示は、customで出します。状態に入れると、チェックポイントに載り、再開するときに古い表示が復活します。
採点ツールは、書かれた説明を信用しません。書かれたモジュールを実際に読み込んで、毎回異なる議事録で動かしてみて、イベントの個数・順序・形を、採点ツールが別に数えた値と突き合わせます。時間は測りません。
ステップ
- /root/work/agstream/stream.pyに
State・4つのノード(split・keypoints・actions・compose)・build_graph()・start(text)・run_values(text)を作成してください。splitのあとでkeypointsとactionsが並んで動き、両方がcomposeに集まります。 run_updates(text)を追加して、stream_mode="updates"のイベントを、そのままリストとして返すようにしてください。REDUCERS・rebuild(text)・rebuild_report(text)を追加して、更新だけを集めて最後の状態を再構築し、それがvaluesの最後と同じかを突き合わせるようにしてください。run_debug(text)を追加して、{수퍼스텝번호: [노드이름...정렬됨]}(プレースホルダーはスーパーステップ番号、ノード名、ソート済みを意味する語です)を返すようにしてください。run_multi(text)とmode_sequence(text)を追加して、stream_mode=["updates", "values"]の結果を扱うようにしてください。composeがwriter: StreamWriterを受け取って断片を2つ出すようにし、run_custom(text)を追加してください。compare_invoke(text)を追加して、流した最後の状態とinvokeの答えが同じか、最初に届くノードが何か、イベントがいくつかを、一度に返すようにしてください。- 測る議事録を保存し(保存先: /root/work/agstream/minutes.txt)、それで測って記録してください(保存先: /root/work/agstream/stream_report.json、/root/work/agstream/stream_report.md)。
参考
- 実行の契約: 採点ツールは、
/root/work/agstream/stream.pyをPythonモジュールとして読み込んで、上に書いた名前を直接使います。スクリプトとして実行しません。 - 状態のキー:
text・lines・points・todos・summary・trace。traceだけが連結するリデューサーを使い、残りは上書きします。ノード名と状態のキーは、重なってはいけません。重なると、コンパイルのときにValueError: 'x' is already being used as a state keyが出ます。 - ノードの動作:
splitは、textを行に分けて空白の行を捨て、linesに入れます。keypointsは、"결정"(韓国語で「決定」を意味する語です)が入った行をpointsに、actionsは、"하기로"(韓国語で「することにした」を意味する語です)が入った行をtodosに入れます。composeは、summaryに"결정 N건 · 할 일 M건"(韓国語の文字列で、「決定N件・やることM件」という意味です)を入れます。4つのノードすべてが、traceに自分の名前を書きます。 start(text)は、{"text": text, "trace": []}を返します。run_values・run_updates・run_multi・run_customは、stream(...)が返したものを、加工せず、そのままリストとして返します。rebuildは、start(text)から始めて、updatesのイベントを順に適用したディクショナリを返します。REDUCERSは、{"trace": "append"}のように、連結するキーを書いておく表です。rebuild_report(text)の答え:{"same": 참거짓, "stream_order": [...], "merged_order": [...]}(プレースホルダーは真偽値です)。stream_orderは再構築した状態のtrace、merged_orderはvaluesの最後のイベントのtraceです。2つは同じではありません。なぜ違うかは、自分で見て書いてください。run_debugは、event["type"] == "task"のものだけを見て、event["step"]とevent["payload"]["name"]を集めます。値はソートします。mode_sequence(text)は、run_multiのタプルから、モード名だけを順に取り出します。compare_invoke(text)の答え:{"same": 참거짓, "first_node": 문자열, "values_events": 정수, "updates_events": 정수}(プレースホルダーは真偽値、文字列、整数です)。- ステップ8で保存する議事録(/root/work/agstream/minutes.txt)は、レポートを作った議事録です。採点ツールがそのファイルを読んで、同じ入力でもう一度測るので、レポートを作ったあとに内容を変えると、数値がずれます。
- このPodにはインターネットがありません。langgraph 0.2.60がすでに入っています。
- 公式ドキュメント: Streaming・Graph API overview・Typesリファレンス
- よくある間違い:
updatesがスーパーステップ単位だと信じること、再構築するときに連結するキーを上書きすること、進行表示を状態に入れること、stream()の結果を加工して返すこと(採点ツールは元の形を見ます)です。
状態全体がステップごとに出る
/root/work/agstream/stream.pyにState・4つのノード・build_graph()・start(text)・run_values(text)を作成してください。splitのあとでkeypointsとactionsが並んで動き、両方がcomposeに集まります。
stream(입력, stream_mode="values")(プレースホルダーは入力です)が返すものを、そのままリストとして返せばよいです。先頭に入力状態が1回出ることを、目で確認してください。並んで動く2つのノードの結果は、統合されたあとに1つのイベントとして出ます。
更新はノードごとに1つずつ
run_updates(text)を追加して、stream_mode="updates"のイベントを、加工せず、そのままリストとして返すようにしてください。
イベント1つが{노드이름: 갱신}(プレースホルダーはノード名と更新です)です。並んで動くノード2つが同じスーパーステップにあっても、イベントは別々に出ます。そのため、valuesよりイベントの数が多くなります。いくつ出るか、自分で数えてみてください。
更新だけで状態を再構築する
REDUCERS = {"trace": "append"}とrebuild(text)を追加してください。start(text)から始めて、updatesのイベントを順に適用したディクショナリを返します。
グラフの中ではリデューサーがしてくれていたことを、外では自分で行う必要があります。すべてを上書きすると、連結するキーでは最後の1つだけが残ります。ところが、リデューサーを正しく真似ても、完全には同じになりません。2つのtraceを並べて出力して、どこが違うか、なぜそうなのかを、自分で確認してください。どちらも毎回同じ答えを出しますが、お互いに違います。
何番目の段階で、何が一緒に動くか
run_debug(text)を追加して、{수퍼스텝번호: [노드이름...정렬됨]}(プレースホルダーはスーパーステップ番号、ノード名、ソート済みを意味する語です)を返すようにしてください。event["type"] == "task"のものだけを見ます。
debugモードは、ノードごとにtaskとtask_resultの2つのイベントを出します。task側のstepがスーパーステップ番号で、payload["name"]がノード名です。並んで動くノードが同じ番号を持つかを確認してください。それが、このモードを使う理由です。
複数のモードを一緒に聞く
run_multi(text)とmode_sequence(text)を追加してください。stream_mode=["updates", "values"]を渡すと、(모드이름, 값)(プレースホルダーはモード名と値です)のタプルが出ます。
リストでモードを渡すと、イベントごとに、どのモードのものかが前に付きます。2つのモードのイベントが、どんな順序で混ざって出てくるかを、自分で見てください。あるモードをすべて出してから別のモードを出す、というわけではありません。mode_sequenceは、その順序を名前だけ取り出して見せます。
状態に残さずに見せる
composeがwriter: StreamWriterを2番目の引数として受け取り、{"stage": "compose", "points": 개수}と{"stage": "compose", "todos": 개수}(プレースホルダーは個数です)を順に出すようにして、run_custom(text)を追加してください。
ノード関数の2番目の引数の名前をwriterにして、StreamWriterで型注釈すると、LangGraphが埋めてくれます。writer(...)で入れたものは、状態に残らず、聞く側にだけ届きます。進行表示を状態に入れると、チェックポイントに載り、再開するときに復活します。
流した最後のものと、一度に受け取った答え
compare_invoke(text)を追加して、{"same": 참거짓, "first_node": 문자열, "values_events": 정수, "updates_events": 정수}(プレースホルダーは真偽値、文字列、整数です)を返すようにしてください。
valuesの最後のイベントがinvoke()の答えと同じなら、画面はストリーミングで描き、記録は別にinvokeで残すように、2回動かす理由がありません。first_nodeは、updatesの最初のイベントに入っているノード名です。ユーザーに最初に見せられるものが何かを教えてくれます。
何をいつ見せるかを整理する
測る議事録を保存してください(保存先: /root/work/agstream/minutes.txt。決定事項とやることが、それぞれ1行以上、空行も1つ)。その議事録で測って、次のように記録してください。/root/work/agstream/stream_report.jsonには、values_events・updates_events・first_node・supersteps・mode_sequence・custom_chunks・rebuild_same・stream_order・merged_orderを、/root/work/agstream/stream_report.mdには、## 다섯 가지 방식이 각각 무엇을 주나 ## 갱신만으로 상태를 다시 세우려면 ## 화면에 먼저 보여 줄 수 있는 것 ## 현장에서 무엇을 고르나の4つの節を書きます(4つの見出しは、順に韓国語で「5つの方式がそれぞれ何を返すか」「更新だけで状態を再構築するには」「画面に先に見せられるもの」「現場では何を選ぶか」を意味します)。
すべて、下の議事録で動かして得た値です。superstepsは、run_debugの答えを文字列のキーに変えて入れ、rebuild_same・stream_order・merged_orderは、rebuild_reportの答えそのままです。数値をでっち上げず、数えて入れてください。