my-pi-agent--todolist与background
三大组成部分
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27┌─────────────────────────────────────────────────────────────────────────────┐
│ 统一任务系统的三大组成部分 │
├─────────────────────────────────────────────────────────────────────────────┤
│ │
│ 【第一部分:数据与 DAG 依赖层】(对标 s10 Task System / Claude Code) │
│ • 解决:“任务是什么、怎么存、前后依赖关系是什么?” │
│ • 核心组件:TaskItem + TaskStore + 单一标准 todo 工具(整合增量与批量) │
│ • 关键能力:DAG 依赖图、传递性成环检测、单 in_progress 约束、自动解锁、 │
│ Crash-Safe 原子文件持久化、串行写操作因果调度。 │
│ │
│ ───────────────────────────────────────────────────────────────────────── │
│ │
│ 【第二部分:看板与视图投影层】(对标 s05 TodoWrite / Pi Todo) │
│ • 解决:“大模型每轮如何感知全局进度,做到不迷路且 0 Token 额外开销?” │
│ • 核心组件:ToolResult 随路回显投影 + action="write" 批量便签支持 │
│ • 关键能力:工具返回时自动携带最新 <TASK_BOARD>,100% 捍卫 Prompt 前缀 │
│ 缓存(Prefix Cache),0 额外轮次开销,Session 历史零污染。 │
│ │
│ ───────────────────────────────────────────────────────────────────────── │
│ │
│ 【第三部分:后台异步执行与通知收割层】(对标 s11 Background / Pi 通知流) │
│ • 解决:“慢操作(跑测试、编译构建)如何不卡死主 Agent 循环?” │
│ • 核心组件:BackgroundRunner + MessageQueue Follow-up 队列 │
│ • 关键能力:bash(run_in_background=True) 立即返回占位符、后台非阻塞运行、│
│ 跨平台进程树递归强杀(杜绝孤儿进程)、完成后通知自动收割闭环。│
│ │
└─────────────────────────────────────────────────────────────────────────────┘
领域模型分工:Subagent 委派句柄 vs 项目工单卡片
为避免概念混淆,彻底解耦子代理委派与看板工单: -
subagent_tasks.py(SubagentTask /
SubagentTaskManager):代表一次子代理的运行时执行容器/句柄,存活于一次
task() 工具调用期间,状态为
RUNNING / COMPLETED / ERROR; -
task_store.py(TaskItem /
TaskStore):代表工程规划中的持久化待办卡片,存在于磁盘
tasks.json 中,状态为
pending / in_progress / completed / deleted。
Task System
在简单的单步任务中,Agent 可以凭记忆搞定;但当面对大型工程开发(如“重构认证模块、新增 OAuth 接口并补全单元测试”)时, 会面临三大挑战:
- 前后顺序不可乱:你不能在“数据库表结构还没建好”之前就去“写 API 接口代码”,也不能在“接口没写完”之前就去“跑集成测试”;
- 死锁与混乱防范:如果大模型逻辑混乱,让任务 A 等待任务 B,又让任务 B 等待任务 A,就会导致整个 Agent 陷入死锁;
- Token 浪费严重:如果像早期系统那样每次修改都重传整张大表,随着任务增多,每轮都会浪费成千上万的 Token。
Task System 的使命:提供一个带 DAG(有向无环图)依赖拓扑、增量修改、自动解锁、崩溃安全的项目工单引擎!
TaskItem 的数据结构
在
packages/my-agent-core/src/my_agent_core/task_store.py
中,每一个任务卡片都是一个 TaskItem 实例:
1
2
3
4
5
6
7
8
9
10
class TaskItem:
id: str # 服务端分配的唯一标识符,如 "task_1", "task_2"
subject: str # 任务短标题(如 "设计数据库表结构")
description: str = "" # 任务长描述与具体要求
status: TaskStatus = "pending" # 状态机:pending | in_progress | completed | deleted
owner: str | None = None # 负责人(如 "agent" 或 "subagent:reviewer")
active_form: str | None = None # 进行时的动态文案(如 "writing schema.sql")
blocked_by: list[str] = field(default_factory=list) # 该任务所依赖的前置任务 ID 列表
metadata: dict[str, Any] = field(default_factory=dict) # 扩展元数据
### 1. status(工单生命周期状态机)
它记录一个任务“当前处于哪个人生阶段”。一共有 4 种合法状态:
1
2
3
4
5
6
7
8
9
10
11┌──────────┐ 认领开工 ┌─────────────┐ 完工打勾 ┌───────────┐
│ pending │ ────────────────► │ in_progress │ ────────────────► │ completed │
└──────────┘ (必须依赖已清空) └─────────────┘ └───────────┘
│ ▲
│ 软删除 │
└───────────────────────────────────────────────────────────────────┘
│
▼
┌───────────┐
│ deleted │ (归档墓碑,不影响拓扑)
└───────────┘
- pending(排队待办): 刚创建出来的初始状态。它可能正被别人卡住(blocked_by 非空),也可能已经就绪随时可做;
- in_progress(正在干活): 已经被 Agent 或 Subagent 认领开工。铁律:只有当它的 blocked_by == [] 时,才允许进入此状态!
- completed(打勾完工): 任务已经搞定。一旦进入此状态,就会像多米诺骨牌一样,去下游任务的 blocked_by 列表里把自己划掉;
- deleted(软删除/墓碑): 作废的任务,不再参与依赖计算。
────────────────────────────────────────────────────────────────────────────────
### 2. blocked_by(前置依赖卡点名单)
它是一个字符串数组(list[str]),记录着“有哪几个上游任务挡在我的面前”: - blocked_by = [](空列表): 说明当前没有任何卡点!它是一条绿色通道,随时可以被认领变成 in_progress; - blocked_by = [“task_1”, “task_2”](非空列表): 说明这是一个被红灯锁定的任务。即便大模型想认领它,系统也会说:“不行!你必须先把 task_1 和 task_2 都做完!”
TaskStore
核心算法与机制
#### ① 两阶段 DAG 建图(Two-Phase DAG Construction)
大模型在同一轮中发出多个 task_create
调用时,无法提前预知系统分配的 ID。 - 阶段一(节点创建):模型调
task_create 创建 task_1 和
task_2; - 阶段二(依赖绑定):模型拿到 ID 后,调用
task_update(task_id="task_2", add_blocked_by=["task_1"])
建立依赖边。
#### ② 传递性成环检测(Cycle Detection)
在执行 add_blocked_by 时,TaskStore
沿着依赖链进行深度遍历回溯(_depends_on 广度优先算法): -
如果尝试添加
task_1 ➔ task_2 ➔ task_1,系统立即拦截并报错抛出
Cycle detected,从根本上防止任务图死锁。
#### ③ 单 in_progress 聚焦原则(Single In-Progress Invariant)
- 默认开启
enforce_single_in_progress=True; - 当 Agent 尝试把
task_2设为in_progress时,如果task_1还在进行中,系统会拒绝并提示必须先完成或暂停前一个任务,强制大模型保持注意力聚焦。 - 在多 Agent 或后台任务场景下,可设为
False放开限制,允许互不依赖的分支并发进行。
#### ④ 下游任务自动解锁(Auto-Unblocking Feedback)
- 当一个任务调用
task_update(status="completed")时,系统自动遍历所有pending任务; - 只要某个任务的前置依赖因本次完成而全部清空,系统在本次工具返回值中显式附带:
大模型收到回显,下一轮立刻知道
1
{ "task": { "id": "task_1", "status": "completed" }, "unblocked": ["task_2"] }
task_2已就绪,实现自动交接!
#### ⑤ 崩溃安全原子落盘与串行调度
- 持久化写入使用
tempfile+fsync+os.replace写入<workspace>/.my_agent_core/tasks.json,断电永不损坏历史; - 写操作工具(
task_create,task_update,todo_write)诚实声明为is_parallel_safe=False,由ToolRegistry按模型原序串行调度,TaskStore内部无需加冗余锁,保持架构纯粹极简。
大模型操作面:从离散 CRUD 到单一统一 todo 工具
在设计底层数据操作面时,系统涵盖了 5 种核心操作能力:
1
2
3
4
5
6
7
8
9
10
11
12
13
141. create(subject, description="", active_form=None)
➔ 创建工单,返回分配的 ID 与初始看板
2. update(task_id, status=None, add_blocked_by=None, remove_blocked_by=None, ...)
➔ 增量修改字段与依赖,返回最新状态、unblocked 解锁列表与最新看板
3. get(task_id)
➔ 查阅单条工单的完整信息(包含长文本 description)
4. list(include_deleted=False)
➔ 获取看板摘要列表(默认不吐出长 description,极省 Token)
5. write(todos: list[dict])
➔ 批量便签覆盖写入快捷能力(对标 s05 TodoWrite 草稿纸模式)
架构演进注记:在初始探索期,这 5 种操作曾作为独立的 5 个工具(
task_create、task_update等)导出;在最终架构收敛时,为了防止工具 Schema 数量膨胀并严格对标 Pi 与 Hermes-Agent,我们将其统一封装为了单一的标准todo工具(通过action: Literal["create", "update", "list", "get", "clear", "write"]统一分发),既保持了概念清晰,又极大精简了大模型的工具调用认知负担。
内部结构解剖图
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19┌─────────────────────────────────────────────────────────────────────────────────────────┐
│ TaskStore 内部全景 │
├─────────────────────────────────────────────────────────────────────────────────────────┤
│ │
│ 【1. 串行调度保护】: 工具层声明 is_parallel_safe=False,ToolRegistry 原序串行执行 │
│ │
│ 【2. 自增序列号】: self._next_id = 3 (自动派发 task_1, task_2, task_3...) │
│ │
│ 【3. 内存工单字典】: self.tasks: dict[str, TaskItem] │
│ │ │
│ ├── "task_1" ──► TaskItem(id="task_1", subject="设计表结构", status="completed", │
│ │ blocked_by=[]) │
│ │ │
│ └── "task_2" ──► TaskItem(id="task_2", subject="编写 API", status="in_progress", │
│ blocked_by=[], owner="agent") │
│ │
│ 【4. 物理文件】: <workspace>/.my_agent_core/tasks.json (Crash-Safe 原子持久化) │
│ │
└─────────────────────────────────────────────────────────────────────────────────────────┘
内部方法
| 方法名 | 大白话作用 | 现实类比 |
|---|---|---|
create(...) |
创建新工单:服务端自动分配递增编号(如
task_1),填入标题,初始化为待办(pending)并落盘。 |
挂号取号机吐出一张新排队票。 |
update(...) (最核心方法) |
增量修改与流转: 1. 改状态(设为进行中或已完成); 2.
绑依赖(add_blocked_by,自带成环死锁拦截); 3.
自动解锁(若完成上游,主动返回哪些下游任务被解锁了)。 |
工单推进:认领开工、绑定先后顺序、打勾并通知下一位接力。 |
get(task_id) |
查单个任务详情:按编号精确调阅一张工单的全部内容(包含长篇描述和具体要求)。 | 调阅某一份病历或工单的完整档案。 |
list() |
看所有活跃任务:列出当前所有没被删除的任务清单。 | 查看当前看板上的全部工单卡片。 |
batch_write(todos) |
批量便签覆盖:支持大模型一次性传一个列表进行批量新建或覆盖更新(用于兼容便签式
todo_write)。 |
在便签纸上一口气写下一组 3~5 步小计划。 |
render_board() |
渲染排版看板:把所有任务格式化为打勾图标的 Markdown
文本([x] 已完成、[>]
进行中、[ ] 待办),每轮自动投影给大模型看。 |
把整个工单进度整理成大屏幕实时滚动看板。 |
create() 方法(工单创建)
create() 只做一件事:
“分配一个全局唯一的递增 ID(如
task_1),填入标题和要求,生成一张全新的待办工单,存入内存并落盘。”
端到端全链路
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26【人类用户】: "帮我重构认证模块并编写测试"
│
▼
【大模型 (LLM)】: 思考后认为需要拆分任务,决定发起工具调用
│
▼ 输出 JSON: {"name": "todo", "arguments": {"action": "create", "subject": "设计数据库表"}}
【Agent 调度器】: 拦截到工具调用,路由给 todo 工具
│
▼ 传入实参调用: todo(action="create", subject="设计数据库表")
【task_tools.py】: 调用底层仓库 await store.create(...)
│
▼ ═══════════════ 进入 TaskStore.create() 核心执行 ═══════════════
│
│ 1. 校验参数 ➔ 2. 派发 ID (task_1) ➔ 3. 构造工单 ➔
│ 4. 存内存字典 ➔ 5. 刷盘到 tasks.json
│
▼ ═══════════════════════════════════════════════════════════════
【TaskStore】: 返回刚刚建好的 TaskItem 对象
│
▼
【task_tools.py】: 包装成标准工具观察结果 ToolResult:
│ {"task": {"id": "task_1", "subject": "设计数据库表", "status": "pending"},
│ "message": "Created task_1"}
▼
【大模型 (LLM)】: 收到观察结果,心里有了底:
“好的,系统分配的编号是 task_1,我下一步可以去给它绑定依赖或者认领它了!”
内部流水线
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48 [ 调用者发起: await store.create(subject="设计表结构") ]
│
▼
┌───────────────────────────────────────────────────────────────────────────┐
│ 【第 1 步:参数清洗与非空校验】 │
│ • sub = subject.strip() 剔除前后空格; │
│ • 检查 sub 是否为空? │
│ ├── 为空 ──► 抛出 ValueError("Task subject cannot be empty") │
│ └── 合法 ──► 进入下一步。 │
└─────────────────────────────────────┬─────────────────────────────────────┘
│
▼
┌───────────────────────────────────────────────────────────────────────────┐
│ 【第 2 步:自增 ID 签发与计数器推进】 │
│ • 读取当前计数器: task_id = f"task_{self._next_id}" (如 "task_1"); │
│ • 计数器立即累加: self._next_id += 1 (变成 2,留给下一个人)。 │
└─────────────────────────────────────┬─────────────────────────────────────┘
│
▼
┌───────────────────────────────────────────────────────────────────────────┐
│ 【第 3 步:实例化 TaskItem 结构体】 │
│ • id = "task_1" │
│ • subject = "设计表结构" │
│ • status = "pending" (初始状态永远是待办,不能一步登天) │
│ • owner = None (初始无人认领) │
│ • blocked_by = [] (初始依赖为空,等待第二阶段用 update 绑定) │
└─────────────────────────────────────┬─────────────────────────────────────┘
│
▼
┌───────────────────────────────────────────────────────────────────────────┐
│ 【第 4 步:更新内存注册表】 │
│ • self.tasks["task_1"] = task │
│ • 此时内存字典已完成登记,后续的 get/list 方法立刻能查到。 │
└─────────────────────────────────────┬─────────────────────────────────────┘
│
▼
┌───────────────────────────────────────────────────────────────────────────┐
│ 【第 5 步:触发原子安全落盘】 self._save_to_disk() │
│ • 内存全量数据序列化为 JSON 字符串; │
│ • 写入临时文件 tasks_xxx.tmp ➔ 强制 os.fsync 刷盘 ➔ os.replace 覆盖; │
│ • 确保此时即使操作系统崩溃,磁盘上的任务记录也 100% 完整无损。 │
└─────────────────────────────────────┬─────────────────────────────────────┘
│
▼
┌───────────────────────────────────────────────────────────────────────────┐
│ 【第 6 步:返回工单实体】 │
│ • 将构建好的 TaskItem 实体对象 return 返回给调用方。 │
└───────────────────────────────────────────────────────────────────────────┘
#### 疑问 1:为什么任务写操作要严格串行? -
场景:大模型在同一轮并发调用 task_create(subject="建表") 和
task_create(subject="写API"); - 方案:写操作工具诚实声明为
is_parallel_safe=False,由 ToolRegistry
保序串行执行。任务 A 先进拿到 task_1 并累加计数器;任务 B
随后进入拿到 task_2。保证每个任务绝对拥有唯一的身份证号,且
TaskStore 无需内部加锁,代码极简。
#### 疑问 2:为什么任务 ID 不能让大模型自己填? - 如果让大模型自己传
ID,大模型第一轮可能起名叫 "1",第二轮叫
"step_one",第三轮叫
"task_A";后续绑定依赖时极易引用错乱; - 由系统统一派发递增
ID(task_1, task_2),规范统一。
update() 方法(增量流转与自动解锁)
端到端全链路图
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30【大模型 (LLM)】: 任务 1 完成了,调用工具 todo(action="update", task_id="task_1", status="completed")
│
▼
【ToolRegistry 调度器】:
│ • 识别出 todo (action="update") 是 is_parallel_safe=False(写操作);
│ • 按照因果顺序严格串行执行,调用底层 task_tools.py。
▼
【task_tools.py】: 调用底层仓库 await store.update(task_id="task_1", status="completed")
│
▼ ═══════════════ 进入 TaskStore.update() 核心流水线 ═══════════════
│
│ 1. 查工单是否存在
│ 2. 若改状态为 in_progress,检查是否已有在跑任务
│ 3. 若增删依赖,做图遍历成环检测(Cycle Detection)
│ 4. 覆盖新字段(status, subject, owner...)
│ 5. 【自动解锁流水线】:若是 completed,自动帮下游清除依赖并挑出就绪任务
│ 6. 触发原子刷盘 (_save_to_disk)
│
▼ ═════════════════════════════════════════════════════════════════
【TaskStore】: 返回元组 (task_1, unblocked=["task_2"])
│
▼
【task_tools.py】: 包装成标准 ToolResult 给大模型回显:
│ {
│ "task": {"id": "task_1", "status": "completed"},
│ "unblocked": ["task_2"],
│ "message": "Updated task_1"
│ }
▼
【大模型 (LLM)】: 拿到回显:“task_1 已完成,下游 task_2 刚刚解锁!我下一轮去认领 task_2!”
内部 6 步微观流水线详解
内部经过 6 步严格处理:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54 [ 调用发起: await store.update(task_id, status, add_blocked_by...) ]
│
▼
┌───────────────────────────────────────────────────────────────────────────┐
│ 【第 1 步:存在性防御检查】 │
│ • if task_id not in self.tasks: │
│ raise KeyError(f"Task '{task_id}' not found") │
│ • 拦截大模型凭空捏造不存在的 ID,保证只操作真实已登记的工单。 │
└─────────────────────────────────────┬─────────────────────────────────────┘
│
▼
┌───────────────────────────────────────────────────────────────────────────┐
│ 【第 2 步:单进行中(in_progress)聚焦校验】 │
│ • if status == "in_progress" and self.enforce_single_in_progress: │
│ 扫描其它任务,若发现已有任务处于 in_progress,直接抛出 ValueError! │
│ • 作用:防止单 Agent 一心二用、多线开花导致全部烂尾。 │
└─────────────────────────────────────┬─────────────────────────────────────┘
│
▼
┌───────────────────────────────────────────────────────────────────────────┐
│ 【第 3 步:依赖图增删边与死锁成环拦截】 │
│ • 如果传入了 add_blocked_by=[dep, ...]: │
│ 1. 检查 dep == task_id ➔ 严禁自己依赖自己; │
│ 2. 检查 dep in self.tasks ➔ 依赖的目标必须真实存在; │
│ 3. 调用 self._depends_on(dep, task_id) 广度优先回溯检查; │
│ 若发现 dep 已经在等 task_id,成环!➔ 抛出 Cycle detected 拦截; │
│ 4. 安全 ➔ 追加到 task.blocked_by(自动去重)。 │
│ • 如果传入了 remove_blocked_by ➔ 从依赖列表中剔除。 │
└─────────────────────────────────────┬─────────────────────────────────────┘
│
▼
┌───────────────────────────────────────────────────────────────────────────┐
│ 【第 4 步:更新标量字段】 │
│ • 哪个参数传了就更新哪个: │
│ status, subject, description, active_form, owner, metadata │
└─────────────────────────────────────┬─────────────────────────────────────┘
│
▼
┌───────────────────────────────────────────────────────────────────────────┐
│ 【第 5 步:下游多米诺骨牌自动解锁(最核心逻辑)】 │
│ • if status == "completed": │
│ 遍历内存中所有其它处于 pending 的任务: │
│ 1. 检查该任务的 blocked_by 是否包含当前完成的 task_id? │
│ 2. 如果包含,自动把 task_id 从依赖中移出; │
│ 3. 移出后检查: len(blocked_by) == 0? │
│ ➔ 如果前置障碍已全清空,追加进 unblocked 列表! │
└─────────────────────────────────────┬─────────────────────────────────────┘
│
▼
┌───────────────────────────────────────────────────────────────────────────┐
│ 【第 6 步:原子安全落盘并返回双结果】 │
│ • 调用 self._save_to_disk() (临时文件 ➔ fsync 刷盘 ➔ replace 覆盖) │
│ • 返回元组: (更新后的 TaskItem, unblocked 列表) │
└───────────────────────────────────────────────────────────────────────────┘
任务收尾守卫(Early-Exit Guard 自动提醒)
大模型做完任务经常容易忘记调用 todo
打勾,准备直接输出纯文本说“我干完了”溜走。 对标 Pi
官方扩展事件哲学,核心 Agent 循环 (Agent.run) 保持
100% 纯净通用,绝不硬编码任何任务状态检查。系统通过挂载在
TurnEnd 上的独立 TaskGuardHook
守卫进行优雅拦截: -
当大模型未发起任何工具调用(not event.tool_results)准备退出本轮时,钩子检查
TaskStore; - 若发现看板上仍有处于 in_progress
的任务未结清,钩子调用 agent.steer(nudge)
注入一条强制提醒:
Task '{task_id}' is still marked as 'in_progress'. If you have completed it, please call todo(action='update', task_id='{task_id}', status='completed') before concluding.
- Agent.run 在安全点检测到 Steering
消息,自动继续推进下一轮,大模型被精准拉回乖乖调 todo
打勾结清,保证任务流 100% 闭环! -
钩子在单次运行中对同一任务最多敲打一次(nudged_ids 集合,在
AgentStart 时重置),杜绝死循环。
添加前置依赖时的检查
检查 1:自依赖防范(禁止自己等自己)
源码:
1
2if dep == task_id:
raise ValueError("Task cannot depend on itself")
为什么必须检查?
- 大模型常见幻觉:模型在更新 task_1 时,误把参数写成了 add_blocked_by=[“task_1”];
- 严重后果:这叫“逻辑自噬”。task_1 必须等 task_1 完成后才能开始,但它不开始就永远无法完成!
- 防御:只要发现依赖的目标编号和自己一模一样,当场拒绝!
────────────────────────────────────────────────────────────────────────────────
二、检查 2:前置存在性检查(禁止依赖幽灵任务)
源码:
1
2if dep not in self.tasks:
raise KeyError(f"Dependency task '{dep}' not found")
为什么必须检查?
- 大模型常见幻觉:模型随手编造了一个不存在的前置编号,比如 add_blocked_by=[“task_888”];
- 严重后果:系统里根本没有 task_888,也就永远没有人能把 task_888 标记为 completed;
- 结果就是该任务的 blocked_by 永远无法清空,变成一张永远被冻结的死卡片!
- 防御:只有当前 tasks.json 中真实存在且已分配的工单,才允许被绑定为前置依赖。
────────────────────────────────────────────────────────────────────────────────
三、检查 3:传递性成环死锁检测(最硬核的图论防御)
这是最关键、技术含量最高的一道关卡!
源码:
1
2if self._depends_on(dep, task_id):
raise ValueError(f"Cycle detected: {task_id} -> {dep} -> {task_id}")
什么是“传递性闭环”?
成环往往不是直接的,而是跨越了多个层级:
1
2
3
4
5已有的依赖链条:
task_3 ──依赖──► task_2 ──依赖──► task_1
此时模型试图让:
task_1 ──依赖──► task_3 ??
_depends_on 是怎么查出这个死锁的?
1
2
3
4
5
6
7
8
9
10
11
12
13def _depends_on(self, task_id: str, target_id: str) -> bool:
visited = set()
queue = [task_id] # 初始把 dep 放进队列
while queue:
curr = queue.pop(0)
if curr == target_id: # 顺藤摸瓜居然顺回了自己!
return True # 发现环死锁!
if curr in visited:
continue
visited.add(curr)
if curr in self.tasks:
queue.extend(self.tasks[curr].blocked_by)
return False
- 核心逻辑: 在把 task_1 ➔ 依赖 ➔ task_3 写入前,先从 task_3 开始往上溯源。 结果系统发现:task_3 的上游是 task_2,task_2 的上游正是 task_1! 系统立刻得出结论:“如果我现在允许 task_1 依赖 task_3,就会形成 1 ➔ 3 ➔ 2 ➔ 1 的闭环死锁!”
- 防御:直接掐死,并在报错信息中把完整的回路打印出来引导大模型。
render_board() 方法(看板与上下文投影)
作用与使用位置
- 调用时机与架构演进:
- 早期探索:曾考虑在
Agent.run()的模型视图构建阶段(BeforeModelCall决策点前)将看板注入系统提示词; - 生产架构定型:为了 100% 捍卫大模型供应商的
Prompt Prefix Cache(避免动态改动前缀导致每轮 KV Cache
全量击穿),系统最终将看板投影全面升级为“在
todo工具的ToolResult.data['board']中随路回显(In-Band Echo)”!
- 早期探索:曾考虑在
- 执行效果: 若存在活跃任务,自动渲染紧凑 Markdown
块随工具执行结果回显:
1
2
3[x] task_1: 设计数据库表结构 (completed)
[>] task_2: 编写 API 接口 (in_progress - writing endpoints)
[ ] task_3: 编写单元测试 (pending, blocked by: ['task_2']) - 核心定位: 这不是面向人类的 UI 可视化(人类 UI 属于产品层/TUI/Web 的职责),而是面向大模型的上下文工程(Prompt Projection)。大模型每轮执行完修改后即可在结果中看到当前焦点与阻塞关系,0 工具额外往返消耗,前缀缓存 100% 稳定,且 Session 磁盘历史绝对零污染。
_save_to_disk
真实磁盘文件内容长这样(完整样例)
打开你的 .my_agent_core/tasks.json,它的真实内容如下:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39{
"next_id": 4,
"tasks": [
{
"id": "task_1",
"subject": "设计数据库表结构",
"description": "定义 users 表与 orders 表,使用 PostgreSQL 语法",
"status": "completed",
"owner": "agent",
"active_form": "writing schema.sql",
"blocked_by": [],
"metadata": {
"priority": "high"
}
},
{
"id": "task_2",
"subject": "编写 API 接口",
"description": "基于 FastAPI 实现用户注册与查询接口",
"status": "in_progress",
"owner": "agent",
"active_form": "implementing endpoints",
"blocked_by": [],
"metadata": {}
},
{
"id": "task_3",
"subject": "编写单元测试",
"description": "覆盖 100% 的边界测试用例",
"status": "pending",
"owner": null,
"active_form": null,
"blocked_by": [
"task_2"
],
"metadata": {}
}
]
}
#### 1. 顶层自增序号:“next_id”(极其重要!)
1
"next_id": self._next_id
- 持久化内容:一个递增整数(比如 4);
- 为什么必须持久化它?
- 这是业界标准的 高水位线(Highwatermark)机制;
- 如果程序中途退出或者机器重启,重新加载时如果丢失了这个数字,计数器就会重置为 1;
- 计数器一旦重置为 1,下次创建新任务又会分配出重复的 “task_1”,导致新旧任务编号严重撞车冲突!
- 存下 “next_id”,保证重启后分配的一定是全新的 task_4。
todolist
整体流转流程图
1 | ┌────────────────────────────────────────────────────────────────────────────────────────┐ |
为什么一定要用“工具返回值回传看板”?
我们对比一下两种不同做法,您就会立刻感受到当前实现的优越性:
### 做法 A:早期粗暴做法(把看板动态拼在 System Prompt 里)
- 实现方式:每轮大模型说话前,框架在 messages[0](System
Prompt)尾部动态追加当前的
… 。 - 致命缺陷:
- 击穿前缀缓存(Cache Busting):Anthropic、OpenAI、DeepSeek 的 KV Cache 是按前缀从前往后匹配的。如果 messages[0] 每一轮因为任务状态变化而改变,整段前缀缓存每轮全量失效!首字延迟(TTFT)从 200ms 飙升到 2~3 秒,Token 账单直接 翻倍;
- 会话历史污染:写入 Session 文件的 System 提示词每一轮都在变,导致回溯历史极其混乱。
### 做法 B:我们当前的做法(利用 ToolResult 随路回传)
- 实现方式: 在
packages/my-agent-core/src/my_agent_core/tools/builtin/task_tools.py
中:
1
2
3
4
5
6
7
8
9
10# 当模型调 todo 工具修改状态时:
return ToolResult(
ok=True,
data={
"action": "update",
"task": {...},
"unblocked": unblocked,
"board": store.render_board(), # 看板随路返回!
},
) - 核心优势:
- 100% 保护前缀缓存:开头的 System Prompt 永久冻结、一字不改;看板作为最新工具调用的客观返回结果自然追加在消息队 列最末端,前序的所有对话前缀完美命中缓存;
- 0 额外查询往返:大模型执行了更新后,不需要再傻傻地调一次 todo(action=“list”),在当前轮次就当场看清了全盘最新格 局;
- 数据流绝对真实纯净:Session JSONL 磁盘文件只记录正常的 ToolCall 和 ToolResult,没有任何人为伪造的系统幽灵消息。
底层三大模块是如何协同工作的?
这套机制由 3 个模块像钟表齿轮一样紧密咬合:
### 1. 排版中枢:TaskStore.render_board()(task_store.py)
它负责把内存里的任务字典“脱水”压缩成大模型最容易解析的紧凑 Markdown: - [x]:已完工; - [>]:当前唯一在跑的焦点(带 active_form 如 writing schema); - [ ]:待办(若有依赖,带 blocked by: [‘task_1’])。 (刻意丢掉几十上百字的长篇描述,保证看板只有几十个 Token,极度节省上下文空间)。
### 2. 通信载体:todo 工具(task_tools.py)
- 单一标准入口,支持 create、update、list、get、clear、write;
- 所有可能改变任务状态的操作,返回时统统附带 store.render_board();
- 声明 is_parallel_safe=False,由底层 ToolRegistry 保证按顺序串行执行,防止大模型同一轮派发多个修改导致看板数据错乱。
### 3. 纪律委员:TaskGuardHook(task_tools.py)
- 守在 TurnEnd 事件出口;
- 如果模型把代码写完了,觉得自己做完了想直接输出文本交差,但忘记调工具去把任务状态从 in_progress 改为 completed;
- 钩子会立刻从外部调用 agent.steer(…),在安全点把模型拦截拉回:“你还有一个任务在 in_progress,请先更新状态再交差 !”,实现管理闭环。
task_tools.py
它扮演着极其关键的角色:它是连接底层数据层(TaskStore)与大模型 ReAct 认知循环的“桥梁与执行终端”。
它在工程上主要交付了两件核心产物: 1. make_todo_tool(store):为大模型提供一个单一、多功能、防并发竞态的 todo 工具; 2. TaskGuardHook:为框架装配一个基于生命周期事件的后台守卫,防止模型早退。
在哪进行注册
一、Agent.__init__ 中的装配流水线
在 Agent.__init__ 中(约第 135~145 行),系统按清晰的顺序执行初始化:
1
2
3
4
5
6
7
8# 1. 解析 task_store(检查工作区或用户传入配置)
self.task_store = self._init_task_store(task_store)
# 2. ① 注册工具(把 todo 工具注入 ToolRegistry)
self._register_tools(tools)
# 3. ② 注册 Hooks(把 TaskGuardHook 注入 HookRegistry)
self._register_hooks(hooks)
只要 self.task_store 处于启用状态,这两个组件就会分别注册进工具注册表与事件钩子注册表。
────────────────────────────────────────────────────────────────────────────────
### 二、1. todo 工具是在哪注册的?
在 agent.py 的 _register_tools() 方法中(第 195~215 行):
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15def _register_tools(self, tools: list[Tool]) -> None:
"""注册用户工具 + 内置 task 工具 + 内置 memory 工具 + 内置 todo 工具。"""
for t in tools:
self.registry.register(t)
...
# 关键:如果启用了 task_store,自动把 make_task_tools 产出的 todo 工具注册进注册表!
if self.task_store:
for task_tool in make_task_tools(self.task_store):
if self.registry.get(task_tool.name) is not None:
raise ValueError(
f"Tool name '{task_tool.name}' conflicts with built-in task tool"
)
self.registry.register(task_tool) # 注册到 self.registry!
- 注册到哪:self.registry(即 ToolRegistry,框架的统一工具库);
- 注册效果:大模型调用 tools = self.registry.get_schemas() 时,就能看到 todo 工具的定义,从而可以在 ReAct 循环中主动调 用它。
────────────────────────────────────────────────────────────────────────────────
### 三、2. TaskGuardHook 守卫是在哪注册的?
在 agent.py 的 _register_hooks() 方法中(第 250~260 行):
1
2
3
4
5
6
7
8
9
10
11def _register_hooks(self, hooks) -> None:
"""构造时批量注册 hooks(对称 _register_tools)。"""
# 关键:如果启用了 task_store,自动实例化守卫并挂载两个生命周期事件!
if self.task_store:
guard = TaskGuardHook(self.task_store, self.steer)
self.hooks.register(AgentStart, guard.on_agent_start)
self.hooks.register(TurnEnd, guard.on_turn_end)
# 注册用户显式传入的其他自定义 hooks
for event_cls, callback in hooks or []:
self.hooks.register(event_cls, callback)
- 注册到哪:self.hooks(即 HookRegistry,框架的事件总线);
- 注册效果:
- 会话启动时触发 AgentStart → 执行 guard.on_agent_start,清空防死循环集合;
- 每轮模型说完话触发 TurnEnd → 执行 guard.on_turn_end,检查是否有在跑任务,若有则调用 self.steer 抓回大模型!
统一工具工厂:make_todo_tool(store: TaskStore)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21def make_todo_tool(store: TaskStore) -> Tool:
"""生成单一统一的标准 todo 工具(对标 Pi & Hermes-Agent)。"""
)
async def todo(
action: Literal["create", "update", "list", "get", "clear", "write"],
subject: str | None = None,
task_id: str | None = None,
status: Literal["pending", "in_progress", "completed", "deleted"] | None = None,
description: str | None = None,
active_form: str | None = None,
owner: str | None = None,
add_blocked_by: list[str] | None = None,
remove_blocked_by: list[str] | None = None,
todos: list[dict[str, Any]] | None = None,
include_deleted: bool = False,
) -> ToolResult:
6 个 Action
todo 工具的 6 个 Action 分点精简说明:
- 1.
create(单项新建)- 作用:传入
subject新建任务,系统分配自增编号(如task_1),初始为pending。 - 返回:新任务基本信息 +
最新随路看板(
board)。
- 作用:传入
- 2.
update(状态流转与自动解锁)- 作用:凭
task_id推进任务(开工标为in_progress、完工标为completed)或调整依赖。 - 返回:任务状态 +
自动解锁的下游列表(
unblocked) + 最新随路看板(board)。 - 规则:强制同一时间只准有 1 个任务在跑;完工打勾时自动解除下游的前置阻塞。
- 作用:凭
- 3.
write(批量草稿覆写)- 作用:一次性传入任务数组,像草稿纸一样一键批量初始化全部待办(对标 s05 TodoWrite)。
- 返回:批量任务清单 +
最新随路看板(
board)。
- 4.
list(轻量全局查)- 作用:主动查阅全局进展。
- 返回:极简摘要列表 +
紧凑看板(
board)。刻意滤掉长篇大论的描述,极省 Token。
- 5.
get(深度单卡查)- 作用:凭
task_id查看某项任务的详细要求。 - 返回:该任务的完整字段(包含长文本
description与metadata),专为深入阅读具体要求设计(唯一不随带看板的分支)。
- 作用:凭
- 6.
clear(一键清盘)- 作用:清空所有工单,并将自增计数器重置为 1。
- 返回:清空确认文案 + 空看板
(No active tasks)。
事件守卫类:TaskGuardHook(防止模型早退)
这是刚才重构的核心亮点(位于文件尾部第 193 行开始):
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31 class TaskGuardHook:
"""任务收尾早退守卫钩子(对标 Pi 扩展架构):在 TurnEnd 时检查未结清工单,通过 steer 提醒大模型。"""
def __init__(self, task_store: TaskStore, steer_fn: Callable[[str], None]) -> None:
self.task_store = task_store
self.steer_fn = steer_fn
self.nudged_ids: set[str] = set()
def on_agent_start(self, event: AgentStart) -> None:
"""会话开始时重置已提醒集合。"""
self.nudged_ids.clear()
def on_turn_end(self, event: TurnEnd) -> None:
"""Turn 结束时检查:若无工具调用且仍有 in_progress 任务,发起 steer 提醒。"""
if event.tool_results:
return # 模型这一轮还在调工具干活(比如在写代码),绝不打扰!
# 检查看板上是否有没结清的 in_progress 工单
in_progress = [t for t in self.task_store.list() if t.status == "in_progress"]
for t in in_progress:
if t.id not in self.nudged_ids:
self.nudged_ids.add(t.id) # 防死循环:单次会话只提醒一次
# 关键:调用 steer_fn 注入纠偏指令,在下一轮安全点打断大模型!
self.steer_fn(
f"Task '{t.id}' ({t.subject}) is still marked as 'in_progress'. "
f"If you have completed it, please call todo(action='update', task_id='{t.id}',
status='completed') "
f"to update your progress before concluding."
)
break
它的运作逻辑:
- on_agent_start:每次用户发起新的任务运行,清空 nudged_ids,保证新一轮有完整的提醒机会;
- on_turn_end:
- 过滤工具轮:if event.tool_results: return,模型在写文件、跑命令时,不触发;
- 捕获早退:当大模型没有调用任何工具,准备说“我做完了”直接交差退出时;
- 状态核实:去 task_store 查一眼——“咦?你还有 in_progress 的任务挂着呢!”;
- 自动干预:调用 steer_fn(直通底层 agent.steer),底层消息队列在内层循环安全点自动拉起下一轮:“别走!你还没打勾更 新状态,先调 todo 更新!”;
- 防死循环(nudged_ids):每个任务只催一次,如果模型执意不听劝,第二次不再重复催促,保证系统绝不陷入死循环。
后台异步执行 (Background Tasks)
后台任务跑完了,何时通知大模型?
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17┌───────────────────────────┐
│ 1. 生产者 (任务完成时) │ BackgroundRunner 发现子进程结束,
│ │ 直接调用:message_queue.add_followup(通知)
└─────────────┬─────────────┘
│
▼
┌───────────────────────────┐
│ 2. 邮箱缓冲区 (暂存排队) │ 通知安静躺在 followup 队列里,
│ │ 绝不粗暴打断大模型正在说的话或正在调的工具
└─────────────┬─────────────┘
│
▼
┌───────────────────────────┐
│ 3. 消费者 (安全点自唤醒) │ 大模型把当前手头的活干完、准备退出的那一瞬间,
│ │ 外层循环检查 if has_followup(),
│ │ 自动把通知取出来,拉起新一轮让大模型总结!
└───────────────────────────┘
并发调度抉择:用多线程还是原生异步?
真正的耗时命令(如 pytest),是在操作系统里作为一个独立的「子进程(Subprocess)」在跑; 而我们在 Python 内部,专门为它创建了一个极其轻量的「监工协程(_worker)」去默默盯着它!
1 | ┌─────────────────────────────────────────────────────────────┐ |
我们当前实现(原生异步)的硬核原理
在 packages/my-agent-core/src/my_agent_core/background.py 中,我们彻底抛弃了多线程,采用了纯粹的 Native Asyncio 协程调度:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20async def run_process(self, command: str, cwd: Path | str, description: str = "") -> str:
...
# 【定义一个后台协程】:它是一个可以随时暂停、恢复的任务
async def _worker() -> None:
# 1. 异步启动子进程,让操作系统内核接管
proc = await asyncio.create_subprocess_shell(...)
# 2. 关键点:【让出 CPU!】
# 子进程在跑慢测试时,这个 _worker 协程在这里原地“暂停挂起”,
# 整个 Python 线程瞬间腾出空来,去继续服务大模型和用户!
stdout, stderr = await proc.communicate()
# 3. 任务彻底跑完了,内核唤醒此协程,继续往下走:
self.message_queue.add_followup(notification)
# 【关键调度动作】:把 _worker 协程交给事件循环,立即放飞!
asyncio.create_task(_worker())
# 毫秒级瞬间返回任务 ID 给大模型!大模型完全感受不到任何卡顿!
return job_id
如何管理后台task
1
2
3
4
5
6
7 【1. 诞生】 【2. 运行中】 【3. 终结】
bash(background=True) self.jobs 花名册实时监控
│ │ │
▼ ▼ ▼
登记进 BackgroundJob ───► 可通过 bg_status 探活 正常完成 ──► 投递 follow_up 邮件
分配唯一 ID (bg_000001) 可通过 bg_logs 查尾部日志 用户按 Ctrl+C ──► taskkill 整树强杀
毫秒级返回凭证给模型 可通过 bg_kill 手工强杀 程序崩溃退出 ──► atexit 自动兜底收尸
内存状态花名册(BackgroundJob)
所有后台任务在启动那一刻,就会被登记在一张唯一的“任务花名册”中。
在我们的 packages/my-agent-core/src/my_agent_core/background.py 中:
1
2
3
4
5
6
7
8
9
10
class BackgroundJob:
"""单个后台作业的完整档案。"""
id: str # 唯一任务编号,如 bg_000001
description: str # 命令描述或命令本身
status: Literal["running", "completed", "failed", "cancelled"] = "running"
result: str | None = None # 执行结果/报错截断
exit_code: int | None = None # 操作系统退出码(0为正常,非0为异常)
process: asyncio.subprocess.Process | None = None # 绑定的系统子进程对象(含 PID)
started_at: float = field(default_factory=time.time) # 启动时间戳
BackgroundRunner 内部维护了 self.jobs: dict[str, BackgroundJob] 字典。 通过这个字典,系统在任何时候都可以秒级获知: - 当前一共有多少个任务在跑? - 每一个任务已经跑了多少秒(time.time() - job.started_at)? - 对应的操作系统物理 PID 是多少?
生命周期治理:如何主动取消与强杀(Cancellation)
当大模型或用户发现一个任务卡死、陷入死循环、或者不再需要时,如何管理它的退出?
#### 1. 框架级联动强杀(agent.abort())
当用户按 Ctrl+C 中止会话,或者调用 agent.abort() 时,管理调度引擎会立刻执行全量清理:
1
2
3
4
5
6
7# background.py 中的 cancel_all 实现:
async def cancel_all(self) -> None:
"""取消所有正在运行的后台子进程并递归杀死进程树。"""
for job in self.jobs.values():
if job.status == "running":
job.status = "cancelled"
_kill_process_tree(job.process) # 整树强杀!
#### 2. 操作系统整树根除(杜绝孤儿)
调用 _kill_process_tree: - Windows:taskkill /F /T /PID
#### 3. 解释器退出兜底(atexit)
BackgroundRunner.__init__ 中注册了 atexit.register(self._sync_cleanup)。 无论用户是直接关掉终端窗口,还是 Python 发生严重崩溃,解释器在退出前都会强制把存活的子进程全部扫地出门,保证操作系统的 干净。
场景全流程演示
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22Turn 1: 大模型规划
➔ 调 todo(action="create", subject="设计表结构") ➔ 获得 "task_1"
➔ 调 todo(action="create", subject="编写API") ➔ 获得 "task_2"
➔ 调 todo(action="update", task_id="task_2", add_blocked_by=["task_1"])
➔ 调 todo(action="update", task_id="task_1", status="in_progress", active_form="writing schema")
Turn 2: 大模型干活
➔ 调 write("schema.sql", ...) 写好了表结构
Turn 3: 大模型打勾并自动解锁
➔ 调 todo(action="update", task_id="task_1", status="completed")
➔ 工具返回: {"unblocked": ["task_2"]}
➔ 大模型收到反馈,立即调 todo(action="update", task_id="task_2", status="in_progress") 开始写 API!
Turn 4: 启动后台测试
➔ 调 bash("pytest tests/ -q", run_in_background=True)
➔ 立即返回: "[Background task bg_000001 started]"
➔ 主 Agent 释放控制权,等待后台跑完
Turn 5: 两层循环自动收割通知并总结
➔ 后台测试完成,MessageQueue 自动注入 <task_notification>
➔ 大模型自动获得测试通过日志,调用 todo(action="update", task_id="task_3", status="completed"),项目圆满完成!
1 | 初始状态 |