Agent Harness 工程:定时任务与 Scheduler
Agent Harness 工程:定时任务与 Scheduler
早上 7:59,服务器一切安静。
8:00,巡检应该开始:检查磁盘、进程和错误日志,把异常整理成报告。
如果这件事依赖用户准时发一条消息,它就不算自动化;如果让 Agent 执行 sleep(86400),它也迟早会失约。进程重启会让等待消失,每轮执行耗时还会让触发时间一点点漂移。
上一篇的后台 Job 解决了“命令启动后如何不阻塞”,却没有解决“谁在约定时刻按下启动键”。
这一章加入一个确定性的 Scheduler:
时间到达
-> 记录这一次触发
-> 创建持久化 Task
-> 唤醒 Agent
-> 必要时启动后台 Job
模型仍然负责完成工作,但“现在是不是八点”不需要模型判断。
先把四个容易混淆的名字分开
一句“每天晚上跑测试”,在系统里其实对应四种不同记录:
Schedule
什么时候触发:每天 23:00,Asia/Shanghai
ScheduleRun
哪一次触发:2026-07-15 23:00 这一轮
Task
这次要完成什么:运行夜间测试并汇总失败
Job
某条命令执行得怎样:pytest 的进程与日志
它们连起来是:
Schedule -> ScheduleRun -> Task -> Agent -> Job
何时 哪一次 目标 决策 执行
Scheduler 只负责可靠地产生工作;Task System 保存目标,后台 Job 记录具体执行。
不要反复“重开”同一个 Task。Task 一旦完成,就应该永远代表那次已经完成的工作。每天的计划时刻都创建独立 ScheduleRun,再生成一个新的 Task,历史才不会被覆盖。
Schedule 保存规则,ScheduleRun 保存现场快照
先定义最小数据结构:
from dataclasses import asdict, dataclass
from typing import Literal
MisfirePolicy = Literal["skip", "run_once", "catch_up"]
OverlapPolicy = Literal["skip", "allow"]
ScheduleRunStatus = Literal[
"pending",
"dispatched",
"skipped",
]
@dataclass
class Schedule:
id: str
name: str
cron: str
timezone: str
task_title: str
task_description: str
next_run_at: str
misfire_policy: MisfirePolicy = "run_once"
misfire_grace_seconds: int = 300
overlap_policy: OverlapPolicy = "skip"
max_catch_up: int = 3
enabled: bool = True
created_at: str = ""
updated_at: str = ""
schema_version: int = 1
def to_dict(self) -> dict:
return asdict(self)
@dataclass
class ScheduleRun:
id: str
schedule_id: str
scheduled_for: str
task_title: str
task_description: str
status: ScheduleRunStatus = "pending"
task_id: str | None = None
reason: str | None = None
created_at: str = ""
dispatched_at: str | None = None
schema_version: int = 1
def to_dict(self) -> dict:
return asdict(self)
Schedule 保存原始 Cron 和 IANA 时区,例如:
cron: 0 8 * * *
timezone: Asia/Shanghai
内部的 next_run_at 统一保存为 UTC,但原时区不能丢。下一次执行仍要按当地日历计算,不能简单地给上一次 UTC 时间加 24 小时。
ScheduleRun 会复制当时的标题和任务描述。这样即使 Schedule 稍后被修改,已经发生的那次触发也不会悄悄换成新含义。
它只记录分发状态。Task 最后成功、失败还是取消,应从关联 Task 读取,不要在 ScheduleRun 里再维护一份重复状态。
日历时间比倒计时麻烦得多
Cron 只有五个字段,看起来很适合自己解析。真正实现后却会遇到月份天数、闰年、星期语义和夏令时。
这里使用成熟库 croniter,时区使用 Python 标准库的 zoneinfo:
pip install croniter tzdata
第一版只接受标准五字段 Cron:
from datetime import datetime, timedelta, timezone
from zoneinfo import ZoneInfo, ZoneInfoNotFoundError
from croniter import croniter
UTC = timezone.utc
MIN_SCHEDULE_INTERVAL = timedelta(minutes=1)
def parse_utc(value: str) -> datetime:
result = datetime.fromisoformat(value)
if result.tzinfo is None:
raise ValueError("timestamp must include timezone")
return result.astimezone(UTC)
def next_fire_time(
expression: str,
timezone_name: str,
after_utc: datetime,
) -> datetime:
if after_utc.tzinfo is None:
raise ValueError("after_utc must include timezone")
zone = ZoneInfo(timezone_name)
local_base = after_utc.astimezone(zone)
local_next = croniter(
expression,
local_base,
day_or=True,
).get_next(datetime)
return local_next.astimezone(UTC)
def validate_schedule_expression(
expression: str,
timezone_name: str,
now_utc: datetime,
) -> datetime:
if len(expression.split()) != 5:
raise ValueError("cron must contain five fields")
if not croniter.is_valid(expression, strict=True):
raise ValueError("invalid cron expression")
try:
ZoneInfo(timezone_name)
except ZoneInfoNotFoundError as exc:
raise ValueError("unknown IANA timezone") from exc
first = next_fire_time(
expression,
timezone_name,
now_utc,
)
second = next_fire_time(
expression,
timezone_name,
first,
)
if second - first < MIN_SCHEDULE_INTERVAL:
raise ValueError("schedule interval is too short")
return first
day_or=True 采用常见 Cron 语义:当“每月第几天”和“星期几”同时受限时,任意一个匹配就会触发。
这条规则很容易被误读。所以创建 Schedule 时,除了保存表达式,还应该向用户展示接下来几次本地执行时间:
每天 08:00,Asia/Shanghai
接下来三次:
- 2026-07-15 08:00 +08:00
- 2026-07-16 08:00 +08:00
- 2026-07-17 08:00 +08:00
Scheduler 的内部比较全部使用带时区的 UTC datetime。不要混用 naive datetime、服务器本地时间和字符串排序。
夏令时让“每天两点半”不再简单
Asia/Shanghai 当前没有夏令时,但 Scheduler 可能服务世界各地的用户。
在采用夏令时的地区:
- 春季有一段本地时间根本不存在;
- 秋季有一段本地时间会出现两次;
- 时区规则还可能随着法规改变。
因此,“每天 02:30”不一定每天都有且只有一个对应时刻。
IANA 时区和 zoneinfo 能处理 UTC 偏移,但产品仍然要明确自己的策略:不存在的时间是跳过还是顺延,重复时间是执行一次还是两次。
ScheduleRun 使用 UTC scheduled_for 生成 ID,所以两个相同的本地钟表时间如果拥有不同偏移,仍然可以区分。
对只关心固定间隔的内部维护任务,直接使用 UTC 能避开大部分歧义;对面向用户的日历时间,则必须为 DST 切换日编写测试,不能依赖一句“库会处理”。
机器停了一夜,错过的任务怎么办
假设夜间测试原本 23:00 运行,但 Harness 从 22:30 停机到第二天 01:00。
重启后有三种合理选择:
skip
过去就过去,不再补跑
run_once
立即补最近一次,忽略更早遗漏
catch_up
补最近若干次,但严格限制数量
还需要 misfire_grace_seconds。Scheduler 晚几秒醒来是正常调度延迟,不应该立刻算作停机恢复。
Misfire 决定错过的触发是否补跑;Overlap 决定上一轮未结束时能否开始下一轮。
核心规划函数可以写成:
def plan_due_times(
schedule: Schedule,
now_utc: datetime,
) -> tuple[list[datetime], datetime]:
scheduled_for = parse_utc(schedule.next_run_at)
if scheduled_for > now_utc:
return [], scheduled_for
lateness = (now_utc - scheduled_for).total_seconds()
is_misfire = (
lateness > schedule.misfire_grace_seconds
)
if not is_misfire:
due_times = [scheduled_for]
next_base = scheduled_for
elif schedule.misfire_policy == "skip":
due_times = []
next_base = now_utc
elif schedule.misfire_policy == "run_once":
due_times = latest_due_times(
schedule,
now_utc,
limit=1,
)
next_base = now_utc
elif schedule.misfire_policy == "catch_up":
due_times = latest_due_times(
schedule,
now_utc,
limit=max(
1,
min(schedule.max_catch_up, 20),
),
)
next_base = now_utc
else:
raise ValueError("unknown misfire policy")
return (
due_times,
next_fire_time(
schedule.cron,
schedule.timezone,
next_base,
),
)
这里的 latest_due_times() 使用同一个 Cron 与时区向前寻找最近计划时刻,并以当前 next_run_at 为下界。
catch_up 一定要有限。停机几个月后无限回放,很可能瞬间制造成千上万个 Task。
某一次执行失败后的重试,也不能通过倒退 Schedule 游标实现。计划周期属于 Schedule,失败重试属于那一次 Task 或 Job。
同一个计划时刻只能有一个身份
Scheduler 每秒都可能扫描一次,也可能在写入后、返回前崩溃。
如果每次触发都用随机 ID,同一个“早上八点”会被创建很多遍。
稳定幂等键应该是:
schedule_id + scheduled_for_utc
import uuid
def schedule_run_id(
schedule_id: str,
scheduled_for: datetime,
) -> str:
canonical = scheduled_for.astimezone(UTC).isoformat(
timespec="seconds"
)
digest = uuid.uuid5(
uuid.NAMESPACE_URL,
f"{schedule_id}:{canonical}",
).hex[:16]
return f"schedule-run-{digest}"
def scheduled_task_id(run_id: str) -> str:
suffix = run_id.removeprefix("schedule-run-")
return f"task-scheduled-{suffix}"
相同 Schedule 在相同 UTC 时刻,无论扫描多少次,都得到同一个 ScheduleRun ID 和 Task ID。
这使一次触发可以安全地走过崩溃窗口。
先留下脚印,再把时钟拨向下一次
一次 Scheduler tick 的写入顺序非常关键:
1. 在锁内重新读取 Schedule
2. 计算本轮到期时刻
3. 写入稳定 ID 的 ScheduleRun
4. 最后更新 next_run_at
不能先推进游标。
如果 next_run_at 已经越过八点,ScheduleRun 却还没落盘,进程此时崩溃,这次巡检会永远消失。
反过来,如果 ScheduleRun 已经写入,游标还没推进,重启后虽然会再次算到八点,却只会看到同一个稳定 ID,然后补写游标,不会重复创建。
这是一种很实用的恢复思路:
先持久化“发生过什么”,再推进“下一次是什么时候”。
文件型 Store 依靠锁和稳定 ID 缩小风险;数据库实现则应在 (schedule_id, scheduled_for) 上建立唯一约束,并把两次写入放进事务。
ScheduleRun 是一份持久化 Outbox
新建的 ScheduleRun 先处于 pending。Dispatcher 负责把它转换为 Task:
from queue import Queue
def dispatch_schedule_run(
schedule_store: ScheduleStore,
task_store: TaskStore,
runtime_events: Queue[dict],
run: ScheduleRun,
) -> None:
if run.status != "pending":
return
task = create_task(
store=task_store,
task_id=scheduled_task_id(run.id),
title=run.task_title,
description=run.task_description,
)
with schedule_store.locked():
current = schedule_store.read_run_unlocked(run.id)
if current.status != "pending":
return
current.status = "dispatched"
current.task_id = task.id
current.dispatched_at = utc_now_text()
schedule_store.write_run_unlocked(current)
runtime_events.put(
{
"event_id": f"scheduled-task-ready:{run.id}",
"kind": "scheduled_task_ready",
"schedule_id": run.schedule_id,
"schedule_run_id": run.id,
"task_id": task.id,
}
)
如果 Task 已创建、ScheduleRun 还没更新时进程崩溃,Dispatcher 重试仍使用同一个 Task ID。上一篇的 create_task() 遇到相同定义会返回原 Task,不会复制工作。
内存事件只是唤醒提示,Task Store 才是最终事实。如果 Task 已经落盘、事件却丢了,Runtime 重启时扫描未完成的定时 Task,就能重新补发通知。
上一次还没结束,八点又到了
这不是 Misfire,而是 Overlap。
第一版只需要两种策略:
skip
记录本次 skipped,不创建新 Task
allow
创建独立 Task,允许两轮并行
巡检、构建、部署等可能争抢资源或产生副作用的工作,默认更适合 skip。
判断上一轮是否结束时,要读取关联 Task 的真实状态:
- ScheduleRun 仍是
pending,说明尚未完成分发; - Task 是
pending或in_progress,说明工作仍未结束; - 只有
completed、failed、cancelled才算终态。
不要在 ScheduleRun 里复制 Task 的完成状态,也不要用 replace 强行杀掉上一轮。取消 Task 不代表它已经发出的外部请求可以回滚。
Scheduler 主循环不需要模型
Scheduler 只做三件事:收集到期时刻、分发 pending run、等待下一次检查。
def scheduler_loop(state: HarnessState) -> None:
while not state.shutdown_requested:
collect_due_runs(
schedule_store=state.schedule_store,
task_store=state.task_store,
now_utc=utc_now(),
)
for run in (
state.schedule_store.list_pending_runs()
):
dispatch_schedule_run(
state.schedule_store,
state.task_store,
state.runtime_event_queue,
run,
)
state.shutdown_event.wait(timeout=1.0)
实际实现可以根据最早的 next_run_at 计算等待时间,但要设置最大等待上限,以便新增、修改和恢复 Schedule 时及时醒来。
启动顺序也不能颠倒:
- 加载并校验 Schedule 与 ScheduleRun;
- 先恢复已有
pendingrun; - 重建已经分发但未结束的 Task 唤醒;
- 再根据当前时间处理 Misfire;
- 最后进入持续循环。
定时执行是一份长期授权
创建 Schedule 和调用一次工具不同。
用户实际上是在说:
以后即使我不在线,也允许系统按这个时间反复开始工作。
所以 schedule_create、修改、恢复和删除都应该经过明确确认。确认界面至少展示:
- 本地时区;
- 接下来几次执行时间;
- 任务内容;
- Misfire 和 Overlap 策略;
- 无人值守权限 Profile;
- 时间、Token 和费用预算。
Schedule 只保存任务模板和稳定 Profile ID,不保存 Token、密码或任意权限文本。
每次触发仍然要经过当时生效的工具权限、路径和网络策略。一次性的工具审批不能自动升级成永久批准,后台和定时执行也不能成为绕过安全检查的暗门。
小结
这一章为 Harness 加入了一只可恢复的时钟:
- Schedule 保存 Cron、时区和任务模板;
- ScheduleRun 为每个计划时刻留下独立记录;
- Task 与 Job 继续分别管理目标和具体执行;
- IANA 时区和 UTC 避免服务器本地时间混乱;
- Misfire 与 Overlap 使用不同策略;
- 稳定 ID 抵抗重复扫描和崩溃重试;
- 先记录触发、再推进游标,避免永久漏跑;
- Scheduler 只产生工作,不参与模型推理;
- 定时执行被视为长期授权,继续接受权限和预算约束。
现在,每天八点的巡检不需要任何人准时敲一下回车。
但随着 Schedule 不断产生新的 Task,一个 Agent 很快会成为瓶颈。前端、后端和测试明明互不依赖,却仍然排队等待同一个执行者。
下一章加入 Agent Team:让多个拥有稳定身份和专业能力的 Agent 共享任务图,同时又不共享一团越来越混乱的上下文。