AG Day 30 多代理通訊:狀態交接與訊息協議
執行需求:CPU 可跑。AG Day 29(原文連結)把 research-agent 拆成了 supervisor 與三個 worker,但 worker 之間交換資料的方式很粗糙:大家共用同一份 TeamState,靠欄位名稱互相約定「誰該寫哪裡、誰該讀哪裡」,沒有明確的訊息格式,也完全沒處理「worker 失敗了該怎麼回報」。今天要把這件事做得更嚴謹:設計一套結構化的訊息協議,讓每個 worker 完成工作後回報的不是隨手塞的字典欄位,而是一份有固定結構、有狀態、有清楚語意的交接封包(handoff payload)。同時我們會補上一個工作佇列(task queue)的概念,讓 supervisor 不只是每次臨場決定下一步,而是能先把整個任務拆成一串子任務,逐一指派。這些改動都在本機 CPU 上就能驗證,不需要外部服務或 API 金鑰。
引言
昨天的架構跑得起來,但如果你認真檢查 search_worker、retrieval_worker、writer_worker 的回傳值,會發現它們各自回傳不同形狀的字典:有的帶 title、url、snippet,有的帶 text、source_url。這在只有三個 worker、流程固定的示範裡還能將就,但只要專案繼續長大——加入第四個 worker、讓某個 worker 可能失敗需要重試、或想知道「這個 worker 花了幾步、呼叫了哪些工具」——目前這種「各自約定欄位」的做法就會迅速失控。沒有共同的協議,除錯時只能一個一個印出來看,也無法寫出通用的記錄或監控邏輯。
今天的目標是引入兩個結構化的核心概念。第一是 任務規格(TaskSpec):把「要做什麼」明確寫成一個有 ID、有描述、有指派目標、有狀態的物件,取代目前隱含在流程順序裡的任務概念。第二是 交接封包(HandoffPayload):把「worker 做完了、結果是什麼」統一成一種格式,不管是哪個 worker 回報,supervisor 都能用同一套邏輯處理成功、失敗、需要重試三種結果。這套協議讓多代理系統從「靠命名默契運作」變成「靠明確契約運作」,跟 AG Day 8 學過的結構化輸出是同一種思路,只是這次套用在代理與代理之間,而不是模型與人類之間。
原理/觀念
為什麼共用大狀態不等於有協議
把所有東西塞進一份 TypedDict 共用狀態,確實是 LangGraph 多代理系統最常見的起手式,但「共用狀態」跟「有訊息協議」是兩件不同的事。共用狀態只解決了「資料放哪裡」的問題,協議要解決的是「這份資料的語意、格式、生命週期是什麼」。舉例來說,如果 search_worker 失敗了(例如網路逾時),它應該回傳什麼?昨天的實作完全沒有處理這個情境,一旦失敗就會讓後面的 retrieval_worker 拿到空資料,卻不知道這是「本來就沒有資料」還是「上一步出錯了」。有協議之後,每個 worker 回報時都要明確指出 status:成功、失敗,或需要主管介入,supervisor 才能依此做出正確的下一步決策,而不是被空值誤導。
集中式路由 vs. 去中心化交接
AG Day 29 用的是集中式路由:所有控制權都經過 supervisor。今天引入的交接封包,其實也能支援另一種常見的多代理拓樸:去中心化交接(decentralized handoff),也就是某個 worker 做完事之後,可以直接在交接封包裡指定「下一步交給誰」,不一定要繞回 supervisor 才能決定。這在 LangGraph 裡通常透過 Command 物件實現:節點回傳 Command(goto=..., update=...),同時完成「更新狀態」與「指定下一個節點」兩件事,圖裡就不需要額外的條件邊。今天我們仍然採用集中式路由(維持昨天的圖拓樸),但會把交接封包設計成兩種拓樸都能沿用的通用格式,方便日後升級。
任務佇列:把流程從「隱含順序」變成「明確清單」
昨天的規則路由靠檢查狀態欄位是否為空來決定下一步,本質上是把任務順序寫死在條件判斷裡。今天引入 TaskSpec 清單後,supervisor 可以先把整個研究問題拆成幾個子任務(例如「查詢主題背景」「查詢近期進展」),放進 task_queue,再逐一指派給對應的 worker,執行完的任務移到 completed_tasks。這個設計為 AG Day 31(原文連結)要做的長任務檢查點鋪路:任務清單本身就是可以序列化、可以中途存檔、可以在恢復執行時重新讀取的狀態,比起原本隱含在條件判斷裡的順序更容易持久化。
完整實作
我們在既有的 agents/ 套件裡新增 protocol.py,定義 TaskSpec 與 HandoffPayload 兩個 Pydantic 模型,並調整 state.py、workers.py、supervisor.py 來套用這套協議:
touch research-agent/src/research_agent/agents/protocol.py
第一步:定義訊息協議本身。TaskSpec 描述「要做的一件事」,HandoffPayload 描述「做完之後的回報」,兩者都是明確型別的 Pydantic 模型,而不是隨意形狀的字典:
# research-agent/src/research_agent/agents/protocol.py
from enum import Enum
from typing import Optional, Any
from pydantic import BaseModel, Field
class TaskStatus(str, Enum):
PENDING = "pending"
IN_PROGRESS = "in_progress"
DONE = "done"
FAILED = "failed"
class TaskSpec(BaseModel):
"""supervisor 派給某個 worker 的一件具體任務。"""
task_id: str
description: str = Field(description="這個任務要完成什麼,用人類可讀的一句話描述")
assigned_to: str = Field(description="負責的 worker 名稱,例如 search_worker")
status: TaskStatus = TaskStatus.PENDING
class HandoffPayload(BaseModel):
"""worker 完成(或失敗)一個任務後,回報給 supervisor 的統一格式。"""
task_id: str
worker_name: str
status: TaskStatus
summary: str = Field(description="這個 worker 做了什麼、結果如何的簡短摘要")
result: Optional[Any] = Field(default=None, description="真正的產出資料,例如搜尋結果串列或報告全文")
error: Optional[str] = Field(default=None, description="status 為 failed 時的錯誤說明")
retry_count: int = 0
第二步:更新 state.py,把 task_queue 與 completed_tasks 換成型別明確的 TaskSpec 與 HandoffPayload 串列,取代昨天分散在多個欄位的做法:
# research-agent/src/research_agent/agents/state.py(更新版)
from typing import TypedDict, Optional
from typing_extensions import Annotated
from langgraph.graph.message import add_messages
from research_agent.agents.protocol import TaskSpec, HandoffPayload
class TeamState(TypedDict):
"""多代理共用狀態:以任務佇列與交接封包取代散裝欄位。"""
messages: Annotated[list, add_messages]
topic: str
task_queue: list[TaskSpec]
completed_tasks: list[HandoffPayload]
final_report: Optional[str]
step_count: int
第三步:改寫 supervisor,讓它先把研究主題拆成任務清單,之後每一輪只需要從 task_queue 拿出下一筆待處理任務指派出去,判斷完成的依據也改成檢查 completed_tasks 裡是否已經有對應的成功回報:
# research-agent/src/research_agent/agents/supervisor.py(更新版)
from research_agent.agents.state import TeamState
from research_agent.agents.protocol import TaskSpec, TaskStatus
def plan_tasks(topic: str) -> list[TaskSpec]:
"""把研究主題拆成固定的三步任務清單;日後可換成模型動態拆解。"""
return [
TaskSpec(task_id="t1", description=f"搜尋「{topic}」的背景資料", assigned_to="search_worker"),
TaskSpec(task_id="t2", description="把搜尋結果整理成可引用片段", assigned_to="retrieval_worker"),
TaskSpec(task_id="t3", description="根據可引用片段撰寫報告", assigned_to="writer_worker"),
]
def supervisor_node(state: TeamState) -> dict:
"""每一輪從 task_queue 挑出下一筆待處理任務,指派給對應的 worker。"""
queue = state.get("task_queue") or plan_tasks(state["topic"])
pending = [t for t in queue if t.status == TaskStatus.PENDING]
if not pending:
return {"task_queue": queue, "step_count": state.get("step_count", 0) + 1}
next_task = pending[0]
next_task.status = TaskStatus.IN_PROGRESS
return {
"task_queue": queue,
"step_count": state.get("step_count", 0) + 1,
"messages": [
{"role": "system", "content": f"supervisor 指派任務 {next_task.task_id} 給 {next_task.assigned_to}"}
],
}
def route_to_worker(state: TeamState) -> str:
"""依 task_queue 目前狀態決定走哪一條邊,全部完成則收尾。"""
queue = state.get("task_queue", [])
for task in queue:
if task.status == TaskStatus.IN_PROGRESS:
return task.assigned_to
return "FINISH"
第四步:改寫 worker,讓每個 worker 收到任務後,一律回傳一份 HandoffPayload,不管成功或失敗都走同一種格式,並在失敗時明確標示 error 而不是安靜地回傳空結果:
# research-agent/src/research_agent/agents/workers.py(改用 HandoffPayload)
from research_agent.config import load_settings
from research_agent.agents.state import TeamState
from research_agent.agents.protocol import TaskStatus, HandoffPayload
def _find_in_progress(state: TeamState, worker_name: str):
for task in state.get("task_queue", []):
if task.assigned_to == worker_name and task.status == TaskStatus.IN_PROGRESS:
return task
return None
def search_worker(state: TeamState) -> dict:
"""負責對外搜尋,失敗時回傳 status=failed 而不是安靜地給空結果。"""
task = _find_in_progress(state, "search_worker")
if task is None:
return {}
settings = load_settings()
try:
if settings.is_dry_run:
data = [{"title": f"{state['topic']} 示範資料", "url": "https://example.com/demo", "snippet": "示範搜尋結果,非真實網路資料。"}]
else:
from research_agent.tools import web_search # AG Day 24 建立的搜尋工具
data = web_search(state["topic"], max_results=5)
payload = HandoffPayload(task_id=task.task_id, worker_name="search_worker", status=TaskStatus.DONE,
summary=f"取得 {len(data)} 筆搜尋資料", result=data)
task.status = TaskStatus.DONE
except Exception as exc:
payload = HandoffPayload(task_id=task.task_id, worker_name="search_worker", status=TaskStatus.FAILED,
summary="搜尋失敗", error=str(exc))
task.status = TaskStatus.FAILED
return {
"task_queue": state["task_queue"],
"completed_tasks": state.get("completed_tasks", []) + [payload],
"messages": [{"role": "system", "content": f"search_worker 回報:{payload.status.value}"}],
}
第五步:writer_worker 收到任務時,不再假設前面一定順利,而是先確認 completed_tasks 裡對應的搜尋與檢索任務是不是都成功,只要有一個失敗,就用失敗訊息組出一份「部分完成」的報告並明確標示,而不是假裝一切正常:
# research-agent/src/research_agent/agents/workers.py(writer_worker:檢查前置任務狀態)
def _lookup(state: TeamState, task_id: str) -> HandoffPayload | None:
for payload in state.get("completed_tasks", []):
if payload.task_id == task_id:
return payload
return None
def writer_worker(state: TeamState) -> dict:
"""寫報告前先檢查前置任務是否全部成功,失敗時明確標示部分完成。"""
task = _find_in_progress(state, "writer_worker")
if task is None:
return {}
search_payload = _lookup(state, "t1")
retrieval_payload = _lookup(state, "t2")
failed = [p for p in (search_payload, retrieval_payload) if p and p.status == TaskStatus.FAILED]
if failed:
report = "# 報告(部分完成)\n\n" + "\n".join(f"- {p.worker_name} 失敗:{p.error}" for p in failed)
status = TaskStatus.FAILED
else:
chunks = retrieval_payload.result if retrieval_payload else []
lines = [f"# {state['topic']}(示範報告,離線模擬模式)", ""]
lines += [f"{i}. {c['text']}[來源:{c['source_title']}]" for i, c in enumerate(chunks, start=1)]
report = "\n".join(lines)
status = TaskStatus.DONE
payload = HandoffPayload(task_id=task.task_id, worker_name="writer_worker", status=status, summary="報告已產生", result=report)
task.status = status
return {
"task_queue": state["task_queue"],
"completed_tasks": state.get("completed_tasks", []) + [payload],
"final_report": report,
"messages": [{"role": "system", "content": f"writer_worker 回報:{status.value}"}],
}
這樣一來,即使搜尋階段真的遇到逾時或速率限制而失敗,使用者看到的也是一份誠實標示「部分完成、哪一步失敗、原因是什麼」的報告,而不是一份看起來完整卻其實資料短缺的報告,這正是結構化交接封包帶來的直接好處:失敗資訊有明確的管道可以一路傳到最後產出的內容裡。
第六步:寫一個小型的驗證腳本,模擬 supervisor 派工、worker 回報成功與失敗兩種情境,確認協議本身在沒有真正跑完整張圖的情況下也能被獨立測試:
# research-agent/scripts/check_protocol.py
from research_agent.agents.protocol import TaskSpec, TaskStatus, HandoffPayload
from research_agent.agents.supervisor import plan_tasks
def main():
tasks = plan_tasks("多代理通訊協議")
print("=== 初始任務清單 ===")
for t in tasks:
print(f"{t.task_id}: {t.description}(指派給 {t.assigned_to},狀態 {t.status.value})")
success = HandoffPayload(task_id="t1", worker_name="search_worker", status=TaskStatus.DONE,
summary="取得 2 筆搜尋資料", result=[{"title": "示範資料"}])
failure = HandoffPayload(task_id="t1", worker_name="search_worker", status=TaskStatus.FAILED,
summary="搜尋逾時", error="httpx.TimeoutException: 讀取逾時")
print("\n=== 成功回報範例 ===")
print(success.model_dump())
print("\n=== 失敗回報範例 ===")
print(failure.model_dump())
if __name__ == "__main__":
main()
執行驗證:
cd research-agent
uv run python scripts/check_protocol.py
示範輸出:
=== 初始任務清單 ===
t1: 搜尋「多代理通訊協議」的背景資料(指派給 search_worker,狀態 pending)
t2: 把搜尋結果整理成可引用片段(指派給 retrieval_worker,狀態 pending)
t3: 根據可引用片段撰寫報告(指派給 writer_worker,狀態 pending)
=== 成功回報範例 ===
{'task_id': 't1', 'worker_name': 'search_worker', 'status': <TaskStatus.DONE: 'done'>, 'summary': '取得 2 筆搜尋資料', 'result': [{'title': '示範資料'}], 'error': None, 'retry_count': 0}
=== 失敗回報範例 ===
{'task_id': 't1', 'worker_name': 'search_worker', 'status': <TaskStatus.FAILED: 'failed'>, 'summary': '搜尋逾時', 'result': None, 'error': 'httpx.TimeoutException: 讀取逾時', 'retry_count': 0}
有了這套協議,supervisor 收到 status=failed 的回報時,就能決定要重試、換另一種搜尋方式,或直接把錯誤寫進最終報告告知使用者,而不是像昨天那樣讓空結果悄悄流到下一個 worker。
常見錯誤與踩雷
第一個常見錯誤是協議定得太細,反而讓每個 worker 都要花很多力氣去符合格式,得不償失。HandoffPayload 的 result 欄位刻意用寬鬆的 Any 型別,是因為不同 worker 的實際產出形狀本來就不同(搜尋結果是串列、報告是字串),協議要統一的是「外層的信封」(狀態、摘要、錯誤),不需要連信封裡的內容也統一格式,過度規範反而會讓協議變得難用。
第二個常見錯誤是把 TaskSpec 物件直接放進 LangGraph 的狀態卻忘了它是可變的 Pydantic 物件,多個節點如果各自拿到同一份清單的參照並直接修改,容易在並行執行時互相干擾。今天的範例都是序列執行、單一 worker 在跑,所以直接修改 task.status 沒有問題;但如果之後要讓多個 worker 平行處理不同任務,就必須改成「回傳一份新的任務清單」而不是原地修改,才符合 LangGraph 狀態更新以回傳值為準的設計原則。
第三個常見錯誤是失敗處理只做半套:只在 worker 內部捕捉例外並記錄,卻沒有讓 supervisor 真正檢查 completed_tasks 裡的失敗狀態。今天 route_to_worker 只看 IN_PROGRESS,還沒有處理「上一個任務失敗了該不該重試」的邏輯,這是刻意留白,重試策略會在正式串接真實 API 呼叫、真正遇到逾時與速率限制時再一併設計,跟 AG Day 7 學過的重試機制銜接。
效能與實務提醒
結構化協議帶來的成本是多了一層 Pydantic 驗證與序列化開銷,但這個開銷相對於一次模型呼叫的延遲幾乎可以忽略,換來的除錯效率提升非常值得。實務上建議把 HandoffPayload 直接寫進 AG Day 9 建立的觀測紀錄裡,因為它本身就是結構化資料,可以直接丟進資料庫或記錄檔,不需要額外轉換格式。
另外,任務清單的拆解方式(今天用 plan_tasks 寫死三步)在真實專案裡通常會隨主題複雜度變化:簡單問題可能只需要一步搜尋,複雜問題可能要拆成五六個子任務。建議先讓拆解邏輯保持簡單、可預測,等到後面章節的評估機制建立起來,再考慮讓模型動態決定要拆成幾步,避免一開始就把「動態拆解」跟「協議設計」這兩件事混在一起除錯。
最後,retry_count 這個欄位今天只是先放著,還沒有實際邏輯使用它,但這是刻意的設計:先把欄位留在協議裡,之後要加重試邏輯時就不需要再改一次資料結構,這是設計訊息協議時值得養成的習慣——預留好欄位,比事後在到處都要改的地方補欄位輕鬆得多。同樣的道理也適用於 summary 欄位:即使今天的範例只是簡單拼字串,把「給人看的摘要」與「給程式用的結構化資料」分開存放,日後接上 AG Day 35 的追蹤平台時,這個欄位可以直接當作追蹤紀錄裡的標題,不需要另外重新產生一份人類可讀的說明。
小結
今天我們把多代理系統的通訊方式從「共用大狀態、靠命名默契」升級成「有明確契約的訊息協議」:TaskSpec 描述要做的一件事,HandoffPayload 描述做完之後的統一回報格式,涵蓋成功、失敗兩種情境。supervisor 從臨場判斷改成先規劃任務清單、逐一指派,worker 一律回傳結構化的交接封包,不再讓失敗被空結果掩蓋。
新增的術語:交接封包(handoff payload,代理之間回報結果的統一格式)、任務佇列(task queue,待處理任務的清單)、去中心化交接(decentralized handoff,worker 直接指定下一步而不必繞回 supervisor)。
結語
有了任務清單與交接封包這兩個結構化狀態,我們已經具備把任務執行進度持久化的基礎——任務清單本身就可以序列化存檔、可以在下次啟動時重新讀取。但目前整條流程還是「跑完就結束」,如果程式在跑到一半時當掉,所有進度都會消失,必須從頭再來一次。
明天,我們會進入「AG Day 31 長任務:檢查點、恢復與背景執行」,把今天的任務佇列接上 LangGraph 的檢查點機制,讓多代理流程可以中途存檔、意外中斷後從斷點恢復,也會談談怎麼把長時間執行的研究任務丟到背景執行,不必讓使用者一直守著終端機等結果。
延伸資源
- LangGraph 官方文件的
Command與 handoff 章節:https://langchain-ai.github.io/langgraph/。去中心化交接的具體寫法與參數請以官方文件當次版本為準。 - Pydantic 官方文件的
Enum欄位與model_dump()用法:https://docs.pydantic.dev/。今天的TaskStatus與序列化範例以此為準。 - Python 官方文件的
enum.Enum:說明字串列舉(str, Enum)的行為與比較規則。
留言
張貼留言