包裹分揀任務(左傳送帶 → 條碼朝上 → 右傳送帶)

文件版本:v1.0(2026-05-12) 對應源碼基線:fmc3-robotics-main/projects/RoboOS(stand-alone) 你的職責範圍(蒲):任務規劃 / 任務調度 / 任務監控 / 任務診斷 / 人機交互(部分) 隊友:尹(任務開關)、朱(Skill Unit / 主策略 / 後處理 / 策略診斷)、李(技能編排 / 技能調度 / 技能監控 / 人機協作)、謝(前處理 / 結構化技能感知 / 視覺定位)


目錄


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,會增加維護成本),透過內部 asyncio event 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.pyslaver/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,新增 TaskQueueManagerTaskMonitorTaskDiagnostics 三個 mixin,再讓 GlobalAgent 持有它們。


5. 詳細模塊設計

以下所有模塊都假設單一 Python process,共享同一個 asyncio event 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_PLANNING prompt。
  • 注入「場景 snapshot」到 prompt:在呼叫 LLM 前,先從 Redis 撈 current_scene key。
  • 長度約束:隊列長度 ≤ 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)— 本系統最核心的模塊

職責:實時整合視覺結果與技能級回傳,回答以下問題:

  1. 包裹類型是什麼?尺寸多大?
  2. 該選哪一個包裹?
  3. 當前包裹的位姿 / 區域?
  4. 條碼面朝上了嗎?
  5. 包裹放置位置對嗎?
  6. 子任務是否完成?是否成功?

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_PLANNING prompt
  • [ ] 跑通既有 RoboOS(master + 一個假 slaver)做 smoke test

Sprint 1:規劃 + 調度 MVP(第二週)

  • [ ] 實作 TaskQueueManagerTaskScheduler 狀態機
  • [ ] 把 GlobalAgent.publish_global_task 包進 scheduler
  • [ ] HMI 增加 /control/* 端點
  • [ ] 用「假視覺事件」(手動 redis-cli publish)驗證端到端

Sprint 2:監控閉環(第三週)

  • [ ] 實作 TaskMonitor + TaskRuntime FSM
  • [ ] 訂閱 vision_eventsskill_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_FAULTTM-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 期先用 curlredis-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_events schema(謝實作、你定義)
  • [ ] skill_events schema(李實作、你定義)
  • [ ] task_faults schema(你實作)
  • [ ] 任務開關 API:startup() / shutdown()(尹實作,你調用)
  • [ ] MCP 工具清單(朱維護的 skill.py,你 review)
  • [ ] 結構化技能語料(謝、李、朱共同維護的技能庫,你只讀)

附錄 B:本文件對應的 PDF 章節

本文件 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 不執行動作,只決定下一步做什麼、有沒有做成、做壞了怎麼辦。」