包裹分揀任務(左傳送帶 → 條碼朝上 → 右傳送帶)
文件版本:v1.0(2026-05-12) 對應源碼基線:
fmc3-robotics-main/projects/RoboOS(stand-alone) 你的職責範圍(蒲):任務規劃 / 任務調度 / 任務監控 / 任務診斷 / 人機交互(部分) 隊友:尹(任務開關)、朱(Skill Unit / 主策略 / 後處理 / 策略診斷)、李(技能編排 / 技能調度 / 技能監控 / 人機協作)、謝(前處理 / 結構化技能感知 / 視覺定位)
目錄
- 1. 任務理解(What & Why)
- 2. 系統總體架構(分層智能)
- 3. Task Manager 模塊地圖
- 4. 從 RoboOS 沿用什麼、改動什麼
- 5. 詳細模塊設計
- 6. 視覺接入:GroundingDINO 與任務監控閉環
- 7. 數據流與接口契約
- 8. 推薦目錄結構
- 9. 開發 Roadmap(建議 4 個 sprint)
- 10. 核心程式碼骨架
- 11. 常見問題(FAQ)
1. 任務理解(What & Why)
1.1 業務目標
機器人佇立在兩條傳送帶之間(左:來料;右:出貨),需要持續處理包裹:
| 步驟 | 動作 | 關鍵約束 |
|---|---|---|
| ① 拿取 | 從左傳送帶抓取包裹 | 包裹姿態、抓取點選擇、防滑落 |
| ② 翻面 / 擺正 | 將條碼面朝上 | 視覺判斷條碼面、可能需翻轉或旋轉 |
| ③ 推送 | 把包裹推到右傳送帶 | 平穩、不掉落、放置區域對齊 |
定義對齊(依據《功能逻辑架构及模块分工202605v1.0.pdf》): - Task:一次「拿放一個包裹」=完成一個 Task,連續拿包裹=多個 Task。 - Skill:完成一個 Task 需要的「原子操作」,例如「視覺定位」「左手抓取」「左手旋轉包裹」「拨移包裹」。 - Policy:執行某個 Skill 時採用的運動規劃 / 控制算法(VLA、VLA+IK、VLA+視覺伺服、…)。
1.2 為什麼需要「Task Manager」這一層?
端到端的 VLA 模型(直接從 pixel 推到 torque)在長時序、多包裹、需錯誤恢復的情境下表現脆弱。因此我們把系統做成「分層智能」:
┌───────────────────────────────────────────────────────────────┐
│ Task Manager(任務級) │
│ - 自主拆解:「下一個抓哪個包裹?順序?」 │
│ - 監控、診斷、重規劃、人機交互 │
└──────────────┬────────────────────────────────────────────────┘
│ 子任務指令 (subtask)
▼
┌───────────────────────────────────────────────────────────────┐
│ Skill Manager(技能級) │
│ - 編排:「用哪隻手?分解成哪幾個技能?」 │
│ - 從技能庫匹配 API、行為樹調度、越界/碰撞檢測 │
└──────────────┬────────────────────────────────────────────────┘
│ 技能 API (skill primitive)
▼
┌───────────────────────────────────────────────────────────────┐
│ Skill Unit(執行級 / 策略級) │
│ - 前處理 → 主策略推理 → 後處理 (IK / 平滑 / 限位) │
│ - 策略監控(力 / 位 / 視覺混合判斷)、冗餘策略切換 │
└───────────────────────────────────────────────────────────────┘
你(蒲)的工作主要在最上層:把人類給的高層指令(例如「持續處理左傳送帶上的包裹直到清空」)拆成可被 Skill Manager 執行的子任務、實時監控執行狀態、發生異常時觸發重規劃。
1.3 具體要對接的「外部信號」
| 來源 | 提供什麼 | 提供者 |
|---|---|---|
| 視覺 Grounding(GroundingDINO + 後處理) | 包裹 bbox、條碼面是否朝上、包裹在傳送帶上的位置區域 | 謝 |
| Skill Manager | 技能執行狀態、技能故障碼、技能可行性回饋 | 李 |
| Skill Unit(策略診斷) | 動作異常、本體故障、傳感器故障 | 朱 |
| 任務開關程序 | 系統初始化完成、任務結束流程觸發 | 尹 |
| User pad UI | 啟動 / 停止 / 暫停 / 復位、模式切換(遙控 / 自主) | 蒲(你) |
2. 系統總體架構(分層智能)
2.1 三個 Harness 的協作
flowchart LR
UI[User pad UI<br/>自然語言 + 按鈕] --> TM
Cam[頭部 / 腕部攝像頭] --> Algo
Algo[Algo-Ground<br/>VLM / GroundingDINO / 深度 / 3D 定位] --> TM
Algo --> SM
Algo --> SU
TM[Task Manager<br/>蒲] -->|子任務| SM
SM[Skill Manager<br/>李 / 朱] -->|技能 API| SU
SU[Skill Unit / Policy Set<br/>朱 / 謝] -->|關節 / 末端指令| MC
SU -.故障碼.-> SM
SM -.故障碼.-> TM
TM -.重規劃 / 復位.-> SM
MC[Motion Control<br/>末端位姿 / 關節角] --> Robot[本體 + 執行器]
2.2 三類「上下文記憶」
每一層都維護自己的上下文,這對於重規劃與故障定位至關重要:
- 任務級上下文(你的):任務提示詞、任務日誌、任務關鍵幀、標定數據、任務約束、質量要求、安全要求、故障報警、任務統計、上下文壓縮。
- 技能級上下文(李):技能日誌、技能關鍵幀、技能約束、操作要求…
- 策略級上下文(朱):策略關鍵幀、策略約束、故障報警…
RoboOS 已經提供了 Collaborator(Redis)作為共享記憶基礎設施,可以直接複用作為上下文存儲後端,詳見第 4 節。
3. Task Manager 模塊地圖
PDF 中明確列出 6 個子模塊。下表加上負責人、輸入、輸出、首版完成標準,方便你排期:
| # | 模塊 | 負責人 | 輸入 | 輸出 | MVP 完成標準 |
|---|---|---|---|---|---|
| 1 | 人機交互程序 | 蒲 | 多模態指令(語音 / 文本 / 按鈕) | 結構化任務、模式切換、回顯 | 能接收 "start / stop / pause / reset" 並回傳當前子任務名稱 |
| 2 | 任務開關程序 | 尹 | 啟動 / 結束信號 | 模型加載、參數加載、任務報告 | 提供 startup() / shutdown() 給你調用 |
| 3 | 任務規劃程序 | 蒲 | 高層指令 + 場景信息 | 子任務隊列(長度 ≤ 3) | 給定「清空左傳送帶」可輸出有效子任務序列 |
| 4 | 任務調度程序 | 蒲 | 子任務隊列 + 故障碼 | 子任務派發、模式切換、復位 | 按隊列順序派發,支持中斷、復位 |
| 5 | 任務監控程序 | 蒲 | 視覺 Grounding 結果、技能級狀態 | 任務進度、跳轉條件、成功 / 失敗 | 能識別「包裹條碼是否朝上」、「是否到達放置區」 |
| 6 | 任務診斷程序 | 蒲 | 規劃失敗信號、軟體故障 | 任務故障碼、寫入上下文 | 列出至少 5 個故障碼,並能寫入上下文供重規劃讀取 |
設計原則:6 個模塊跑在同一個 Python process 內(不要拆 micro-services,會增加維護成本),透過內部
asyncioevent loop + 一個共享的TaskContext對象通信;對外(給 Skill Manager / UI)走 Redis pub/sub,沿用 RoboOS 既有的Collaborator。
4. 從 RoboOS 沿用什麼、改動什麼
4.1 可以直接沿用的
| RoboOS 組件 | 路徑 | 對應到新架構的 | 沿用理由 |
|---|---|---|---|
Collaborator(Redis pub/sub + agent 註冊) |
flag_scale.flagscale.agent.collaboration |
任務級 / 技能級上下文記憶的後端、跨層消息匯流排 | Pub/Sub + Hash 已足夠,不必重造 |
GlobalAgent(master agent.py) |
master/agents/agent.py |
任務調度程序(部分)+ 任務級事件監聽 | 已實現 robot 註冊、subtask 派發、async dispatch |
GlobalTaskPlanner |
master/agents/planner.py |
任務規劃程序(核心)+ 提示詞模板 | 已封裝 LLM 連接 + 重試機制 |
Flask /publish_task /system_status /robot_status |
master/run.py |
人機交互程序(HTTP API 後端) | 直接擴展即可 |
MASTER_PLANNING_PLANNING 提示詞 |
master/agents/prompts.py |
任務規劃的 prompt | 改寫成包裹分揀領域即可 |
Slaver 的 ToolCallingAgent + MCP tool 機制 |
slaver/agents/slaver_agent.py、slaver/demo_robot_local/skill.py |
Skill Unit 的調用接口 | Skill Manager 對接 MCP 工具即可 |
SceneMemory(動作對場景的副作用模型) |
slaver/tools/memory.py |
任務監控的「狀態跳轉判定」雛形 | 已支援 add/remove/move 三類動作 |
4.2 需要改動 / 新增的
| 模塊 | 改動類型 | 說明 |
|---|---|---|
| 任務隊列管理 | 新增 | 現在 GlobalAgent 一次處理一個 task;包裹分揀需要 隊列長度=3 的滑動隊列(出隊、入隊、重排、插隊、剔除)。 |
| 任務監控閉環 | 新增 | 現在 master 只在收到 slaver 回傳時被動更新;新架構需要主動訂閱視覺結果,實時判斷任務進度。 |
| 任務故障碼 + 重規劃 | 新增 | 沿用 _handle_result 但加入故障碼解析與「重規劃 vs 重試 vs 中止」決策。 |
| 場景 profile | 改寫 | 把 scene/profile.yaml 從廚房場景改為傳送帶場景(左傳送帶、右傳送帶、抓取區、條碼面狀態列舉)。 |
| 任務規劃 prompt | 改寫 | 領域改成包裹分揀,輸出加上「包裹優先級」「條碼姿態目標」等欄位。 |
| 人機交互 | 改寫 | UI 需要顯示:當前子任務、包裹隊列、條碼朝向、進度條、故障碼。 |
| 任務開關程序 | 新增(尹) | 系統初始化、模型熱加載、任務報告生成。 |
重點:先不要重寫,而是用「裝飾/擴展」的方式包住
GlobalAgent,新增TaskQueueManager、TaskMonitor、TaskDiagnostics三個 mixin,再讓GlobalAgent持有它們。
5. 詳細模塊設計
以下所有模塊都假設單一 Python process,共享同一個
asyncioevent loop 與一個TaskContext對象。對外通信走 Redis(Collaborator)。
5.1 任務規劃程序(Task Planner)
職責:把高層自然語言指令(例如「持續處理左傳送帶上的包裹直到清空」)拆解為有優先級、有依賴關係的子任務隊列。
5.1.1 輸入
task_text: str:人類指令scene_snapshot: dict:當前場景 snapshot(從 GroundingDINO 拿到的包裹列表、條碼朝向)robot_caps: dict:已註冊機器人的能力(沿用collaborator.read_all_agents_info())
5.1.2 輸出
{
"reasoning_explanation": "left conveyor 上偵測到 3 個包裹,依據離夾爪距離排序…",
"subtask_list": [
{"subtask_id": "t1", "subtask": "pick package P1 from left_conveyor", "subtask_order": 0, "robot_name": "humanoid_01", "priority": 0.92},
{"subtask_id": "t2", "subtask": "rotate P1 so barcode_face is up", "subtask_order": 1, "robot_name": "humanoid_01", "priority": 0.92, "depends_on": "t1"},
{"subtask_id": "t3", "subtask": "push P1 onto right_conveyor", "subtask_order": 2, "robot_name": "humanoid_01", "priority": 0.92, "depends_on": "t2"}
]
}
5.1.3 實作要點
- 沿用
GlobalTaskPlanner.forward(),但替換MASTER_PLANNING_PLANNINGprompt。 - 注入「場景 snapshot」到 prompt:在呼叫 LLM 前,先從 Redis 撈
current_scenekey。 - 長度約束:隊列長度 ≤ 3。如果規劃出更多,截斷後存到
pending_buffer。 - 可行性校驗:每個 subtask 都要包含
robot_name,並驗證該 robot 已註冊(沿用reasoning_and_subtasks_is_right)。
5.1.4 新版 Prompt(範本,可放到 master/agents/prompts.py)
PACKAGE_SORTING_PLANNING = """
你是一台工業分揀人形機器人的任務規劃器。
場景描述:機器人面前有兩條傳送帶,左邊是來料、右邊是出貨。
任務目標:把左傳送帶上的所有包裹依序搬到右傳送帶,**並確保條碼面朝上**。
## 當前場景快照
{scene_snapshot}
## 可用機器人與技能
{robot_tools_info}
## 約束
1. 一個 Task = 拿放一個包裹(pick → reorient → push)。
2. 同時最多規劃 3 個 Task;其餘留在 pending buffer。
3. 若包裹的 barcode_face 已經朝上,跳過 reorient 步驟。
4. 包裹優先級 = 1 / (距離夾爪的距離) × 視覺置信度。
## 輸出 JSON
{{
"reasoning_explanation": "...",
"subtask_list": [
{{"subtask_id": "tX", "subtask": "...", "subtask_order": N, "robot_name": "...", "priority": F, "depends_on": "tY"}}
]
}}
## 人類指令
{task}
"""
5.2 任務調度程序(Task Scheduler)
職責:按隊列順序把子任務派發給 Skill Manager;處理中斷、復位、模式切換。
5.2.1 與 RoboOS 既有 _dispath_subtasks_async 的差異
| 既有 | 新需求 |
|---|---|
| 一次性把所有 subtasks 派完 | 隊列式調度,可中途插入 / 剔除 |
| 沒有故障碼處理 | 收到故障碼 → 暫停、重規劃或中止 |
| 沒有模式切換 | 支持「遙控」「自主」「半自主」三種模式 |
5.2.2 狀態機
┌────────┐
│ IDLE │
└───┬────┘
start_task │
▼
┌────────┐
┌──────► │ RUNNING│ ◄────┐
│ └───┬────┘ │
resume│ │ pause │ resume
│ ▼ │
│ ┌────────┐ │
│ │ PAUSED │ ─────┘
│ └───┬────┘
│ │ fault
│ ▼
│ ┌────────┐
└──────── │ FAULT │
└───┬────┘
│ reset
▼
┌────────┐
│ RESET │
└────────┘
5.2.3 接口(給其他模塊調用)
class TaskScheduler:
def dispatch_next(self) -> Optional[str]: ...
def pause(self) -> None: ...
def resume(self) -> None: ...
def abort(self, reason: str) -> None: ...
def reset(self) -> None: ...
def switch_mode(self, mode: Literal["teleop", "auto", "semi"]) -> None: ...
5.3 任務監控程序(Task Monitor)— 本系統最核心的模塊
職責:實時整合視覺結果與技能級回傳,回答以下問題:
- 包裹類型是什麼?尺寸多大?
- 該選哪一個包裹?
- 當前包裹的位姿 / 區域?
- 條碼面朝上了嗎?
- 包裹放置位置對嗎?
- 子任務是否完成?是否成功?
5.3.1 設計成「事件驅動的狀態機」
每個運行中的 Task 都是一個 TaskRuntime 物件,內部有一個小型 FSM:
PICK ──ok──► REORIENT ──ok──► PUSH ──ok──► DONE
│ │ │
└fail─────────┴────────────────┴─► FAULT
每次收到視覺事件或技能完成事件,跑一次 evaluate(),決定是否跳轉。
5.3.2 監控數據來源
| 數據 | 來源 | 頻率 |
|---|---|---|
| 包裹 bbox + class | GroundingDINO(謝) | ~10 Hz |
| 條碼面朝向(up / down / side) | 後處理(謝) | ~10 Hz |
| 抓取狀態(grasped / not) | Skill Unit 回傳 | event |
| 放置成功 | Skill Unit 回傳 | event |
| 技能故障碼 | Skill Manager | event |
5.3.3 跳轉條件範例
def evaluate(self, evt: VisionEvent | SkillEvent) -> Transition | None:
if self.state == State.PICK and evt.type == "skill_done" and evt.name == "G1_left_grasp":
return Transition(to=State.REORIENT)
if self.state == State.REORIENT and evt.type == "vision":
if evt.barcode_face == "up":
return Transition(to=State.PUSH)
elif evt.barcode_face == "down":
return Transition(to=State.REORIENT_RETRY, reason="still down")
if self.state == State.PUSH and evt.type == "vision":
if evt.position_zone == "right_conveyor":
return Transition(to=State.DONE)
...
5.4 任務診斷程序(Task Diagnostics)
職責:把不同類型的故障歸類為任務級故障碼並寫入上下文,供調度程序、規劃程序作為下一步決策依據。
5.4.1 故障碼分類(首版至少這 7 個)
| 故障碼 | 名稱 | 觸發條件 | 默認處理 |
|---|---|---|---|
TM-001 |
PLAN_FAIL | LLM 連續 5 次無法輸出合法 JSON | 中止 + 告警 |
TM-002 |
SCENE_EMPTY | 視覺連續 N 幀無包裹 | 任務完成 |
TM-003 |
GRASP_TIMEOUT | PICK 階段超時 | 重試一次 |
TM-004 |
REORIENT_FAIL | REORIENT 重試 ≥ 3 次仍非朝上 | 重規劃(換手 / 換策略) |
TM-005 |
PLACE_OUT_OF_ZONE | PUSH 後不在 right_conveyor | 重試或人工介入 |
TM-006 |
SKILL_FAULT | 收到 Skill Manager 故障碼 | 透傳 + 重規劃 |
TM-007 |
SOFTWARE_FAULT | 內部 exception / 超時 | 記日誌 + 暫停 |
5.4.2 上下文寫入
context.faults.append({
"code": "TM-004",
"ts": time.time(),
"subtask_id": "t2",
"details": {"retry": 3, "last_barcode_face": "side"},
})
collaborator.record_environment("task_context", json.dumps(context.dict()))
5.5 人機交互程序(HMI)
職責:作為 Task Manager 對使用者的唯一接觸點。
5.5.1 對外 HTTP API(沿用 Flask)
| 端點 | 方法 | 用途 |
|---|---|---|
/publish_task |
POST | 接收自然語言任務(已存在,擴展即可) |
/control/{start,stop,pause,resume,reset} |
POST | 控制指令 |
/mode |
POST | 模式切換 |
/task_status |
GET | 當前隊列 + 進度 + 子任務名稱 |
/system_status |
GET | CPU/RAM + 機器人狀態(已存在) |
/faults |
GET | 當前故障碼列表 |
/logs/stream |
WebSocket | 實時日誌 |
5.5.2 對外 WebSocket 推送(沿用 flask-socketio)
socketio.emit("task_progress", {
"task_id": "...", "subtask": "REORIENT P1",
"progress": 0.66, "barcode_face": "down",
})
5.6 任務級上下文記憶
用 Redis hash 存儲,鍵設計如下:
| Redis Key | 內容 |
|---|---|
task:{task_id}:meta |
task_text, created_at, status |
task:{task_id}:queue |
子任務 JSON 列表 |
task:{task_id}:faults |
故障碼歷史 |
task:{task_id}:keyframes |
視覺關鍵幀(圖像路徑 + 標註) |
task:current_scene |
最新的場景快照 |
task:stats |
累計成功 / 失敗 / 處理時間 |
利用 RoboOS 既有的
Collaborator.record_environment()/read_environment()直接存 JSON 字串即可。
6. 視覺接入:GroundingDINO 與任務監控閉環
6.1 你(Task Manager)不負責訓 / 跑 GroundingDINO
訓練、推理、後處理(含條碼朝向判斷)由謝完成。你只負責訂閱結果。
6.2 約定的視覺事件 schema
{
"event": "vision_update",
"ts": 1747040000.123,
"camera": "head",
"frame_id": "img_000123",
"objects": [
{
"id": "P1",
"class": "package_box",
"bbox": [x1, y1, x2, y2],
"score": 0.93,
"pose_3d": {"position": [0.42, 0.10, 0.85], "yaw": 0.3},
"barcode_face": "down",
"barcode_score": 0.81,
"zone": "left_conveyor"
}
]
}
6.3 視覺與監控的數據流
GroundingDINO (謝) ── publish ──► Redis channel "vision_events"
│
▼
TaskMonitor.on_vision()
│
┌─────────────┼─────────────┐
▼ ▼ ▼
scene_snapshot TaskRuntime.evaluate()
update (FSM 跳轉)
│
▼
發布 task_progress / 故障碼
6.4 如何在 Task Manager 接收
沿用 Collaborator.listen():
threading.Thread(
target=lambda: self.collaborator.listen("vision_events", self._on_vision_event),
daemon=True,
).start()
6.5 「條碼面朝上」如何判斷(給謝參考)
- 對偵測到的
package_box做二級 grounding:用 prompt"barcode on the box top face"拿條碼 bbox。 - 用深度 / 表面法線(謝負責)算條碼面與相機 z 軸夾角:
- 夾角 < 30° →
up - 30°–60° →
side -
60° →
down - 把結果塞進
barcode_face欄位。
重點:這部分不是你寫,但你要在介面文件(你維護的
interfaces.md)裡把 schema 寫死,這樣謝才知道怎麼往 Redis 推。
7. 數據流與接口契約
7.1 Redis Pub/Sub channels(你維護的合約)
| Channel | 方向 | Schema |
|---|---|---|
AGENT_REGISTRATION |
Slaver → Master | (沿用 RoboOS) |
roboos_to_{robot_name} |
Master → Slaver | (沿用 RoboOS) |
{robot_name}_to_RoboOS |
Slaver → Master | (沿用 RoboOS) |
vision_events |
Algo-Ground → TaskMonitor | 見 §6.2 |
skill_events |
Skill Manager → TaskMonitor | {event, ts, skill_name, status, fault_code?} |
task_progress |
TaskMonitor → UI | {task_id, subtask, progress, state} |
task_faults |
TaskDiagnostics → UI | {code, ts, details} |
7.2 與 Skill Manager(李)的接口
兩種風格擇一即可,建議首版用 (A) 直接沿用 RoboOS:
- (A) 完全沿用:Task Scheduler 把子任務字串
"pick package P1 from left_conveyor"透過roboos_to_{robot_name}發給 slaver,slaver 內部的 ToolCallingAgent + LLM 自行選擇 MCP tool。 - (B) 結構化指令:發送結構化 JSON
{action: "pick", target: "P1", from: "left_conveyor"},由 Skill Manager 做技能組裝。第二期再演進到這個。
7.3 場景 profile 改寫(你負責 master/scene/profile.yaml)
scene:
- name: left_conveyor
type: conveyor
role: source
contains: [] # 動態,由 GroundingDINO 填充
- name: right_conveyor
type: conveyor
role: sink
contains: []
- name: workspace
type: zone
bounds: [[0.2, -0.4, 0.7], [0.8, 0.4, 1.2]] # x_min,y_min,z_min .. x_max,y_max,z_max
properties:
task_queue_max_length: 3
barcode_face_required: up
vision_event_channel: vision_events
8. 推薦目錄結構
在
master/底下擴展,不要新開一個新 repo(會跟 slaver 的依賴脫鉤)。
projects/RoboOS/master/
├── run.py # Flask 入口(沿用 + 擴充端點)
├── config.yaml # 沿用 + 加 task_manager 區段
├── scene/
│ └── profile.yaml # 改寫為傳送帶場景
├── agents/
│ ├── agent.py # 沿用 GlobalAgent(最小改動)
│ ├── planner.py # 沿用 GlobalTaskPlanner
│ └── prompts.py # 加入 PACKAGE_SORTING_PLANNING
└── task_manager/ # ★ 你新增的目錄
├── __init__.py
├── context.py # TaskContext 對象(持久化到 Redis)
├── queue.py # TaskQueueManager(長度 3 滑窗)
├── scheduler.py # TaskScheduler(狀態機 + 派發)
├── monitor.py # TaskMonitor(FSM + 視覺訂閱)
├── diagnostics.py # TaskDiagnostics(故障碼)
├── hmi.py # HMI(Flask blueprint)
├── runtime.py # TaskRuntime(單一任務的 FSM 實例)
├── events.py # dataclass: VisionEvent / SkillEvent
└── interfaces.md # ★ 寫給隊友看的合約文件
9. 開發 Roadmap(建議 4 個 sprint)
Sprint 0:架構落地(本週)
- [ ] 寫完本文件 +
interfaces.md,與隊友對齊接口(最重要!) - [ ] 建立
task_manager/目錄與空骨架 - [ ] 改寫
scene/profile.yaml - [ ] 改寫
PACKAGE_SORTING_PLANNINGprompt - [ ] 跑通既有 RoboOS(master + 一個假 slaver)做 smoke test
Sprint 1:規劃 + 調度 MVP(第二週)
- [ ] 實作
TaskQueueManager、TaskScheduler狀態機 - [ ] 把
GlobalAgent.publish_global_task包進 scheduler - [ ] HMI 增加
/control/*端點 - [ ] 用「假視覺事件」(手動 redis-cli publish)驗證端到端
Sprint 2:監控閉環(第三週)
- [ ] 實作
TaskMonitor+TaskRuntimeFSM - [ ] 訂閱
vision_events、skill_events - [ ] 與謝對接真實 GroundingDINO 輸出
- [ ] HMI 顯示進度條 + barcode_face
Sprint 3:診斷 + 重規劃(第四週)
- [ ] 實作 7 個故障碼
- [ ] 重規劃路徑(PLAN_FAIL / REORIENT_FAIL)
- [ ] 任務級上下文壓縮(避免 Redis 撐爆)
- [ ] 集成測試:連續處理 10 個包裹,故意製造故障,看是否能恢復
10. 核心程式碼骨架
僅展示骨架,重點是接口與責任邊界;細節留給實作階段。
10.1 task_manager/context.py
from dataclasses import dataclass, field
from typing import List, Dict, Any, Optional
import json, time
@dataclass
class FaultRecord:
code: str
ts: float
subtask_id: Optional[str]
details: Dict[str, Any] = field(default_factory=dict)
@dataclass
class TaskContext:
task_id: str
task_text: str
created_at: float = field(default_factory=time.time)
state: str = "IDLE"
queue: List[Dict] = field(default_factory=list)
faults: List[FaultRecord] = field(default_factory=list)
stats: Dict[str, Any] = field(default_factory=dict)
def to_redis(self, collaborator) -> None:
collaborator.record_environment(
f"task:{self.task_id}", json.dumps(self.__dict__, default=str)
)
@classmethod
def from_redis(cls, collaborator, task_id: str) -> "TaskContext":
raw = collaborator.read_environment(f"task:{task_id}")
return cls(**json.loads(raw)) if raw else cls(task_id=task_id, task_text="")
10.2 task_manager/queue.py
from collections import deque
from typing import Dict, List, Optional
class TaskQueueManager:
"""滑動窗口長度 = 3 的任務隊列;多餘的塞到 pending_buffer。"""
def __init__(self, max_active: int = 3):
self.active: deque[Dict] = deque(maxlen=max_active)
self.pending: deque[Dict] = deque()
def push(self, subtask: Dict) -> None:
if len(self.active) < self.active.maxlen:
self.active.append(subtask)
else:
self.pending.append(subtask)
def pop_next(self) -> Optional[Dict]:
if not self.active:
return None
nxt = self.active.popleft()
if self.pending:
self.active.append(self.pending.popleft())
return nxt
def remove(self, subtask_id: str) -> bool: ...
def reorder(self, new_order: List[str]) -> None: ...
def insert(self, idx: int, subtask: Dict) -> None: ...
10.3 task_manager/runtime.py(單一 Task 的 FSM)
from enum import Enum
from dataclasses import dataclass
from typing import Optional
class TaskState(Enum):
PICK = "PICK"
REORIENT = "REORIENT"
PUSH = "PUSH"
DONE = "DONE"
FAULT = "FAULT"
@dataclass
class VisionEvent:
barcode_face: str # "up" / "down" / "side"
zone: str # "left_conveyor" / "workspace" / "right_conveyor"
score: float
@dataclass
class SkillEvent:
skill_name: str
status: str # "ok" / "fail" / "timeout"
fault_code: Optional[str] = None
class TaskRuntime:
def __init__(self, task_id: str, package_id: str):
self.task_id = task_id
self.package_id = package_id
self.state = TaskState.PICK
self.retry_count = {s: 0 for s in TaskState}
def evaluate(self, evt) -> Optional[TaskState]:
if isinstance(evt, SkillEvent) and evt.status == "fail":
self.state = TaskState.FAULT
return self.state
if self.state == TaskState.PICK and isinstance(evt, SkillEvent) and evt.skill_name.startswith("G") and evt.status == "ok":
self.state = TaskState.REORIENT
return self.state
if self.state == TaskState.REORIENT and isinstance(evt, VisionEvent):
if evt.barcode_face == "up":
self.state = TaskState.PUSH
return self.state
self.retry_count[TaskState.REORIENT] += 1
if self.retry_count[TaskState.REORIENT] >= 3:
self.state = TaskState.FAULT
return self.state
if self.state == TaskState.PUSH and isinstance(evt, VisionEvent):
if evt.zone == "right_conveyor":
self.state = TaskState.DONE
return self.state
return None
10.4 task_manager/monitor.py
import json, threading, logging
from typing import Dict
from .runtime import TaskRuntime, VisionEvent, SkillEvent
log = logging.getLogger(__name__)
class TaskMonitor:
def __init__(self, collaborator, diagnostics, hmi):
self.collaborator = collaborator
self.diagnostics = diagnostics
self.hmi = hmi
self.runtimes: Dict[str, TaskRuntime] = {}
def attach(self, task_id: str, package_id: str) -> None:
self.runtimes[task_id] = TaskRuntime(task_id, package_id)
def start_listeners(self) -> None:
for ch, handler in [
("vision_events", self._on_vision),
("skill_events", self._on_skill),
]:
threading.Thread(
target=lambda c=ch, h=handler: self.collaborator.listen(c, h),
daemon=True, name=f"mon-{ch}"
).start()
def _on_vision(self, raw: str) -> None:
data = json.loads(raw)
for obj in data.get("objects", []):
evt = VisionEvent(
barcode_face=obj.get("barcode_face", "unknown"),
zone=obj.get("zone", "unknown"),
score=obj.get("score", 0.0),
)
for tid, rt in self.runtimes.items():
if rt.package_id == obj.get("id"):
self._tick(tid, evt)
def _on_skill(self, raw: str) -> None:
data = json.loads(raw)
evt = SkillEvent(**data)
for tid in list(self.runtimes.keys()):
self._tick(tid, evt)
def _tick(self, tid: str, evt) -> None:
rt = self.runtimes[tid]
new_state = rt.evaluate(evt)
if not new_state:
return
self.hmi.push_progress(tid, new_state.value)
if new_state.name == "FAULT":
self.diagnostics.raise_fault(tid, evt)
10.5 task_manager/diagnostics.py
import time, json
from dataclasses import asdict
FAULT_TABLE = {
"TM-001": "PLAN_FAIL",
"TM-002": "SCENE_EMPTY",
"TM-003": "GRASP_TIMEOUT",
"TM-004": "REORIENT_FAIL",
"TM-005": "PLACE_OUT_OF_ZONE",
"TM-006": "SKILL_FAULT",
"TM-007": "SOFTWARE_FAULT",
}
class TaskDiagnostics:
def __init__(self, collaborator, scheduler):
self.collaborator = collaborator
self.scheduler = scheduler
def raise_fault(self, task_id: str, evt, code: str = "TM-006") -> None:
record = {
"code": code, "name": FAULT_TABLE.get(code, "UNKNOWN"),
"ts": time.time(), "task_id": task_id,
"details": getattr(evt, "__dict__", {}),
}
self.collaborator.send("task_faults", json.dumps(record))
# 依故障碼決定後續處理
if code in ("TM-001", "TM-004"):
self.scheduler.request_replan(task_id, reason=code)
elif code in ("TM-003",):
self.scheduler.retry_current(task_id)
else:
self.scheduler.pause()
10.6 task_manager/scheduler.py(與既有 GlobalAgent 整合)
import json, threading
from enum import Enum
from typing import Optional
class SchedState(Enum):
IDLE = "IDLE"; RUNNING = "RUNNING"; PAUSED = "PAUSED"
FAULT = "FAULT"; RESET = "RESET"
class TaskScheduler:
def __init__(self, global_agent, queue_mgr):
self.global_agent = global_agent
self.queue = queue_mgr
self.state = SchedState.IDLE
self.lock = threading.Lock()
def start_task(self, task_text: str) -> str:
reasoning = self.global_agent.publish_global_task(task_text, refresh=True, task_id=None)
for sub in reasoning.get("subtask_list", []):
self.queue.push(sub)
self.state = SchedState.RUNNING
return reasoning.get("task_id", "")
def pause(self): self.state = SchedState.PAUSED
def resume(self): self.state = SchedState.RUNNING
def abort(self, reason: str): self.state = SchedState.FAULT
def reset(self): self.queue = type(self.queue)(); self.state = SchedState.RESET
def request_replan(self, task_id: str, reason: str) -> None:
self.queue.active.clear()
self.global_agent.publish_global_task(
f"Replan task {task_id} due to {reason}", refresh=False, task_id=task_id
)
def retry_current(self, task_id: str) -> None: ...
10.7 task_manager/hmi.py(Flask blueprint)
from flask import Blueprint, jsonify, request
bp = Blueprint("hmi", __name__)
def register(app, scheduler, monitor):
@bp.post("/control/<action>")
def control(action):
getattr(scheduler, action)()
return jsonify({"state": scheduler.state.value})
@bp.get("/task_status")
def task_status():
return jsonify({
"state": scheduler.state.value,
"queue": list(scheduler.queue.active),
"pending": list(scheduler.queue.pending),
"runtimes": {tid: rt.state.value for tid, rt in monitor.runtimes.items()},
})
app.register_blueprint(bp)
10.8 改動 master/run.py
# 在 master_agent 構造之後新增:
from task_manager.queue import TaskQueueManager
from task_manager.scheduler import TaskScheduler
from task_manager.diagnostics import TaskDiagnostics
from task_manager.monitor import TaskMonitor
from task_manager.hmi import register as register_hmi
queue_mgr = TaskQueueManager(max_active=3)
scheduler = TaskScheduler(global_agent=master_agent, queue_mgr=queue_mgr)
hmi_pusher = type("Pusher", (), {"push_progress": lambda self, tid, st: socketio.emit("task_progress", {"task_id": tid, "state": st})})()
diagnostics = TaskDiagnostics(master_agent.collaborator, scheduler)
monitor = TaskMonitor(master_agent.collaborator, diagnostics, hmi_pusher)
monitor.start_listeners()
register_hmi(app, scheduler, monitor)
11. 常見問題(FAQ)
Q1:為什麼還要做監控?讓 LLM 自己看著辦不行嗎? A:LLM 的延遲(百毫秒級)和 token cost 都無法承受 10 Hz 的視覺更新。監控用 FSM + 規則寫死「物理可判定」的條件(朝向、區域),LLM 只在重規劃這種低頻決策上介入。
Q2:條碼朝向的判斷一定要在謝那邊做嗎?我能不能自己再判斷一次?
A:建議單一源頭,避免兩邊結果不一致。你最多在 TaskMonitor 做 confidence 的平滑(例如連續 3 幀都判 up 才算)。
Q3:要不要替換掉 GlobalAgent 換一套?
A:不要。GlobalAgent 已經處理了 robot 註冊、async dispatch、JSON 重試這些工程細節,重寫得不償失。把它當「規劃 + 派發底層」,你的 TaskScheduler 是它的上層 wrapper。
Q4:故障碼會不會跟 Skill Manager 的故障碼撞?
A:不會,命名空間分開:
- 任務級:TM-xxx
- 技能級:SK-xxx(李定)
- 策略級:PL-xxx(朱定)
透過 task_faults channel 只發 TM-xxx;上游的 SK-xxx / PL-xxx 一律由 diagnostics 模塊翻譯為 TM-006: SKILL_FAULT 或 TM-007: SOFTWARE_FAULT。
Q5:如果 GroundingDINO 還沒整合好,我能怎麼先開發? A:寫一個 mock publisher:
# tools/mock_vision.py
import json, time, redis
r = redis.Redis()
while True:
r.publish("vision_events", json.dumps({
"event": "vision_update", "ts": time.time(), "objects": [
{"id": "P1", "barcode_face": "down", "zone": "workspace", "score": 0.9}
]
}))
time.sleep(0.3)
Q6:UI 要前端嗎?
A:MVP 期先用 curl 或 redis-cli 觸發即可。前端可以等 Sprint 3 再做(或對接公司既有的 User pad UI)。
Q7:和「分層智能」其他層的開發節奏怎麼對齊?
A:建議三層並行 + 雙向 mock:
- 你給李、朱一個 mock 的「子任務 publisher」,讓他們先寫 Skill Manager。
- 李、朱給你一個 mock 的 skill_events publisher,讓你先寫 Monitor。
- 謝給你一個 mock 的 vision_events publisher。
Sprint 2 末做第一次 e2e 拉通。
附錄 A:與隊友的接口 Checklist
請在每週同步會議時,逐項確認:
- [ ] 場景 profile:
scene/profile.yaml(你維護) - [ ] 規劃 prompt:
master/agents/prompts.py::PACKAGE_SORTING_PLANNING(你維護) - [ ]
vision_eventsschema(謝實作、你定義) - [ ]
skill_eventsschema(李實作、你定義) - [ ]
task_faultsschema(你實作) - [ ] 任務開關 API:
startup()/shutdown()(尹實作,你調用) - [ ] MCP 工具清單(朱維護的
skill.py,你 review) - [ ] 結構化技能語料(謝、李、朱共同維護的技能庫,你只讀)
附錄 B:本文件對應的 PDF 章節
| 本文件 | |
|---|---|
| §1 任務理解 | 第 1 頁「重要概念 / 總體思路」、第 3 頁「Task manager 功能點」 |
| §2 系統總體架構 | 第 1 頁「系統設計思路」、第 2 頁「功能邏輯架構」 |
| §3 模塊地圖 | 第 3 頁「Task manager」全圖 |
| §5.3 監控程序 | 第 3 頁「任務監控程序」條目 |
| §5.4 診斷程序 | 第 3 頁「任務診斷程序」條目 |
| §6 視覺接入 | 第 2 頁「Algo-Ground / 視覺 Grounding」 |
| §10 代碼骨架 | (新增,PDF 未涵蓋) |
結語
這份文件的目的是讓你在 Sprint 0 結束前,能對著它把 task_manager/ 目錄敲完骨架,並且讓謝、李、朱、尹四位隊友看完就知道接口在哪。後續真實開發中遇到的接口微調,請直接更新本文件 + interfaces.md,作為團隊唯一事實來源(single source of truth)。
任何時候卡住,回到一句話: 「Task Manager 不執行動作,只決定下一步做什麼、有沒有做成、做壞了怎麼辦。」