Agent开发

Agent loop

Agent Loop (智能体循环)指AI Agent为完成目标而反复执行的一套闭环过程:理解目标 -> 获取信息 -> 制定下一步 -> 调用工具执行 -> 检查结果 -> 决定继续、调整或结束

普通聊天通常是一问一答,Agent Loop则允许模型根据执行结果持续采取行动,直到任务完成、达到次数/时间限制,或者遇到必须由用户处理的阻塞。

One loop & Bash is all you need 一个工具 + 一个循环 = 一个Agent

当提出一个问题给大模型:”帮我读取下我的目录下有哪些文件并且执行XXX.py”,这个时候模型能输出一条bash命令,输出完就停止,不会自己跑因此不会有结果也不会有后续的推理。为了解决这个问题,必须自己手动跑一遍bash命令然后把输出粘贴回去,大模型才能接着分析输出下一个bash命令然后再手动执行bash命令粘贴输出结果循环往复。

这个过程人去手动执行都是在做中间层,为了实现Agent Loop就是要将其自动化:Harness层:循环-模型与真实世界的第一道连接

Agent Loop

一个while True循环,模型调用工具就继续,不调用就停止,整个过程只有两个信号:

信号含义循环动作
stop_reason == "tool_use"模型举手说”我要用工具”执行 → 结果喂回去 → 继续
stop_reason != "tool_use"模型说”我做完了”退出循环

工作原理

第一步:把用户的问题作为第一条消息

message = [{"role":"user","content": query}]

第二步:将消息和工具定义一起发给LLM

response = client.messages.create(
    model=MODEL, system=SYSTEM , message=messages,
    tools=TOOLS, max_tokens=8000,
)

第三步:追加模型回答,检查它是否调用了工具。没有调用就结束

messages.append({"role":"assistant","content":response.content})
if response.stop_reason != "tool_use":
    return

第四步:执行模型要求的工具,收集结果

results = []
for block in response.content:
    if block.type == "tool_use":
        output = run_bash(block.input["command"])
        results.append({
            "type": "tool_result",
            "tool_use_id": block.id,
            "content": output,
        })

第五步:把工具结果作为新消息追加回到第二步开始循环

messages.append({"role": "user", "content": results})

组成一个完成函数:

def agent_loop(messages):
    while True:
        response = client.messages.create(
            model=MODEL, system=SYSTEM, messages=messages,
            tools=TOOLS, max_tokens=8000,
        )
        messages.append({"role":"assistant","content": response.content})

        if response.stop_reason != "tool_use":
            return False
        results = []
        for block in response.content:
            if block.type == "tool_use":
                output = run_bash(block.input["command"])
                results.append({
                    "type":"tool_result",
                    "tool_use_id":block.id,
                    "content":output,
                })
        messages.append({"role":"user","content": results})

不到30行,这就是最小可运行的agent harness内核。它不是智能体本身,而是让模型能持续行动的最小运行框架,模型负责决策,harness负责执行。

Tool Use

前面完成了Agent loop,这个Agent中只有一个bash工具。这个时候如果要读取一个文件/path/to/file

如果有专门读取工具就可以直接生成

{
  "name": "read_file",
  "input": {
    "path": "/path/to/file"
  }
}

只有Bash工具时,则要转换成

{
  "name": "bash",
  "input": {
    "command": "cat /path/to/file"
  }
}

很明显的如果有专用的read_file工具,模型思考和生成shell命令所需的token就会明显减少也能够避免模型选择cat、head等命令的认知负担。举个例子,这里用户让读取10w个文件,那么纯bash工具相比有read_file+bash工具的Agent就会多消耗大量的token。同时有read_file工具还能够避免路径引用、空格和特殊字符的转义错误。

Tool Dispatch

对比前面的Agent唯一的变动在工具执行,run_bash()替换为TOOL_HANDLERS[block.name]()查表分发。

给Agent加一个工具只需要做两件事:
1.定义工具:在TOOLS数组里加一条描述
2.注册处理函数:在TOOL_HANDLERS字典里加一个映射

从一个工具到五个工具

前面只有一个工具bash

TOOLS = [{"name":"bash",...}]
def run_bash(command):...

现在追加到五个工具,每个工具都是独立定义

TOOLS = [
    {"name": "bash",       "description": "Run a shell command.", ...},
    {"name": "read_file",  "description": "Read file contents.",  ...},
    {"name": "write_file", "description": "Write content to file.", ...},
    {"name": "edit_file",  "description": "Replace text in file once.", ...},
    {"name": "glob",       "description": "Find files by pattern.", ...},
]

每个工具都有自己的实现函数

def run_read(path, limit=None):
    lines = safe_path(path).read_text().splitlines()
    if limit:
        lines = lines[:limit]
    return "\n".join(lines)

def run_write(path, content):
    safe_path(path).write_text(content)
    return f"Wrote {len(content)} bytes to {path}"

def run_edit(path, old_text, new_text):
    text = safe_path(path).read_text()
    if old_text not in text:
        return "Error: text not found"
    safe_path(path).write_text(text.replace(old_text, new_text, 1))
    return f"Edited {path}"

def run_glob(pattern):
    import glob as g
    return "\n".join(g.glob(pattern, root_dir=WORKDIR))

工具分发

TOOL_HANDLERS = {
    "bash":       run_bash,
    "read_file":  run_read,
    "write_file": run_write,
    "edit_file":  run_edit,
    "glob":       run_glob,
}

# 循环里只改了一行——从硬编码 run_bash 变成查表:
for block in response.content:
    if block.type == "tool_use":
        handler = TOOL_HANDLERS[block.name]    # 查表
        output = handler(**block.input)         # 调用
        results.append(...)

加一个工具 = 在TOOLS数组加一条 + 在TOOL_HANDLERS字典加一行。循环不变

执行前做权限判断

前面的Agent有五个工具。file tools受safe_path保护,但bash不受限制。让它清理一下项目,可能会执行rm -rf /,安全不能靠信任模型,要靠代码-在工具执行之前做判断。

Permission Overview

核心逻辑完整保留,唯一的变动在工具执行前插入check_permission()每个工具调用经过三道闸门,顺序固定:硬拒绝优先,软询问次之,都没命中就放行。

三道闸门对应三种策略

