code.py 28 KB

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