| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808 |
- #!/usr/bin/env python3
- """
- s14: Cron 调度器 — 独立守护线程 + 队列处理器。
- 运行: python s14_cron_scheduler/code.py
- 需要: pip install anthropic python-dotenv + .env 中配置 ANTHROPIC_API_KEY
- 相对 s13 的变化:
- - Cron任务 dataclass(id、cron、prompt、recurring、durable)
- - cron_matches:带 DOM/DOW OR 语义的 5 字段 cron 表达式匹配
- - schedule_job / cancel_job:注册/移除 cron 任务(带校验)
- - cron_scheduler_loop:独立守护线程,每 1 秒轮询
- - cron_queue:线程安全队列,调度器写入,队列处理器负责投递
- - queue_processor_loop:当 cron_queue 有任务时自动运行 agent_loop
- - 持久化存储:.scheduled_tasks.json(重启后仍保留)
- - 3 个新工具:schedule_cron、list_crons、cancel_cron
- 四层结构:
- 1. 调度器:守护线程检查时间 → 触发匹配任务
- 2. 队列:cron_queue 将调度器与 Agent 循环解耦
- 3. 队列处理器:当队列中有任务且 Agent 空闲时唤醒 Agent
- 4. 消费者:agent_loop 消费队列任务,并注入到 messages
- """
- import os, subprocess, json, time, random, threading
- from pathlib import Path
- from datetime import datetime
- from dataclasses import dataclass, asdict
- try:
- import readline
- readline.parse_and_bind('set bind-tty-special-chars off')
- except ImportError:
- pass
- from anthropic import Anthropic
- from dotenv import load_dotenv
- load_dotenv(override=True)
- if os.getenv("ANTHROPIC_BASE_URL"):
- os.environ.pop("ANTHROPIC_AUTH_TOKEN", None)
- WORKDIR = Path.cwd()
- MEMORY_DIR = WORKDIR / ".memory"
- MEMORY_INDEX = MEMORY_DIR / "MEMORY.md"
- client = Anthropic(base_url=os.getenv("ANTHROPIC_BASE_URL"))
- MODEL = os.environ["MODEL_ID"]
- # ── 任务系统 (来自 s12,已同步) ──
- TASKS_DIR = WORKDIR / ".tasks"
- TASKS_DIR.mkdir(exist_ok=True)
- @dataclass
- class Task:
- id: str
- subject: str
- description: str
- status: str # pending | in_progress | completed
- owner: str | None
- blockedBy: list[str]
- def _task_path(task_id: str) -> Path:
- return TASKS_DIR / f"{task_id}.json"
- def create_task(subject: str, description: str = "",
- blockedBy: list[str] | None = None) -> Task:
- task = Task(
- id=f"task_{int(time.time())}_{random.randint(0, 9999):04d}",
- subject=subject, description=description,
- status="pending", owner=None,
- blockedBy=blockedBy or [],
- )
- save_task(task)
- return task
- def save_task(task: Task):
- _task_path(task.id).write_text(json.dumps(asdict(task), indent=2))
- def load_task(task_id: str) -> Task:
- return Task(**json.loads(_task_path(task_id).read_text()))
- def list_tasks() -> list[Task]:
- return [Task(**json.loads(p.read_text()))
- for p in sorted(TASKS_DIR.glob("task_*.json"))]
- def get_task(task_id: str) -> str:
- """以 JSON 返回完整任务详情。"""
- task = load_task(task_id)
- return json.dumps(asdict(task), indent=2)
- def can_start(task_id: str) -> bool:
- """检查所有 blockedBy 依赖是否已完成;缺失依赖视为阻塞。"""
- task = load_task(task_id)
- for dep_id in task.blockedBy:
- if not _task_path(dep_id).exists():
- return False
- if load_task(dep_id).status != "completed":
- return False
- return True
- def claim_task(task_id: str, owner: str = "agent") -> str:
- task = load_task(task_id)
- if task.status != "pending":
- status = {"pending": "待处理", "in_progress": "进行中",
- "completed": "已完成"}.get(task.status, task.status)
- return f"任务 {task_id} 当前状态为 {status},无法认领"
- if not can_start(task_id):
- deps = [d for d in task.blockedBy
- if not _task_path(d).exists() or load_task(d).status != "completed"]
- return f"被阻塞于:{deps}"
- task.owner = owner
- task.status = "in_progress"
- save_task(task)
- print(f" \033[36m[认领] {task.subject} → in_progress(负责人:{owner})\033[0m")
- return f"已认领 {task.id} ({task.subject})"
- def complete_task(task_id: str) -> str:
- task = load_task(task_id)
- if task.status != "in_progress":
- status = {"pending": "待处理", "in_progress": "进行中",
- "completed": "已完成"}.get(task.status, task.status)
- return f"任务 {task_id} 当前状态为 {status},无法完成"
- task.status = "completed"
- save_task(task)
- unblocked = [t.subject for t in list_tasks()
- if t.status == "pending" and t.blockedBy and can_start(t.id)]
- print(f" \033[32m[完成] {task.subject} ✓\033[0m")
- msg = f"已完成 {task.id} ({task.subject})"
- if unblocked:
- msg += f"\n已解除阻塞:{', '.join(unblocked)}"
- print(f" \033[33m[已解除阻塞] {', '.join(unblocked)}\033[0m")
- return msg
- # ── 提示词组装 (来自 s10,已同步) ──
- PROMPT_SECTIONS = {
- "identity": "你是一个编码 Agent。直接行动,不要只解释。",
- "tools": "可用工具:bash, read_file, write_file, "
- "create_task, list_tasks, get_task, claim_task, complete_task, "
- "schedule_cron, list_crons, cancel_cron.",
- "workspace": f"工作目录:{WORKDIR}",
- "memory": "有可用的相关记忆时,会在下方注入。",
- }
- def assemble_system_prompt(context: dict) -> str:
- sections = [PROMPT_SECTIONS["identity"],
- PROMPT_SECTIONS["tools"],
- PROMPT_SECTIONS["workspace"]]
- memories = context.get("memories", "")
- if memories:
- sections.append(f"相关记忆:\n{memories}")
- return "\n\n".join(sections)
- _last_context_key, _last_prompt = None, None
- def get_system_prompt(context: dict) -> str:
- global _last_context_key, _last_prompt
- key = json.dumps(context, sort_keys=True, ensure_ascii=False, default=str)
- if key == _last_context_key and _last_prompt:
- return _last_prompt
- _last_context_key = key
- _last_prompt = assemble_system_prompt(context)
- return _last_prompt
- # ── 工具 ──
- def safe_path(p: str) -> Path:
- path = (WORKDIR / p).resolve()
- if not path.is_relative_to(WORKDIR):
- raise ValueError(f"路径逃逸出工作区:{p}")
- return path
- def run_bash(command: str, run_in_background: bool = False) -> str:
- # run_in_background 由 agent_loop 分发处理,不在这里处理
- try:
- r = subprocess.run(command, shell=True, cwd=WORKDIR,
- capture_output=True, text=True, timeout=120)
- out = (r.stdout + r.stderr).strip()
- return out[:50000] if out else "(无输出)"
- except subprocess.TimeoutExpired:
- return "错误:执行超时(120 秒)"
- def run_read(path: str, limit: int | None = None) -> str:
- try:
- lines = safe_path(path).read_text().splitlines()
- if limit and limit < len(lines):
- lines = lines[:limit] + [f"... ({len(lines) - limit} 行更多内容)"]
- return "\n".join(lines)
- except Exception as e:
- return f"错误:{e}"
- def run_write(path: str, content: str) -> str:
- try:
- fp = safe_path(path)
- fp.parent.mkdir(parents=True, exist_ok=True)
- fp.write_text(content)
- return f"已写入 {len(content)} 字节到 {path}"
- except Exception as e:
- return f"错误:{e}"
- # 任务工具
- def run_create_task(subject: str, description: str = "",
- blockedBy: list[str] | None = None) -> str:
- task = create_task(subject, description, blockedBy)
- deps = f"(依赖:{', '.join(blockedBy)})" if blockedBy else ""
- print(f" \033[34m[创建] {task.subject}{deps}\033[0m")
- return f"已创建 {task.id}: {task.subject}{deps}"
- def run_list_tasks() -> str:
- 个任务 = list_tasks()
- if not 个任务:
- return "暂无任务。请使用 create_task 添加任务。"
- lines = []
- for t in 个任务:
- icon = {"pending": "○", "in_progress": "●",
- "completed": "✓"}.get(t.status, "?")
- deps = f"(依赖:{', '.join(t.blockedBy)})" if t.blockedBy else ""
- owner = f" [负责人:{t.owner}]" if t.owner else ""
- status = {"pending": "待处理", "in_progress": "进行中",
- "completed": "已完成"}.get(t.status, t.status)
- lines.append(f" {icon} {t.id}: {t.subject} "
- f"[{status}]{owner}{deps}")
- return "\n".join(lines)
- def run_get_task(task_id: str) -> str:
- try:
- return get_task(task_id)
- except FileNotFoundError:
- return f"错误:任务 {task_id} 未找到"
- def run_claim_task(task_id: str) -> str:
- return claim_task(task_id, owner="agent")
- def run_complete_task(task_id: str) -> str:
- return complete_task(task_id)
- # ── 后台任务 (来自 s13,已同步) ──
- _bg_计数器 = 0
- background_tasks: dict[str, dict] = {}
- background_results: dict[str, str] = {}
- background_lock = threading.Lock()
- def is_slow_operation(tool_name: str, tool_input: dict) -> bool:
- """兜底启发式:判断命令是否可能超过 30 秒。"""
- if tool_name != "bash":
- return False
- cmd = tool_input.get("command", "").lower()
- slow_keywords = ["install", "build", "test", "deploy", "compile",
- "docker build", "pip install", "npm install",
- "cargo build", "pytest", "make"]
- return any(kw in cmd for kw in slow_keywords)
- def should_run_background(tool_name: str, tool_input: dict) -> bool:
- """模型的显式请求优先;否则使用启发式兜底。"""
- if tool_input.get("run_in_background"):
- return True
- return is_slow_operation(tool_name, tool_input)
- def execute_tool(block) -> str:
- """执行工具调用块并返回输出。"""
- handler = {
- "bash": run_bash, "read_file": run_read, "write_file": run_write,
- "create_task": run_create_task, "list_tasks": run_list_tasks,
- "get_task": run_get_task, "claim_task": run_claim_task,
- "complete_task": run_complete_task,
- "schedule_cron": run_schedule_cron, "list_crons": run_list_crons,
- "cancel_cron": run_cancel_cron,
- }.get(block.name)
- if handler:
- return handler(**block.input)
- return f"未知工具:{block.name}"
- def start_background_task(block) -> str:
- """在守护线程中运行工具,并返回后台任务 ID。"""
- global _bg_计数器
- _bg_计数器 += 1
- bg_id = f"bg_{_bg_计数器:04d}"
- cmd = block.input.get("command", block.name)
- def worker():
- result = execute_tool(block)
- with background_lock:
- background_tasks[bg_id]["status"] = "completed"
- background_results[bg_id] = result
- with background_lock:
- background_tasks[bg_id] = {
- "tool_use_id": block.id,
- "command": cmd,
- "status": "running",
- }
- threading.Thread(target=worker, daemon=True).start()
- print(f" \033[33m[后台] 已分发 {bg_id}: {cmd[:40]}\033[0m")
- return bg_id
- def collect_background_results() -> list[str]:
- """将已完成的后台结果收集为 task_notification 消息。"""
- with background_lock:
- ready_ids = [bid for bid, task in background_tasks.items()
- if task["status"] == "completed"]
- notifications = []
- for bg_id in ready_ids:
- with background_lock:
- task = background_tasks.pop(bg_id)
- output = background_results.pop(bg_id, "")
- summary = output[:200] if len(output) > 200 else output
- notifications.append(
- f"<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[后台完成] {bg_id}: "
- f"{task['command'][:40]} ({len(output)} 个字符)\033[0m")
- return notifications
- # ── Cron 调度器 (s14 新增) ──
- DURABLE_PATH = WORKDIR / ".scheduled_tasks.json"
- @dataclass
- class CronJob:
- id: str
- cron: str # "0 9 * * *"
- prompt: str # 触发时要注入的消息
- recurring: bool # 为真 = 重复,为假 = 一次性
- durable: bool # 为真 = 持久化到磁盘
- scheduled_jobs: dict[str, CronJob] = {}
- cron_queue: list[CronJob] = []
- cron_lock = threading.Lock()
- agent_lock = threading.Lock()
- _last_fired: dict[str, str] = {} # job_id → "YYYY-MM-DD HH:MM"
- def _cron_field_matches(field: str, value: int) -> bool:
- """将单个 cron 字段与一个值进行匹配。"""
- if field == "*":
- return True
- if field.startswith("*/"):
- step = int(field[2:])
- return step > 0 and value % step == 0
- if "," in field:
- return any(_cron_field_matches(f.strip(), value)
- for f in field.split(","))
- if "-" in field:
- lo, hi = field.split("-", 1)
- return int(lo) <= value <= int(hi)
- return value == int(field)
- def cron_matches(cron_expr: str, dt: datetime) -> bool:
- """检查 5 字段 cron 表达式是否匹配给定时间。
- 标准 cron 语义:月内日和周内日同时受限时,两者使用 OR。"""
- fields = cron_expr.strip().split()
- if len(fields) != 5:
- return False
- minute, hour, dom, month, dow = fields
- dow_val = (dt.weekday() + 1) % 7 # Python Monday=0 → cron Sunday=0
- m = _cron_field_matches(minute, dt.minute)
- h = _cron_field_matches(hour, dt.hour)
- dom_ok = _cron_field_matches(dom, dt.day)
- month_ok = _cron_field_matches(month, dt.month)
- dow_ok = _cron_field_matches(dow, dow_val)
- # 分钟、小时、月份必须全部匹配
- if not (m and h and month_ok):
- return False
- # DOM 和 DOW:如果两者都有限制,任一匹配即可(OR)
- dom_unconstrained = dom == "*"
- dow_unconstrained = dow == "*"
- if dom_unconstrained and dow_unconstrained:
- return True
- if dom_unconstrained:
- return dow_ok
- if dow_unconstrained:
- return dom_ok
- return dom_ok or dow_ok
- def _validate_cron_field(field: str, lo: int, hi: int) -> str | None:
- """校验单个 cron 字段值是否位于 [lo, hi] 范围内。"""
- if field == "*":
- return None
- if field.startswith("*/"):
- step_str = field[2:]
- if not step_str.isdigit():
- return f"无效步长:{field}"
- step = int(step_str)
- if step <= 0:
- return f"步长必须 > 0:{field}"
- return None
- if "," in field:
- for part in field.split(","):
- err = _validate_cron_field(part.strip(), lo, hi)
- if err: return err
- return None
- if "-" in field:
- parts = field.split("-", 1)
- if not parts[0].isdigit() or not parts[1].isdigit():
- return f"无效范围:{field}"
- a, b = int(parts[0]), int(parts[1])
- if a < lo or a > hi or b < lo or b > hi:
- return f"范围 {field} 超出边界 [{lo}-{hi}]"
- if a > b:
- return f"范围起点大于终点:{field}"
- return None
- if not field.isdigit():
- return f"无效字段:{field}"
- val = int(field)
- if val < lo or val > hi:
- return f"值 {val} 超出边界 [{lo}-{hi}]"
- return None
- def validate_cron(cron_expr: str) -> str | None:
- """校验 cron 表达式。返回错误消息或 None。"""
- fields = cron_expr.strip().split()
- if len(fields) != 5:
- return f"期望 5 个字段,实际得到 {len(fields)}"
- bounds = [(0, 59), (0, 23), (1, 31), (1, 12), (0, 6)]
- names = ["分", "时", "月内日", "月", "周内日"]
- for i, (field, (lo, hi), name) in enumerate(zip(fields, bounds, names)):
- err = _validate_cron_field(field, lo, hi)
- if err:
- return f"{name}: {err}"
- return None
- def save_durable_jobs():
- """将持久任务保存到 .scheduled_tasks.json。"""
- durable = [asdict(j) for j in scheduled_jobs.values() if j.durable]
- DURABLE_PATH.write_text(json.dumps(durable, indent=2))
- def load_durable_jobs():
- """启动时从磁盘加载持久任务。"""
- if not DURABLE_PATH.exists():
- return
- try:
- jobs = json.loads(DURABLE_PATH.read_text())
- for j in jobs:
- job = CronJob(**j)
- err = validate_cron(job.cron)
- if err:
- print(f" \033[31m[cron] 跳过无效任务 {job.id}: {err}\033[0m")
- continue
- scheduled_jobs[job.id] = job
- valid = [j for j in jobs if j["id"] in scheduled_jobs]
- if valid:
- print(f" \033[35m[cron] 已加载 {len(valid)} 个持久任务\033[0m")
- except Exception:
- pass
- def schedule_job(cron: str, prompt: str, recurring: bool = True,
- durable: bool = True) -> CronJob | str:
- """注册一个新的 cron 任务。返回 CronJob 或错误字符串。"""
- err = validate_cron(cron)
- if err:
- return err
- job = CronJob(
- id=f"cron_{random.randint(0, 999999):06d}",
- cron=cron, prompt=prompt,
- recurring=recurring, durable=durable,
- )
- with cron_lock:
- scheduled_jobs[job.id] = job
- if durable:
- save_durable_jobs()
- print(f" \033[35m[cron 注册] {job.id} '{cron}' → {prompt[:40]}\033[0m")
- return job
- def cancel_job(job_id: str) -> str:
- """取消一个 cron 任务。"""
- with cron_lock:
- job = scheduled_jobs.pop(job_id, None)
- if not job:
- return f"任务 {job_id} 未找到"
- if job.durable:
- save_durable_jobs()
- print(f" \033[31m[cron 取消] {job_id}\033[0m")
- return f"已取消 {job_id}"
- def cron_scheduler_loop():
- """独立守护线程:每秒轮询一次并触发匹配任务。
- 单个任务出错会被捕获,避免整个调度线程退出。"""
- while True:
- time.sleep(1)
- now = datetime.now()
- # 带日期感知的标记,防止每日任务从第 2 天起被跳过
- minute_marker = now.strftime("%Y-%m-%d %H:%M")
- with cron_lock:
- for job in list(scheduled_jobs.values()):
- try:
- if cron_matches(job.cron, now):
- if _last_fired.get(job.id) != minute_marker:
- cron_queue.append(job)
- _last_fired[job.id] = minute_marker
- print(f" \033[35m[cron 触发] {job.id} → "
- f"{job.prompt[:40]}\033[0m")
- if not job.recurring:
- scheduled_jobs.pop(job.id, None)
- if job.durable:
- save_durable_jobs()
- except Exception as e:
- print(f" \033[31m[cron 错误] {job.id}: {e}\033[0m")
- def consume_cron_queue() -> list[CronJob]:
- """消费 cron_queue 中已触发的任务(由 agent_loop 调用)。"""
- with cron_lock:
- fired = list(cron_queue)
- cron_queue.clear()
- return fired
- def has_cron_queue() -> bool:
- """返回是否有已触发但尚未投递的 cron 任务。"""
- with cron_lock:
- return bool(cron_queue)
- # 启动时加载持久任务,然后启动调度线程
- load_durable_jobs()
- threading.Thread(target=cron_scheduler_loop, daemon=True).start()
- print(" \033[35m[cron] 调度线程已启动\033[0m")
- # ── Cron 工具 ──
- def run_schedule_cron(cron: str, prompt: str,
- recurring: bool = True, durable: bool = True) -> str:
- result = schedule_job(cron, prompt, recurring, durable)
- if isinstance(result, str):
- return f"错误:{result}"
- return f"已调度 {result.id}: '{cron}' → {prompt}"
- def run_list_crons() -> str:
- with cron_lock:
- jobs = list(scheduled_jobs.values())
- if not jobs:
- return "暂无 cron 任务。请使用 schedule_cron 添加一个。"
- lines = []
- for j in jobs:
- tag = "重复" if j.recurring else "一次性"
- dur = "持久化" if j.durable else "仅本会话"
- lines.append(f" {j.id}: '{j.cron}' → {j.prompt[:40]} "
- f"[{tag}, {dur}]")
- return "\n".join(lines)
- def run_cancel_cron(job_id: str) -> str:
- return cancel_job(job_id)
- # ── 工具定义 ──
- TOOLS = [
- {"name": "bash", "description": "运行一条 shell 命令。",
- "input_schema": {"type": "object",
- "properties": {
- "command": {"type": "string"},
- "run_in_background": {"type": "boolean"}},
- "required": ["command"]}},
- {"name": "read_file", "description": "读取文件内容。",
- "input_schema": {"type": "object",
- "properties": {"path": {"type": "string"},
- "limit": {"type": "integer"}},
- "required": ["path"]}},
- {"name": "write_file", "description": "向文件写入内容。",
- "input_schema": {"type": "object",
- "properties": {"path": {"type": "string"},
- "content": {"type": "string"}},
- "required": ["path", "content"]}},
- {"name": "create_task",
- "description": "创建一个新任务,可选 blockedBy 依赖。",
- "input_schema": {"type": "object",
- "properties": {
- "subject": {"type": "string"},
- "description": {"type": "string"},
- "blockedBy": {"type": "array",
- "items": {"type": "string"}}},
- "required": ["subject"]}},
- {"name": "list_tasks",
- "description": "列出所有任务及其状态、负责人和依赖。",
- "input_schema": {"type": "object", "properties": {},
- "required": []}},
- {"name": "get_task",
- "description": "按 ID 获取指定任务的完整详情。",
- "input_schema": {"type": "object",
- "properties": {"task_id": {"type": "string"}},
- "required": ["task_id"]}},
- {"name": "claim_task",
- "description": "认领一个待处理任务。设置 owner,并将状态改为 in_progress。",
- "input_schema": {"type": "object",
- "properties": {"task_id": {"type": "string"}},
- "required": ["task_id"]}},
- {"name": "complete_task",
- "description": "完成一个进行中的任务。报告被解除阻塞的下游任务。",
- "input_schema": {"type": "object",
- "properties": {"task_id": {"type": "string"}},
- "required": ["task_id"]}},
- {"name": "schedule_cron",
- "description": "调度一个 cron 任务。cron 为 5 字段:分 时 月内日 月 周内日。",
- "input_schema": {"type": "object",
- "properties": {
- "cron": {"type": "string",
- "description": "5 字段 cron 表达式"},
- "prompt": {"type": "string",
- "description": "触发时要注入的消息"},
- "recurring": {"type": "boolean",
- "description": "为真表示重复,为假表示一次性"},
- "durable": {"type": "boolean",
- "description": "为真表示持久化到磁盘"}},
- "required": ["cron", "prompt"]}},
- {"name": "list_crons",
- "description": "列出所有已注册的 cron 任务。",
- "input_schema": {"type": "object", "properties": {},
- "required": []}},
- {"name": "cancel_cron",
- "description": "按 ID 取消 cron 任务。",
- "input_schema": {"type": "object",
- "properties": {"job_id": {"type": "string"}},
- "required": ["job_id"]}},
- ]
- # ── 上下文 ──
- def update_context(context: dict, messages: list) -> dict:
- """从真实状态推导上下文。"""
- memories = ""
- if MEMORY_INDEX.exists():
- content = MEMORY_INDEX.read_text().strip()
- if content:
- memories = content
- return {
- "enabled_tools": [t["name"] for t in TOOLS],
- "workspace": str(WORKDIR),
- "memories": memories,
- }
- # ── Agent 循环(简化版,聚焦 cron 调度器) ──
- # 教学代码保留基础 Agent 循环。省略 S11 的完整错误恢复。
- # cron_scheduler_loop 产出任务;当
- # 存在排队任务且没有其他 Agent 轮次运行时,queue_processor_loop 会唤醒该循环。
- def agent_loop(messages: list, context: dict) -> dict:
- system = get_system_prompt(context)
- while True:
- # 第 4 层:消费已触发的 cron 任务 → 作为消息注入
- fired = consume_cron_queue()
- for job in fired:
- messages.append({"role": "user",
- "content": f"[已调度] {job.prompt}"})
- print(f" \033[35m[注入 cron] {job.prompt[:50]}\033[0m")
- try:
- response = client.messages.create(
- model=MODEL, system=system, messages=messages,
- tools=TOOLS, max_tokens=8000)
- except Exception as e:
- messages.append({"role": "assistant", "content": [
- {"type": "text",
- "text": f"[错误] {type(e).__name__}: {e}"}]})
- return context
- messages.append({"role": "assistant", "content": response.content})
- if response.stop_reason != "tool_use":
- return context
- results = []
- for block in response.content:
- if block.type != "tool_use":
- continue
- print(f"\033[36m> {block.name}\033[0m")
- if should_run_background(block.name, block.input):
- bg_id = start_background_task(block)
- results.append({"type": "tool_result",
- "tool_use_id": block.id,
- "content": f"[后台任务 {bg_id} 已启动] "
- f"完成后结果将可用。"})
- else:
- output = execute_tool(block)
- print(str(output)[:300])
- results.append({"type": "tool_result",
- "tool_use_id": block.id,
- "content": output})
- # 将后台工具结果 + 通知合并成一条用户消息
- user_content = list(results)
- bg_notifications = collect_background_results()
- if bg_notifications:
- for notif in bg_notifications:
- user_content.append({"type": "text", "text": notif})
- messages.append({"role": "user", "content": user_content})
- context = update_context(context, messages)
- system = get_system_prompt(context)
- session_history: list = []
- session_context = update_context({}, [])
- def print_latest_assistant_text(messages: list):
- """打印最新 assistant 消息中的文本块。"""
- if not messages:
- return
- msg = messages[-1]
- if not isinstance(msg, dict) or msg.get("role") != "assistant":
- return
- content = msg.get("content", "")
- if isinstance(content, str):
- print(content)
- return
- for block in content:
- if getattr(block, "type", None) == "text":
- print(block.text)
- elif isinstance(block, dict) and block.get("type") == "text":
- print(block.get("text", ""))
- def run_agent_turn_locked(user_query: str | None = None):
- """运行一轮 Agent。调用方必须持有 agent_lock。"""
- global session_context
- if user_query is not None:
- session_history.append({"role": "user", "content": user_query})
- session_context = agent_loop(session_history, session_context)
- session_context = update_context(session_context, session_history)
- print_latest_assistant_text(session_history)
- print()
- def queue_processor_loop():
- """当 Agent 空闲时自动投递已触发的 cron 任务。"""
- global session_context
- while True:
- time.sleep(0.2)
- if not has_cron_queue():
- continue
- if not agent_lock.acquire(blocking=False):
- continue
- try:
- if not has_cron_queue():
- continue
- print("\n \033[35m[队列处理器] 正在投递调度任务\033[0m")
- run_agent_turn_locked()
- finally:
- agent_lock.release()
- if __name__ == "__main__":
- print("s14: cron 调度器")
- print("输入问题后按回车发送。输入 q 退出。\n")
- threading.Thread(target=queue_processor_loop, daemon=True).start()
- print(" \033[35m[队列处理器] 已启动\033[0m")
- while True:
- try:
- query = input("\033[36ms14 >> \033[0m")
- except (EOFError, KeyboardInterrupt):
- break
- if query.strip().lower() in ("q", "exit", ""):
- break
- with agent_lock:
- run_agent_turn_locked(query)
|