闸门作用命中后
1. 拒绝列表永远禁止的操作(rm -rf /sudo直接拒绝,不执行
2. 规则匹配取决于上下文的操作(读/写工作区外、rm 文件)交给闸门 3
3. 用户审批闸门 2 命中后,暂停等用户确认用户决定允许或拒绝

三道都没命中就直接执行

Permission Pipeline

闸门1:一张硬拒绝表,先查,命中就返回阻止信息

DENY_LIST = [
    "rm -rf /", "sudo", "shutdown", "reboot",
    "mkfs", "dd if=", "> /dev/sda",
]

def check_deny_list(command: str) -> str | None:
    for pattern in DENY_LIST:
        if pattern in command:
            return f"Blocked: '{pattern}' is on the deny list"
    return None

闸门2:规则匹配,描述什么时候需要问用户。每条规则指定工具和检查条件

PERMISSION_RULES = [
    {
        "tools": ["read_file", "write_file", "edit_file"],
        "check": lambda args: not (WORKDIR / args.get("path", "")).resolve().is_relative_to(WORKDIR),
        "message": "Access outside workspace",
    },
    {
        "tools": ["bash"],
        "check": lambda args: any(kw in args.get("command", "") for kw in ["rm ", "> /etc/", "chmod 777"]),
        "message": "Potentially destructive command",
    },
]

def check_rules(tool_name: str, args: dict) -> str | None:
    for rule in PERMISSION_RULES:
        if tool_name in rule["tools"] and rule["check"](args):
            return rule["message"]
    return None

闸门3:规则命中后,暂停等用户输入

def ask_user(tool_name: str, args: dict, reason: str) -> str:
    print(f"\n⚠  {reason}")
    print(f"   Tool: {tool_name}({args})")
    choice = input("   Allow? [y/N] ").strip().lower()
    return "allow" if choice in ("y", "yes") else "deny"

三个闸门串在一起,插在工具执行之前

def check_permission(block) -> bool:
    # 闸门 1: 硬拒绝
    if block.name == "bash":
        reason = check_deny_list(block.input.get("command", ""))
        if reason:
            print(f"\n⛔ {reason}")
            return False

    # 闸门 2 + 3: 规则匹配 → 用户审批
    reason = check_rules(block.name, block.input)
    if reason:
        decision = ask_user(block.name, block.input, reason)
        if decision == "deny":
            return False

    return True

# 在 agent_loop 中——s02 的循环只加了一行:
for block in response.content:
    if block.type == "tool_use":
        if not check_permission(block):           # ← 新增
            results.append({... "content": "Permission denied."})
            continue
        output = TOOL_HANDLERS[block.name](**block.input)  # s02 原有
        results.append(...)

这里删除成功是因为环境问题源码中写的是bash + rm而Windows下的删除应该是对应Remove-Item test.txt

扩展检查

前面已经添加了权限的检查。但是每次加一个新的检查,都要修改agent_loop函数,就会导致循环内部挤压

def agent_loop(messages):
    while True:
        # ... LLM call ...
        for block in response.content:
            if block.type != "tool_use":
                continue
            log_to_file(block)          # 加一行
            check_permission(block)     # 加一行
            notify_slack(block)         # 又加一行
            output = execute(block)
            auto_git_add(block)         # 再加一行
            # ... 很快循环就认不出来了

想要扩展的是Agent的行为,但是却改了循环本身,循环应该是一个稳定的核心,扩展应该挂在外面

“挂在循环上,不写进循环里”hook在工具执行前后注入扩展逻辑

Hooks Overview

之前的循环逻辑完全保留。唯一的变动是把check_permission()从循环体内转移到了hook上 ,循环不再直接调用任何检查函数,改为trigger_hooks(“PreToolUse”,block),由注册表决定跑什么

四个事件覆盖一个完整的agent cycle

事件触发时机典型用途
UserPromptSubmit用户输入提交后、进入 LLM 前输入验证、注入上下文
PreToolUse工具执行前权限检查、日志记录
PostToolUse工具执行后副作用(自动 git add 等)、输出检查
Stop循环即将退出时收尾清理(CC 还支持强制续跑)

扩展通过register_hook()添加,循环只调用trigger_hooks()

hook注册表:一个字典,事件名映射到回调列表

HOOKS = {
    "UserPromptSubmit": [],
    "PreToolUse": [],
    "PostToolUse": [],
    "Stop": [],
}

def register_hook(event: str, callback):
    HOOKS[event].append(callback)

def trigger_hooks(event: str, *args):
    for callback in HOOKS[event]:
        result = callback(*args)
        if result is not None:   # 返回值 ≠ None → hook 说"停"
            return result
    return None

PreToolUse的非None返回值会阻止本次工具执行,Stop的非None返回值会强制继续跑。UserPromptSubmit和PostToolUse的返回值未被使用

UserPromptSubmit,用户输入提交后、进入LLM前触发

def context_inject_hook(query: str) -> str | None:
    """Inject current working directory info into every prompt."""
    print(f"\033[90m[HOOK] UserPromptSubmit: working in {WORKDIR}\033[0m")
    return None   # return None = no modification, let prompt through

register_hook("UserPromptSubmit", context_inject_hook)

在主循环中,用户输入后立即触发

query = input("s04 >> ")
trigger_hooks("UserPromptSubmit", query)   # ← 进入 LLM 之前
history.append({"role": "user", "content": query})
agent_loop(history)

PreToolUse/PostToolUse,工具执行前后的hook。将前面的权限检查逻辑现在包装成PreToolUse hook,再加一个日志hook和一个输出提醒

# PreToolUse: 权限检查(s03 的逻辑,从循环移到 hook)
def permission_hook(block):
    if block.name == "bash":
        for pattern in DENY_LIST:
            if pattern in block.input.get("command", ""):
                return "Permission denied by deny list"
    if block.name in ("read_file", "write_file", "edit_file"):
        path = block.input.get("path", "")
        if not (WORKDIR / path).resolve().is_relative_to(WORKDIR):
            choice = input("   Allow? [y/N] ").strip().lower()
            if choice not in ("y", "yes"):
                return "Permission denied by user"
    return None

# PreToolUse: 日志
def log_hook(block):
    print(f"[HOOK] {block.name}(...)")

# PostToolUse: 大文件提醒
def large_output_hook(block, output):
    if len(str(output)) > 100000:
        print(f"[HOOK] ⚠ Large output from {block.name}")

register_hook("PreToolUse", permission_hook)
register_hook("PreToolUse", log_hook)
register_hook("PostToolUse", large_output_hook)

stop,循环即将退出时触发(stop_reason != “tool_use”)

def summary_hook(messages: list) -> str | None:
    """Print a summary when the loop is about to stop."""
    tool_count = sum(1 for m in messages
                     for b in (m.get("content") if isinstance(m.get("content"), list) else [])
                     if isinstance(b, dict) and b.get("type") == "tool_result")
    print(f"\033[90m[HOOK] Stop: session used {tool_count} tool calls\033[0m")
    return None   # return None = allow stop, return string = force continuation

register_hook("Stop", summary_hook)

在agent_loop中退出前触发

if response.stop_reason != "tool_use":
    force = trigger_hooks("Stop", messages)   # ← 退出之前
    if force:
        # hook returned a message → inject it and continue
        messages.append({"role": "user", "content": force})
        continue
    return

循环里只改了一处,从直接调用check_permission(block)改为trigger_hooks(“PreToolUse”,block):

for block in response.content:
    if block.type != "tool_use":
        continue

    # s03: if not check_permission(block): ...
    # s04: hook 替代硬编码
    blocked = trigger_hooks("PreToolUse", block)
    if blocked:
        results.append({"type": "tool_result", "tool_use_id": block.id,
                        "content": str(blocked)})
        continue

    handler = TOOL_HANDLERS.get(block.name)
    output = handler(**block.input) if handler else f"Unknown: {block.name}"

    trigger_hooks("PostToolUse", block, output)

    results.append({"type": "tool_result", "tool_use_id": block.id,
                    "content": output})

四个hook覆盖了agent cycle的关键节点:输入 –> 执行前 –> 执行后 –> 退出。循环只负责调用trigger_hooks(),具体逻辑全在hook回调里。

Agent Loop 是智能体的核心执行骨架,负责维持“模型决策—工具执行—结果反馈”的基本流程,因此不应因附加功能而频繁修改。权限控制、日志、监控、缓存等扩展能力,应通过 Hook、Middleware 或插件接口以可插拔方式接入,从而降低耦合,并保证核心循环的稳定性。

TodoWrite

给Agent一个复杂任务:”把所有Python文件改成snake_case命名,然后跑测试,修好失败”。Agent开始干活,改了3个文件,跑了个测试,发现2个失败,开始修。修着修着就忘了最初是”改成snake_case”,测试失败把注意力全吸走了。对话越长越严重:工具结果不断填满上下文,系统提示的影响力被稀释。一个10步重构,做完1-3步就开始即兴发挥,因为4-10步已经被挤出注意力了。

Todo Overview

保留上一章的最小hook结构,重点看新增的todo_write工具和reminder机制。todo_write本身不做任何实际工作,不能读文件、不能跑命令,只是让Agent在动手之前先理清思路。

dispatch机制(工具调用的分发机制)不变,新工具仍然走TOOL_HANDLERS[block.name]分发。为了演示todo reminder,循环里加了一个计数器:连续3轮没调todo_write就注入一条提醒

todo_write工具,接收一个带状态的列表,保存在当前进程内存中,同时在终端显示进度

CURRENT_TODOS: list[dict] = []

def run_todo_write(todos: list) -> str:
    global CURRENT_TODOS
    CURRENT_TODOS = todos

    lines = ["\n## Current Tasks"]
    for t in CURRENT_TODOS:
        icon = {"pending": " ", "in_progress": "▸", "completed": "✓"}[t["status"]]
        lines.append(f"  [{icon}] {t['content']}")
    print("\n".join(lines))
    return f"Updated {len(CURRENT_TODOS)} tasks"

工具定义和其他5个工具一起加入dispatch map

TOOLS = [
    {"name": "bash",       ...},
    {"name": "read_file",  ...},
    {"name": "write_file", ...},
    {"name": "edit_file",  ...},
    {"name": "glob",       ...},
    # s05: 新增一条
    {"name": "todo_write", "description": "Create and manage a task list ...",
     "input_schema": {
         "type": "object",
         "properties": {
             "todos": {
                 "type": "array",
                 "items": {
                     "type": "object",
                     "properties": {
                         "content": {"type": "string"},
                         "status": {"type": "string", "enum": ["pending", "in_progress", "completed"]},
                     },
                 },
             },
         },
     },
    },
]

TOOL_HANDLERS["todo_write"] = run_todo_write

Nag reminder,模型连续3轮没调todo_write时,自动注入一条提醒

if rounds_since_todo >= 3 and messages:
    messages.append({
        "role": "user",
        "content": "<reminder>Update your todos.</reminder>",
    })
    rounds_since_todo = 0

Agent收到任务后的典型流程:先调todo_write列出所有步骤 -> 做一个步骤,改成in_progress -> 做完改成completed -> 看下一个pending -> 继续。连续3轮没有调用todo_write时,循环会在下一次LLM调用前追加一条reminder。

todo_write不给Agent增加任何执行能力,它增加的时规划能力

Subagent

Agent在修一个bug。它读了30个文件来追踪调用链,中间聊了60轮。messages列表涨到了120条,其中大部分是”追踪调用链”的中间过程,和”修bug”这个最终目标无关,而这些中间过程占着上下文位置,让Agent越来越健忘,记不住最初的问题是什么了

因此Agent需要一个能力:开一个独立的子进程,给它一个独立的消息列表,让它专心做一件事

Subagent Overview

新增一个task工具,调用它时,spawn一个子Agent,拥有全新的messages[],跑自己的循环,结束后只把摘要文本传回给主Agent。对话上下文被丢弃,但文件系统的副作用保留在工作目录中

子Agent的工具受限:有bash/read/write/edit/glob,但没有task,不能递归spawn新的子Agent。子Agent的工具调用仍经过权限hook,安全策略不因上下文隔离而跳过

spawn_subagent,给子Agent一个全新的messages列表,跑自己的循环,只回传结论。这里与agent_loop明显的区别是没有采用while true的模式而是设置了轮次防止subagent无限自主运行

def spawn_subagent(description: str) -> str:
    # 子 Agent 的工具:基础工具,但没有 task(禁止递归)
    sub_tools = [
        {"name": "bash", ...}, {"name": "read_file", ...},
        {"name": "write_file", ...}, {"name": "edit_file", ...},
        {"name": "glob", ...},
    ]
    messages = [{"role": "user", "content": description}]  # 全新 messages[]

    for _ in range(30):  # safety limit
        response = client.messages.create(
            model=MODEL, system=SUB_SYSTEM,
            messages=messages, tools=sub_tools, max_tokens=8000,
        )
        messages.append({"role": "assistant", "content": response.content})
        if response.stop_reason != "tool_use":
            break
        results = []
        for block in response.content:
            if block.type == "tool_use":
                blocked = trigger_hooks("PreToolUse", block)
                if blocked:
                    results.append({... "content": str(blocked)})
                    continue
                handler = SUB_HANDLERS.get(block.name)
                output = handler(**block.input) if handler else f"Unknown"
                trigger_hooks("PostToolUse", block, output)
                results.append({... "content": output})
        messages.append({"role": "user", "content": results})

    # 只返回最后的文本结论,中间过程全部丢弃
    return extract_text(messages[-1]["content"])

主Agent调用时,跟调用其他工具一样

TOOLS = [
    {"name": "bash", ...},
    {"name": "read_file", ...},
    {"name": "write_file", ...},
    {"name": "edit_file", ...},
    {"name": "glob", ...},
    {"name": "todo_write", ...},
    # s06: 新增 task 工具
    {"name": "task",
     "description": "Launch a subagent to handle a complex subtask. Returns only the final conclusion.",
     "input_schema": {"type": "object", "properties": {"description": {"type": "string"}}, "required": ["description"]}},
]

TOOL_HANDLERS["task"] = spawn_subagent

dispatch机制不变,task工具通过TOOL_HANDLERS[block.name]分发。子Agent有独立的SUB_SYSTEM提示,明确要求直接完成任务不要再委派

Skill Loading

项目有一套React组件规范、一份SQL风格指南、一份API设计文档。希望Agent自动遵守这些规范。最直接的想法是全塞进system prompt

SYSTEM = (
    f"You are a coding agent. "
    + open("docs/react-style.md").read()       # 2000 行
    + open("docs/sql-style.md").read()         # 1500 行
    + open("docs/api-design.md").read()        # 3000 行
)

6500行system prompt。Agent每次调用LLM都带着这些文档不管是在改CSS颜色还是修SQL查询。99%的内容和当前任务无关,白白消耗token

Skill Overview

新增load_skill工具,启动时把技能目录注入SYSTEM prompt,运行时多注册一个工具加载完整内容,用到才花token

两层设计

位置时机代价
1. 目录system prompt启动时注入(harness 扫描 skills/)~100 tokens/skill,每轮都带
2. 内容tool_resultAgent 调用 load_skill 时;SKILL.md 可指引后续的 read_file/bash 调用,用于按需访问额外资源~2000 tokens/skill,按需

skills/目录,每个技能一个子目录,包含SKILL.md文件

skills/
  agent-builder/SKILL.md
  code-review/SKILL.md
  mcp-builder/SKILL.md
  pdf/SKILL.md

第一级:启动时注入目录:harness启动时调用_scan_skills()扫描skills/目录,解析每个SKILL.md的YAML frontmatter(name、description),存入SKILL_REGISTRY字典。list_skills()从注册表生成目录,注入SYSTEM prompt。Agent每轮都能看到有哪些技能可用,不花额外的API调用

SKILL_REGISTRY: dict[str, dict] = {}

def _scan_skills():
    if not SKILLS_DIR.exists():
        return
    for d in sorted(SKILLS_DIR.iterdir()):
        if not d.is_dir():
            continue
        manifest = d / "SKILL.md"
        if manifest.exists():
            raw = manifest.read_text()
            meta, body = _parse_frontmatter(raw)
            name = meta.get("name", d.name)
            desc = meta.get("description", raw.split("\n")[0].lstrip("#").strip())
            SKILL_REGISTRY[name] = {"name": name, "description": desc, "content": raw}

_scan_skills()  # runs once at startup

def list_skills() -> str:
    return "\n".join(f"- **{s['name']}**: {s['description']}" for s in SKILL_REGISTRY.values())

def build_system() -> str:
    catalog = list_skills()
    return (
        f"You are a coding agent at {WORKDIR}. "
        f"Skills available:\n{catalog}\n"
        "Use load_skill to get full details when needed."
    )

SYSTEM = build_system()

第二级:load_skill:Agent决定”我需要SQL风格指南”,调用load_skill(“sql-style”)通过注册表查找,不走文件路径,没有路径遍历风险。SKILL.md内容通过tool_result注入,并可通过现有的file和bash工具进一步访问引用的reference/、scripts/或assets/

def load_skill(name: str) -> str:
    skill = SKILL_REGISTRY.get(name)
    if not skill:
        return f"Skill not found: {name}"
    return skill["content"]

SKILL和System Prompt,SKILL不是System Prompt的一部分。System Prompt是永久性的规则而SKILL按当前任务临时加载的专项说明。SKILL会作为一次工具结果进入当前的message后续调用会随历史一起携带,直到上下文压缩、截断或者会话结束

Context Compact上下文压缩

Agent跑着跑着就不动了,有bash、有read、有write能力是够的,但是它读了一个1000行的文件,又读了30个文件,跑了20条命令。每条命令的输出、每个文件的内容,全部堆在messages列表里。上下文窗口是有限的,满了之后API直接拒绝:prompt_too_long,不压缩上下文,Agent根本没法在大项目里干活

Compact Overview

核心设计逻辑:便宜的先跑,贵的后跑

每轮LLM调用前插入三层预处理器(0 API),token仍超过阈值时触发LLM摘要(1 API),API报错时应急裁剪

四层压缩管线

L1:snip_compact-裁掉无关的旧对话

Agent跑了80轮对话,message攒了160条。最前面的”帮我创建hello.py”和当前工作几乎无关了,但全占着位置。

消息树超过50条 –> 保留头部3条(初始上下文)和尾部47条(当前工作),中间裁掉;唯一额外边界条件是,不能把assistant(tool_use)和后面的user(tool_result)拆开

def snip_compact(messages, max_messages=50):
    if len(messages) <= max_messages:
        return messages
    head_end, tail_start = 3, len(messages) - (max_messages - 3)
    if head_end > 0 and _message_has_tool_use(messages[head_end - 1]):
        while head_end < len(messages) and _is_tool_result_message(messages[head_end]):
            head_end += 1
    if (tail_start > 0 and tail_start < len(messages)
            and _is_tool_result_message(messages[tail_start])
            and _message_has_tool_use(messages[tail_start - 1])):
        tail_start -= 1
    snipped = tail_start - head_end
    placeholder = {"role": "user", "content": f"[snipped {snipped} messages from conversation middle]"}
    return messages[:head_end] + [placeholder] + messages[tail_start:]

这里还特别处理了工具的调用
Agent:我要读取文件 <– tool_use
工具:这是文件内容 <– tool_result
这两条必须一起保留,不能只保留其中一条,如果截断位置刚好落在它们中间,函数会移动截断位置。防止拆开头部附近的工具调用配对和防止尾部从孤立的工具结果开始有了一下操作。

if head_end > 0 and _message_has_tool_use(messages[head_end - 1]):
        while head_end < len(messages) and _is_tool_result_message(messages[head_end]):
            head_end += 1

头部以工具调用为起点补工具结果这里就是头部判断如果前面一条是工具调用并且这一条是调用结果那么两者要在一起进行保留。这里假设后面的消息全部是 tool_resulthead_end 可能一直增加这会超出列表范围导致异常所以为了保护循环后续的每次判断设置了head_end < len(messages)

if (tail_start > 0 and tail_start < len(messages)
            and _is_tool_result_message(messages[tail_start])
            and _message_has_tool_use(messages[tail_start - 1])):
        tail_start -= 1

尾部是以工具结果为尾点补工具调用因此只需要判断前位有一个是工具调用即可

裁掉的是消息本身,但是消息数量变少,不一定代表上下文真正变小很多

L2:micro_compact-旧工具结果占位

旧结果占位

Agent连续读了10个文件。第1-7次的完整内容还躺在上下文里,早就不需要了,但占着大量空间

只保留最近3条tool_result的完整内容,更旧的替换为一行占位符

KEEP_RECENT_TOOL_RESULTS = 3

def micro_compact(messages):
    tool_results = collect_tool_result_blocks(messages)
    if len(tool_results) <= KEEP_RECENT_TOOL_RESULTS:
        return messages
    for _, _, block in tool_results[:-KEEP_RECENT_TOOL_RESULTS]:
        if len(block.get("content", "")) > 120:
            block["content"] = "[Earlier tool result compacted. Re-run if needed.]"
    return messages

旧结果清掉了,但单条新结果可能就有500KB,一个cat大文件的输出就能打满上下文

L3:tool_result_budget-大结果落盘

大结果落盘

模型一次读了5个大文件,单条user消息里所有tool_result加起来500KB。统计最后一条user消息里所有tool_result的总大小。超过200KB -> 按大小排序,从最大的开始落盘到.task_outputs/tool-results/,上下文只留<persisted-output>标记 + 前2000字符预览。模型看到标记后知道完整内容在磁盘上,需要时可以重新读。

def tool_result_budget(messages, max_bytes=200_000):
    last = messages[-1]
    blocks = [(i, b) for i, b in enumerate(last["content"])
              if b.get("type") == "tool_result"]
    total = sum(len(str(b.get("content", ""))) for _, b in blocks)
    if total <= max_bytes:
        return messages
    ranked = sorted(blocks, key=lambda p: len(str(p[1].get("content", ""))), reverse=True)
    for idx, block in ranked:
        if total <= max_bytes:
            break
        block["content"] = persist_large_output(block["tool_use_id"], str(block["content"]))
        total = recalculate_total(blocks)
    return messages

L4:compact_history-LLM全量摘要

前三层都是纯文件/结构操作,0 API调用,但也无法理解对话内容,上下文可能仍然太大

LLM 全量摘要

前三层全跑完了,但在超大项目中连续工作30min后,token仍然超过阈值

三步流程:

1.保存transcript:完整对话写入.transcripts/,JSONL格式。transcript保留了可恢复记录,但模型的活跃上下文里只剩摘要。对模型当下推理来说,细节已经不在上下文中了。

2.LLM生成摘要:把对话历史发给LLM,要求保留当前目标、重要发现、已改文件、剩余工作、用户约束等关键信息。

3.替换消息列表:所有旧消息被替换为一条摘要。

def compact_history(messages):
    transcript_path = write_transcript(messages)  # 先保存完整对话
    summary = summarize_history(messages)          # LLM 生成摘要
    return [{"role": "user",
             "content": f"[Compacted]\n\n{summary}"}]

应急:reactive_compact

有时候API还是返回prompt_too_long(413),上下文增长速度快于压缩触发速度时。

这时触发reactive_compact触发方式比compact_history更激进,但压缩策略更温和,保留最近约5条原始消息,只总结较早历史。同样避免留下孤立tool_result

def reactive_compact(messages):
    transcript = write_transcript(messages)
    tail_start = max(0, len(messages) - 5)
    if (tail_start > 0 and tail_start < len(messages)
            and _is_tool_result_message(messages[tail_start])
            and _message_has_tool_use(messages[tail_start - 1])):
        tail_start -= 1
    summary = summarize_history(messages[:tail_start])
    return [{"role": "user",
             "content": f"[Reactive compact]\n\n{summary}"}, *messages[tail_start:]]

reactive compact有重试上限。再失败就抛出异常,不无限循环

def agent_loop(messages):
    reactive_retries = 0
    while True:
        # 三个预处理器(0 API 调用)
        # 顺序:budget 先跑,确保大内容落盘后再做占位和裁剪
        messages[:] = tool_result_budget(messages)    # L3: 大结果落盘
        messages[:] = snip_compact(messages)          # L1: 裁中间
        messages[:] = micro_compact(messages)         # L2: 旧结果占位

        # 还不够?LLM 摘要(1 API 调用)
        if estimate_token_count(messages) > THRESHOLD:
            messages[:] = compact_history(messages)

        try:
            response = client.messages.create(...)
        except PromptTooLongError:
            if reactive_retries < MAX_REACTIVE_RETRIES:
                messages[:] = reactive_compact(messages)  # 应急
                reactive_retries += 1
                continue
            raise  # 超过重试上限,抛出异常
        # ... 工具执行 ...

        # compact 工具:模型主动调用时触发 compact_history
        if block.name == "compact":
            messages[:] = compact_history(messages)
            results.append({..., "content": "[Compacted. History summarized.]"})
            messages.append({"role": "user", "content": results})
            break  # 结束当前 turn,用压缩后的上下文开始新一轮

落地下来L3在L2前面,因为micro会把旧的大tool_result替换成一行占位符,budget必须在那之前把完整内容落盘。

Memory

前面的autoCompact(L4)会把当前目标、剩余工作、用户约束写进摘要,但细节会丢失:”用tab缩进不要用空格”可能被简化成”用户有代码风格偏好”。新开一个会话,连摘要也没了。

LLM没有持久状态,所有信息都在上下文窗口里。上下文满了要压缩,压缩就有损。需要一层不参与压缩、跨会话保留的存储。

Memory Overview

前面的压缩管线保留,存储选文件系统:.memory/目录下,每个记忆一个.md文件,带YAML frontmatter(name / description / type)文件多了需要索引:MEMORY.md一行一个链接注入SYSTEM

关键设计:索引常驻SYSTEM prompt(可被 prompt cache缓存),文件内容按需注入到当前user turn(按filename/description匹配当前对话,不破坏cache)。写入由每轮结束后的提取器完成:用户显示说”记住”或表达稳定偏好时,提取器会保存为记忆。文件积累多了,定期整理去重。

四类记忆,各有用途

类型回答什么示例
user你是谁“用 tab 不用空格”
feedback怎么做事“别 mock 数据库”
project正在发生什么“auth 重写是合规驱动”
reference东西在哪找“pipeline bug 在 Linear INGEST”
Memory Subsystems

存储md文件+索引

每个记忆是一个.md文件,YAML frontmatter记录元数据

---
name: user-preference-tabs
description: User prefers tabs for indentation
type: user
---

User prefers using tabs, not spaces, for indentation.
**Why:** Consistency with existing codebase conventions.
**How to apply:** Always use tabs when writing or editing files.

MEMORY.md是索引,一行一个链接

- [user-preference-tabs](user-preference-tabs.md) — User prefers tabs for indentation

写入新记忆时自动重建索引

def write_memory_file(name, mem_type, description, body):
    slug = name.lower().replace(" ", "-")
    filepath = MEMORY_DIR / f"{slug}.md"
    filepath.write_text(
        f"---\nname: {name}\ndescription: {description}\ntype: {mem_type}\n---\n\n{body}\n"
    )
    _rebuild_index()

两条路径加载

路径一:索引常驻SYSTEM。build_system()在每次用户请求开始时读取MEMORY.md,把记忆清单注入。记忆提取和整理只在本轮结束时触发,因此同一轮用户请求中不需要重复重建SYSTEM

路径二:相关记忆按需注入。每次用户请求开始时,load_memories()把最近对话和记忆目录(name + description)一起发给LLM做一次轻量side-query,选出相关的文件名,再读文件内容临时注入到当前user turn。最多五条控制开销。

def select_relevant_memories(messages, max_items=5):
    files = list_memory_files()
    if not files:
        return []

    # Build catalog: "0: user-preference-tabs — User prefers tabs..."
    catalog = "\n".join(f"{i}: {f['name']} — {f['description']}" for i, f in enumerate(files))

    response = client.messages.create(model=MODEL, messages=[{"role": "user",
        "content": f"Select relevant memory indices. Return JSON array.\n\n"
                   f"Recent conversation:\n{recent}\n\nMemory catalog:\n{catalog}"}],
        max_tokens=200)
    text = extract_text(response.content).strip()
    indices = json.loads(re.search(r'\[.*?\]', text).group())
    return [files[i]["filename"] for i in indices if 0 <= i < len(files)]

如果side-query失败(在正式回答用户之前,harness 额外向 LLM 发起的一次内部小查询。让LLM选择相关记忆),降级到关键词匹配name + description

写入:每轮结束后提取

用户不会每次都说”记住这个”。偏好通常散落在正常对话中:”用tab比空格好”、”以后都用单引号”。

extract_memories()在每轮结束时运行,条件是模型停止且没有tool_use(说明对话告一段落)

# In agent_loop:
if response.stop_reason != "tool_use":
    extract_memories(pre_compress)   # 从压缩前快照提取新记忆
    consolidate_memories()       # 检查是否需要整理
    return

提取前先检查已有记忆,避免重复。提取prompt要求LLM返回{name,type,description,body}的JSON数组,只有确实由新信息时才写文件

def extract_memories(messages):
    dialogue = format_recent_messages(messages[-10:])
    existing = "\n".join(f"- {m['name']}: {m['description']}" for m in list_memory_files())

    prompt = (
        "Extract user preferences, constraints, or project facts.\n"
        "Return JSON array: [{name, type, description, body}].\n"
        "If nothing new or already covered, return [].\n\n"
        f"Existing memories:\n{existing}\n\nDialogue:\n{dialogue[:4000]}"
    )
    # ... parse response, write files ...

整理:低频合并去重

记忆文件会积累。consolidate_memories()在文件数达到阈值时触发,让LLM去重、合并矛盾、淘汰过时记忆:

CONSOLIDATE_THRESHOLD = 10

def consolidate_memories():
    files = list_memory_files()
    if len(files) < CONSOLIDATE_THRESHOLD:
        return  # 太少,不值得整理
    # Send all memories to LLM, get back deduplicated list
    # Replace all files with consolidated results

CC把这个过程叫Dream,实际有四门门控:时间间隔、扫描节流、会话数、文件锁。这里简化为文件数阈值

Memory适合保存什么

Memory保存跨会话仍有用的信息:用户偏好、反复出现的反馈、项目背景、常用入口和排查线索。它关注”以后还会用到什么”,并通过索引+按需加载把这些信息带回当前对话。

session memory关注同一会话内的连续性:compact之后,当前会话还需要保留哪些上下文。两者配合使用:Memory管长期知识,session memory管当前会话的压缩连续。

System Prompt运行时组装

前面system prompt都是一行硬编码,当只有bash、read、write三个工具。但到后面Agent已经有记忆、有压缩、有技能加载。prompt该提的能力越来越多,就会面临三个问题

1.换项目要重写整个prompt,不知道哪些该改、哪些该留
2.修改一处可能影响全局,加一段工具描述可能跟前面的指令冲突
3.每次请求都带全部内容,即使当前对话用不到某些段落也浪费token

System prompt应该是运行时根据当前状态组装的配置:哪些工具启用、哪些上下文可见、哪些记忆相关、哪些内容必须保持稳定以命中prompt cache。

System Prompt Overview

把硬编码的SYSTEM拆成独立段落(section),运行时根据真实状态按需拼接,缓存结果避免重复组装。

四个section,两种加载策略

Section加载策略内容判断依据
identity始终你是谁、怎么做事始终存在
tools始终可用工具列表enabled_tools
workspace始终工作目录始终存在
memory按需相关记忆内容.memory/MEMORY.md 是否存在

关键设计:section是否加载取决于真实状态(工具是否存在、文件是否存在),不是消息里的关键词。

PROMPT_SECTIONS:分段定义

把一大段字符串拆成字典,每个key是一个主题

PROMPT_SECTIONS = {
    "identity": "You are a coding agent. Act, don't explain.",
}

每个section独立维护。修改tools不影响identity,新增memory不动workspace。

assemble_system_prompt:按需拼接

不是所有section每次都需要。当前没有记忆文件,加载memory section只是浪费token。根据context的真实状态决定加载哪些

def assemble_system_prompt(context: dict) -> str:
    sections = []

    # 始终加载
    sections.append(PROMPT_SECTIONS["identity"])

    # 从 context 动态获取 tools 和 workspace
    tools = ", ".join(context.get("enabled_tools", []))
    if tools:
        sections.append(f"Available tools: {tools}.")
    sections.append(f"Working directory: {context.get("workspace", WORKDIR)}")

    # 按需加载 — 基于真实状态,不是关键词
    memories = context.get("memories", "")
    if memories:
        sections.append(f"Relevant memories:\n{memories}")

    return "\n\n".join(sections)

“始终加载”的是每轮都需要的:身份、工具、工作目录。“按需加载”的只是在特定条件下才有用(为什么不全加载?token有成本,信息越少LLM越专注,一些无关指令是噪音)

get_system_prompt:缓存避免重复拼接

上下文没变时,重新拼接是浪费。用确定性序列化检测变化,命中缓存直接返回

def get_system_prompt(context: dict) -> str:
    global _last_context_key, _last_prompt
    key = json.dumps(context, sort_keys=True, ensure_ascii=False, default=str)
    if key == _last_context_key and _last_prompt:
        return _last_prompt
    _last_context_key = key
    _last_prompt = assemble_system_prompt(context)
    return _last_prompt

通过把context这个字典,转换成一个尽量稳定、可比较的字符串,然后用这个字符串判断,这次的context和上一次是不是一样

context = {
    "tools": ["read", "write"],
    "workspace": "/project"
}
序列化后
{"tools": ["read", "write"], "workspace": "/project"}
设置了sort_keys=True
a = {
    "tools": ["read"],
    "workspace": "/project"
}

b = {
    "workspace": "/project",
    "tools": ["read"]
}
序列化后a=b

用json.dumps而不是hash():Python内置hash()有进程随机化,不适合做稳定cache key,而且遇到list/dict会报错

context:真实状态,不是关键词猜测

context反映当前运行态的真实状态

def update_context(context: dict, messages: list) -> dict:
    memories = ""
    if MEMORY_INDEX.exists():
        content = MEMORY_INDEX.read_text().strip()
        if content:
            memories = content
    return {
        "enabled_tools": list(TOOL_HANDLERS.keys()),
        "workspace": str(WORKDIR),
        "memories": memories,
    }

enabled_tools列出实际注册的工具。memories检查.memory/MEMORY.md是否存在。section加载基于这些真实状态,不在消息里搜关键词。

def agent_loop(messages: list, context: dict):
    system = get_system_prompt(context)
    while True:
        response = client.messages.create(
            model=MODEL, system=system, messages=messages,
            tools=TOOLS, max_tokens=8000)
        # ... 工具执行 ...
        context = update_context(context, messages)
        system = get_system_prompt(context)

每轮循环开头拿一次system prompt。context变了就重新组装,没变就返回缓存。

Error Recovery

Agent跑着跑着报错了。Agent崩溃了,它没有重试,没有换模型,没有减少上下文-直接崩溃

生产环境中API错误是常态。三种最常见的故障模式:输出被截断(模型话说一半token用完了)、上下文超限(压缩后还是太长)、临时故障(429限流/529过载)。一个不处理错误的Agent就像一个一碰就熄火的车。

Error Recovery Overview

将LLM调用包裹在try/except里,根据错误类型走不同的恢复路径。恢复后continue回到循环开头重新调用LLM

三种最常见的恢复模式

模式触发恢复动作
输出截断max_tokens升级 8K→64K / 续写提示
上下文超限prompt_too_longreactive compact → 重试
临时故障429 / 529指数退避 + 抖动,连续 529 可切换备用模型

路径1:输出被截断

模型话说一半,max_tokens用完了。默认8000 token不够它输出完整回答。第一次发生时,直接把max_tokens从8K升级到64K,重试同一 请求-此时不追加截断输出到messages,保持原始请求不变。如果64K还是不够,才保存截断输出并注入续写提示让模型接着刚才的话继续说最多3次

if response.stop_reason == "max_tokens":
    # First escalation: don't append truncated output, retry same request
    if not state.has_escalated:
        max_tokens = ESCALATED_MAX_TOKENS
        state.has_escalated = True
        continue  # messages unchanged, same request with more tokens
    # 64K still truncated: save output + continuation prompt
    messages.append({"role": "assistant", "content": response.content})
    if state.recovery_count < MAX_RECOVERY_RETRIES:
        messages.append({"role": "user", "content":
            "Output token limit hit. Resume directly — "
            "no apology, no recap. Pick up mid-thought."})
        state.recovery_count += 1
        continue
    return  # still truncated after 3 continuations
# Normal: append after max_tokens check
messages.append({"role": "assistant", "content": response.content})

路径2:上下文超限

LLM说”你的上下文太长了”,虽然通过前面的压缩技术全跑过了,还是超了。

触发reactive compact-比auto compact更激进。只保留最后5条消息模拟压缩效果。但如果压缩过一次还是超限,只能退出-再压缩也不会变小

except PromptTooLongError:
    if not state.has_attempted_reactive_compact:
        messages[:] = reactive_compact(messages)
        state.has_attempted_reactive_compact = True
        continue
    return  # 压缩过了还是超限,只能退出

路径3:临时故障

网络抖动、429限流、529过载-这些不是bug,是分布式系统的常态。

429和529统一走指数退避+抖动:第一次等0.5秒,第二次等1秒,第三次等2秒,最多10次。加随机抖动让并发请求不在同一时刻重试。连续3次529过载->切换到备用模型

def retry_delay(attempt, retry_after=None):
    if retry_after:
        return retry_after
    base = min(500 * (2 ** attempt), 32000) / 1000
    return base + random.uniform(0, base * 0.25)

def with_retry(fn, state, max_retries=10):
    for attempt in range(max_retries):
        try:
            return fn()
        except (RateLimitError, OverloadedError):
            delay = retry_delay(attempt)
            time.sleep(delay)
            if is_overloaded:
                state.consecutive_529 += 1
                if state.consecutive_529 >= 3 and FALLBACK_MODEL:
                    state.current_model = FALLBACK_MODEL
    raise MaxRetriesExceeded()

退避公式:min(500 * 2 ** attempt,32000) + random(0~25%)。如果服务器返回Retry-After header,优先用那个值。

def agent_loop(messages, context):
    system = get_system_prompt(context)
    state = RecoveryState()
    max_tokens = 8000

    while True:
        try:
            response = with_retry(
                lambda: client.messages.create(
                    model=state.current_model, system=system,
                    messages=messages, tools=TOOLS,
                    max_tokens=max_tokens),
                state)
        except Exception as e:
            if is_prompt_too_long_error(e):
                if not state.has_attempted_reactive_compact:
                    messages[:] = reactive_compact(messages)
                    state.has_attempted_reactive_compact = True
                    continue
                return
            log_error(e)
            return

        # max_tokens check BEFORE appending to messages
        if response.stop_reason == "max_tokens":
            if not state.has_escalated:
                max_tokens = 64000
                state.has_escalated = True
                continue  # retry same request, messages unchanged
            # save truncated output + continuation prompt
            messages.append({"role": "assistant", "content": response.content})
            messages.append({"role": "user", "content": CONTINUATION_PROMPT})
            continue
        # Normal completion
        messages.append({"role": "assistant", "content": response.content})

        if response.stop_reason != "tool_use":
            return
        # ... tool execution ...

外层try/except捕获API异常(prompt_too_long等),with_retry处理瞬态错误(429/529),stop_reason检查处理截断。三种恢复机制各管各的错误类型。

Task System

Agent接到一个项目:搭建数据库、写API、加测试,用TodoWrite列了一张清单,然后开始写API,写到一半发现没数据库表,回头补;加测试时发现API接口签名又变了

盖房子不能先盖屋顶再打地基。任务之间有先后。任务依赖应该形成有向无环图。前面的TodoWrite是当前任务的执行清单,保存在会话内存中。这里需要的是任务系统:每个任务是一个JSON文件,任务之间有blockedBy依赖,跨会话持久化在磁盘上。

这里的依赖指的是在运行当前任务之前,该任务被哪些任务阻塞

Task(
    id="T4",
    subject="部署",
    description="部署最终项目",
    status="pending",
    owner="agent_deploy",
    blockedBy=["T2", "T3"]
)
      T2 ──┐
           ├──→ T4
      T3 ──┘
只有当T2、T3都完成T4才会开始
Task System Overview

工作原理

Task DAG

Task:数据结构

每个任务是一个JSON文件,存于.tasks/目录

@dataclass
class Task:
    id: str
    subject: str
    description: str
    status: str          # pending | in_progress | completed
    owner: str | None    # Agent 名(多 Agent 场景)
    blockedBy: list[str] # 依赖的任务 ID 列表

依赖的任务ID 用 timestamp + random hex生成,简单够用。

create_task:创建任务

def create_task(subject: str, description: str = "",
                blockedBy: list[str] | None = None) -> Task:
    task = Task(
        id=f"task_{int(time.time())}_{random_hex(4)}",
        subject=subject, description=description,
        status="pending", owner=None,
        blockedBy=blockedBy or [],
    )
    save_task(task)
    return task

创建时自动save_task到.tasks/{id}.json。blockedBy声明依赖,比如”写API”的blockedBy是[“task_schema”]

can_start:依赖检查

一个任务只能在它的blockedBy全部completed之后才能开始

def can_start(task_id: str) -> bool:
    task = load_task(task_id)
    for dep_id in task.blockedBy:
        if not _task_path(dep_id).exists():
            return False  # missing dependency = blocked
        dep = load_task(dep_id)
        if dep.status != "completed":
            return False
    return True

can_start是claim_task的前置检查:blockedBy里有任何一个不是completed,就不能认领。不存在的依赖视为blocked,避免引用错误ID时崩溃

claim_task:认领任务

Agent开始做一个任务时,调用claim_task:设置owner,状态从pending -> in_progress。owner字段记录谁在做个任务,多Agent场景下防止重复认领

def claim_task(task_id: str, owner: str = "agent") -> str:
    task = load_task(task_id)
    if task.status != "pending":
        return f"Task {task_id} is {task.status}, cannot claim"
    if not can_start(task_id):
        deps = [d for d in task.blockedBy
                if load_task(d).status != "completed"]
        return f"Blocked by: {deps}"
    task.owner = owner
    task.status = "in_progress"
    save_task(task)
    return f"Claimed {task_id} ({task.subject})"

如果任务已被别人认领(status != “pending”),或者依赖没完成(can_start返回False),拒绝认领。

complete_task:完成与解锁

任务做完后,设为completed。同时扫描所有其他任务,找出刚刚被解锁的下游任务:

def complete_task(task_id: str) -> str:
    task = load_task(task_id)
    task.status = "completed"
    save_task(task)
    # 找出被解锁的下游任务
    unblocked = [t.subject for t in list_tasks()
                 if t.status == "pending" and t.blockedBy
                 and can_start(t.id)]
    msg = f"Completed {task_id} ({task.subject})"
    if unblocked:
        msg += f"\nUnblocked: {', '.join(unblocked)}"
    return msg

完成”schema”后,”endpoints”和”docs”的can_start返回True,它们可以开始。

get_task:查看完整细节

list_tasks只显示一行摘要。get_task返回完整的任务JSON,包括description和依赖细节。跨会话恢复时,Agent需要读取完整描述才能继续工作

def get_task(task_id: str) -> str:
    task = load_task(task_id)
    return json.dumps(asdict(task), indent=2)

状态机:两个动作,三个状态

pending ──claim──→ in_progress ──complete──→ completed

这里的claim/complete是动作,pending/in_progress/completed是状态:

  • claim_task:pending -> in_progress。设置owner,开始工作。
  • complete_task:in_progress -> completed。把任务标记为完成,并解锁下游。

CC没有in_progress -> pending的release路径。如果teammate终止或shutdown,CC会把它未完成的任务unassign(清除owner),并将status重置为pending,方便其他agent重新认领。

# 创建有依赖的任务
schema = create_task("setup database schema")
endpoints = create_task("create API endpoints", blockedBy=[schema.id])
tests = create_task("write tests", blockedBy=[endpoints.id])
docs = create_task("write docs", blockedBy=[schema.id])

# Agent 认领第一个可做的任务
claim_task(schema.id)       # ✓ Claimed (无依赖)
complete_task(schema.id)    # ✓ Completed → 解锁 endpoints, docs

claim_task(endpoints.id)    # ✓ Claimed (schema 已完成)
complete_task(endpoints.id) # ✓ Completed → 解锁 tests

claim_task(docs.id)         # ✓ Claimed (schema 已完成)
complete_task(docs.id)      # ✓ Completed

claim_task(tests.id)        # ✓ Claimed (endpoints 已完成)
complete_task(tests.id)     # ✓ Completed

单Agent只能顺序完成task,但是当多Agent的情况下endpoints和docs可以同时进行。每个create_task写一个JOSN文件,每个claim_task/complete_task更新文件。跨会话时,.tasks/目录还在,Agent读文件就能恢复进度。

Background Tasks

类比洗衣机,把衣服扔进去,按下启动,然后去别的事情。30min后洗衣机提醒洗好了衣服,这30分钟不需要干等而是可以干自己的事情。

Agent的bash工具也一样。pip install torch要等10min,npm run build 要3min。这些命令一跑,Agent就在等bash工具返回,没法利用这段时间处理别的任务。读文件是毫秒级,不需要等待。git status一秒钟内返回,不需要等待。但是npm install?分钟级。Agent等10min什么都不做,而LLM按token计费,空转就是浪费。

Background Tasks Overview

采用慢操作丢后台,agent继续处理-后台线程跑命令,完成后注入通知。

should_run_background:显示请求优先,启发式兜底

模型通过bash工具的run_in_background参数显示请求后台执行。如果模型没指定,这里用关键词启发式兜底

def is_slow_operation(tool_name: str, tool_input: dict) -> bool:
    """Fallback heuristic: commands likely to take > 30s."""
    if tool_name != "bash":
        return False
    cmd = tool_input.get("command", "").lower()
    slow_keywords = ["install", "build", "test", "deploy", "compile",
                     "docker build", "pip install", "npm install",
                     "cargo build", "pytest", "make"]
    return any(kw in cmd for kw in slow_keywords)

def should_run_background(tool_name: str, tool_input: dict) -> bool:
    """Model explicit request takes priority; fallback to heuristic."""
    if tool_input.get("run_in_background"):
        return True
    return is_slow_operation(tool_name, tool_input)

CC的bash工具schema里有run_in_background:boolean参数。模型自己决定哪些命令丢后台,不靠关键词猜。

start_background_task:后台执行与生命周期

把工具调用包装成worker函数,扔到daemon线程里执行。每个后台执行任务有唯一ID,状态存在background_tasks字典里

_bg_counter = 0
background_tasks: dict[str, dict] = {}   # bg_id → {tool_use_id, command, status}
background_results: dict[str, str] = {}   # bg_id → output
background_lock = threading.Lock()

def start_background_task(block) -> str:
    """Run tool in a daemon thread. Returns background task ID."""
    global _bg_counter
    _bg_counter += 1
    bg_id = f"bg_{_bg_counter:04d}"

    def worker():
        result = execute_tool(block)
        with background_lock:
            background_tasks[bg_id]["status"] = "completed"
            background_results[bg_id] = result

    with background_lock:
        background_tasks[bg_id] = {
            "tool_use_id": block.id,
            "command": block.input.get("command", ""),
            "status": "running",
        }
    thread = threading.Thread(target=worker, daemon=True)
    thread.start()
    return bg_id

返回bg_id而不是只返回[Running in background…]。daemon=True确保Agent进程退出时线程跟着退出。

collect_background_results:通知收集

后台任务完成后,收集结果并格式化为 <task_notification>通知:

def collect_background_results() -> list[str]:
    """Collect completed results as task_notification messages."""
    with background_lock:
        ready_ids = [bid for bid, task in background_tasks.items()
                     if task["status"] == "completed"]
    notifications = []
    for bg_id in ready_ids:
        with background_lock:
            task = background_tasks.pop(bg_id)
            output = background_results.pop(bg_id, "")
        notifications.append(
            f"<task_notification>\n"
            f"  <task_id>{bg_id}</task_id>\n"
            f"  <status>completed</status>\n"
            f"  <command>{task['command']}</command>\n"
            f"  <summary>{output[:200]}</summary>\n"
            f"</task_notification>")
    return notifications

通知不复用原始tool_use_id。原始tool call已经用占位tool_result回复了,后台完成是独立事件,用task_notification格式注入。这符合Message API的工具配对语义:一个tool_use只对应一个tool_result。

agent_loop里,工具执行分两条路,通知和结果合并为一条user消息

results = []
for block in response.content:
    if block.type != "tool_use":
        continue
    if should_run_background(block.name, block.input):
        bg_id = start_background_task(block)
        results.append({"type": "tool_result",
            "tool_use_id": block.id,
            "content": f"[Background task {bg_id} started] "
                       f"Result will be available when complete."})
    else:
        output = execute_tool(block)
        results.append({"type": "tool_result",
            "tool_use_id": block.id, "content": output})

# 通知和工具结果合入同一条 user 消息
user_content = []
bg_notifications = collect_background_results()
if bg_notifications:
    for notif in bg_notifications:
        user_content.append({"type": "text", "text": notif})
user_content.extend(results)
messages.append({"role": "user", "content": user_content})

慢操作先回一个带bg_id的占位tool_result,LLM知道这个命令还在跑,可以先做别的事。后台完成后,通知作为独立text block和当前轮的tool_result一起组成user消息。

Cron Scheduler-按时间表生产工作

前面让Agent能后台执行慢操作,但所有操作仍然是你手动触发的。你说一句,Agent动一下。”每天早上9点跑测试”、”每30分钟检查CI状态”,这些周期性任务不该需要人每次来推。

Cron Scheduler Overview

通过独立的cron调度线程,每秒检查一次,时间到了把让任务塞进cron_queue;再由queue processor在Agent空闲时自动交付。

手动 vs 定时

手动触发定时触发
触发者用户输入调度线程
触发时机随时cron 表达式指定
需要人参与否(调度器自动入队,空闲时自动交付)
持久性durable 跨重启

四层模型

Cron调度分四层:

1.Scheduler:daemon线程,每秒轮询吗,判断时间到了没有
2.Queue:cron_queue,调度线程写入已触发任务
3.Queue Processor:发现队列非空且Agent空闲,启动一轮agent_loop
4.Consumer:agent_loop从队列消费,注入到messages

CronJob:数据结构

每个cron任务是一个CronJob对象

@dataclass
class CronJob:
    id: str
    cron: str        # "0 9 * * *" (五段式 cron 表达式)
    prompt: str      # 触发时注入给 Agent 的消息
    recurring: bool  # True=周期性,False=一次性
    durable: bool    # True=写磁盘,跨会话保留

Cron表达式,五段式

分钟  小时  日  月  星期
  *    *   *   *   *      每分钟
  0    9   *   *   *      每天早上 9:00
 */5    *   *   *   *      每 5 分钟
  0    9   *   *  1-5     工作日早上 9:00

cron_matches:五段式匹配

标准cron语义:分钟、小时、月必须全部匹配;日(DOM)和星期(DOW)同时被约束任一匹配即可

def cron_matches(cron_expr: str, dt: datetime) -> bool:
    fields = cron_expr.strip().split()
    if len(fields) != 5:
        return False
    minute, hour, dom, month, dow = fields
    dow_val = (dt.weekday() + 1) % 7  # Python Monday=0 → cron Sunday=0

    m = _cron_field_matches(minute, dt.minute)
    h = _cron_field_matches(hour, dt.hour)
    dom_ok = _cron_field_matches(dom, dt.day)
    month_ok = _cron_field_matches(month, dt.month)
    dow_ok = _cron_field_matches(dow, dow_val)

    if not (m and h and month_ok):
        return False
    # DOM and DOW: both constrained → either matching is enough (OR)
    dom_unconstrained = dom == "*"
    dow_unconstrained = dow == "*"
    if dom_unconstrained and dow_unconstrained:
        return True
    if dom_unconstrained:
        return dow_ok
    if dow_unconstrained:
        return dom_ok
    return dom_ok or dow_ok

独立调度线程:每秒轮询

调度器跑在独立的daemon线程里,不依赖agent_loop是否在执行。单个job异常不会杀掉整个线程

def cron_scheduler_loop():
    while True:
        time.sleep(1)
        now = datetime.now()
        minute_marker = now.strftime("%Y-%m-%d %H:%M")
        with cron_lock:
            for job in list(scheduled_jobs.values()):
                try:
                    if cron_matches(job.cron, now):
                        if _last_fired.get(job.id) != minute_marker: #cron只精确到分钟,防止一分钟重复执行60次
                            cron_queue.append(job)
                            _last_fired[job.id] = minute_marker
                        if not job.recurring:
                            scheduled_jobs.pop(job.id, None)
                            if job.durable:
                                save_durable_jobs()
                except Exception as e:
                    print(f"[cron error] {job.id}: {e}")
  • 独立于agent_loop:即使agent_loop没在跑,调度器也在后台检查时间
  • date-aware minute_marker:用”YYYY-MM-DD HH:MM”防止同一分钟重复触发,同时不会在第二天跳过
  • 单 job try/except:一个坏job不会拖垮整个调度线程
  • 一次性任务:触发后自动从scheduled_jobs里删除

Queue Processor + agent_loop:交付端

queue processor不检查时间,只负责在队列有任务且Agent空闲时拉起一轮执行

def queue_processor_loop():
    while True:
        time.sleep(0.2)
        if not has_cron_queue():
            continue
        if not agent_lock.acquire(blocking=False):
            continue
        try:
            if has_cron_queue():
                run_agent_turn_locked()
        finally:
            agent_lock.release()

agent_loop也不负责检查时间,它只从cron_queue里拿已触发的任务,注入到messages里

fired = consume_cron_queue()
for job in fired:
    messages.append({"role": "user",
                     "content": f"[Scheduled] {job.prompt}"})

生产者(调度线程)、交付者(queue processor)和消费者(agent_loop)通过cron_queue、cron_lock、agent_lock解耦

校验:防止坏cron杀掉调度器

schedule_job在注册前校验cron表达式,非法的直接返回错误

def schedule_job(cron, prompt, recurring=True, durable=True):
    err = validate_cron(cron)
    if err:
        return err
    # ... register job

从磁盘加载durable job时也会跳过非法表达式,避免单个坏任务拖垮启动

Durable vs Session-only

  • Durable:任务定义写进 .scheduled_tasks.json。Agent重启后加载文件,恢复任务
  • Session-only:只在内存里。Agent关闭就没了

重要前提:cron 调度器必须在 Agent 进程内跑。进程关闭,调度也停。Durable 只意味着任务定义跨重启保留,下次 Agent 启动时调度器才会发现”该触发了”并触发。如果需要”即使应用关闭也能定时跑”,请用系统 crontab 或 systemd timer。

  1. 启动时:
    load_durable_jobs() → 从 .scheduled_tasks.json 恢复持久化任务
    Thread(cron_scheduler_loop, daemon=True).start() → 调度线程开始轮询
    Thread(queue_processor_loop, daemon=True).start() → 队列处理器等待交付
  2. 注册任务:
    schedule_cron(cron=”*/2 * * * *”, prompt=”run date”, durable=True)
    → CronJob 写入 scheduled_jobs + .scheduled_tasks.json
  3. 每 2 分钟:
    调度线程检查 → cron_matches 返回 True → cron_queue.append(job)
    → queue processor 发现 Agent 空闲 → agent_loop consume_cron_queue
    → 注入 “[Scheduled] run date”
    → LLM 收到消息,执行 date 命令
  4. 关闭进程:
    调度线程跟着停(daemon=True)
    .scheduled_tasks.json 还在磁盘上
    下次启动 → load_durable_jobs → 任务恢复

Agent Teams

“重构整个后端”涉及认证模块、数据库层、API路由、测试。一个Agent在修API路由时,认证模块的细节已经不在上下文里了。上下文窗口就那么大,单个Agent的注意力覆盖不了所有模块。前面的子Agent是临时工,叫来干一件事就走了。但有些任务需要能通信、能协作的队友。

新增三洋:MessageBus(文件收件箱)、spawn_teammate_thread(启动队友线程)、inbox注入(Lead接收队友消息并注入history)

子Agent vs 队友

子 Agent队友
生命周期一次性,用完销毁多轮(教学版限 10 轮,真实 CC 用 idle loop)
通信只回传结论异步收件箱,随时通信
上下文完全隔离通过消息共享信息
数量一个主 Agent + 偶尔子 Agent一个 Lead + 多个队友

MessageBus:文件收件箱

每个Agent(包括Lead和队友)有一个.jsonl邮箱。发消息=往对方的文件里append一行JSON。读消息 = 读文件 + 删除(消费式)

class MessageBus:
    def send(self, from_agent: str, to_agent: str,
             content: str, msg_type: str = "message"):
        msg = {"from": from_agent, "to": to_agent,
               "content": content, "type": msg_type,
               "ts": time.time()}
        inbox = MAILBOX_DIR / f"{to_agent}.jsonl"
        with open(inbox, "a") as f:
            f.write(json.dumps(msg) + "\n")

    def read_inbox(self, agent: str) -> list[dict]:
        inbox = MAILBOX_DIR / f"{agent}.jsonl"
        if not inbox.exists():
            return []
        msgs = [json.loads(line) for line in inbox.read_text().splitlines()]
        inbox.unlink()  # 消费式:读完删除
        return msgs

为什么用文件而不是内存队列,这里选文件是因为直观、跨线程可观察。真实CC也用文件收件箱(~/.claude/teams/{team}/inboxes/),但加了proper-lockfile防并发写冲突。这里的read_inbox有read+unlink竟态,多线程同时读可能丢消息,但对目前场景可以接受。

spawn_teammate_thread:启动队友

Lead调用spawn_teammate工具启动一个队友。队友跑在自己的daemon线程里,有自己的system prompt、自己的messages、自己的简化工具集

def spawn_teammate_thread(name: str, role: str, prompt: str) -> str:
    system = f"You are '{name}', a {role}. Use tools to complete tasks."

    def run():
        messages = [{"role": "user", "content": prompt}]
        sub_tools = [bash, read_file, write_file, send_message]
        for _ in range(10):           # 最多 10 轮
            inbox = BUS.read_inbox(name)
            if inbox:
                messages.append({"role": "user",
                    "content": f"<inbox>{json.dumps(inbox)}</inbox>"})
            response = client.messages.create(
                model=MODEL, system=system, messages=messages[-20:],
                tools=sub_tools, max_tokens=8000)
            # ... 执行工具、处理结果
        # 完成后发 summary 给 Lead
        BUS.send(name, "lead", summary, "result")

    threading.Thread(target=run, daemon=True).start()

关键设计:

  • 队友有简化工具集:bash、read、write、send_message。真实的CC队友也有TaskCreate、TaskUpdate等工具,任务系统是团队共享的这里简洁版省略。
  • 10轮限制:防止队友无限循环。真实CC用idle loop:跑完一轮后发idle_notification,等inbox消息,收到消息后,知道shutdown_request才退出
  • 完成后自动汇报:BUS.send(name,”lead”,summary)把最终结果发到Lead的收件箱

Lead的inbox注入

Lead在每轮主循环结束后检查收件箱。队友发来的消息注入到history里,让LLM能看到并作出反应

# 主循环结束后
inbox = BUS.read_inbox("lead")
if inbox:
    inbox_text = "\n".join(
        f"From {m['from']}: {m['content'][:200]}" for m in inbox)
    history.append({"role": "user",
                    "content": f"[Inbox]\n{inbox_text}"})

这里采用在用户输入循环外注入。CC更精细,Lead的useInboxPoller每1s检查一次,有消息就提交为新的turn,不需要等用户输入。

权限冒泡

真实的CC流程:

  • 队友遇到需要审批的操作 -> 发 permission_request到Lead收件箱
  • Lead的useInboxPoller检测到请求 -> 路由到审批队列
  • 用户审批后 -> Lead发permission_response回队友
  • 队友的useSwarmPermissionPoller(每500ms轮询)收到回复 -> 继续或拒绝
1. Lead: "搭建后端:一个人搞不定,组队吧"
2. Lead → spawn_teammate("alice", "backend dev", "创建数据库 schema")
3. Lead → spawn_teammate("bob", "frontend dev", "写 API 客户端")
4. alice 线程启动 → 自己的 LLM 调用 → bash "python manage.py migrate"
5. bob 线程启动 → 自己的 LLM 调用 → write_file("client.ts", ...)
6. alice 完成 → BUS.send("alice", "lead", "Schema done: users, orders tables")
7. bob 完成 → BUS.send("bob", "lead", "Client written with types")
8. Lead 下次循环 → inbox 注入 history → LLM 看到 alice 和 bob 的结果

Team Protocols-队友之间要有约定

有了队友,队友都能干活了,但是协调是松散的:Lead发消息,队友回复,没有结构化的协议。

关机:Lead想让Alice关机。直接杀线程,Alice泄写了一半的文件留在磁盘上。需要握手:Lead发请求,Alice确认收尾后关机

计划审批:Bob想重构认证模块,属于高风险操作。应该先让Lead看Bob的计划,审批通过后再动手。

这两个场景结构完全一样:一方发请求,另一方给回复,请求和回复通过同一个ID关联。有状态机追踪:pending -> approved / rejected

新增三样:ProtocolState(请求状态追踪)、dispatch_message(按消息类型路由到处理器)、match_response(通过request_id关联回复与请求,含类型校验)

两种协议,一套机制:

协议方向用途
shutdown_request / responseLead → 队友体面关机握手
plan_approval_request / response队友 → Lead计划审批协议示例

工作原理

ProtocolState:请求状态

每个协议请求创建一条状态记录,记录谁发的、发给谁、当前状态、附带内容

@dataclass
class ProtocolState:
    request_id: str      # 唯一 ID,如 "req_004281"
    type: str            # "shutdown" | "plan_approval"
    sender: str          # 发起方
    target: str          # 接收方
    status: str          # pending | approved | rejected
    payload: str         # 计划文本或关机原因
    created_at: float    # 时间戳

pending_requests: dict[str, ProtocolState] = {}

发送请求时创建记录,收回复时通过request_id找到对应记录,更新状态。

四步协议流程

四步协议流程

以关机为例,完整链路

① Lead 发请求
   req_id = new_request_id()           # "req_004281"
   pending_requests[req_id] = ProtocolState(type="shutdown", status="pending", ...)
   BUS.send("lead", "alice", "shutdown_request", metadata={"request_id": req_id})

② 队友收到 → dispatch
   inbox = BUS.read_inbox("alice")
   msg_type = msg["type"]              # "shutdown_request"
   → 路由到 handle_shutdown_request()

③ 队友回复
   BUS.send("alice", "lead", "shutdown_response",
            metadata={"request_id": req_id, "approve": True})

④ Lead 收响应 → match
   match_response("shutdown_response", req_id, approve=True)
   pending_requests[req_id].status = "approved"

request_id是贯穿全链路的关键键,请求带着它出去,回复带着它回来。

dispatch_messages:按类型路由

队友的inbox不只收普通消息,还收协议消息。handle_inbox_messge按消息类型分发:

def handle_inbox_message(name, msg, messages):
    msg_type = msg.get("type", "message")
    req_id = msg.get("metadata", {}).get("request_id", "")

    if msg_type == "shutdown_request":
        BUS.send(name, "lead", "Shutting down.", "shutdown_response",
                 {"request_id": req_id, "approve": True})
        return True   # 停止循环

    if msg_type == "plan_approval_response":
        approve = msg["metadata"].get("approve", False)
        messages.append({"role": "user",
            "content": "[Plan approved]" if approve else "[Plan rejected]"})
    return False       # 继续循环

新增协议类型只需加新的if分支

match_response:类型校验

match_response不只按request_id找状态,还会校验响应类型是否匹配请求类型

def match_response(response_type, request_id, approve):
    state = pending_requests.get(request_id)
    if not state:
        return
    if state.type == "shutdown" and response_type != "shutdown_response":
        return  # type mismatch, skip
    if state.type == "plan_approval" and response_type != "plan_approval_response":
        return
    if state.status != "pending":
        return  # already resolved, skip duplicate
    state.status = "approved" if approve else "rejected"

一个shutdown_response不会意外approve一个plan_approval请求

统一inbox消费:consume_lead_inbox

check_inbox工具和主循环末尾都调用同一个consume_lead_inbox()函数,先路由协议消息再返回剩余内容,避免消息被读走但协议状态没更新

def consume_lead_inbox(route_protocol=True) -> list[dict]:
    msgs = BUS.read_inbox("lead")
    if route_protocol:
        for msg in msgs:
            meta = msg.get("metadata", {})
            req_id = meta.get("request_id", "")
            msg_type = msg.get("type", "")
            if req_id and msg_type.endswith("_response"):
                match_response(msg_type, req_id, meta.get("approve", False))
    return msgs

主循环末尾还会把inbox消息注入到history,让LLM能看到并作出反应

队友idle loop:等待而不是退出

前面队友跑完10轮就退出,现在队友在LLM返回非tool_use后进入idle等待:轮询inbox,收到shutdown_request就响应退出,收到新消息就继续工作。

LLM 返回非 tool_use
  → idle: 每秒轮询 inbox
  → 收到 shutdown_request → 回复 shutdown_response → 退出
  → 收到新消息 → 注入 messages → 继续 LLM turn

整合后就是

1. Lead: "让 Alice 创建一个文件,然后关机"
2. Lead → spawn_teammate("alice", "backend", "创建 config.py")
3. alice 线程启动 → write_file("config.py", "...") → 完成 → idle
4. Lead → request_shutdown("alice")
   → BUS.send("shutdown_request", {request_id: "req_000142"})
5. alice idle 轮询收到 → handle_shutdown_request
   → BUS.send("shutdown_response", {request_id: "req_000142", approve: True})
6. Lead consume_lead_inbox → match_response("req_000142", approve=True)
   → pending_requests["req_000142"].status = "approved"
   → inbox 消息注入 history,LLM 看到关机结果

Autonomouts Agents-自己看板自己认领

在前面的基础上进行扩展,队友应该自己看任务看板,发现没人做的任务就认领,做完再找下一个

新增:idle_poll(空闲时每5s轮询一次)、scan_unclaimed_tasks(扫描看板上可认领的任务)、自动认领(找到任务就claim,不用Lead操心)

队友生命周期从两阶段变成三阶段

阶段行为退出条件
WORKinbox → LLM → 工具循环stop_reason != tool_use
IDLE每 5s 轮询 inbox + 任务板60s 超时
SHUTDOWN发 summary,退出

工作原理

idle_poll:空闲轮询

队友完成当前任务后不退出,进入IDLE阶段–每5s检查一次有没有新工作

IDLE_POLL_INTERVAL = 5   # seconds
IDLE_TIMEOUT = 60         # seconds

def idle_poll(name, messages, role) -> str:
    """Return 'work', 'shutdown', or 'timeout'."""
    for _ in range(IDLE_TIMEOUT // IDLE_POLL_INTERVAL):
        time.sleep(IDLE_POLL_INTERVAL)

        # ① 检查收件箱(优先)
        inbox = BUS.read_inbox(name)
        if inbox:
            # shutdown_request 立即处理
            for msg in inbox:
                if msg.get("type") == "shutdown_request":
                    # ... 回复 shutdown_response
                    return "shutdown"
            # 普通消息注入上下文,回到 WORK
            messages.append(...)
            return "work"

        # ② 扫描任务看板
        unclaimed = scan_unclaimed_tasks()
        if unclaimed:
            task = unclaimed[0]
            result = claim_task(task["id"], name)
            if "Claimed" in result:
                messages.append(...)
                return "work"
    return "timeout"

inbox优先(可能包含shutdown_request等协议消息),任务板其次。IDLE阶段收到shutdown_request会直接回复并退出,不等到下一轮WORK。

scan_unclaimed_tasks:扫描任务看板

找pending状态,无owner、所有依赖已经完成(can_start)的任务

def scan_unclaimed_tasks() -> list[dict]:
    unclaimed = []
    for f in sorted(TASKS_DIR.glob("task_*.json")):
        task = json.loads(f.read_text())
        if (task.get("status") == "pending"
                and not task.get("owner")
                and can_start(task["id"])):
            unclaimed.append(task)
    return unclaimed

三个条件:必须是pending、没有owner、所有blockedBy依赖已完成。can_start检查依赖任务的状态–有依赖不代表不能做,只有被未完成的任务阻塞才不能做。这里按文件名排序取第一个;CC用文件锁防止多个队友同时认领同一个任务。

claim_task:owner检查

自动认领时检查claim结果,不把失败当成功

def claim_task(task_id: str, owner: str = "agent") -> str:
    task = load_task(task_id)
    if task.status != "pending":
        return f"Task {task_id} is {task.status}, cannot claim"
    if task.owner:
        return f"Task {task_id} already owned by {task.owner}"
    if not can_start(task_id):
        return f"Blocked by: {deps}"
    task.owner = owner
    task.status = "in_progress"
    save_task(task)
    return f"Claimed {task.id} ({task.subject})"

为什么这里还需要进行一次owner检查?

前面空闲寻论的时候扫描任务的时候进行了一次owner检查。那么这个时候可能出现两个Agent空闲轮询到同一个无owner的任务,为了防止复写所以再次检查owner。

这里没有文件锁,并发认领可能出现竞争。并发认领可能出现竞争。但至少task_owner检查避免了最明显的”后写覆盖”问题。CC用proper-lockfile保护任务文件,claimTask在文件锁内完成读-改-写

队友生命周期:WORK –> IDLE –> SHUTDOWN

添加IDLE阶段,队友在外层循环中反复WORK –> IDLE

# Outer loop: WORK → IDLE cycle
while True:
    # WORK phase: 内层循环(最多 10 轮 LLM 调用)
    for _ in range(10):
        # 检查 inbox、处理协议消息、调 LLM、执行工具
        ...
        if response.stop_reason != "tool_use":
            break  # WORK 阶段结束

    # IDLE phase
    idle_result = idle_poll(name, messages, role)
    if idle_result == "shutdown":
        break
    if idle_result == "timeout":
        break  # 60s 超时 → SHUTDOWN

# SHUTDOWN: 发 summary 给 Lead
BUS.send(name, "lead", summary, "result")

关键设计

  • 外层while True:WORK和IDLE交替进行,直到超时或收到关机请求
  • 内层 for 10:WORK阶段最多10轮LLM调用(防止无限循环)
  • IDLE超时 60 s:12次轮询 * 5 s = 60s。超时后发送summary并推出
  • shutdown_request两阶段都能响应:WORK阶段通过handle_inbox_message分发;IDLE阶段idle_poll直接检查并回复

身份重注入

autoCompact之后,队友的message列表可能被压缩成一段摘要。每次进入新的WORK阶段时检查

if len(messages) <= 3:
    messages.insert(0, {"role": "user",
        "content": f"<identity>You are '{name}', role: {role}. "
                   f"Continue your work.</identity>"})

消息过短说明发生了压缩,此时重新注入身份信息。真实CC中context compaction会保留system prompt,这里简化实现需要手动处理。

consume_lead_inbox:统一inbox消费

check_inbox工具和主循环末尾都调用一个consume_lead_inbox()函数:先路由协议response更新状态,再把所有消息注入Lead的对话历史。队友发来的summary/result不会只打印在终端,Lead的LLM能看到并协调下一步。

1. Lead: "搭建后端——任务太多,让队友自己认领"
2. Lead → create_task("创建数据库 schema")
3. Lead → create_task("写 API 路由")
4. Lead → create_task("写单元测试")
5. Lead → spawn_teammate("alice", "backend", "你是后端开发者")
6. Lead → spawn_teammate("bob", "backend", "你是后端开发者")

7. alice 线程启动 → WORK: 没有初始 inbox → 空转 → IDLE
8. bob 线程启动 → WORK: 没有初始 inbox → 空转 → IDLE

9. alice IDLE 第 1 次轮询 → scan_unclaimed → 发现"创建数据库 schema"
10. alice → claim_task → "创建数据库 schema" → 回到 WORK
11. bob IDLE 第 1 次轮询 → scan_unclaimed → 发现"写 API 路由"
12. bob → claim_task → "写 API 路由" → 回到 WORK

13. alice WORK: write_file("schema.sql", ...) → complete_task → WORK 结束
14. alice IDLE → scan → "写单元测试" → claim → WORK
15. alice WORK: write_file("test_api.py", ...) → complete_task → WORK 结束
16. alice IDLE → 60s 无新任务 → SHUTDOWN

17. bob 类似流程 → 做完 → SHUTDOWN
18. Lead consume_lead_inbox → 看到 alice 和 bob 的 summary

两个队友并行认领、并行工作。Lead只需要创建任务和启动队友,不需要手动分配。

Worktree Isolation-各干各的,互不干扰

前面Alice和Bob都在同一个目录下工作。Alice的任务是”重构认证模块”,Bob的任务是”重构UI登录页”

Alice write_file(“config.py”,…)Bob也 write_file(“config.py”,…)两个人改同一个文件,互相覆盖。而且无法干净地回滚–分不清哪些改动是谁的。

可以通过Git worktree让其在同一仓库中创建多个独立的工作目录,每个有自己的分支。Alice在.worktrees/auth-refactor/下工作,Bob在.worktrees/ui-login/下工作–互不干扰。

能力作用
create_worktree为任务创建独立目录 + 独立分支
bind_task_to_worktree把任务和工作目录绑定(不改状态)
remove_worktree / keep_worktree完成后清理或保留
validate_worktree_name拒绝路径穿越和非法字符

工作原理

创建:任务-Worktree绑定

def create_worktree(name: str, task_id: str = "") -> str:
    validate_worktree_name(name)       # 只允许 [A-Za-z0-9._-]{1,64}
    path = WORKTREES_DIR / name
    ok, result = run_git(["worktree", "add", str(path), "-b", f"wt/{name}", "HEAD"])
    if not ok:
        return f"Git error: {result}"
    if task_id:
        bind_task_to_worktree(task_id, name)
    log_event("create", name, task_id)
    return f"Worktree '{name}' created at {path}"

def bind_task_to_worktree(task_id: str, worktree_name: str):
    task = load_task(task_id)
    task.worktree = worktree_name       # 只写 worktree 字段
    save_task(task)                     # 状态保持 pending,等队友 claim

绑定规则:一个任务绑定一个worktree。绑定不改任务状态–任务仍是pending,队友自动认领时才推进到in_progress。这样Lead可以提前创建任务和worktree,队友idle是自然认领带worktree的任务

队友工具的cwd切换

这里给每个队友维护一个wt_ctx字典,记录当前worktree路径。队友认领带worktree的任务时,wt_ctx自动设置为worktree路径;队友的bash、read_file、write_file在worktree目录下执行

# 队友线程内部
wt_ctx = {"path": None}

def _run_claim_task(task_id):
    result = claim_task(task_id, owner=name)
    if "Claimed" in result:
        task = load_task(task_id)
        if task.worktree:
            wt_ctx["path"] = str(WORKTREES_DIR / task.worktree)
    return result

def _run_bash(command):
    return run_bash(command, cwd=wt_ctx["path"])  # 在 worktree 下执行

真实CC的EnterWorktree用process.chdir()切换整个进程目录,AgentTool isolation用cwdOverride包住子Agent执行。

收尾:Keep还是Remove

任务完成后,两个选择

def remove_worktree(name: str, discard_changes: bool = False) -> str:
    # 安全检查:有改动时默认拒绝
    if not discard_changes:
        files, commits = _count_worktree_changes(path)
        if files > 0 or commits > 0:
            return "有未提交改动,使用 discard_changes=true 强制删除,或 keep_worktree 保留"
    ok, _ = run_git(["worktree", "remove", str(path), "--force"])
    if not ok:
        return "删除失败"
    run_git(["branch", "-D", f"wt/{name}"])
    log_event("remove", name)

def keep_worktree(name: str) -> str:
    log_event("keep", name)
    return f"Worktree '{name}' kept for review (branch: wt/{name})"

Keep = 留着分支,等人工review后合并到主分支。Remove = 有改动时默认拒绝,需要discard_changes=true确认。不自动complete task –任务完成由队友的complete_task显式触发。

事件流:可审计

每次生命周期操作写入日志,方便排查

def log_event(event_type: str, worktree_name: str, task_id: str = ""):
    event = {"type": event_type, "worktree": worktree_name,
             "task_id": task_id, "ts": time.time()}
    # append to .worktrees/events.jsonl

事件类型:create(创建)、remove(删除)、keep(保留)。这里只记录事件用于人工排查;完整恢复还需要index或git worktree list扫描

run_git:返回成功/失败

def run_git(args: list[str]) -> tuple[bool, str]:
    r = subprocess.run(["git"] + args, cwd=WORKDIR, ...)
    return r.returncode == 0, output

create_worktree和remove_worktree只在git命令成功后才写事件日志,保证日志反映真实状态。

MCP Tools-外接工具标准协议

前面Agent的所有工具都是手写的–bash、read、write、task、worktree。每个工具的输入验证、执行逻辑、错误处理,都是一行一行写的。

现在有3个外部服务想接入:公司的Jira API(查 issue、建 ticket)、自建的部署系统(触发deploy、看日志)、团队的Notion知识库(搜文档、建页面),不想为每个服务重写一套工具代码,此时就需要一个标准协议–外部只要实现它,Agent就能直接调用。

MCP定义了Agent如何发现和调用外部工具

概念作用
MCPClientAgent 端的客户端,连接 server、发现工具、调用工具
MCP Server外部服务,实现 tools/list + tools/call
assemble_tool_pool把内置工具和 MCP 工具组装成一个工具池
mcp__server__tool 命名避免不同 server 的工具名冲突

这里用mock handler模拟外部server。真实版会启动子进程,通过stdin/stdout发送JSON-RPC请求。mock的好处是不依赖外部服务就能跑完整流程;代价是你看不到真正 网络通信和进程管理。

工作原理

MCPClient:发现+调用

class MCPClient:
    def __init__(self, name: str):
        self.name = name
        self.tools: list[dict] = []
        self._handlers: dict[str, callable] = {}

    def register(self, tool_defs, handlers):
        """Simulates tools/list discovery."""
        self.tools = tool_defs
        self._handlers = handlers

    def call_tool(self, tool_name: str, args: dict) -> str:
        """Simulates tools/call."""
        handler = self._handlers.get(tool_name)
        if not handler:
            return f"MCP error: unknown tool '{tool_name}'"
        return handler(**args)

这里用python函数模拟server的工具实现。真实版通过stdio JSON-RPC与子进程通信。

connect_mcp:连接+发现

def connect_mcp(name: str) -> str:
    if name in mcp_clients:
        return f"MCP server '{name}' already connected"
    factory = MOCK_SERVERS.get(name)
    if not factory:
        return f"Unknown server '{name}'. Available: ..."
    mcp_client = factory()
    mcp_clients[name] = mcp_client
    return f"Connected to '{name}'. Discovered: ..."

连接后,server提供的工具立即可用

normalize_mcp_name:名称规范化

_DISALLOWED_CHARS = re.compile(r'[^a-zA-Z0-9_-]')

def normalize_mcp_name(name: str) -> str:
    return _DISALLOWED_CHARS.sub('_', name)

所有非[a-zA-Z0-9_-]的字符替换为_防止server名或工具名中包含特殊字符导致命名冲突或注入问题

assemble_tool_pool:组装工具池

def assemble_tool_pool() -> tuple[list[dict], dict]:
    tools = list(BUILTIN_TOOLS)
    handlers = dict(BUILTIN_HANDLERS)
    for server_name, mcp_client in mcp_clients.items():
        safe_server = normalize_mcp_name(server_name)
        for tool_def in mcp_client.tools:
            safe_tool = normalize_mcp_name(tool_def["name"])
            prefixed = f"mcp__{safe_server}__{safe_tool}"
            tools.append(...)
            handlers[prefixed] = (
                lambda *, c=mcp_client, t=tool_def["name"], **kw:
                    c.call_tool(t, kw))
    return tools, handlers

前缀mcp_{server}_{tool}避免不同server的工具名冲突。名称经过normalize_mcp_name规范化

MCP工具的description带(readOnly)或(destructive)标注

无缓存:工具池变了,prompt也变

前面agent_loop用prompt cache避免重复序列化,这里去掉缓存重新生成

def agent_loop(messages, context):
    tools, handlers = assemble_tool_pool()     # 每次重新构建
    system = assemble_system_prompt(context)    # 每次重新生成
    ...
    if any(b.name == "connect_mcp" ...):
        tools, handlers = assemble_tool_pool()  # 连接后重建
        system = assemble_system_prompt(context)

原因:connect_mcp之后工具池变化了–新增了mcp_docs_search等工具。缓存中的工具列表是旧的,继续用会导致模型调用不了新工具。

MCP工具只有Lead可用

connect_mcp 是 Lead 工具,assemble_tool_pool 也只服务于 Lead 的 agent_loop。Teammate 仍使用固定的 8 个子集工具(bash、read_file、write_file、send_message、submit_plan、list_tasks、claim_task、complete_task)。

真实 CC 中,MCP 工具对主 agent 和子 agent 都可用——子 agent 继承父级的 MCP 配置。

Comprehensive Agent

一个长期工作的coding agent需要同时拥有

  • 工具分发和权限边界
  • hooks扩展点
  • todo计划和任务图
  • 技能、记忆、系统prompt组装
  • 压缩和错误恢复
  • 后台任务和cron调度
  • 团队、协议、自治认领
  • worktree隔离
  • MCP外部工具接入

难点不是把功能堆起来,而是看清楚它们都挂在循环的哪个位置

一个完整的harness

用户输入
  → UserPromptSubmit hooks
  → cron/background 通知注入
  → context compact
  → memory + skills + MCP 状态组装 system prompt
  → LLM
  → has tool_use block?
      否 → Stop hooks → 返回
      是 → PreToolUse hooks + permission
          → TOOL_HANDLERS / MCP handlers / background dispatch
          → PostToolUse hooks
          → tool_result / task_notification 回 messages
          → 下一轮

循环本身仍然是同一个结构:调用模型,检查响应里是否出现tool_useblock,执行工具,把结果追加回messages。

组件在循环中的位置

位置组件作用
用户输入前后UserPromptSubmit hooks记录、注入、审计用户输入
LLM 前cron queue把定时触发的 prompt 注入 messages
LLM 前background notifications后台任务完成后以 <task_notification> 注入
LLM 前compaction pipeline先压大输出,再裁历史,再压旧 tool_result,必要时摘要
LLM 前memory / skills / MCP state组装 system prompt,让模型看到当前能力和长期上下文
LLM 调用error recovery429/529 重试,max_tokens 升级,prompt too long 触发 reactive compact
工具执行前PreToolUse hooks + permission拦截危险命令、写越界、破坏性 MCP 工具
工具分发assemble_tool_pool组装内置工具和 MCP 动态工具
工具执行时background dispatch慢 bash 操作放 daemon thread,主循环先返回占位结果
工具执行后PostToolUse hooks大输出告警、日志等后处理
返回循环tool_result每个 tool_use 对应一个 tool_result,再回到下一轮
本轮没有 tool_use / 停止时Stop hooks统计、清理、审计

工具与分发

内置工具池包含27个工具

bash, read_file, write_file, edit_file, glob
todo_write, task, load_skill, compact
create_task, list_tasks, get_task, claim_task, complete_task
schedule_cron, list_crons, cancel_cron
spawn_teammate, send_message, check_inbox
request_shutdown, request_plan, review_plan
create_worktree, remove_worktree, keep_worktree
connect_mcp

assemble_tool_pool()每轮组装

BUILTIN_TOOLS + connected MCP tools
BUILTIN_HANDLERS + mcp__server__tool handlers

所以connect_mcp(“docs”)后,下一轮工具池里会出现mcp_docs_search

权限和hooks

权限不写死在工具执行行里,而是作为PreToolUse hook

blocked = trigger_hooks("PreToolUse", block)
if blocked:
    results.append(tool_result(block.id, blocked))
    continue

这样permission、log、审计都可以挂在同一个hook点上。执行后再触发PostToolUse

计划与任务

两层计划

  • todo_write:当前会话的轻量计划,保存在内存中
  • task graph:跨会话、可依赖、可认领的任务文件,写入.tasks/task_*.json

前者帮助单个Agent不飘移;后者支撑团队协作

子agent与团队

  • task:一次性subagent。独立messages[],中间过程丢弃,只返回最终摘要。
  • spawn_teammate:持久队友线程。通过MessageBus收发消息,能idle轮询任务板并自动认领

一次性subagent解决”上下文隔离”;持久队友解决”长期并行协作”

记忆、技能与prompt

assemble_system_prompt(context)每轮组装

  • 身份和工具说明
  • workspace
  • skills catalog
  • .memory/MEMORY.md
  • 已连接MCP server

技能只在system prompt里放目录。完整内容通过load_skill(name)按需加载

压缩和恢复

LLM前先跑压缩管线

tool_result_budget → snip_compact → micro_compact → compact_history

调用模型时再包一层恢复

调用模型时再包一层恢复:

  • 429:指数退避重试
  • 529:指数退避,连续失败可切fallback model
  • max_tokens:先提高max_tokens,再要求continuation
  • prompt too long:reactive compact后重试

后台和cron

慢bash操作不会阻塞主循环

should_run_background → start_background_task → placeholder tool_result
后台完成 → task_notification → 下一轮注入 messages

cron调度器独立daemon thread每秒检查一次。CLI会监听cron_queue,命中后主动把[Scheduled]注入并运行一轮Agent

worktree与MCP

worktree负责隔离目录

  • create_worktree(name,task_id)创建独立分支和目录
  • task的worktree字段绑定目录
  • 队友claim到带worktree的task后,bash/read/write自动在对应目录下执行

MCP负责外部能力

  • connect_mcp(name)连接mock server
  • assemble_tool_pool()把MCP工具组装进工具池
  • 工具名统一为mcp_server_tool

文末附加内容
暂无评论

发送评论 编辑评论


				
|´・ω・)ノ
ヾ(≧∇≦*)ゝ
(☆ω☆)
(╯‵□′)╯︵┴─┴
 ̄﹃ ̄
(/ω\)
∠( ᐛ 」∠)_
(๑•̀ㅁ•́ฅ)
→_→
୧(๑•̀⌄•́๑)૭
٩(ˊᗜˋ*)و
(ノ°ο°)ノ
(´இ皿இ`)
⌇●﹏●⌇
(ฅ´ω`ฅ)
(╯°A°)╯︵○○○
φ( ̄∇ ̄o)
ヾ(´・ ・`。)ノ"
( ง ᵒ̌皿ᵒ̌)ง⁼³₌₃
(ó﹏ò。)
Σ(っ °Д °;)っ
( ,,´・ω・)ノ"(´っω・`。)
╮(╯▽╰)╭
o(*////▽////*)q
>﹏<
( ๑´•ω•) "(ㆆᴗㆆ)
😂
😀
😅
😊
🙂
🙃
😌
😍
😘
😜
😝
😏
😒
🙄
😳
😡
😔
😫
😱
😭
💩
👻
🙌
🖕
👍
👫
👬
👭
🌚
🌝
🙈
💊
😶
🙏
🍦
🍉
😣
Source: github.com/k4yt3x/flowerhd
颜文字
Emoji
小恐龙
花!
上一篇