#!/usr/bin/env python3 """ s15: Agent 团队 — MessageBus + spawn_teammate_thread + 收件箱注入。 运行: python s15_agent_teams/code.py 需要: pip install anthropic python-dotenv + .env 中配置 ANTHROPIC_API_KEY 相对 s14 的变化: - MessageBus 类:基于文件的邮箱(.mailboxes/*.jsonl) - spawn_teammate_thread:在后台线程中创建队友 - 队友运行自己的简化 agent_loop(bash、read、write、send_message) - Lead 工具:spawn_teammate、send_message、check_inbox(3 个新增) - Lead 收件箱:队友消息会注入历史(不只是打印) - 教学版本:队友限制为 10 轮(真实 CC 使用空闲循环) ASCII 流程: Lead: cron_queue → messages → prompt → LLM → TOOLS ────→ 循环 ↑ ↓ | └── inbox ← MessageBus ← teammate.send_message ←┘ Teammate: inbox → LLM → bash/read/write/send → 循环(最多 10 轮) """ import os, subprocess, json, time, random, threading, queue from pathlib import Path from datetime import datetime 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"] # ── 任务系统 (来自 s12,已同步) ── 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 blockedBy: list[str] 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, " "get_task, create_task, list_tasks, claim_task, complete_task, " "schedule_cron, list_crons, cancel_cron, " "spawn_teammate, send_message, check_inbox.", "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, run_in_background: bool = False) -> str: # run_in_background 由 agent_loop 分发处理,不在这里处理 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) # ── 后台任务 (来自 s13,已同步) ── _bg_计数器 = 0 background_tasks: dict[str, dict] = {} background_results: dict[str, str] = {} background_lock = threading.Lock() def is_slow_operation(tool_name: str, tool_input: dict) -> bool: """兜底启发式:判断命令是否可能超过 30 秒。""" 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: """模型的显式请求优先;否则使用启发式兜底。""" if tool_input.get("run_in_background"): return True return is_slow_operation(tool_name, tool_input) def execute_tool(block) -> str: """执行工具调用块并返回输出。""" handler = { "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, "schedule_cron": run_schedule_cron, "list_crons": run_list_crons, "cancel_cron": run_cancel_cron, "spawn_teammate": run_spawn_teammate, "send_message": run_send_message, "check_inbox": run_check_inbox, }.get(block.name) if handler: return handler(**block.input) return f"未知工具:{block.name}" def start_background_task(block) -> str: """在守护线程中运行工具,并返回后台任务 ID。""" global _bg_计数器 _bg_计数器 += 1 bg_id = f"bg_{_bg_计数器:04d}" cmd = block.input.get("command", block.name) 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": cmd, "status": "running", } threading.Thread(target=worker, daemon=True).start() print(f" \033[33m[后台] 已分发 {bg_id}: {cmd[:40]}\033[0m") return bg_id def collect_background_results() -> list[str]: """将已完成的后台结果收集为 task_notification 消息。""" 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, "") summary = output[:200] if len(output) > 200 else output notifications.append( f"\n" f" {bg_id}\n" f" completed\n" f" {task['command']}\n" f" {summary}\n" f"") print(f" \033[32m[后台完成] {bg_id}: " f"{task['command'][:40]} ({len(output)} 个字符)\033[0m") return notifications def has_pending_background() -> bool: """非破坏性检查:是否有后台任务已完成并等待收集。 收件箱轮询器会把它作为唤醒条件之一。""" with background_lock: return any(t["status"] == "completed" for t in background_tasks.values()) # ── Cron 调度器 (来自 s14,已同步) ── DURABLE_PATH = WORKDIR / ".scheduled_tasks.json" @dataclass class CronJob: id: str cron: str # "0 9 * * *" prompt: str # 触发时要注入的消息 recurring: bool # 为真 = 重复,为假 = 一次性 durable: bool # 为真 = 持久化到磁盘 scheduled_jobs: dict[str, CronJob] = {} cron_queue: list[CronJob] = [] cron_lock = threading.Lock() _last_fired: dict[str, str] = {} # job_id → "YYYY-MM-DD HH:MM" def _cron_field_matches(field: str, value: int) -> bool: """将单个 cron 字段与一个值进行匹配。""" if field == "*": return True if field.startswith("*/"): step = int(field[2:]) return step > 0 and value % step == 0 if "," in field: return any(_cron_field_matches(f.strip(), value) for f in field.split(",")) if "-" in field: lo, hi = field.split("-", 1) return int(lo) <= value <= int(hi) return value == int(field) def cron_matches(cron_expr: str, dt: datetime) -> bool: """检查 5 字段 cron 表达式是否匹配给定时间。 标准 cron 语义:月内日和周内日同时受限时,两者使用 OR。""" 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 和 DOW:如果两者都有限制,任一匹配即可(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 def _validate_cron_field(field: str, lo: int, hi: int) -> str | None: """校验单个 cron 字段值是否位于 [lo, hi] 范围内。""" if field == "*": return None if field.startswith("*/"): step_str = field[2:] if not step_str.isdigit(): return f"无效步长:{field}" step = int(step_str) if step <= 0: return f"步长必须 > 0:{field}" return None if "," in field: for part in field.split(","): err = _validate_cron_field(part.strip(), lo, hi) if err: return err return None if "-" in field: parts = field.split("-", 1) if not parts[0].isdigit() or not parts[1].isdigit(): return f"无效范围:{field}" a, b = int(parts[0]), int(parts[1]) if a < lo or a > hi or b < lo or b > hi: return f"范围 {field} 超出边界 [{lo}-{hi}]" if a > b: return f"范围起点大于终点:{field}" return None if not field.isdigit(): return f"无效字段:{field}" val = int(field) if val < lo or val > hi: return f"值 {val} 超出边界 [{lo}-{hi}]" return None def validate_cron(cron_expr: str) -> str | None: """校验 cron 表达式。返回错误消息或 None。""" fields = cron_expr.strip().split() if len(fields) != 5: return f"期望 5 个字段,实际得到 {len(fields)}" bounds = [(0, 59), (0, 23), (1, 31), (1, 12), (0, 6)] names = ["分", "时", "月内日", "月", "周内日"] for i, (field, (lo, hi), name) in enumerate(zip(fields, bounds, names)): err = _validate_cron_field(field, lo, hi) if err: return f"{name}: {err}" return None def save_durable_jobs(): """将持久任务保存到 .scheduled_tasks.json。""" durable = [asdict(j) for j in scheduled_jobs.values() if j.durable] DURABLE_PATH.write_text(json.dumps(durable, indent=2)) def load_durable_jobs(): """启动时从磁盘加载持久任务。""" if not DURABLE_PATH.exists(): return try: jobs = json.loads(DURABLE_PATH.read_text()) for j in jobs: job = CronJob(**j) err = validate_cron(job.cron) if err: print(f" \033[31m[cron] 跳过无效任务 {job.id}: {err}\033[0m") continue scheduled_jobs[job.id] = job valid = [j for j in jobs if j["id"] in scheduled_jobs] if valid: print(f" \033[35m[cron] 已加载 {len(valid)} 个持久任务\033[0m") except Exception: pass def schedule_job(cron: str, prompt: str, recurring: bool = True, durable: bool = True) -> CronJob | str: """注册一个新的 cron 任务。返回 CronJob 或错误字符串。""" err = validate_cron(cron) if err: return err job = CronJob( id=f"cron_{random.randint(0, 999999):06d}", cron=cron, prompt=prompt, recurring=recurring, durable=durable, ) with cron_lock: scheduled_jobs[job.id] = job if durable: save_durable_jobs() print(f" \033[35m[cron 注册] {job.id} '{cron}' → {prompt[:40]}\033[0m") return job def cancel_job(job_id: str) -> str: """取消一个 cron 任务。""" with cron_lock: job = scheduled_jobs.pop(job_id, None) if not job: return f"任务 {job_id} 未找到" if job.durable: save_durable_jobs() print(f" \033[31m[cron 取消] {job_id}\033[0m") return f"已取消 {job_id}" def cron_scheduler_loop(): """独立守护线程:每秒轮询一次并触发匹配任务。 单个任务出错会被捕获,避免整个调度线程退出。""" while True: time.sleep(1) now = datetime.now() # 带日期感知的标记,防止每日任务从第 2 天起被跳过 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_queue.append(job) _last_fired[job.id] = minute_marker print(f" \033[35m[cron 触发] {job.id} → " f"{job.prompt[:40]}\033[0m") if not job.recurring: scheduled_jobs.pop(job.id, None) if job.durable: save_durable_jobs() except Exception as e: print(f" \033[31m[cron 错误] {job.id}: {e}\033[0m") def consume_cron_queue() -> list[CronJob]: """消费 cron_queue 中已触发的任务(由 agent_loop 调用)。""" with cron_lock: fired = list(cron_queue) cron_queue.clear() return fired # 启动时加载持久任务,然后启动调度线程 load_durable_jobs() threading.Thread(target=cron_scheduler_loop, daemon=True).start() print(" \033[35m[cron] 调度线程已启动\033[0m") # Cron 工具处理器 def run_schedule_cron(cron: str, prompt: str, recurring: bool = True, durable: bool = True) -> str: result = schedule_job(cron, prompt, recurring, durable) if isinstance(result, str): return f"错误:{result}" return f"已调度 {result.id}: '{cron}' → {prompt}" def run_list_crons() -> str: with cron_lock: jobs = list(scheduled_jobs.values()) if not jobs: return "暂无 cron 任务。请使用 schedule_cron 添加一个。" lines = [] for j in jobs: tag = "重复" if j.recurring else "一次性" dur = "持久化" if j.durable else "仅本会话" lines.append(f" {j.id}: '{j.cron}' → {j.prompt[:40]} " f"[{tag}, {dur}]") return "\n".join(lines) def run_cancel_cron(job_id: str) -> str: return cancel_job(job_id) # ── MessageBus (s15 新增) ── # 教学版本使用简单的文件追加 + 删除。 # 真实 CC 使用 proper-lockfile 保证并发写入安全。 MAILBOX_DIR = WORKDIR / ".mailboxes" MAILBOX_DIR.mkdir(exist_ok=True) class MessageBus: """基于文件的消息总线。每个 Agent 都有一个 .jsonl 收件箱。 读取是破坏性的:read_text + unlink 会消费消息。 教学版本不做文件锁;真实 CC 使用 proper-lockfile。""" 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") print(f" \033[33m[bus] {from_agent} → {to_agent}: " f"{content[:50]}\033[0m") 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() if line.strip()] inbox.unlink() # 消费:读取 + 删除 return msgs def peek(self, agent: str) -> bool: """非破坏性检查:Agent 是否有未读收件箱消息。 Lead 的收件箱轮询器用它决定是否唤醒一轮,同时不消费邮箱。""" inbox = MAILBOX_DIR / f"{agent}.jsonl" return inbox.exists() and inbox.stat().st_size > 0 BUS = MessageBus() # 跟踪已启动的队友 active_teammates: dict[str, bool] = {} # ── 队友线程 (s15 新增) ── def spawn_teammate_thread(name: str, role: str, prompt: str) -> str: """在后台线程中启动队友 Agent。 教学版本中每个队友最多运行 10 轮。 真实 CC 的队友使用空闲循环:等待收件箱、工作、重复,直到收到 shutdown_request。""" if name in active_teammates: return f"队友 '{name}' 已存在" system = (f"你是 '{name}',角色是 {role}。" f"使用工具完成任务。" f"通过 send_message 将结果发送给 'lead'。") def run(): messages = [{"role": "user", "content": prompt}] sub_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"}}, "required": ["path"]}}, {"name": "write_file", "description": "向文件写入内容。", "input_schema": {"type": "object", "properties": {"path": {"type": "string"}, "content": {"type": "string"}}, "required": ["path", "content"]}}, {"name": "send_message", "description": "向另一个 Agent 发送消息。", "input_schema": {"type": "object", "properties": {"to": {"type": "string"}, "content": {"type": "string"}}, "required": ["to", "content"]}}, ] sub_handlers = { "bash": run_bash, "read_file": run_read, "write_file": run_write, "send_message": lambda to, content: (BUS.send(name, to, content), "已发送")[1], } for _ in range(10): inbox = BUS.read_inbox(name) if inbox: messages.append({"role": "user", "content": f"{json.dumps(inbox)}"}) try: response = client.messages.create( model=MODEL, system=system, messages=messages[-20:], tools=sub_tools, max_tokens=8000) except Exception: break 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": handler = sub_handlers.get(block.name) output = handler(**block.input) if handler else "未知" results.append({"type": "tool_result", "tool_use_id": block.id, "content": str(output)}) messages.append({"role": "user", "content": results}) # 向 Lead 发送最终摘要 summary = "已完成。" for msg in reversed(messages): if msg["role"] == "assistant" and isinstance(msg["content"], list): for b in msg["content"]: if getattr(b, "type", None) == "text": summary = b.text break else: continue break BUS.send(name, "lead", summary, "result") active_teammates.pop(name, None) print(f" \033[32m[队友] {name} 已完成\033[0m") active_teammates[name] = True threading.Thread(target=run, daemon=True).start() print(f" \033[36m[队友] {name} 已作为 {role} 启动\033[0m") return f"队友 '{name}' 已启动为 {role}" # ── 团队工具处理器 (s15 新增) ── def run_spawn_teammate(name: str, role: str, prompt: str) -> str: return spawn_teammate_thread(name, role, prompt) def run_send_message(to: str, content: str) -> str: BUS.send("lead", to, content) return f"已发送给 {to}" def run_check_inbox() -> str: msgs = BUS.read_inbox("lead") if not msgs: return "(收件箱为空)" lines = [] for m in msgs: lines.append(f" [{m['from']}] {m['content'][:200]}") return "\n".join(lines) # ── 工具定义 ── TOOLS = [ {"name": "bash", "description": "运行一条 shell 命令。", "input_schema": {"type": "object", "properties": { "command": {"type": "string"}, "run_in_background": {"type": "boolean"}}, "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"]}}, {"name": "schedule_cron", "description": "调度一个 cron 任务。cron 为 5 字段:分 时 月内日 月 周内日。", "input_schema": {"type": "object", "properties": { "cron": {"type": "string", "description": "5 字段 cron 表达式"}, "prompt": {"type": "string", "description": "触发时要注入的消息"}, "recurring": {"type": "boolean", "description": "为真表示重复,为假表示一次性"}, "durable": {"type": "boolean", "description": "为真表示持久化到磁盘"}}, "required": ["cron", "prompt"]}}, {"name": "list_crons", "description": "列出所有已注册的 cron 任务。", "input_schema": {"type": "object", "properties": {}, "required": []}}, {"name": "cancel_cron", "description": "按 ID 取消 cron 任务。", "input_schema": {"type": "object", "properties": {"job_id": {"type": "string"}}, "required": ["job_id"]}}, {"name": "spawn_teammate", "description": "在后台线程中启动一个队友 Agent。", "input_schema": {"type": "object", "properties": { "name": {"type": "string"}, "role": {"type": "string"}, "prompt": {"type": "string"}}, "required": ["name", "role", "prompt"]}}, {"name": "send_message", "description": "通过 MessageBus 向队友发送消息。", "input_schema": {"type": "object", "properties": {"to": {"type": "string"}, "content": {"type": "string"}}, "required": ["to", "content"]}}, {"name": "check_inbox", "description": "检查 Lead 收件箱中的队友消息。", "input_schema": {"type": "object", "properties": {}, "required": []}}, ] # ── 上下文 ── 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": [t["name"] for t in TOOLS], "workspace": str(WORKDIR), "memories": memories, } # ── Agent 循环 ── # 教学代码保留基础 Agent 循环。省略 S11 的完整错误恢复。 # 调用 agent_loop 时消费 Cron 队列;真实 CC 会通过 # 队列处理器 (useQueueProcessor.ts) 在条目到达时。 def agent_loop(messages: list, context: dict): system = get_system_prompt(context) while True: # 消费已触发的 cron 任务 → 作为消息注入 fired = consume_cron_queue() for job in fired: messages.append({"role": "user", "content": f"[已调度] {job.prompt}"}) print(f" \033[35m[注入 cron] {job.prompt[:50]}\033[0m") 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") 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"[后台任务 {bg_id} 已启动] " f"完成后结果将可用。"}) else: output = execute_tool(block) print(str(output)[:300]) results.append({"type": "tool_result", "tool_use_id": block.id, "content": output}) # 将后台工具结果 + 通知合并成一条用户消息 user_content = list(results) bg_notifications = collect_background_results() if bg_notifications: for notif in bg_notifications: user_content.append({"type": "text", "text": notif}) messages.append({"role": "user", "content": user_content}) context = update_context(context, messages) system = get_system_prompt(context) if __name__ == "__main__": print("s15: Agent 团队") print("输入问题后按回车发送。输入 q 退出。\n") history = [] context = update_context({}, []) # input() 和 1 秒轮询器(队友收件箱或后台结果)共同写入一个 # 事件队列(问题 #291、#46)。 events = queue.Queue() def input_reader(): while True: try: line = input("\033[36ms15 >> \033[0m") except (EOFError, KeyboardInterrupt): events.put(("quit", None)) return events.put(("user", line)) def inbox_poller(): # 每约 1 秒轮询一次;当异步结果就绪时唤醒 Lead:队友 # 收件箱消息或已完成的后台任务。不要依赖 # active_teammates:队友发送结果后会移除自身, # 因此最终消息可能比注册表条目存在得更久。 while True: time.sleep(1) if BUS.peek("lead") or has_pending_background(): events.put(("唤醒", None)) threading.Thread(target=input_reader, daemon=True).start() threading.Thread(target=inbox_poller, daemon=True).start() had_teammates = False while True: kind, payload = events.get() if kind == "quit": break if kind == "user": if payload.strip().lower() in ("q", "exit", ""): break history.append({"role": "user", "content": payload}) else: # "唤醒": 队友收件箱或后台结果已就绪 parts = [] inbox = BUS.read_inbox("lead") if inbox: parts.append("[收件箱]\n" + "\n".join( f"来自 {m['from']}: {m['content'][:200]}" for m in inbox)) bg = collect_background_results() parts.extend(bg) if not parts: continue # 已被更早的唤醒消费(幂等) history.append({"role": "user", "content": "\n".join(parts)}) print(f"\n\033[33m[唤醒: {len(inbox)} 条收件箱消息 + {len(bg)} 条后台通知 " f"-> 新一轮]\033[0m") # 为唤醒来源执行一轮。 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", "")) # 当所有队友完成且输出已消费后,只公告一次。 if active_teammates: had_teammates = True elif had_teammates and not BUS.peek("lead") and not has_pending_background(): print("\033[32m[所有队友已完成]\033[0m") had_teammates = False print()