code.py 29 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808
  1. #!/usr/bin/env python3
  2. """
  3. s14: Cron 调度器 — 独立守护线程 + 队列处理器。
  4. 运行: python s14_cron_scheduler/code.py
  5. 需要: pip install anthropic python-dotenv + .env 中配置 ANTHROPIC_API_KEY
  6. 相对 s13 的变化:
  7. - Cron任务 dataclass(id、cron、prompt、recurring、durable)
  8. - cron_matches:带 DOM/DOW OR 语义的 5 字段 cron 表达式匹配
  9. - schedule_job / cancel_job:注册/移除 cron 任务(带校验)
  10. - cron_scheduler_loop:独立守护线程,每 1 秒轮询
  11. - cron_queue:线程安全队列,调度器写入,队列处理器负责投递
  12. - queue_processor_loop:当 cron_queue 有任务时自动运行 agent_loop
  13. - 持久化存储:.scheduled_tasks.json(重启后仍保留)
  14. - 3 个新工具:schedule_cron、list_crons、cancel_cron
  15. 四层结构:
  16. 1. 调度器:守护线程检查时间 → 触发匹配任务
  17. 2. 队列:cron_queue 将调度器与 Agent 循环解耦
  18. 3. 队列处理器:当队列中有任务且 Agent 空闲时唤醒 Agent
  19. 4. 消费者:agent_loop 消费队列任务,并注入到 messages
  20. """
  21. import os, subprocess, json, time, random, threading
  22. from pathlib import Path
  23. from datetime import datetime
  24. from dataclasses import dataclass, asdict
  25. try:
  26. import readline
  27. readline.parse_and_bind('set bind-tty-special-chars off')
  28. except ImportError:
  29. pass
  30. from anthropic import Anthropic
  31. from dotenv import load_dotenv
  32. load_dotenv(override=True)
  33. if os.getenv("ANTHROPIC_BASE_URL"):
  34. os.environ.pop("ANTHROPIC_AUTH_TOKEN", None)
  35. WORKDIR = Path.cwd()
  36. MEMORY_DIR = WORKDIR / ".memory"
  37. MEMORY_INDEX = MEMORY_DIR / "MEMORY.md"
  38. client = Anthropic(base_url=os.getenv("ANTHROPIC_BASE_URL"))
  39. MODEL = os.environ["MODEL_ID"]
  40. # ── 任务系统 (来自 s12,已同步) ──
  41. TASKS_DIR = WORKDIR / ".tasks"
  42. TASKS_DIR.mkdir(exist_ok=True)
  43. @dataclass
  44. class Task:
  45. id: str
  46. subject: str
  47. description: str
  48. status: str # pending | in_progress | completed
  49. owner: str | None
  50. blockedBy: list[str]
  51. def _task_path(task_id: str) -> Path:
  52. return TASKS_DIR / f"{task_id}.json"
  53. def create_task(subject: str, description: str = "",
  54. blockedBy: list[str] | None = None) -> Task:
  55. task = Task(
  56. id=f"task_{int(time.time())}_{random.randint(0, 9999):04d}",
  57. subject=subject, description=description,
  58. status="pending", owner=None,
  59. blockedBy=blockedBy or [],
  60. )
  61. save_task(task)
  62. return task
  63. def save_task(task: Task):
  64. _task_path(task.id).write_text(json.dumps(asdict(task), indent=2))
  65. def load_task(task_id: str) -> Task:
  66. return Task(**json.loads(_task_path(task_id).read_text()))
  67. def list_tasks() -> list[Task]:
  68. return [Task(**json.loads(p.read_text()))
  69. for p in sorted(TASKS_DIR.glob("task_*.json"))]
  70. def get_task(task_id: str) -> str:
  71. """以 JSON 返回完整任务详情。"""
  72. task = load_task(task_id)
  73. return json.dumps(asdict(task), indent=2)
  74. def can_start(task_id: str) -> bool:
  75. """检查所有 blockedBy 依赖是否已完成;缺失依赖视为阻塞。"""
  76. task = load_task(task_id)
  77. for dep_id in task.blockedBy:
  78. if not _task_path(dep_id).exists():
  79. return False
  80. if load_task(dep_id).status != "completed":
  81. return False
  82. return True
  83. def claim_task(task_id: str, owner: str = "agent") -> str:
  84. task = load_task(task_id)
  85. if task.status != "pending":
  86. status = {"pending": "待处理", "in_progress": "进行中",
  87. "completed": "已完成"}.get(task.status, task.status)
  88. return f"任务 {task_id} 当前状态为 {status},无法认领"
  89. if not can_start(task_id):
  90. deps = [d for d in task.blockedBy
  91. if not _task_path(d).exists() or load_task(d).status != "completed"]
  92. return f"被阻塞于:{deps}"
  93. task.owner = owner
  94. task.status = "in_progress"
  95. save_task(task)
  96. print(f" \033[36m[认领] {task.subject} → in_progress(负责人:{owner})\033[0m")
  97. return f"已认领 {task.id} ({task.subject})"
  98. def complete_task(task_id: str) -> str:
  99. task = load_task(task_id)
  100. if task.status != "in_progress":
  101. status = {"pending": "待处理", "in_progress": "进行中",
  102. "completed": "已完成"}.get(task.status, task.status)
  103. return f"任务 {task_id} 当前状态为 {status},无法完成"
  104. task.status = "completed"
  105. save_task(task)
  106. unblocked = [t.subject for t in list_tasks()
  107. if t.status == "pending" and t.blockedBy and can_start(t.id)]
  108. print(f" \033[32m[完成] {task.subject} ✓\033[0m")
  109. msg = f"已完成 {task.id} ({task.subject})"
  110. if unblocked:
  111. msg += f"\n已解除阻塞:{', '.join(unblocked)}"
  112. print(f" \033[33m[已解除阻塞] {', '.join(unblocked)}\033[0m")
  113. return msg
  114. # ── 提示词组装 (来自 s10,已同步) ──
  115. PROMPT_SECTIONS = {
  116. "identity": "你是一个编码 Agent。直接行动,不要只解释。",
  117. "tools": "可用工具:bash, read_file, write_file, "
  118. "create_task, list_tasks, get_task, claim_task, complete_task, "
  119. "schedule_cron, list_crons, cancel_cron.",
  120. "workspace": f"工作目录:{WORKDIR}",
  121. "memory": "有可用的相关记忆时,会在下方注入。",
  122. }
  123. def assemble_system_prompt(context: dict) -> str:
  124. sections = [PROMPT_SECTIONS["identity"],
  125. PROMPT_SECTIONS["tools"],
  126. PROMPT_SECTIONS["workspace"]]
  127. memories = context.get("memories", "")
  128. if memories:
  129. sections.append(f"相关记忆:\n{memories}")
  130. return "\n\n".join(sections)
  131. _last_context_key, _last_prompt = None, None
  132. def get_system_prompt(context: dict) -> str:
  133. global _last_context_key, _last_prompt
  134. key = json.dumps(context, sort_keys=True, ensure_ascii=False, default=str)
  135. if key == _last_context_key and _last_prompt:
  136. return _last_prompt
  137. _last_context_key = key
  138. _last_prompt = assemble_system_prompt(context)
  139. return _last_prompt
  140. # ── 工具 ──
  141. def safe_path(p: str) -> Path:
  142. path = (WORKDIR / p).resolve()
  143. if not path.is_relative_to(WORKDIR):
  144. raise ValueError(f"路径逃逸出工作区:{p}")
  145. return path
  146. def run_bash(command: str, run_in_background: bool = False) -> str:
  147. # run_in_background 由 agent_loop 分发处理,不在这里处理
  148. try:
  149. r = subprocess.run(command, shell=True, cwd=WORKDIR,
  150. capture_output=True, text=True, timeout=120)
  151. out = (r.stdout + r.stderr).strip()
  152. return out[:50000] if out else "(无输出)"
  153. except subprocess.TimeoutExpired:
  154. return "错误:执行超时(120 秒)"
  155. def run_read(path: str, limit: int | None = None) -> str:
  156. try:
  157. lines = safe_path(path).read_text().splitlines()
  158. if limit and limit < len(lines):
  159. lines = lines[:limit] + [f"... ({len(lines) - limit} 行更多内容)"]
  160. return "\n".join(lines)
  161. except Exception as e:
  162. return f"错误:{e}"
  163. def run_write(path: str, content: str) -> str:
  164. try:
  165. fp = safe_path(path)
  166. fp.parent.mkdir(parents=True, exist_ok=True)
  167. fp.write_text(content)
  168. return f"已写入 {len(content)} 字节到 {path}"
  169. except Exception as e:
  170. return f"错误:{e}"
  171. # 任务工具
  172. def run_create_task(subject: str, description: str = "",
  173. blockedBy: list[str] | None = None) -> str:
  174. task = create_task(subject, description, blockedBy)
  175. deps = f"(依赖:{', '.join(blockedBy)})" if blockedBy else ""
  176. print(f" \033[34m[创建] {task.subject}{deps}\033[0m")
  177. return f"已创建 {task.id}: {task.subject}{deps}"
  178. def run_list_tasks() -> str:
  179. 个任务 = list_tasks()
  180. if not 个任务:
  181. return "暂无任务。请使用 create_task 添加任务。"
  182. lines = []
  183. for t in 个任务:
  184. icon = {"pending": "○", "in_progress": "●",
  185. "completed": "✓"}.get(t.status, "?")
  186. deps = f"(依赖:{', '.join(t.blockedBy)})" if t.blockedBy else ""
  187. owner = f" [负责人:{t.owner}]" if t.owner else ""
  188. status = {"pending": "待处理", "in_progress": "进行中",
  189. "completed": "已完成"}.get(t.status, t.status)
  190. lines.append(f" {icon} {t.id}: {t.subject} "
  191. f"[{status}]{owner}{deps}")
  192. return "\n".join(lines)
  193. def run_get_task(task_id: str) -> str:
  194. try:
  195. return get_task(task_id)
  196. except FileNotFoundError:
  197. return f"错误:任务 {task_id} 未找到"
  198. def run_claim_task(task_id: str) -> str:
  199. return claim_task(task_id, owner="agent")
  200. def run_complete_task(task_id: str) -> str:
  201. return complete_task(task_id)
  202. # ── 后台任务 (来自 s13,已同步) ──
  203. _bg_计数器 = 0
  204. background_tasks: dict[str, dict] = {}
  205. background_results: dict[str, str] = {}
  206. background_lock = threading.Lock()
  207. def is_slow_operation(tool_name: str, tool_input: dict) -> bool:
  208. """兜底启发式:判断命令是否可能超过 30 秒。"""
  209. if tool_name != "bash":
  210. return False
  211. cmd = tool_input.get("command", "").lower()
  212. slow_keywords = ["install", "build", "test", "deploy", "compile",
  213. "docker build", "pip install", "npm install",
  214. "cargo build", "pytest", "make"]
  215. return any(kw in cmd for kw in slow_keywords)
  216. def should_run_background(tool_name: str, tool_input: dict) -> bool:
  217. """模型的显式请求优先;否则使用启发式兜底。"""
  218. if tool_input.get("run_in_background"):
  219. return True
  220. return is_slow_operation(tool_name, tool_input)
  221. def execute_tool(block) -> str:
  222. """执行工具调用块并返回输出。"""
  223. handler = {
  224. "bash": run_bash, "read_file": run_read, "write_file": run_write,
  225. "create_task": run_create_task, "list_tasks": run_list_tasks,
  226. "get_task": run_get_task, "claim_task": run_claim_task,
  227. "complete_task": run_complete_task,
  228. "schedule_cron": run_schedule_cron, "list_crons": run_list_crons,
  229. "cancel_cron": run_cancel_cron,
  230. }.get(block.name)
  231. if handler:
  232. return handler(**block.input)
  233. return f"未知工具:{block.name}"
  234. def start_background_task(block) -> str:
  235. """在守护线程中运行工具,并返回后台任务 ID。"""
  236. global _bg_计数器
  237. _bg_计数器 += 1
  238. bg_id = f"bg_{_bg_计数器:04d}"
  239. cmd = block.input.get("command", block.name)
  240. def worker():
  241. result = execute_tool(block)
  242. with background_lock:
  243. background_tasks[bg_id]["status"] = "completed"
  244. background_results[bg_id] = result
  245. with background_lock:
  246. background_tasks[bg_id] = {
  247. "tool_use_id": block.id,
  248. "command": cmd,
  249. "status": "running",
  250. }
  251. threading.Thread(target=worker, daemon=True).start()
  252. print(f" \033[33m[后台] 已分发 {bg_id}: {cmd[:40]}\033[0m")
  253. return bg_id
  254. def collect_background_results() -> list[str]:
  255. """将已完成的后台结果收集为 task_notification 消息。"""
  256. with background_lock:
  257. ready_ids = [bid for bid, task in background_tasks.items()
  258. if task["status"] == "completed"]
  259. notifications = []
  260. for bg_id in ready_ids:
  261. with background_lock:
  262. task = background_tasks.pop(bg_id)
  263. output = background_results.pop(bg_id, "")
  264. summary = output[:200] if len(output) > 200 else output
  265. notifications.append(
  266. f"<task_notification>\n"
  267. f" <task_id>{bg_id}</task_id>\n"
  268. f" <status>completed</status>\n"
  269. f" <command>{task['command']}</command>\n"
  270. f" <summary>{summary}</summary>\n"
  271. f"</task_notification>")
  272. print(f" \033[32m[后台完成] {bg_id}: "
  273. f"{task['command'][:40]} ({len(output)} 个字符)\033[0m")
  274. return notifications
  275. # ── Cron 调度器 (s14 新增) ──
  276. DURABLE_PATH = WORKDIR / ".scheduled_tasks.json"
  277. @dataclass
  278. class CronJob:
  279. id: str
  280. cron: str # "0 9 * * *"
  281. prompt: str # 触发时要注入的消息
  282. recurring: bool # 为真 = 重复,为假 = 一次性
  283. durable: bool # 为真 = 持久化到磁盘
  284. scheduled_jobs: dict[str, CronJob] = {}
  285. cron_queue: list[CronJob] = []
  286. cron_lock = threading.Lock()
  287. agent_lock = threading.Lock()
  288. _last_fired: dict[str, str] = {} # job_id → "YYYY-MM-DD HH:MM"
  289. def _cron_field_matches(field: str, value: int) -> bool:
  290. """将单个 cron 字段与一个值进行匹配。"""
  291. if field == "*":
  292. return True
  293. if field.startswith("*/"):
  294. step = int(field[2:])
  295. return step > 0 and value % step == 0
  296. if "," in field:
  297. return any(_cron_field_matches(f.strip(), value)
  298. for f in field.split(","))
  299. if "-" in field:
  300. lo, hi = field.split("-", 1)
  301. return int(lo) <= value <= int(hi)
  302. return value == int(field)
  303. def cron_matches(cron_expr: str, dt: datetime) -> bool:
  304. """检查 5 字段 cron 表达式是否匹配给定时间。
  305. 标准 cron 语义:月内日和周内日同时受限时,两者使用 OR。"""
  306. fields = cron_expr.strip().split()
  307. if len(fields) != 5:
  308. return False
  309. minute, hour, dom, month, dow = fields
  310. dow_val = (dt.weekday() + 1) % 7 # Python Monday=0 → cron Sunday=0
  311. m = _cron_field_matches(minute, dt.minute)
  312. h = _cron_field_matches(hour, dt.hour)
  313. dom_ok = _cron_field_matches(dom, dt.day)
  314. month_ok = _cron_field_matches(month, dt.month)
  315. dow_ok = _cron_field_matches(dow, dow_val)
  316. # 分钟、小时、月份必须全部匹配
  317. if not (m and h and month_ok):
  318. return False
  319. # DOM 和 DOW:如果两者都有限制,任一匹配即可(OR)
  320. dom_unconstrained = dom == "*"
  321. dow_unconstrained = dow == "*"
  322. if dom_unconstrained and dow_unconstrained:
  323. return True
  324. if dom_unconstrained:
  325. return dow_ok
  326. if dow_unconstrained:
  327. return dom_ok
  328. return dom_ok or dow_ok
  329. def _validate_cron_field(field: str, lo: int, hi: int) -> str | None:
  330. """校验单个 cron 字段值是否位于 [lo, hi] 范围内。"""
  331. if field == "*":
  332. return None
  333. if field.startswith("*/"):
  334. step_str = field[2:]
  335. if not step_str.isdigit():
  336. return f"无效步长:{field}"
  337. step = int(step_str)
  338. if step <= 0:
  339. return f"步长必须 > 0:{field}"
  340. return None
  341. if "," in field:
  342. for part in field.split(","):
  343. err = _validate_cron_field(part.strip(), lo, hi)
  344. if err: return err
  345. return None
  346. if "-" in field:
  347. parts = field.split("-", 1)
  348. if not parts[0].isdigit() or not parts[1].isdigit():
  349. return f"无效范围:{field}"
  350. a, b = int(parts[0]), int(parts[1])
  351. if a < lo or a > hi or b < lo or b > hi:
  352. return f"范围 {field} 超出边界 [{lo}-{hi}]"
  353. if a > b:
  354. return f"范围起点大于终点:{field}"
  355. return None
  356. if not field.isdigit():
  357. return f"无效字段:{field}"
  358. val = int(field)
  359. if val < lo or val > hi:
  360. return f"值 {val} 超出边界 [{lo}-{hi}]"
  361. return None
  362. def validate_cron(cron_expr: str) -> str | None:
  363. """校验 cron 表达式。返回错误消息或 None。"""
  364. fields = cron_expr.strip().split()
  365. if len(fields) != 5:
  366. return f"期望 5 个字段,实际得到 {len(fields)}"
  367. bounds = [(0, 59), (0, 23), (1, 31), (1, 12), (0, 6)]
  368. names = ["分", "时", "月内日", "月", "周内日"]
  369. for i, (field, (lo, hi), name) in enumerate(zip(fields, bounds, names)):
  370. err = _validate_cron_field(field, lo, hi)
  371. if err:
  372. return f"{name}: {err}"
  373. return None
  374. def save_durable_jobs():
  375. """将持久任务保存到 .scheduled_tasks.json。"""
  376. durable = [asdict(j) for j in scheduled_jobs.values() if j.durable]
  377. DURABLE_PATH.write_text(json.dumps(durable, indent=2))
  378. def load_durable_jobs():
  379. """启动时从磁盘加载持久任务。"""
  380. if not DURABLE_PATH.exists():
  381. return
  382. try:
  383. jobs = json.loads(DURABLE_PATH.read_text())
  384. for j in jobs:
  385. job = CronJob(**j)
  386. err = validate_cron(job.cron)
  387. if err:
  388. print(f" \033[31m[cron] 跳过无效任务 {job.id}: {err}\033[0m")
  389. continue
  390. scheduled_jobs[job.id] = job
  391. valid = [j for j in jobs if j["id"] in scheduled_jobs]
  392. if valid:
  393. print(f" \033[35m[cron] 已加载 {len(valid)} 个持久任务\033[0m")
  394. except Exception:
  395. pass
  396. def schedule_job(cron: str, prompt: str, recurring: bool = True,
  397. durable: bool = True) -> CronJob | str:
  398. """注册一个新的 cron 任务。返回 CronJob 或错误字符串。"""
  399. err = validate_cron(cron)
  400. if err:
  401. return err
  402. job = CronJob(
  403. id=f"cron_{random.randint(0, 999999):06d}",
  404. cron=cron, prompt=prompt,
  405. recurring=recurring, durable=durable,
  406. )
  407. with cron_lock:
  408. scheduled_jobs[job.id] = job
  409. if durable:
  410. save_durable_jobs()
  411. print(f" \033[35m[cron 注册] {job.id} '{cron}' → {prompt[:40]}\033[0m")
  412. return job
  413. def cancel_job(job_id: str) -> str:
  414. """取消一个 cron 任务。"""
  415. with cron_lock:
  416. job = scheduled_jobs.pop(job_id, None)
  417. if not job:
  418. return f"任务 {job_id} 未找到"
  419. if job.durable:
  420. save_durable_jobs()
  421. print(f" \033[31m[cron 取消] {job_id}\033[0m")
  422. return f"已取消 {job_id}"
  423. def cron_scheduler_loop():
  424. """独立守护线程:每秒轮询一次并触发匹配任务。
  425. 单个任务出错会被捕获,避免整个调度线程退出。"""
  426. while True:
  427. time.sleep(1)
  428. now = datetime.now()
  429. # 带日期感知的标记,防止每日任务从第 2 天起被跳过
  430. minute_marker = now.strftime("%Y-%m-%d %H:%M")
  431. with cron_lock:
  432. for job in list(scheduled_jobs.values()):
  433. try:
  434. if cron_matches(job.cron, now):
  435. if _last_fired.get(job.id) != minute_marker:
  436. cron_queue.append(job)
  437. _last_fired[job.id] = minute_marker
  438. print(f" \033[35m[cron 触发] {job.id} → "
  439. f"{job.prompt[:40]}\033[0m")
  440. if not job.recurring:
  441. scheduled_jobs.pop(job.id, None)
  442. if job.durable:
  443. save_durable_jobs()
  444. except Exception as e:
  445. print(f" \033[31m[cron 错误] {job.id}: {e}\033[0m")
  446. def consume_cron_queue() -> list[CronJob]:
  447. """消费 cron_queue 中已触发的任务(由 agent_loop 调用)。"""
  448. with cron_lock:
  449. fired = list(cron_queue)
  450. cron_queue.clear()
  451. return fired
  452. def has_cron_queue() -> bool:
  453. """返回是否有已触发但尚未投递的 cron 任务。"""
  454. with cron_lock:
  455. return bool(cron_queue)
  456. # 启动时加载持久任务,然后启动调度线程
  457. load_durable_jobs()
  458. threading.Thread(target=cron_scheduler_loop, daemon=True).start()
  459. print(" \033[35m[cron] 调度线程已启动\033[0m")
  460. # ── Cron 工具 ──
  461. def run_schedule_cron(cron: str, prompt: str,
  462. recurring: bool = True, durable: bool = True) -> str:
  463. result = schedule_job(cron, prompt, recurring, durable)
  464. if isinstance(result, str):
  465. return f"错误:{result}"
  466. return f"已调度 {result.id}: '{cron}' → {prompt}"
  467. def run_list_crons() -> str:
  468. with cron_lock:
  469. jobs = list(scheduled_jobs.values())
  470. if not jobs:
  471. return "暂无 cron 任务。请使用 schedule_cron 添加一个。"
  472. lines = []
  473. for j in jobs:
  474. tag = "重复" if j.recurring else "一次性"
  475. dur = "持久化" if j.durable else "仅本会话"
  476. lines.append(f" {j.id}: '{j.cron}' → {j.prompt[:40]} "
  477. f"[{tag}, {dur}]")
  478. return "\n".join(lines)
  479. def run_cancel_cron(job_id: str) -> str:
  480. return cancel_job(job_id)
  481. # ── 工具定义 ──
  482. TOOLS = [
  483. {"name": "bash", "description": "运行一条 shell 命令。",
  484. "input_schema": {"type": "object",
  485. "properties": {
  486. "command": {"type": "string"},
  487. "run_in_background": {"type": "boolean"}},
  488. "required": ["command"]}},
  489. {"name": "read_file", "description": "读取文件内容。",
  490. "input_schema": {"type": "object",
  491. "properties": {"path": {"type": "string"},
  492. "limit": {"type": "integer"}},
  493. "required": ["path"]}},
  494. {"name": "write_file", "description": "向文件写入内容。",
  495. "input_schema": {"type": "object",
  496. "properties": {"path": {"type": "string"},
  497. "content": {"type": "string"}},
  498. "required": ["path", "content"]}},
  499. {"name": "create_task",
  500. "description": "创建一个新任务,可选 blockedBy 依赖。",
  501. "input_schema": {"type": "object",
  502. "properties": {
  503. "subject": {"type": "string"},
  504. "description": {"type": "string"},
  505. "blockedBy": {"type": "array",
  506. "items": {"type": "string"}}},
  507. "required": ["subject"]}},
  508. {"name": "list_tasks",
  509. "description": "列出所有任务及其状态、负责人和依赖。",
  510. "input_schema": {"type": "object", "properties": {},
  511. "required": []}},
  512. {"name": "get_task",
  513. "description": "按 ID 获取指定任务的完整详情。",
  514. "input_schema": {"type": "object",
  515. "properties": {"task_id": {"type": "string"}},
  516. "required": ["task_id"]}},
  517. {"name": "claim_task",
  518. "description": "认领一个待处理任务。设置 owner,并将状态改为 in_progress。",
  519. "input_schema": {"type": "object",
  520. "properties": {"task_id": {"type": "string"}},
  521. "required": ["task_id"]}},
  522. {"name": "complete_task",
  523. "description": "完成一个进行中的任务。报告被解除阻塞的下游任务。",
  524. "input_schema": {"type": "object",
  525. "properties": {"task_id": {"type": "string"}},
  526. "required": ["task_id"]}},
  527. {"name": "schedule_cron",
  528. "description": "调度一个 cron 任务。cron 为 5 字段:分 时 月内日 月 周内日。",
  529. "input_schema": {"type": "object",
  530. "properties": {
  531. "cron": {"type": "string",
  532. "description": "5 字段 cron 表达式"},
  533. "prompt": {"type": "string",
  534. "description": "触发时要注入的消息"},
  535. "recurring": {"type": "boolean",
  536. "description": "为真表示重复,为假表示一次性"},
  537. "durable": {"type": "boolean",
  538. "description": "为真表示持久化到磁盘"}},
  539. "required": ["cron", "prompt"]}},
  540. {"name": "list_crons",
  541. "description": "列出所有已注册的 cron 任务。",
  542. "input_schema": {"type": "object", "properties": {},
  543. "required": []}},
  544. {"name": "cancel_cron",
  545. "description": "按 ID 取消 cron 任务。",
  546. "input_schema": {"type": "object",
  547. "properties": {"job_id": {"type": "string"}},
  548. "required": ["job_id"]}},
  549. ]
  550. # ── 上下文 ──
  551. def update_context(context: dict, messages: list) -> dict:
  552. """从真实状态推导上下文。"""
  553. memories = ""
  554. if MEMORY_INDEX.exists():
  555. content = MEMORY_INDEX.read_text().strip()
  556. if content:
  557. memories = content
  558. return {
  559. "enabled_tools": [t["name"] for t in TOOLS],
  560. "workspace": str(WORKDIR),
  561. "memories": memories,
  562. }
  563. # ── Agent 循环(简化版,聚焦 cron 调度器) ──
  564. # 教学代码保留基础 Agent 循环。省略 S11 的完整错误恢复。
  565. # cron_scheduler_loop 产出任务;当
  566. # 存在排队任务且没有其他 Agent 轮次运行时,queue_processor_loop 会唤醒该循环。
  567. def agent_loop(messages: list, context: dict) -> dict:
  568. system = get_system_prompt(context)
  569. while True:
  570. # 第 4 层:消费已触发的 cron 任务 → 作为消息注入
  571. fired = consume_cron_queue()
  572. for job in fired:
  573. messages.append({"role": "user",
  574. "content": f"[已调度] {job.prompt}"})
  575. print(f" \033[35m[注入 cron] {job.prompt[:50]}\033[0m")
  576. try:
  577. response = client.messages.create(
  578. model=MODEL, system=system, messages=messages,
  579. tools=TOOLS, max_tokens=8000)
  580. except Exception as e:
  581. messages.append({"role": "assistant", "content": [
  582. {"type": "text",
  583. "text": f"[错误] {type(e).__name__}: {e}"}]})
  584. return context
  585. messages.append({"role": "assistant", "content": response.content})
  586. if response.stop_reason != "tool_use":
  587. return context
  588. results = []
  589. for block in response.content:
  590. if block.type != "tool_use":
  591. continue
  592. print(f"\033[36m> {block.name}\033[0m")
  593. if should_run_background(block.name, block.input):
  594. bg_id = start_background_task(block)
  595. results.append({"type": "tool_result",
  596. "tool_use_id": block.id,
  597. "content": f"[后台任务 {bg_id} 已启动] "
  598. f"完成后结果将可用。"})
  599. else:
  600. output = execute_tool(block)
  601. print(str(output)[:300])
  602. results.append({"type": "tool_result",
  603. "tool_use_id": block.id,
  604. "content": output})
  605. # 将后台工具结果 + 通知合并成一条用户消息
  606. user_content = list(results)
  607. bg_notifications = collect_background_results()
  608. if bg_notifications:
  609. for notif in bg_notifications:
  610. user_content.append({"type": "text", "text": notif})
  611. messages.append({"role": "user", "content": user_content})
  612. context = update_context(context, messages)
  613. system = get_system_prompt(context)
  614. session_history: list = []
  615. session_context = update_context({}, [])
  616. def print_latest_assistant_text(messages: list):
  617. """打印最新 assistant 消息中的文本块。"""
  618. if not messages:
  619. return
  620. msg = messages[-1]
  621. if not isinstance(msg, dict) or msg.get("role") != "assistant":
  622. return
  623. content = msg.get("content", "")
  624. if isinstance(content, str):
  625. print(content)
  626. return
  627. for block in content:
  628. if getattr(block, "type", None) == "text":
  629. print(block.text)
  630. elif isinstance(block, dict) and block.get("type") == "text":
  631. print(block.get("text", ""))
  632. def run_agent_turn_locked(user_query: str | None = None):
  633. """运行一轮 Agent。调用方必须持有 agent_lock。"""
  634. global session_context
  635. if user_query is not None:
  636. session_history.append({"role": "user", "content": user_query})
  637. session_context = agent_loop(session_history, session_context)
  638. session_context = update_context(session_context, session_history)
  639. print_latest_assistant_text(session_history)
  640. print()
  641. def queue_processor_loop():
  642. """当 Agent 空闲时自动投递已触发的 cron 任务。"""
  643. global session_context
  644. while True:
  645. time.sleep(0.2)
  646. if not has_cron_queue():
  647. continue
  648. if not agent_lock.acquire(blocking=False):
  649. continue
  650. try:
  651. if not has_cron_queue():
  652. continue
  653. print("\n \033[35m[队列处理器] 正在投递调度任务\033[0m")
  654. run_agent_turn_locked()
  655. finally:
  656. agent_lock.release()
  657. if __name__ == "__main__":
  658. print("s14: cron 调度器")
  659. print("输入问题后按回车发送。输入 q 退出。\n")
  660. threading.Thread(target=queue_processor_loop, daemon=True).start()
  661. print(" \033[35m[队列处理器] 已启动\033[0m")
  662. while True:
  663. try:
  664. query = input("\033[36ms14 >> \033[0m")
  665. except (EOFError, KeyboardInterrupt):
  666. break
  667. if query.strip().lower() in ("q", "exit", ""):
  668. break
  669. with agent_lock:
  670. run_agent_turn_locked(query)