| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382 |
- #!/usr/bin/env python3
- """
- s12: 任务系统 — 用文件持久化带 blockedBy 依赖的任务图。
- 运行: python s12_task_system/code.py
- 需要: pip install anthropic python-dotenv + .env 中配置 ANTHROPIC_API_KEY
- 相对 s11 的变化:
- - 任务 dataclass(id、subject、description、status、owner、blockedBy)
- - TASKS_DIR = .tasks/,用于持久化 JSON 存储
- - create_task / save_task / load_task / list_tasks / get_task
- - can_start:检查 blockedBy 是否全部完成(缺失依赖 = 被阻塞)
- - claim_task:设置 owner + pending -> in_progress
- - complete_task:设置 completed + 报告下游解锁任务
- - 5 个新工具:create_task、list_tasks、get_task、claim_task、complete_task
- 说明:教学代码保留一个基础 Agent 循环,以便聚焦任务系统。
- S11 的完整错误恢复(RecoveryState、退避、升级、reactive compact、备用模型)被省略;
- 在真实 CC 中,tasks.ts 和 withRetry 是可以自然组合的独立层。
- """
- import os, subprocess, json, time, random
- from pathlib import Path
- from dataclasses import dataclass, asdict
- try:
- import readline
- readline.parse_and_bind('set bind-tty-special-chars off')
- except ImportError:
- pass
- from anthropic import Anthropic
- from dotenv import load_dotenv
- load_dotenv(override=True)
- if os.getenv("ANTHROPIC_BASE_URL"):
- os.environ.pop("ANTHROPIC_AUTH_TOKEN", None)
- WORKDIR = Path.cwd()
- MEMORY_DIR = WORKDIR / ".memory"
- MEMORY_INDEX = MEMORY_DIR / "MEMORY.md"
- client = Anthropic(base_url=os.getenv("ANTHROPIC_BASE_URL"))
- MODEL = os.environ["MODEL_ID"]
- # ── 任务系统 ──
- TASKS_DIR = WORKDIR / ".tasks"
- TASKS_DIR.mkdir(exist_ok=True)
- @dataclass
- class Task:
- id: str
- subject: str
- description: str
- status: str # pending | in_progress | completed
- owner: str | None # Agent 名称(多 Agent 场景)
- blockedBy: list[str] # 依赖任务 ID
- def _task_path(task_id: str) -> Path:
- return TASKS_DIR / f"{task_id}.json"
- def create_task(subject: str, description: str = "",
- blockedBy: list[str] | None = None) -> Task:
- task = Task(
- id=f"task_{int(time.time())}_{random.randint(0, 9999):04d}",
- subject=subject,
- description=description,
- status="pending",
- owner=None,
- blockedBy=blockedBy or [],
- )
- save_task(task)
- return task
- def save_task(task: Task):
- _task_path(task.id).write_text(json.dumps(asdict(task), indent=2))
- def load_task(task_id: str) -> Task:
- return Task(**json.loads(_task_path(task_id).read_text()))
- def list_tasks() -> list[Task]:
- return [Task(**json.loads(p.read_text()))
- for p in sorted(TASKS_DIR.glob("task_*.json"))]
- def get_task(task_id: str) -> str:
- """以 JSON 返回完整任务详情。"""
- task = load_task(task_id)
- return json.dumps(asdict(task), indent=2)
- def can_start(task_id: str) -> bool:
- """检查所有 blockedBy 依赖是否已完成;缺失依赖视为阻塞。"""
- task = load_task(task_id)
- for dep_id in task.blockedBy:
- if not _task_path(dep_id).exists():
- return False
- if load_task(dep_id).status != "completed":
- return False
- return True
- def claim_task(task_id: str, owner: str = "agent") -> str:
- task = load_task(task_id)
- if task.status != "pending":
- status = {"pending": "待处理", "in_progress": "进行中",
- "completed": "已完成"}.get(task.status, task.status)
- return f"任务 {task_id} 当前状态为 {status},无法认领"
- if not can_start(task_id):
- deps = [d for d in task.blockedBy
- if not _task_path(d).exists() or load_task(d).status != "completed"]
- return f"被阻塞于:{deps}"
- task.owner = owner
- task.status = "in_progress"
- save_task(task)
- print(f" \033[36m[认领] {task.subject} → in_progress(负责人:{owner})\033[0m")
- return f"已认领 {task.id} ({task.subject})"
- def complete_task(task_id: str) -> str:
- task = load_task(task_id)
- if task.status != "in_progress":
- status = {"pending": "待处理", "in_progress": "进行中",
- "completed": "已完成"}.get(task.status, task.status)
- return f"任务 {task_id} 当前状态为 {status},无法完成"
- 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)]
- print(f" \033[32m[完成] {task.subject} ✓\033[0m")
- msg = f"已完成 {task.id} ({task.subject})"
- if unblocked:
- msg += f"\n已解除阻塞:{', '.join(unblocked)}"
- print(f" \033[33m[已解除阻塞] {', '.join(unblocked)}\033[0m")
- return msg
- # ── 提示词组装 (来自 s10,已同步) ──
- PROMPT_SECTIONS = {
- "identity": "你是一个编码 Agent。直接行动,不要只解释。",
- "tools": "可用工具:bash, read_file, write_file, "
- "create_task, list_tasks, get_task, claim_task, complete_task.",
- "workspace": f"工作目录:{WORKDIR}",
- "memory": "有可用的相关记忆时,会在下方注入。",
- }
- def assemble_system_prompt(context: dict) -> str:
- sections = [PROMPT_SECTIONS["identity"],
- PROMPT_SECTIONS["tools"],
- PROMPT_SECTIONS["workspace"]]
- memories = context.get("memories", "")
- if memories:
- sections.append(f"相关记忆:\n{memories}")
- return "\n\n".join(sections)
- _last_context_key, _last_prompt = None, None
- 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
- # ── 工具 ──
- def safe_path(p: str) -> Path:
- path = (WORKDIR / p).resolve()
- if not path.is_relative_to(WORKDIR):
- raise ValueError(f"路径逃逸出工作区:{p}")
- return path
- def run_bash(command: str) -> str:
- try:
- r = subprocess.run(command, shell=True, cwd=WORKDIR,
- capture_output=True, text=True, timeout=120)
- out = (r.stdout + r.stderr).strip()
- return out[:50000] if out else "(无输出)"
- except subprocess.TimeoutExpired:
- return "错误:执行超时(120 秒)"
- def run_read(path: str, limit: int | None = None) -> str:
- try:
- lines = safe_path(path).read_text().splitlines()
- if limit and limit < len(lines):
- lines = lines[:limit] + [f"... ({len(lines) - limit} 行更多内容)"]
- return "\n".join(lines)
- except Exception as e:
- return f"错误:{e}"
- def run_write(path: str, content: str) -> str:
- try:
- fp = safe_path(path)
- fp.parent.mkdir(parents=True, exist_ok=True)
- fp.write_text(content)
- return f"已写入 {len(content)} 字节到 {path}"
- except Exception as e:
- return f"错误:{e}"
- # 任务工具
- def run_create_task(subject: str, description: str = "",
- blockedBy: list[str] | None = None) -> str:
- task = create_task(subject, description, blockedBy)
- deps = f"(依赖:{', '.join(blockedBy)})" if blockedBy else ""
- print(f" \033[34m[创建] {task.subject}{deps}\033[0m")
- return f"已创建 {task.id}: {task.subject}{deps}"
- def run_list_tasks() -> str:
- 个任务 = list_tasks()
- if not 个任务:
- return "暂无任务。请使用 create_task 添加任务。"
- lines = []
- for t in 个任务:
- icon = {"pending": "○", "in_progress": "●",
- "completed": "✓"}.get(t.status, "?")
- deps = f"(依赖:{', '.join(t.blockedBy)})" if t.blockedBy else ""
- owner = f" [负责人:{t.owner}]" if t.owner else ""
- status = {"pending": "待处理", "in_progress": "进行中",
- "completed": "已完成"}.get(t.status, t.status)
- lines.append(f" {icon} {t.id}: {t.subject} "
- f"[{status}]{owner}{deps}")
- return "\n".join(lines)
- def run_get_task(task_id: str) -> str:
- try:
- return get_task(task_id)
- except FileNotFoundError:
- return f"错误:任务 {task_id} 未找到"
- def run_claim_task(task_id: str) -> str:
- return claim_task(task_id, owner="agent")
- def run_complete_task(task_id: str) -> str:
- return complete_task(task_id)
- TOOLS = [
- {"name": "bash", "description": "运行一条 shell 命令。",
- "input_schema": {"type": "object",
- "properties": {"command": {"type": "string"}},
- "required": ["command"]}},
- {"name": "read_file", "description": "读取文件内容。",
- "input_schema": {"type": "object",
- "properties": {"path": {"type": "string"},
- "limit": {"type": "integer"}},
- "required": ["path"]}},
- {"name": "write_file", "description": "向文件写入内容。",
- "input_schema": {"type": "object",
- "properties": {"path": {"type": "string"},
- "content": {"type": "string"}},
- "required": ["path", "content"]}},
- {"name": "create_task",
- "description": "创建一个新任务,可选 blockedBy 依赖。",
- "input_schema": {"type": "object",
- "properties": {
- "subject": {"type": "string"},
- "description": {"type": "string"},
- "blockedBy": {"type": "array",
- "items": {"type": "string"}}},
- "required": ["subject"]}},
- {"name": "list_tasks",
- "description": "列出所有任务及其状态、负责人和依赖。",
- "input_schema": {"type": "object", "properties": {},
- "required": []}},
- {"name": "get_task",
- "description": "按 ID 获取指定任务的完整详情。",
- "input_schema": {"type": "object",
- "properties": {"task_id": {"type": "string"}},
- "required": ["task_id"]}},
- {"name": "claim_task",
- "description": "认领一个待处理任务。设置 owner,并将状态改为 in_progress。",
- "input_schema": {"type": "object",
- "properties": {"task_id": {"type": "string"}},
- "required": ["task_id"]}},
- {"name": "complete_task",
- "description": "完成一个进行中的任务。报告被解除阻塞的下游任务。",
- "input_schema": {"type": "object",
- "properties": {"task_id": {"type": "string"}},
- "required": ["task_id"]}},
- ]
- TOOL_HANDLERS = {
- "bash": run_bash, "read_file": run_read, "write_file": run_write,
- "create_task": run_create_task, "list_tasks": run_list_tasks,
- "get_task": run_get_task, "claim_task": run_claim_task,
- "complete_task": run_complete_task,
- }
- # ── 上下文 ──
- 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,
- }
- # ── Agent 循环(简化版,聚焦任务系统) ──
- def agent_loop(messages: list, context: dict):
- system = get_system_prompt(context)
- while True:
- try:
- response = client.messages.create(
- model=MODEL, system=system, messages=messages,
- tools=TOOLS, max_tokens=8000)
- except Exception as e:
- messages.append({"role": "assistant", "content": [
- {"type": "text",
- "text": f"[错误] {type(e).__name__}: {e}"}]})
- return
- 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":
- continue
- print(f"\033[36m> {block.name}\033[0m")
- handler = TOOL_HANDLERS.get(block.name)
- output = handler(**block.input) if handler else f"未知工具:{block.name}"
- print(str(output)[:300])
- results.append({"type": "tool_result",
- "tool_use_id": block.id, "content": output})
- messages.append({"role": "user", "content": results})
- context = update_context(context, messages)
- system = get_system_prompt(context)
- if __name__ == "__main__":
- print("s12: 任务系统")
- print("输入问题后按回车发送。输入 q 退出。\n")
- history = []
- context = update_context({}, [])
- while True:
- try:
- query = input("\033[36ms12 >> \033[0m")
- except (EOFError, KeyboardInterrupt):
- break
- if query.strip().lower() in ("q", "exit", ""):
- break
- history.append({"role": "user", "content": query})
- agent_loop(history, context)
- context = update_context(context, history)
- for block in history[-1]["content"]:
- if getattr(block, "type", None) == "text":
- print(block.text)
- elif isinstance(block, dict) and block.get("type") == "text":
- print(block.get("text", ""))
- print()
|