|
@@ -0,0 +1,985 @@
|
|
|
|
|
+#!/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:
|
|
|
|
|
+ """Check if all blockedBy dependencies are completed.
|
|
|
|
|
+ Missing dependencies are treated as blocked."""
|
|
|
|
|
+ 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":
|
|
|
|
|
+ return f"任务 {task_id} 当前状态为 {task.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"Blocked by: {deps}"
|
|
|
|
|
+ task.owner = owner
|
|
|
|
|
+ task.status = "in_progress"
|
|
|
|
|
+ save_task(task)
|
|
|
|
|
+ print(f" \033[36m[claim] {task.subject} → in_progress (owner: {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":
|
|
|
|
|
+ return f"任务 {task_id} 当前状态为 {task.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[complete] {task.subject} ✓\033[0m")
|
|
|
|
|
+ msg = f"已完成 {task.id} ({task.subject})"
|
|
|
|
|
+ if unblocked:
|
|
|
|
|
+ msg += f"\n已解除阻塞:{', '.join(unblocked)}"
|
|
|
|
|
+ print(f" \033[33m[unblocked] {', '.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"Relevant memories:\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" (blockedBy: {', '.join(blockedBy)})" if blockedBy else ""
|
|
|
|
|
+ print(f" \033[34m[create] {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" (blockedBy: {', '.join(t.blockedBy)})" if t.blockedBy else ""
|
|
|
|
|
+ owner = f" [{t.owner}]" if t.owner else ""
|
|
|
|
|
+ lines.append(f" {icon} {t.id}: {t.subject} "
|
|
|
|
|
+ f"[{t.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[background] dispatched {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"<task_notification>\n"
|
|
|
|
|
+ f" <task_id>{bg_id}</task_id>\n"
|
|
|
|
|
+ f" <status>completed</status>\n"
|
|
|
|
|
+ f" <command>{task['command']}</command>\n"
|
|
|
|
|
+ f" <summary>{summary}</summary>\n"
|
|
|
|
|
+ f"</task_notification>")
|
|
|
|
|
+ print(f" \033[32m[background done] {bg_id}: "
|
|
|
|
|
+ f"{task['command'][:40]} ({len(output)} 个字符)\033[0m")
|
|
|
|
|
+ return notifications
|
|
|
|
|
+
|
|
|
|
|
+
|
|
|
|
|
+def has_pending_background() -> bool:
|
|
|
|
|
+ """Non-destructive: True if any background task has completed and is
|
|
|
|
|
+ waiting to be collected. The inbox poller uses this in its 唤醒 condition."""
|
|
|
|
|
+ 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 # True = 重复,False = 一次性
|
|
|
|
|
+ durable: bool # True = 持久化到磁盘
|
|
|
|
|
+
|
|
|
|
|
+
|
|
|
|
|
+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:
|
|
|
|
|
+ """Check if a 5-field cron expression matches the given datetime.
|
|
|
|
|
+ Standard cron semantics: DOM and DOW use OR when both are constrained."""
|
|
|
|
|
+ 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 = ["minute", "hour", "day-of-month", "month", "day-of-week"]
|
|
|
|
|
+ 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] skipping invalid job {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] loaded {len(valid)} durable job(s)\033[0m")
|
|
|
|
|
+ except Exception:
|
|
|
|
|
+ pass
|
|
|
|
|
+
|
|
|
|
|
+
|
|
|
|
|
+def schedule_job(cron: str, prompt: str, recurring: bool = True,
|
|
|
|
|
+ durable: bool = True) -> Cron任务 | 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 register] {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 cancel] {job_id}\033[0m")
|
|
|
|
|
+ return f"已取消 {job_id}"
|
|
|
|
|
+
|
|
|
|
|
+
|
|
|
|
|
+def cron_scheduler_loop():
|
|
|
|
|
+ """Independent daemon thread: poll every 1s, fire matching jobs.
|
|
|
|
|
+ Individual job errors are caught to prevent one bad job from
|
|
|
|
|
+ killing the entire scheduler thread."""
|
|
|
|
|
+ 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 fire] {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 error] {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] scheduler thread started\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 = "recurring" if j.recurring else "one-shot"
|
|
|
|
|
+ dur = "durable" if j.durable else "session"
|
|
|
|
|
+ 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:
|
|
|
|
|
+ """File-based message bus. Each agent has a .jsonl inbox.
|
|
|
|
|
+ Read is destructive: read_text + unlink (consumes messages).
|
|
|
|
|
+ Teaching version: no file locking; real CC uses 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:
|
|
|
|
|
+ """Non-destructive: True if the agent has unread inbox messages.
|
|
|
|
|
+ The Lead's inbox poller uses this to decide whether to 唤醒 a turn
|
|
|
|
|
+ without consuming the mailbox."""
|
|
|
|
|
+ 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:
|
|
|
|
|
+ """Spawn a teammate agent in a background thread.
|
|
|
|
|
+ Teaching version: max 10 rounds per teammate.
|
|
|
|
|
+ Real CC: teammates use idle loop (wait for inbox, work, repeat)
|
|
|
|
|
+ until 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),
|
|
|
|
|
+ "Sent")[1],
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ for _ in range(10):
|
|
|
|
|
+ inbox = BUS.read_inbox(name)
|
|
|
|
|
+ if inbox:
|
|
|
|
|
+ messages.append({"role": "user",
|
|
|
|
|
+ "content": f"<inbox>{json.dumps(inbox)}</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[teammate] {name} finished\033[0m")
|
|
|
|
|
+
|
|
|
|
|
+ active_teammates[name] = True
|
|
|
|
|
+ threading.Thread(target=run, daemon=True).start()
|
|
|
|
|
+ print(f" \033[36m[teammate] {name} spawned as {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": "True=重复,False=一次性"},
|
|
|
|
|
+ "durable": {"type": "boolean",
|
|
|
|
|
+ "description": "True=持久化到磁盘"}},
|
|
|
|
|
+ "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:
|
|
|
|
|
+ """Derive context from real state."""
|
|
|
|
|
+ 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"[Scheduled] {job.prompt}"})
|
|
|
|
|
+ print(f" \033[35m[inject 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"[Background task {bg_id} started] "
|
|
|
|
|
+ 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("[Inbox]\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)} background "
|
|
|
|
|
+ f"-> new turn]\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()
|