Replace the dual task/monitor model, NO_REPLY string protocol, and Agent self-delivery with a single Scheduled Run path: claim-time JobRun snapshots, isolated Root/named Agent execution, exactly-once complete_scheduled_run termination, and Scheduler-owned policy delivery through a persistent outbox. - SQLite v11: drop job_kind/model/delete_after_run, add job_runs with status/outcome joint constraints and delivery lease columns; one-shot BEGIN IMMEDIATE migration with atomic rollback. - Non-blocking JoinSet event loop with bounded run/delivery concurrency; terminal commit before any channel I/O; recover unfinished runs as unknown. - ExecutionOrigin::Scheduled propagates to descendants, completion sink is top-level only, background delegation downgrades to foreground. - Typed delivery receipts, fixed target_session_id, idempotent scheduled:<job_run_id> history insert. - New cron_runs read-only tool; cron_add/update drop kind/model; WebUI and Health consume the same JobRun projection. - Bump version to 1.22.0.
1143 lines
57 KiB
Markdown
1143 lines
57 KiB
Markdown
# 统一定时任务执行与投递设计
|
||
|
||
状态:已实施
|
||
|
||
目标数据库版本:v11
|
||
|
||
适用范围:`scheduler`、Scheduled Agent、Cron tools、SQLite、WebUI 任务页
|
||
|
||
## 1. 背景
|
||
|
||
当前定时任务同时存在四组互相耦合的概念:
|
||
|
||
- `JobKind::Task` / `JobKind::Monitor`;
|
||
- `DeliveryPolicy::Direct` / `Always` / `OnAlert` / `Never`;
|
||
- Agent 自己调用 `send_message` 的旧执行路径,以及 Scheduler 托管投递的新执行路径;
|
||
- 通过 `NO_REPLY`、`NO_REPLY[INFO]`、`NO_REPLY[FAIL]`、`NO_REPLY[REFUSE]` 等字符串猜测运行结果的协议。
|
||
|
||
这导致任务“做什么”和“是否投递”没有形成正交模型。特别是 `NO_REPLY[INFO] XXXX`、Markdown 包裹、缺少冒号或模型附加解释时,字符串解析会把本应静默的结果当作普通内容投递。继续放宽正则只会扩大不确定协议,无法从根本上解决问题。
|
||
|
||
本设计把所有 AI 定时任务统一成一种 **Scheduled Run**:任务配置只声明调度、执行 Agent、完整 prompt、目标和投递策略;每次运行必须通过运行时注入的终结工具提交结构化结果;Scheduler 根据结构化结果和投递策略做确定性决策。
|
||
|
||
## 2. 设计目标
|
||
|
||
1. 删除 `task` / `monitor` 运行类型,任务语义只由 prompt 表达。
|
||
2. 删除所有 `NO_REPLY[...]` 魔法字符串和自然语言结果解析。
|
||
3. 删除 Agent 自行投递的 `Direct` 路径,所有投递由 Scheduler 拥有。
|
||
4. 统一普通通知、异常巡检和后台维护的执行路径。
|
||
5. 定时任务可选择 Root 或一个命名 Agent;命名 Agent 的 Provider、模型、工具、Skills 和委托边以 AgentCatalog 为准。
|
||
6. Scheduled Run 不产生脱离当前执行的后台子 Agent;请求后台委托时自动按前台委托执行。
|
||
7. 执行结果先持久化,再投递;进程重启后可以恢复未完成投递。
|
||
8. 无法确认是否完成的执行标记为 `unknown`,不盲目重跑同一 occurrence。
|
||
9. 使用一次性、原子、可失败回滚的 SQLite v11 自动迁移;迁移后运行时代码只认识新 schema 和新枚举。
|
||
10. `never` 等静默任务的每次运行仍可通过管理 API、WebUI 和只读工具审计。
|
||
|
||
## 3. 非目标
|
||
|
||
- 不提供分布式 Scheduler 或跨节点共识。
|
||
- 不承诺外部渠道 exactly-once。渠道发送成功但本地确认前进程崩溃时,仍可能重复投递。
|
||
- 不保存 Scheduled Run 的用户聊天历史;每次执行默认是隔离上下文。
|
||
- 不自动把上一次执行的工具结果带到下一次执行。
|
||
- 不增加 `shell`、Webhook 等第二种 Job 执行类型;本期所有任务仍是 Agent 任务。
|
||
- 不保留旧字段、旧枚举、旧工具参数或旧字符串协议的运行时兼容分支。
|
||
- `never` 不提供“失败时例外通知”;选择该策略意味着 `failed`、`refused`、`timed_out` 和 `unknown` 也只进入审计、WebUI 与 Health。
|
||
|
||
## 4. 参考项目取舍
|
||
|
||
本设计采用以下参考经验:
|
||
|
||
- ZeroClaw:Cron 明确归属于某个 Agent,复用该 Agent 的身份、模型和权限;Cron 本身是一项顶层运行,而不是子 Agent 完成消息。
|
||
- Hermes:有限生命周期调用方不能接收后台完成结果时,委托回退为同步;进程崩溃后把副作用不确定的执行标记为 `unknown`,不自动重放。
|
||
- PicoBot:保留现有 SQLite 租约、Scheduler 集中投递、MessageBus 目标有序发送和命名 AgentCatalog。
|
||
|
||
明确不采用:
|
||
|
||
- Hermes 大量平台特例、JSON Job 主存储和多套 Scheduler provider;
|
||
- ZeroClaw 仅清除过期锁后重新执行的恢复方式;
|
||
- pi 的子进程式、无持久化 Subagent 示例;
|
||
- 把 Cron 结果伪装成后台 Agent Inbox 事件。
|
||
|
||
## 5. 核心模型
|
||
|
||
### 5.1 ScheduledJob
|
||
|
||
ScheduledJob 只回答五个问题:
|
||
|
||
```text
|
||
何时执行:schedule
|
||
由谁执行:agent_id
|
||
执行什么:prompt
|
||
发到哪里:channel + chat_id
|
||
何时投递:delivery_policy
|
||
```
|
||
|
||
目标 Rust 类型:
|
||
|
||
```rust
|
||
pub struct ScheduledJob {
|
||
pub id: String,
|
||
pub name: String,
|
||
pub schedule: Schedule,
|
||
pub prompt: String,
|
||
/// None 表示 Root;Some(id) 表示当前 AgentCatalog 中的命名 Agent。
|
||
pub agent_id: Option<String>,
|
||
pub channel: String,
|
||
pub chat_id: String,
|
||
pub delivery_policy: DeliveryPolicy,
|
||
pub enabled: bool,
|
||
pub next_run_at: i64,
|
||
pub last_run_at: Option<i64>,
|
||
pub last_outcome: Option<ScheduledOutcomeKind>,
|
||
pub created_at: i64,
|
||
pub updated_at: i64,
|
||
// durable claim
|
||
pub locked_at: Option<i64>,
|
||
pub lock_owner: Option<String>,
|
||
pub lease_until: Option<i64>,
|
||
}
|
||
```
|
||
|
||
删除以下字段:
|
||
|
||
- `job_kind`:与 `delivery_policy` 重复;
|
||
- `model`:当前并未实际应用,且会绕开命名 Agent 的 Provider/Model 定义;
|
||
- `delete_after_run`:当前工具不开放且始终写 `false`;`Schedule::At` 在 claim 事务内立即禁用,用户显式删除即可。
|
||
|
||
### 5.2 DeliveryPolicy
|
||
|
||
只保留三个值:
|
||
|
||
```rust
|
||
pub enum DeliveryPolicy {
|
||
Always,
|
||
OnAlert,
|
||
Never,
|
||
}
|
||
```
|
||
|
||
- `always`:任何终态都投递;
|
||
- `on_alert`:仅 `alert`、`failed`、`refused`、`unknown` 投递;
|
||
- `never`:任何终态都不投递,只保留执行记录。
|
||
|
||
删除 `Direct`。Agent 不再有自行完成最终投递的职责。
|
||
|
||
### 5.3 ScheduledOutcome
|
||
|
||
Agent 可以主动提交四种结果,运行时恢复还可以产生 `unknown`:
|
||
|
||
```rust
|
||
pub enum ScheduledOutcome {
|
||
Ok { message: String },
|
||
Alert { message: String },
|
||
Failed { message: String },
|
||
Refused { message: String },
|
||
Unknown { message: String }, // 仅运行时生成,模型不能提交
|
||
}
|
||
```
|
||
|
||
语义:
|
||
|
||
- `ok`:任务成功完成,没有需要用户关注的异常;
|
||
- `alert`:任务成功完成并发现需要用户关注的事实;
|
||
- `failed`:检查或任务没有可靠完成;
|
||
- `refused`:安全策略、授权或 Agent 自身约束拒绝执行;
|
||
- `unknown`:进程在持久化终态前退出,无法判断外部副作用是否发生。
|
||
|
||
投递矩阵:
|
||
|
||
| Outcome | `always` | `on_alert` | `never` |
|
||
|---|---:|---:|---:|
|
||
| `ok` | 投递 | 静默 | 静默 |
|
||
| `alert` | 投递 | 投递 | 静默 |
|
||
| `failed` | 投递 | 投递 | 静默 |
|
||
| `refused` | 投递 | 投递 | 静默 |
|
||
| `unknown` | 投递 | 投递 | 静默 |
|
||
|
||
Agent 只报告事实,不能通过工具参数提供 `notify=true/false`。投递权始终属于 Scheduler。
|
||
|
||
### 5.4 ScheduledRunStatus
|
||
|
||
运行状态只描述生命周期,不表达业务是否正常:
|
||
|
||
```rust
|
||
pub enum ScheduledRunStatus {
|
||
Claimed,
|
||
Running,
|
||
Completed,
|
||
Failed,
|
||
TimedOut,
|
||
Cancelled,
|
||
Interrupted,
|
||
Unknown,
|
||
}
|
||
```
|
||
|
||
成功调用终结工具后,生命周期状态是 `completed`,业务结果由 `outcome` 表达;例如 `completed + failed` 表示 Agent 正常结束并明确报告检查失败。Provider 错误或结果协议缺失则是 `failed + failed`。
|
||
|
||
两者必须满足以下组合约束,不能独立随意取值:
|
||
|
||
| Run status | 合法 Outcome |
|
||
|---|---|
|
||
| `claimed` / `running` | `NULL` |
|
||
| `completed` | `ok` / `alert` / `failed` / `refused` |
|
||
| `failed` / `timed_out` / `cancelled` / `interrupted` | `failed` |
|
||
| `unknown` | `unknown` |
|
||
|
||
因此 `unknown` status 与 `unknown` outcome 恒同现,且只能由平台恢复路径生成:包括启动时恢复未完成运行,以及迁移无法映射的旧状态;模型不能提交它。历史迁移必须先按生命周期状态归一化,再决定 Outcome,不能产生 `unknown + ok` 等非法组合。
|
||
|
||
## 6. 总体架构
|
||
|
||
```mermaid
|
||
flowchart TD
|
||
EventLoop[Scheduler event loop] -->|tick + free slots| Claim[Storage claim occurrence]
|
||
EventLoop -->|run completion| Reap[reap JoinSet]
|
||
EventLoop -->|pending due| Delivery[bounded delivery drain]
|
||
Claim --> RunRow[(job_runs: claimed)]
|
||
Claim --> Advance[提前推进 next_run / 禁用 At]
|
||
RunRow --> JoinSet[bounded JoinSet]
|
||
JoinSet --> Resolve[ScheduledAgentRunner]
|
||
Resolve --> Catalog[Root config / AgentCatalog]
|
||
Resolve --> Agent[AgentLoop]
|
||
Agent --> Tools[受限工具 + complete_scheduled_run]
|
||
Tools --> Outcome[ScheduledOutcome]
|
||
Outcome --> Reap
|
||
Reap --> Commit[原子提交 Run 终态与 delivery_status]
|
||
Commit -->|suppressed / not_requested| Done[完成]
|
||
Commit -->|pending| Delivery
|
||
Delivery --> Bus[MessageBus / OutboundDispatcher]
|
||
Bus --> Channel[Channel]
|
||
Channel --> Ack[delivery_status = delivered]
|
||
```
|
||
|
||
组件职责:
|
||
|
||
| 组件 | 职责 | 明确不负责 |
|
||
|---|---|---|
|
||
| `Scheduler` | 非阻塞事件循环、领取、并发上限、执行超时、恢复、投递 drain | 解释模型自然语言 |
|
||
| `ScheduledAgentRunner` | 解析 Root/命名 Agent、构造隔离上下文、运行 Agent | 渠道发送、next-run 计算 |
|
||
| `AgentLoop` | 模型/工具循环、识别终结工具已提交 | Job 状态或渠道策略 |
|
||
| `complete_scheduled_run` | Schema 校验并提交一次结构化 Outcome | 消息发送、数据库写入 |
|
||
| `Storage` | occurrence、租约、终态、投递状态的原子转换 | Provider 和 Channel I/O |
|
||
| `OutboundDispatcher` | 目标有序、瞬态错误重试、渠道调用、类型化投递回执 | Scheduled Run 业务分类 |
|
||
|
||
Scheduler 主循环维护一个 `JoinSet` 和固定并发上限,不再像旧实现一样等待整批 Job 全部结束后才进入下一次 poll:
|
||
|
||
1. tick 到达时先回收已完成 Run,并按空闲槽位领取新 occurrence;
|
||
2. 每个 Run 终态提交后立即尝试认领并投递自己的 pending 结果;
|
||
3. 启动和后续 tick 另行有界回收遗留 pending delivery;
|
||
4. 长任务只占用一个执行槽,不阻塞其他到期 Job 的领取或已完成结果的投递;
|
||
5. Run task、delivery task 和主循环都由 TaskSupervisor 拥有,关停时停止 admission、取消并有界 join。
|
||
|
||
这仍是一个 Scheduler、一种 Job 和一条投递路径;`JoinSet` 只解决生命周期阻塞,不引入新的任务类型或第二套执行器。
|
||
|
||
## 7. 结构化终结协议
|
||
|
||
### 7.1 工具定义
|
||
|
||
所有 Scheduled Run 都运行时注入同一个工具:
|
||
|
||
```json
|
||
{
|
||
"name": "complete_scheduled_run",
|
||
"parameters": {
|
||
"type": "object",
|
||
"properties": {
|
||
"outcome": {
|
||
"type": "string",
|
||
"enum": ["ok", "alert", "failed", "refused"]
|
||
},
|
||
"message": {
|
||
"type": "string",
|
||
"minLength": 1,
|
||
"maxLength": 16384
|
||
}
|
||
},
|
||
"required": ["outcome", "message"],
|
||
"additionalProperties": false
|
||
}
|
||
}
|
||
```
|
||
|
||
约束:
|
||
|
||
- `runtime_injected() == true`,Agent definition 的 `tools` 不得声明它;
|
||
- 只在 `ToolExecutionContext` 带 Scheduled completion sink 时可用;
|
||
- `exclusive() == true`,不与其他工具并行;
|
||
- 每次运行只接受一次;
|
||
- `unknown` 不出现在模型可见 schema 中;
|
||
- message 去除首尾空白后必须非空,并在字符边界安全截断到上限。
|
||
|
||
### 7.2 AgentLoop 终结行为
|
||
|
||
运行上下文增加两个用途不同的字段:
|
||
|
||
```rust
|
||
pub enum ExecutionOrigin {
|
||
Interactive,
|
||
Scheduled { job_run_id: i64 },
|
||
}
|
||
|
||
pub struct ToolExecutionContext {
|
||
// existing fields...
|
||
pub execution_origin: ExecutionOrigin,
|
||
pub scheduled_completion: Option<Arc<ScheduledCompletionSink>>,
|
||
}
|
||
```
|
||
|
||
- `execution_origin` 是事实标记,普通构造器缺省为 `Interactive`,`execute_scheduled()` 显式设为 Scheduled;每次 Coordinator 构造子 Agent 的 ToolExecutionContext 时从父 context 复制,因此所有子孙都继承。它用来统一禁止 signal/inbox、把 background 委托降级为 foreground;不得通过字符串形式的 session ID 或 Agent ID 推断来源。
|
||
- `scheduled_completion` 是顶层独占能力,只存在于 Scheduled 顶层 Agent,不得传给任何子 Agent。
|
||
|
||
终结工具用一次性 compare-and-set 写入 `ScheduledCompletionSink`。AgentLoop 在每个工具调用之后检查 sink:
|
||
|
||
1. 未提交:继续普通工具循环;
|
||
2. 已提交:把本次终结工具调用和结果写入 Agent transcript;
|
||
3. 同一 Provider 消息中排在终结工具之后的工具调用统一归约为 `Cancelled`,原因是 Scheduled Run 已完成;
|
||
4. 不再调用 Provider,立即返回结构化 Scheduled 结果。
|
||
|
||
终结工具必须是最后的语义动作,但运行时不依赖模型遵守这一提示来保证结束。
|
||
|
||
### 7.3 Fail-closed
|
||
|
||
只有合法的结构化提交才能得到 `ok` 并可能静默。以下情况统一生成 `failed`,不读取文本猜测:
|
||
|
||
- Agent 输出普通最终文本但没有调用终结工具;
|
||
- 输出任何 `NO_REPLY` 变体;
|
||
- 工具参数不合法且模型未修正;
|
||
- 达到最大工具迭代次数;
|
||
- Provider 错误;
|
||
- Agent 返回空结果;
|
||
- 执行超时。
|
||
|
||
普通最终文本可以截断后保存在 `diagnostic`,但不得作为通知正文。用户通知由 Scheduler 生成,例如:
|
||
|
||
```text
|
||
定时任务「生产站点巡检」未能完成:Agent 未提交结构化运行结果。
|
||
```
|
||
|
||
## 8. Agent 解析与权限
|
||
|
||
### 8.1 Root 任务
|
||
|
||
`agent_id = NULL` 表示使用当前 Root Provider、模型、上下文窗口和 Skills。工具从 Root registry 派生,但移除:
|
||
|
||
- `send_message`;
|
||
- `cron_add`、`cron_update`、`cron_remove`、`cron_enable`、`cron_disable`、`cron_list`、`cron_runs`;
|
||
- `reload_config`;
|
||
- 只服务于交互会话或后台收件箱控制的工具。
|
||
|
||
再注入 `complete_scheduled_run` 和 `delegate`。`delegate` 的目标仍受 AgentCatalog `root_can_delegate()` 约束,其 background 参数在 Scheduled origin 下统一降级为 foreground。Root Scheduled Run 不加载用户聊天历史。
|
||
|
||
### 8.2 命名 Agent 任务
|
||
|
||
`agent_id = Some(id)` 必须从当前不可变 AgentCatalog 解析:
|
||
|
||
- Provider profile、模型、上下文窗口来自 Agent definition;
|
||
- 工具和 Skills 完全按 definition allowlist;
|
||
- `complete_scheduled_run` 作为固有运行时工具额外注入,不构成权限扩张;
|
||
- 委托边按 definition 的 `delegates`;
|
||
- Job 不提供 per-job tools 或 model override。
|
||
|
||
创建和更新任务时校验当前 `agent_id`。如果之后重载删除或禁用了该 Agent,Gateway 不因一个 Job 无法启动;该 occurrence 记录为 `failed`,并按投递策略通知。Health 页面同时报告悬空 Agent 引用。
|
||
|
||
### 8.3 AgentRun 记录
|
||
|
||
Scheduled 顶层执行仍复用 `agent_runs` 审计和 transcript,但不增加新的 `AgentRunMode`:
|
||
|
||
- `mode = foreground`,因为 Scheduler 同步等待它;
|
||
- `caller_agent_id = SCHEDULER`;
|
||
- `caller_scope_id = scheduled:<job_id>`;
|
||
- `root_session_id = scheduled-run:<job_run_id>`,它是审计 scope,不是可接收 Inbox 的 Session;
|
||
- 顶层 `ToolExecutionContext.session_id = scheduled-run:<job_run_id>`,供 browser 等 session-scoped 状态工具隔离使用;
|
||
- 顶层 `ToolExecutionContext.agent` 必须指向该 AgentRun 的完整 AgentExecutionContext;Agent 访问授权以其中的 `root_session_id` 和 ancestry 为准,不以 `session_id` 字符串代替;
|
||
- `execution_origin = Scheduled { job_run_id }`;
|
||
- `completion_slot_reserved = false`;
|
||
- `job_runs.agent_run_id` 关联顶层 AgentRun。
|
||
|
||
这是执行来源的区别,不是第三种并发模式,因此不扩展 `foreground/background` 枚举。
|
||
|
||
Scheduled 顶层及其后代不得创建 `agent_session_state`,不得调用 `reserve_completion_slots`,也不得插入 signal/completion inbox event。`recover_agent_state()` 的通用容量对账不应看到合成 Scheduled scope。
|
||
|
||
## 9. 子 Agent 语义
|
||
|
||
Scheduled Agent 可以调用 `delegate`,但所有子任务必须在本次 Scheduled Run 内收敛:
|
||
|
||
```text
|
||
delegate(mode=foreground) → 正常执行
|
||
delegate(mode=background) → 自动改为 foreground,并在工具结果中说明降级
|
||
```
|
||
|
||
理由:
|
||
|
||
- `scheduled-run:<id>` 不是用户 Session,不能消费 durable inbox continuation;
|
||
- Scheduler 必须在一次 run 内得到完整 Outcome;
|
||
- 避免创建无法投递的 `cron:<id>` completion;
|
||
- 前台子 Agent 仍可并行批量执行并由父 Agent 汇总。
|
||
|
||
Scheduled Agent 不能使用 `emit_signal` 向用户 Session 建立旁路。即使命名 Agent definition 声明了 signal contract,Scheduled origin 也不得注入 `emit_signal`,不得预留 completion slot。子 Agent 的最终结果作为普通工具结果返回父 Scheduled Agent,只有父 Agent 可以调用 `complete_scheduled_run`。
|
||
|
||
`ExecutionOrigin::Scheduled` 必须随每一层子 Agent 继承,使嵌套委托也保持同步收敛语义;`ScheduledCompletionSink` 则是顶层运行能力,Coordinator 构造子 Agent context 时必须显式清空,子 Agent registry 也不得注入终结工具。二者不能合并为同一个可选字段,否则清空 sink 后嵌套子 Agent 会丢失 Scheduled 限制。
|
||
|
||
Root Scheduled Agent 的第一次委托必须走 AgentCatalog 的 `root_can_delegate()` 规则,不能把合成的 `SCHEDULER` 或 `ROOT` 身份误当作命名 Agent 传给 `can_delegate()`。命名 Scheduled Agent 的后续委托仍走 `can_delegate(caller, target)`。这一区分只影响委托边校验,不创建第二条 Scheduled 执行路径。
|
||
|
||
## 10. occurrence、领取和崩溃恢复
|
||
|
||
### 10.1 领取事务
|
||
|
||
`claim_due_scheduled_jobs` 改为 `claim_due_scheduled_runs`。一次领取在同一事务中:
|
||
|
||
1. 选择 `enabled=1 AND next_run_at<=now` 且租约为空/过期的 Job;
|
||
2. 条件更新租约;
|
||
3. 插入唯一的 `job_runs(status=claimed, delivery_status=awaiting_result)` 行,并把其自增 `id` 作为本次 occurrence 的稳定 ID,同时快照 Agent、投递策略和目标;
|
||
4. 对 `Every/Cron` 把 `next_run_at` 推进到领取时刻之后的第一个未来时间;
|
||
5. 对 `At` 在本次 claim 事务内立即设置 `enabled=0`;
|
||
6. 提交后返回 `ClaimedScheduledRun { job_snapshot, run_id, owner }`。
|
||
|
||
领取条件、Run 插入和 Job 推进处于同一写事务;SQLite 的写串行化和条件租约保证一次到期状态只能提交一个 Run。无需额外维护 occurrence key 或 schedule revision。推进下次时间发生在执行之前,因此进程崩溃不会使本次 occurrence 被自动重放。错过的多个历史 tick 不逐个补跑,只执行当前到期 occurrence,并计算下一个未来时间。
|
||
|
||
`Every` 的间隔从领取/开始时刻计算,而不是从完成时刻计算;例如每小时任务执行 55 分钟,下一次约 5 分钟后到期。`Cron` 始终表达绝对日历时间。任务仍受 Job 租约约束,不允许重叠;如果执行时间持续超过调度间隔,Health 应报告“执行时长覆盖调度间隔”,由管理员调整周期或拆分任务。
|
||
|
||
过去时间的 `At` 不能通过创建、更新或重新启用隐式重跑:`cron_add` / `cron_update` 拒绝 `at <= now`,`cron_enable` 拒绝启用已经过期的 At,并要求先把 schedule 更新到未来时间。迁移前已经 enabled 且到期的 At 仍保留为一次合法待执行 occurrence。
|
||
|
||
Job 执行租约固定为 `execution_timeout + shutdown_grace`,执行本身必须在 `execution_timeout` 内终止,因此不增加心跳续租任务。终态提交同时校验 `job_run_id + lock_owner + 非终态 status`;租约过期后迟到的旧执行不能修改新 Run 或 Job 摘要。
|
||
|
||
### 10.2 状态转换
|
||
|
||
```text
|
||
claimed → running → completed
|
||
├→ failed
|
||
├→ timed_out
|
||
├→ cancelled
|
||
└→ interrupted
|
||
|
||
claimed/running --process restart--> unknown
|
||
```
|
||
|
||
所有终态不可重写。Job 的租约 owner 必须匹配才能提交终态,旧执行不得覆盖新领取。
|
||
|
||
### 10.3 启动恢复
|
||
|
||
Gateway 启动、Scheduler admission 尚未开放前,Storage 先执行一次 `recover_scheduled_runs(active_generation, now)` 原子事务:
|
||
|
||
1. 找到所有 `job_runs.status IN ('claimed','running')`;
|
||
2. 把 JobRun 标记为 `status=unknown/outcome=unknown`,写入固定诊断;
|
||
3. 若已关联非终态 AgentRun,把 AgentRun 标记为 `interrupted`,说明其执行生命周期随旧进程结束;
|
||
4. 根据 JobRun 快照的 delivery policy 设置 `pending` 或 `not_requested`;
|
||
5. 清除关联 Job 的旧租约;
|
||
6. 不回退已推进的 `next_run_at`,也不重跑 occurrence;
|
||
7. 在同一事务内提交以上跨表变化。
|
||
|
||
随后再执行通用 `AgentCoordinator::recover_on_activation()`:已由 Scheduled 恢复事务终结的 AgentRun 不会再次处理。`job_runs.status/outcome` 是业务结果和投递决策的唯一权威;`agent_runs.status=interrupted` 只表示编排执行被进程切断,不得反向覆盖 JobRun 的 `unknown`。开放 Scheduler admission 后先 drain pending delivery,再领取新任务。
|
||
|
||
优雅关停由 TaskSupervisor 先停止新领取,再取消/限时等待运行;能够得到明确取消结果时记录 `interrupted`,只有硬崩溃才在下次启动归为 `unknown`。
|
||
|
||
## 11. 投递设计
|
||
|
||
### 11.1 先提交再发送
|
||
|
||
执行完成事务负责:
|
||
|
||
1. 提交 AgentRun 终态;
|
||
2. 提交 JobRun status、outcome、message、diagnostic 和 duration;
|
||
3. 更新 ScheduledJob 的 `last_run_at`、`last_outcome`;
|
||
4. 按矩阵把 delivery status 设置为:
|
||
- `pending`:需要投递;
|
||
- `suppressed`:`on_alert + ok`;
|
||
- `not_requested`:`never`。
|
||
5. 释放 Job 租约。
|
||
|
||
任何渠道 I/O 都发生在事务之后。
|
||
|
||
### 11.2 复用 job_runs 作为轻量 outbox
|
||
|
||
不新增通用消息队列表。JobRun 自带目标快照和投递状态:
|
||
|
||
```text
|
||
awaiting_result → pending / suppressed / not_requested
|
||
pending → delivering → delivered
|
||
├→ pending(瞬态失败、退避后重试)
|
||
└→ failed(永久失败或次数耗尽)
|
||
```
|
||
|
||
Scheduler 在 Run 终态提交后立即尝试该 Run,并在每次 tick 和启动恢复后有界 drain 遗留项:
|
||
|
||
- 原子认领 `pending` 或租约过期的 `delivering` 行;
|
||
- 最多 3 次持久化尝试;
|
||
- 仅对 Channel 明确分类为瞬态的错误重试;
|
||
- 重试间隔使用有上限的指数退避;
|
||
- 复用 OutboundDispatcher 的 `(channel, chat_id)` 顺序锁;
|
||
- 使用稳定 `scheduled_delivery_id = job_run_id` 写入 metadata。
|
||
|
||
具体发送复用现有 `MessageBus::deliver_outbound()` 和 OutboundDispatcher,不允许 Scheduler 绕过 Bus 直接调用 Channel。现有 delivery watch 回执从 `Result<(), String>` 收紧成不含敏感信息的类型化结果,至少区分:
|
||
|
||
- `Delivered`;
|
||
- `TransientFailure { summary }`;
|
||
- `PermanentFailure { summary }`;
|
||
- `TimedOut`;
|
||
- `DispatcherClosed`。
|
||
|
||
OutboundDispatcher 保留现有单次调用内的短暂、内存级瞬态重试;该调用最终返回的回执算一次持久化 delivery attempt。`TransientFailure`、`TimedOut` 和 `DispatcherClosed` 在未达到持久化尝试上限时回到 `pending`,`PermanentFailure` 直接进入 `failed`。Scheduler 不根据错误字符串猜测是否可重试,也不把“成功写入 outbound 队列”当作渠道送达。
|
||
|
||
类型映射固定如下,Channel adapter 必须先把平台错误归一化为正确的 ChannelError:
|
||
|
||
| 来源 | Delivery receipt | 行为 |
|
||
|---|---|---|
|
||
| Channel 成功返回 | `Delivered` | 标记 delivered |
|
||
| `ConnectionError` / `SendError` 经 Dispatcher 内部重试耗尽 | `TransientFailure` | 持久化退避后重试 |
|
||
| Dispatcher 单次发送最终超时 | `TimedOut` | 持久化退避后重试 |
|
||
| Bus/Dispatcher 关停、lane 暂时饱和 | `DispatcherClosed` / `TransientFailure` | 保留 pending,等待当前或下一运行代 |
|
||
| channel 不存在、`ConfigError`、明确的永久平台错误 | `PermanentFailure` | 立即标记 failed |
|
||
|
||
未知 `Other` 默认按永久失败处理,除非产生它的调用点明确证明可以安全重试。回执只持久化经过清洗和长度限制的 summary,不保存平台响应正文、凭据或临时 URL。Dispatcher 内部的秒级重试与 Scheduler 的持久化分钟级退避职责不同,不得互相递归调用或把每次内部尝试计入 `delivery_attempts`。
|
||
|
||
外部发送成功但本地 `delivered` 提交前崩溃时可能重复发送。渠道支持幂等键时传递稳定 ID;不支持时允许“至少一次”并在重试消息 metadata 中标出可能重复。不得为了避免重复而丢失告警。
|
||
|
||
### 11.3 会话历史
|
||
|
||
需要投递时先解析一次目标会话,并用 `UPDATE ... WHERE target_session_id IS NULL` 把结果固定到 JobRun。后续重试不得重新选择“当前最近 dialog”,避免用户切换会话后同一个通知写入不同历史。目标会话 ID 同时用于本地历史和 OutboundMessage `_session_id` metadata。
|
||
|
||
会话解析或幂等历史事务失败时不得绕过历史直接发送:明确的无效目标进入永久 delivery failure;SQLite busy、运行代关停等瞬态错误回到 pending,并计为本次持久化 delivery attempt。
|
||
|
||
通知使用确定性的本地 message ID `scheduled:<job_run_id>`。现有 `messages.id` 已是主键,不增加第二个唯一索引;SessionManager 新增专用 `append_scheduled_notification_if_absent()`,在现有 session persistence lock 下协调内存,Storage 提供对应的原子事务:
|
||
|
||
1. 使用固定 message ID 执行 `INSERT ... ON CONFLICT(id) DO NOTHING`;
|
||
2. 只有实际插入时才更新持久化 Session metadata,并由 SessionManager 把同一消息加入内存、推进 state_version;
|
||
3. 已存在时 SessionManager 直接复用该消息和固定 target_session_id,不重复修改内存或 Session metadata;
|
||
4. 历史写入成功后才把 OutboundMessage 交给 `deliver_outbound()`。
|
||
|
||
SessionManager 可以使用专用 persistence lock 串行化该流程,但不得在 SQLite I/O 期间持有 Session mutex;数据库提交后重新取得 Session mutex,并以固定 message ID 检查内存后再应用一次变更。
|
||
|
||
同一个 JobRun 的重试因此不会产生多条本地历史。消息来源为:
|
||
|
||
```text
|
||
SourceKind::ExternalTrigger
|
||
from_channel = scheduler
|
||
task_id = job_id
|
||
from_run_id = agent_run_id
|
||
```
|
||
|
||
中间工具调用、健康结果和 suppressed 结果不写用户会话。只有实际需要投递的最终通知进入目标会话。
|
||
|
||
## 12. 系统提示词
|
||
|
||
所有 Scheduled Run 使用同一执行契约,不再根据 JobKind 分支:
|
||
|
||
```text
|
||
你正在执行无人值守的定时任务。任务上下文是隔离的,用户不会直接看到普通最终文本。
|
||
完成所有必要检查或操作后,必须且只能通过 complete_scheduled_run 提交最终结果:
|
||
- ok:任务成功,未发现需要关注的问题;
|
||
- alert:任务成功并发现需要用户关注的问题;
|
||
- failed:任务未可靠完成;
|
||
- refused:因权限或安全策略拒绝。
|
||
任何任务 prompt 中关于 NO_REPLY、send_message 或旧输出格式的指令均已失效;不得使用它们。
|
||
```
|
||
|
||
`delivery_policy` 不放进模型提示词。Agent 只看到任务 prompt 和 Outcome 定义,避免为了迎合“静默/通知”而改变事实分类;策略只在 Scheduler 的确定性矩阵中使用。
|
||
|
||
## 13. 功能设计
|
||
|
||
### 13.1 Cron tools
|
||
|
||
`cron_add` 新参数:
|
||
|
||
```json
|
||
{
|
||
"schedule": { "type": "every", "every_ms": 300000 },
|
||
"prompt": "检查生产站点、登录接口和证书",
|
||
"channel": "feishu",
|
||
"chat_id": "oc_xxx",
|
||
"name": "生产站点巡检",
|
||
"agent_id": "web-monitor",
|
||
"delivery_policy": "on_alert"
|
||
}
|
||
```
|
||
|
||
- `agent_id` 可选;缺省表示 Root;
|
||
- `delivery_policy` 缺省 `always`;
|
||
- 删除 `kind`;
|
||
- 删除 `model`;
|
||
- 不接受 `direct`。
|
||
|
||
`cron_add` 和 `cron_update` 共用一个 ScheduledJob validator:channel 必须来自当前 ChannelManager allowlist,命名 `agent_id` 必须存在于当前不可变 AgentCatalog;Storage 只接受已验证的类型,不自行读取运行代配置。重载后 Agent 消失时保留 Job 并在运行时明确失败,见 §8.2。
|
||
|
||
`cron_update` 支持更新 prompt、schedule、channel、chat_id、agent_id 和 delivery_policy。`agent_id: null` 明确切回 Root;字段缺失表示不修改。新建或更新 `At` 时要求 `at > now`;`cron_enable` 对已过期 At 返回失败并要求先更新 schedule。
|
||
|
||
`cron_disable` 只阻止后续 occurrence,不取消已经 claimed/running 的 Run。`cron_remove` 在存在非终态 JobRun 或 Job 租约时返回 conflict,要求先 disable 并等待当前 Run 到达终态;不得依赖 `ON DELETE CASCADE` 静默删除正在执行的 occurrence。
|
||
|
||
`cron_list` 展示:
|
||
|
||
```text
|
||
enabled · agent=web-monitor · delivery=on_alert · next=... · last=alert
|
||
```
|
||
|
||
不再展示 kind、model 或截断的 `last_error`;`last` 只显示 `last_outcome`,详细信息由 `cron_runs` 查询。
|
||
|
||
新增只读 `cron_runs` 工具并注册到普通全局 registry;Root 交互 Agent 可用,命名交互 Agent 仍须由 definition allowlist 授权,Scheduled Agent 的收窄规则一律移除它:
|
||
|
||
```json
|
||
{
|
||
"job_id": "job-id",
|
||
"limit": 20,
|
||
"run_id": 123
|
||
}
|
||
```
|
||
|
||
- `job_id` 必填;
|
||
- 不传 `run_id` 时列出最近记录,`limit` 缺省 20、范围 1–100,每条返回 run_id、时间、status、outcome、delivery status/attempts、duration 及有界 message/diagnostic 摘要;
|
||
- 传 `run_id` 时校验它属于该 Job,并返回该次运行的完整结构化字段;message 上限沿用 16 KiB,diagnostic 经过清洗并执行独立上限;
|
||
- `read_only() == true`,不修改 Job/Run、不触发重投;
|
||
- `never`、`suppressed` 和投递失败记录与其他策略一样可查询;
|
||
- 输出不包含 Provider 私有状态、reasoning、凭据或 Agent transcript。
|
||
|
||
### 13.2 WebUI
|
||
|
||
任务页删除“任务/巡检”徽标,展示:
|
||
|
||
- Agent:Root 或命名 Agent;
|
||
- 投递:始终通知 / 异常通知 / 从不通知;
|
||
- 最近 Outcome;
|
||
- 最近运行生命周期状态;
|
||
- 投递状态、尝试次数及失败原因;
|
||
- 最近运行的 message/diagnostic 摘要,并可展开单次详情。
|
||
|
||
`GET /api/jobs/{id}/runs` 明确切换到 v11 shape:`id/job_id/scheduled_for/started_at/finished_at/status/outcome/message/diagnostic/duration_ms/delivery_status/delivery_attempts/delivery_error`;删除 `output/error/result_kind`。API 继续受管理认证保护并保持 limit 上限。
|
||
|
||
创建表单可以提供三个用户友好模板,但模板不进入后端模型:
|
||
|
||
- 定期通知 → `always`;
|
||
- 异常巡检 → `on_alert`;
|
||
- 后台维护 → `never`。
|
||
|
||
### 13.3 Health
|
||
|
||
HealthService 增加只读 Scheduler 配置检查:
|
||
|
||
- ScheduledJob 引用不存在或 disabled 的 Agent;
|
||
- 非法 channel;
|
||
- 长时间停留在 pending/delivering 的投递;
|
||
- 最近一次 `unknown`、`failed` 或 `timed_out`;
|
||
- `never` Job 最近的 unknown(明确提示其不会自动通知);
|
||
- 最近执行时长持续覆盖 `Every` 周期;
|
||
- enabled Job 无法计算 next run。
|
||
|
||
Health 不执行 Job、不连接 Provider、不修复数据。
|
||
|
||
## 14. SQLite v11 schema
|
||
|
||
### 14.1 scheduled_jobs
|
||
|
||
```sql
|
||
CREATE TABLE scheduled_jobs (
|
||
id TEXT PRIMARY KEY,
|
||
name TEXT NOT NULL,
|
||
schedule TEXT NOT NULL,
|
||
prompt TEXT NOT NULL,
|
||
agent_id TEXT,
|
||
channel TEXT NOT NULL,
|
||
chat_id TEXT NOT NULL,
|
||
delivery_policy TEXT NOT NULL
|
||
CHECK (delivery_policy IN ('always','on_alert','never')),
|
||
enabled INTEGER NOT NULL DEFAULT 1 CHECK (enabled IN (0,1)),
|
||
next_run_at INTEGER NOT NULL,
|
||
last_run_at INTEGER,
|
||
last_outcome TEXT CHECK (
|
||
last_outcome IS NULL OR
|
||
last_outcome IN ('ok','alert','failed','refused','unknown')
|
||
),
|
||
locked_at INTEGER,
|
||
lock_owner TEXT,
|
||
lease_until INTEGER,
|
||
created_at INTEGER NOT NULL,
|
||
updated_at INTEGER NOT NULL,
|
||
CHECK (length(trim(id)) > 0),
|
||
CHECK (length(trim(name)) > 0),
|
||
CHECK (length(trim(prompt)) > 0),
|
||
CHECK (length(trim(channel)) > 0),
|
||
CHECK (length(trim(chat_id)) > 0),
|
||
CHECK (agent_id IS NULL OR length(trim(agent_id)) > 0),
|
||
CHECK (
|
||
(locked_at IS NULL AND lock_owner IS NULL AND lease_until IS NULL) OR
|
||
(locked_at IS NOT NULL AND lock_owner IS NOT NULL AND lease_until IS NOT NULL)
|
||
)
|
||
);
|
||
|
||
CREATE INDEX idx_jobs_claimable
|
||
ON scheduled_jobs(enabled, next_run_at, lease_until);
|
||
```
|
||
|
||
### 14.2 job_runs
|
||
|
||
```sql
|
||
CREATE TABLE job_runs (
|
||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||
job_id TEXT NOT NULL
|
||
REFERENCES scheduled_jobs(id) ON DELETE CASCADE,
|
||
scheduled_for INTEGER NOT NULL,
|
||
agent_run_id TEXT UNIQUE
|
||
REFERENCES agent_runs(id) ON DELETE SET NULL,
|
||
|
||
-- claim-time immutable snapshot
|
||
agent_id TEXT,
|
||
delivery_policy TEXT NOT NULL
|
||
CHECK (delivery_policy IN ('always','on_alert','never')),
|
||
target_channel TEXT NOT NULL,
|
||
target_chat_id TEXT NOT NULL,
|
||
-- resolved once when delivery is first prepared; stable across retries
|
||
target_session_id TEXT,
|
||
|
||
started_at INTEGER,
|
||
finished_at INTEGER,
|
||
status TEXT NOT NULL CHECK (status IN (
|
||
'claimed','running','completed','failed',
|
||
'timed_out','cancelled','interrupted','unknown'
|
||
)),
|
||
outcome TEXT CHECK (
|
||
outcome IS NULL OR
|
||
outcome IN ('ok','alert','failed','refused','unknown')
|
||
),
|
||
message TEXT,
|
||
diagnostic TEXT,
|
||
duration_ms INTEGER,
|
||
|
||
delivery_status TEXT NOT NULL CHECK (delivery_status IN (
|
||
'awaiting_result','not_requested','suppressed','pending',
|
||
'delivering','delivered','failed'
|
||
)),
|
||
delivery_attempts INTEGER NOT NULL DEFAULT 0,
|
||
delivery_next_attempt_at INTEGER,
|
||
delivery_lease_owner TEXT,
|
||
delivery_lease_until INTEGER,
|
||
delivery_error TEXT,
|
||
|
||
created_at INTEGER NOT NULL,
|
||
updated_at INTEGER NOT NULL,
|
||
CHECK (agent_id IS NULL OR length(trim(agent_id)) > 0),
|
||
CHECK (length(trim(target_channel)) > 0),
|
||
CHECK (length(trim(target_chat_id)) > 0),
|
||
CHECK (target_session_id IS NULL OR length(trim(target_session_id)) > 0),
|
||
CHECK (delivery_attempts >= 0),
|
||
CHECK (duration_ms IS NULL OR duration_ms >= 0),
|
||
CHECK (
|
||
(status IN ('claimed','running') AND outcome IS NULL
|
||
AND finished_at IS NULL AND delivery_status = 'awaiting_result') OR
|
||
(status = 'completed' AND outcome IS NOT NULL
|
||
AND outcome IN ('ok','alert','failed','refused')
|
||
AND finished_at IS NOT NULL AND delivery_status != 'awaiting_result') OR
|
||
(status IN ('failed','timed_out','cancelled','interrupted')
|
||
AND outcome IS NOT NULL AND outcome = 'failed'
|
||
AND finished_at IS NOT NULL AND delivery_status != 'awaiting_result') OR
|
||
(status = 'unknown' AND outcome IS NOT NULL AND outcome = 'unknown'
|
||
AND finished_at IS NOT NULL AND delivery_status != 'awaiting_result')
|
||
),
|
||
CHECK (
|
||
(delivery_lease_owner IS NULL AND delivery_lease_until IS NULL) OR
|
||
(delivery_lease_owner IS NOT NULL AND delivery_lease_until IS NOT NULL)
|
||
)
|
||
);
|
||
|
||
CREATE INDEX idx_job_runs_job_finished
|
||
ON job_runs(job_id, finished_at DESC);
|
||
|
||
CREATE INDEX idx_job_runs_recovery
|
||
ON job_runs(status, updated_at);
|
||
|
||
CREATE INDEX idx_job_runs_delivery
|
||
ON job_runs(delivery_status, delivery_next_attempt_at, delivery_lease_until);
|
||
```
|
||
|
||
不为 AgentCatalog 建数据库外键;Agent 定义是配置运行代资源,不是 SQLite 行。`agent_id` 和 definition hash 的最终审计值保存在关联 AgentRun 中。`target_session_id` 也不建 Session 外键:它是投递时固定的路由快照,Session 使用软删除且重试不能因当前 dialog 变化而自动改投;写入历史时仍由 SessionManager 验证目标行和作用域。
|
||
|
||
## 15. 自动迁移设计
|
||
|
||
### 15.1 总体原则
|
||
|
||
- `SCHEMA_VERSION` 从 10 增加到 11;
|
||
- Storage 在 Gateway 启动后台任务之前执行迁移;
|
||
- migration 使用一个专用连接和 `BEGIN IMMEDIATE`,在单个 SQLite 事务中完成;
|
||
- 任一步失败则整体回滚,`user_version` 保持原值,Gateway 拒绝启动;
|
||
- v11 运行时代码不读取旧列、不解析旧枚举、不调用旧函数;
|
||
- 旧版本形状探测只存在于冻结的 `storage/migrations/legacy_to_v10.rs` 和 `storage/migrations/v11_scheduled.rs`;
|
||
- v11 数据库由旧版本二进制打开时,沿用现有“数据库版本过新”拒绝策略;
|
||
- 不支持自动降级。
|
||
|
||
`BEGIN IMMEDIATE` 必须在配置的 SQLite busy timeout 内取得写锁;若另一个 PicoBot 进程仍在使用同一数据库且无法取得锁,启动直接报错,不等待后台任务运行后再迁移。部署流程必须先停旧进程再启动新版本。
|
||
|
||
### 15.2 初始化顺序调整
|
||
|
||
当前 `migrate_schema()` 是一段按实际 table shape 累积修补的迁移,并不存在现成的逐版本 `v10::up()`;同时 `init_scheduler_schema()` 在 migration 前创建旧 Scheduler DDL。实施不重写全部 v1–v10 历史,而采用最小版本化:
|
||
|
||
1. 删除 migration 前的 `init_scheduler_schema()` 调用,`migrate_schema()` 成为 Scheduler 表 DDL 的唯一入口;
|
||
2. 把当前累积逻辑冻结、抽取为 `normalize_legacy_to_v10(conn, detected_shape)`,只供 `current < 10` 的启动迁移调用;
|
||
3. 新增 `migrate_v10_to_v11(conn)`,只接受经过预检的 canonical v10 Scheduler 表;
|
||
4. 全新数据库直接创建 latest base/Agent/Scheduler v11 schema,不先建立再重建 v10 Scheduler 表;
|
||
5. `user_version=0` 且已有历史应用表时先运行 legacy normalizer,再运行 v11;若旧数据库从未有 Scheduler 表,normalizer 只处理其他历史表,随后直接创建 v11 Scheduler 表;
|
||
6. `current=10` 只执行 v11,`current=11` 不执行迁移,`current>11` 拒绝;
|
||
7. Agent schema 必须在创建带 `agent_runs` 外键的 v11 `job_runs` 之前就绪;
|
||
8. 整个所需步骤共享同一个专用连接和 `BEGIN IMMEDIATE` 事务,只有最后写 `user_version=11`;
|
||
9. 测试不再直接调用可绕过 migration 的旧初始化函数,统一通过 `Storage::new`、fresh-schema fixture 或明确的 v10 fixture。
|
||
|
||
这只是启动 migration 的版本兼容,不是运行时双路径。无需把历史逻辑拆成 `v1→v2→...→v10`,也不得在 v11 Scheduler 查询中保留 shape probing。
|
||
|
||
### 15.3 v10 预检
|
||
|
||
在创建新表前读取并验证全部旧行:
|
||
|
||
- schedule JSON 必须能反序列化;
|
||
- delivery policy 必须是 `direct/always/on_alert/never`;
|
||
- Job ID、channel、chat_id 和 prompt 必须满足非空约束;
|
||
- 旧 JobRun 外键必须能找到 Job;
|
||
- Job lease 字段必须全部为空或全部非空;
|
||
- 发现损坏数据时返回包含 Job ID 的 migration error,不做默认修复。
|
||
|
||
### 15.4 Job 字段映射
|
||
|
||
| v10 | v11 | 规则 |
|
||
|---|---|---|
|
||
| `job_kind` | 删除 | 不读取其运行语义 |
|
||
| `delivery_policy=direct` | `always` | 新代码统一中央投递 |
|
||
| `always/on_alert/never` | 原值 | 保留 |
|
||
| `model` | 删除 | 当前执行路径未应用该值;迁移日志报告被删除的非空数量 |
|
||
| `delete_after_run` | 删除 | `At` 任务改为 claim 事务内立即禁用 |
|
||
| 无 `agent_id` | `NULL` | 使用 Root Scheduled Agent |
|
||
| `last_status` | `last_outcome` | 优先取迁移后最新 JobRun.outcome;没有历史 Run 时 `ok→ok`、其他非空值→failed |
|
||
| `last_error` | 删除 | 作为最新旧 JobRun diagnostic 的 fallback 后删除,不保留 Job 级兼容列 |
|
||
|
||
内置 `picobot-routine-maintenance` 的 prompt 在迁移中按固定 Job ID 改成不含 `NO_REPLY` 的任务描述。用户自定义 prompt 不做不安全的字符串替换;新的系统级 Scheduled 契约明确忽略其中关于 `NO_REPLY`、`send_message` 和旧输出格式的指令。
|
||
|
||
### 15.5 历史 JobRun 映射
|
||
|
||
历史行只用于审计,绝不因升级重新投递。status 与 outcome 必须联合映射:
|
||
|
||
| v10 条件 | v11 status | v11 outcome |
|
||
|---|---|---|
|
||
| `status=ok` | `completed` | 按下表 result_kind 映射 |
|
||
| `status=delivery_error` 且已有模型结果/result_kind | `completed` | 按下表 result_kind 映射 |
|
||
| `status=delivery_error` 且执行 error 非空、无结果 | `failed` | `failed` |
|
||
| `status=timeout` | `timed_out` | `failed` |
|
||
| `status=error` | `failed` | `failed` |
|
||
| 其他状态 | `unknown` | `unknown` |
|
||
|
||
仅在新 status 为 `completed` 时读取 result_kind:
|
||
|
||
| v10 result_kind | v11 outcome |
|
||
|---|---|
|
||
| `quiet` | `ok` |
|
||
| `content` 且 Job 为 `on_alert` | `alert` |
|
||
| `content` 其他情况 | `ok` |
|
||
| `reported_failure` | `failed` |
|
||
| `refused` | `refused` |
|
||
| `NULL`(旧 Direct) | `ok` |
|
||
| 其他无法映射值 | migration error,不猜测 |
|
||
|
||
delivery 单独映射,不反向改变执行 status/outcome:
|
||
|
||
| v10 delivery_status | v11 delivery_status |
|
||
|---|---|
|
||
| `delivery_status=direct/delivered` | `delivered` |
|
||
| `suppressed/skipped` | `suppressed` |
|
||
| `failed` 或存在 `delivery_error` | `failed` |
|
||
| 其他 | `not_requested` |
|
||
|
||
其他字段:
|
||
|
||
- `scheduled_for = started_at`;
|
||
- `message = output`;
|
||
- `diagnostic = error`;最新历史 Run 的 error 为空时可以用同 Job 的 `last_error` 补全;
|
||
- `target_*`、`delivery_policy` 从迁移时的 Job 行快照;
|
||
- `target_session_id = NULL`,历史行不重新投递,因此不解析会话;
|
||
- `agent_run_id = NULL`;
|
||
- 所有 delivery lease 字段为空。
|
||
|
||
若 Job 存在非空 `last_error` 但没有任何可承载它的历史 Run,migration 记录 Job ID 和计数后丢弃该冗余摘要,不为兼容旧投影伪造一条执行记录。
|
||
|
||
### 15.6 迁移时发现旧执行锁
|
||
|
||
v10 只有 Job lease,没有预先建立的 JobRun。若迁移时发现 `lock_owner` 非空或 lease 尚未清理,说明旧进程可能在终态提交前退出:
|
||
|
||
1. 为它创建一条 `status=unknown/outcome=unknown` 的 JobRun;
|
||
2. message 固定说明升级时发现未完成执行;
|
||
3. `always/on_alert` 设置 delivery `pending`,`never` 设置 `not_requested`;
|
||
4. recurring Job 的 `next_run_at` 推进到迁移时刻之后;
|
||
5. `At` Job 设为 disabled;
|
||
6. 清除旧租约。
|
||
|
||
这条未知通知可能与崩溃前已经送达但未提交的通知重复,但不会静默掩盖不确定状态。
|
||
|
||
### 15.7 SQLite 表重建
|
||
|
||
SQLite 删除列和修改 CHECK 约束采用 canonical table rebuild:
|
||
|
||
1. 确保 latest `agent_runs` schema 已创建;
|
||
2. 把旧 `job_runs`、`scheduled_jobs` 依次 rename 为内部 legacy 临时表;
|
||
3. 使用 fresh database 相同的 v11 DDL 常量创建 canonical `scheduled_jobs` 和 `job_runs`;
|
||
4. 写入转换后的 Job、历史 Run 和旧锁恢复 Run;
|
||
5. 删除 legacy `job_runs`;
|
||
6. 删除 legacy `scheduled_jobs`;
|
||
7. 重建索引;
|
||
8. 执行 `PRAGMA foreign_key_check` 并要求零行;
|
||
9. 最后设置 `PRAGMA user_version = 11` 并提交。
|
||
|
||
全新数据库走 latest schema creation,不经过临时 v10 表和上述数据搬迁;但使用同一组 v11 DDL 常量,避免 fresh schema 与 migrated schema 漂移。
|
||
|
||
表重建必须在同一连接、同一事务中完成。不得用多个 pool connection 分散 DDL,也不得在事务提交前启动 Scheduler。
|
||
|
||
### 15.8 “不在代码层面兼容过去”的准确含义
|
||
|
||
允许且必须存在:
|
||
|
||
- 一次性 migration 对 v10 表和旧枚举的读取与转换;
|
||
- migration tests 的 v10 fixture;
|
||
- 迁移日志和损坏数据诊断。
|
||
|
||
明确禁止:
|
||
|
||
- `row_to_job` 同时尝试新旧列;
|
||
- `DeliveryPolicy::parse("direct")`;
|
||
- 保留 `JobKind` 但在新路径忽略;
|
||
- 保留 `handle_cron_message` 作为 fallback;
|
||
- 继续解析任意 `NO_REPLY` 文本;
|
||
- 新旧工具参数并存;
|
||
- 根据 `user_version` 在 Scheduler 运行期分支。
|
||
|
||
迁移成功后,进程内只有 v11 类型和 v11 SQL。
|
||
|
||
## 16. 旧代码清理清单
|
||
|
||
### `src/scheduler/mod.rs`
|
||
|
||
删除:
|
||
|
||
- `ScheduledDisposition`;
|
||
- `parse_scheduled_disposition()`;
|
||
- `managed` / `Direct` 双路径;
|
||
- `job_kind == Monitor` 分支;
|
||
- Agent 返回普通文本后再分类的代码;
|
||
- 先发送、后写 JobRun 的顺序。
|
||
|
||
替换为:非阻塞 tick/JoinSet 事件循环 → claim occurrence → `ScheduledAgentRunner` → typed outcome → terminal commit → immediate/bounded delivery drain。删除等待整批 Run 完成后才继续 poll 的 `for_each_concurrent(...).await` 结构。
|
||
|
||
### `src/session/session.rs`
|
||
|
||
删除:
|
||
|
||
- `create_cron_agent()`;
|
||
- `create_managed_scheduled_agent()`;
|
||
- `handle_cron_message()`;
|
||
- `handle_managed_scheduled_message()`;
|
||
- Cron 专用 `NO_REPLY` / `send_message` prompt。
|
||
|
||
Scheduled Agent 构造迁移到 Agent/Coordinator 边界,SessionManager 不再执行 Cron Agent。
|
||
|
||
### `src/storage/scheduler.rs`
|
||
|
||
删除:
|
||
|
||
- `JobKind` 及 parser;
|
||
- `DeliveryPolicy::Direct`;
|
||
- `model`、`job_kind`、`delete_after_run` 映射;
|
||
- `set_scheduled_job_behavior()`;
|
||
- 完成时才插入 JobRun 的旧事务。
|
||
|
||
新增严格 typed parser、occurrence claim、status/outcome 联合约束、跨 JobRun/AgentRun 的原子未知恢复、终态提交、固定 target_session_id 和持久化 delivery claim。
|
||
|
||
### `src/storage/mod.rs` 与 `src/storage/migrations/`
|
||
|
||
把 `SCHEMA_VERSION` 一次提升到 11,删除 migration 之前调用 Scheduler 旧 DDL 初始化函数的顺序。把现有累积迁移冻结为 `legacy_to_v10` normalizer,新增唯一的 v11 Scheduler migration、fresh v11 DDL 和跨表终态/恢复事务 API;迁移完成后通用 Storage 查询不得包含任何 v10 列名。
|
||
|
||
### `src/tools/cron.rs`
|
||
|
||
删除 `kind`、`model` 参数与所有默认联动;默认 `delivery_policy=always`。新增可选 `agent_id` 和共享 Agent/Channel validator,拒绝过去的 At;更新 list 输出,并新增只读 `CronRunsTool`。
|
||
|
||
### `src/tools`
|
||
|
||
新增 `complete_scheduled_run.rs`。它是 runtime-injected、exclusive、Scheduled context-only 的无外部副作用控制工具。
|
||
|
||
### `src/agent`
|
||
|
||
新增可继承的 `ExecutionOrigin::Scheduled`、顶层独占的 `ScheduledCompletionSink` 和 `AgentCoordinator::execute_scheduled()`;复用 foreground AgentRun 持久化,不创建 agent_session_state、signal 或 inbox completion。AgentLoop 在终结工具成功后停止。
|
||
|
||
### `src/session/messenger.rs`
|
||
|
||
第一次投递时固定 target_session_id;把 Scheduled 通知改成利用现有 `messages.id` 主键的稳定 message ID/原子 insert-if-absent 写入;只有实际插入时更新 Session 内存和 metadata。删除任何依赖 `cron:<id>` 作为可消费会话的行为。
|
||
|
||
### `src/bus`
|
||
|
||
保留 `MessageBus::deliver_outbound()` 的等待回执入口,把 `OutboundMessage.delivery` 和 `BusError` 的字符串结果改为类型化、可判定 retry class 的安全回执,并补齐 ChannelError→receipt 映射。OutboundDispatcher 仍是唯一 Channel 调用方;普通无回执消息继续使用 `publish_outbound()`。
|
||
|
||
### Gateway / WebUI / Protocol
|
||
|
||
- Gateway 激活时先执行原子 Scheduled 恢复,再执行通用 Agent recovery,最后开放 Scheduler admission;
|
||
- API JSON 删除 `job_kind`、`model`、`delete_after_run`、`output/error/result_kind`;
|
||
- 增加 `agent_id`、`outcome/message/diagnostic` 和 delivery attempts/状态;JobRun 内部另存固定 target Session 路由快照,不向非管理客户端暴露;
|
||
- `webui/src/pages/TasksPage.svelte` 删除“巡检/任务”判断;
|
||
- `/api/jobs/{id}/runs`、TasksPage、Health 与 `cron_runs` 消费同一 JobRun projection。
|
||
|
||
### 文档与内置知识
|
||
|
||
实施时同步更新:
|
||
|
||
- `README.md`;
|
||
- `docs/ARCHITECTURE.md`;
|
||
- `AGENTS.md`;
|
||
- `resources/skills/about-picobot/references/architecture.md`;
|
||
- `resources/skills/about-picobot/references/config.md`;
|
||
- `resources/skills/about-picobot/references/db-schema.md`;
|
||
- `resources/skills/about-picobot/references/tools.md`;
|
||
- 内置维护任务 prompt。
|
||
|
||
全仓库应不存在运行时 `NO_REPLY`、`JobKind`、`DeliveryPolicy::Direct` 或 `handle_cron_message` 引用。
|
||
|
||
## 17. 并发与生命周期不变量
|
||
|
||
1. 一个 Job 同时最多有一个持租约 occurrence;短周期任务不会重叠执行。
|
||
2. occurrence 的 JobRun 在 Agent 启动前持久化;同一个到期状态只能提交一个 Run。
|
||
3. `next_run_at` 在 claim 事务中推进;执行失败和进程崩溃不会重放同一 occurrence。
|
||
4. `At` 在 claim 事务内禁用;过期 At 不能通过 enable 隐式重跑。
|
||
5. disable 只影响未来 occurrence;存在非终态 Run 时禁止删除 Job。
|
||
6. Scheduler 主循环不等待整批执行结束;长任务只占并发槽,不能阻塞其他 Job 或 pending delivery。
|
||
7. 只有 `job_run_id + lease owner + 非终态 status` 同时匹配才可以提交终态。
|
||
8. 所有终态不可重写;迟到 Agent 结果必须丢弃并记录日志。
|
||
9. status/outcome 必须满足 §5.4 联合约束;`unknown` 只能由运行时生成。
|
||
10. Outcome 和投递策略都使用 claim-time snapshot;运行过程中编辑 Job 只影响后续 occurrence。
|
||
11. Outcome、Job 摘要和初始 delivery status 在同一事务提交。
|
||
12. 渠道 I/O 不发生在 SQLite 事务或 Session mutex 内。
|
||
13. Scheduled origin 传递给所有后代;ScheduledCompletionSink 只属于顶层。
|
||
14. Scheduled Run 及其后代不创建 agent_session_state、不预留 Inbox slot、不发 signal/completion event。
|
||
15. 后台委托在模型可见参数不变,但 Scheduled origin 强制变为 foreground。
|
||
16. JobRun 是 Outcome/投递的唯一权威;关联 AgentRun 只提供编排审计,不能反向改写 JobRun。
|
||
17. 第一次投递解析出的 target_session_id、message ID 和目标 metadata 在所有重试中保持不变。
|
||
18. `never` 严格不通知,包括失败、拒绝、超时和 unknown;这些状态仍可通过 `cron_runs`、管理页面和 Health 查看。
|
||
19. 静默只来源于结构化 `ok` 与策略矩阵,不能来源于普通文本。
|
||
|
||
## 18. 错误处理矩阵
|
||
|
||
| 场景 | Scheduled JobRun status | Outcome | 投递内容来源 |
|
||
|---|---|---|---|
|
||
| Agent 提交 `ok` | completed | ok | Agent message |
|
||
| Agent 提交 `alert` | completed | alert | Agent message |
|
||
| Agent 提交 `failed` | completed | failed | Agent message |
|
||
| Agent 提交 `refused` | completed | refused | Agent message |
|
||
| Agent 未调用终结工具 | failed | failed | Scheduler 固定协议错误 |
|
||
| Provider 失败 | failed | failed | 安全归一化错误,不含响应正文 |
|
||
| 运行超时 | timed_out | failed | Scheduler 固定超时说明 |
|
||
| Gateway 优雅取消 | interrupted | failed | Scheduler 固定中断说明 |
|
||
| 硬崩溃恢复 | unknown | unknown | Scheduler 固定 unknown 说明;关联 AgentRun 审计为 interrupted |
|
||
| Agent 定义不存在 | failed | failed | Agent ID 和修复建议 |
|
||
| 渠道发送永久失败 | 原 run 不变 | 原 outcome | delivery_status=failed |
|
||
|
||
投递失败不改变已经确定的执行 Outcome;它只改变 delivery status。
|
||
|
||
## 19. 实施顺序
|
||
|
||
本设计应在一个功能版本中完成,不能长期保留两套路径:
|
||
|
||
1. 抽取 legacy→v10 normalizer,增加 fresh v11/v10→v11 migration、新 Storage 类型与迁移测试;
|
||
2. 增加 `ExecutionOrigin::Scheduled`、Scheduled completion sink、终结工具和 AgentLoop 终止语义;
|
||
3. 增加 `AgentCoordinator::execute_scheduled()`,落实同步委托、signal/inbox 禁止和顶层 AgentRun;
|
||
4. 改写 Scheduler 为非阻塞 JoinSet 事件循环,实现 occurrence claim、跨表恢复、终态提交和持久化 delivery drain;
|
||
5. 实现固定 target_session_id、幂等历史插入和类型化 Bus 回执;
|
||
6. 切换 Cron tools(含 `cron_runs`)、管理 API、WebUI 与 Health;
|
||
7. 更新内置维护任务和文档;
|
||
8. 删除所有旧类型、旧函数、旧字段读取和魔法字符串;
|
||
9. 运行全量验证后再提交,并按功能变化增加产品中段版本号一次。
|
||
|
||
代码合并点只允许新路径。迁移代码可以先写,但最终提交中不允许 Scheduler 通过 feature flag 或 schema 判断走旧路径。
|
||
|
||
## 20. 测试设计
|
||
|
||
### 单元测试
|
||
|
||
- `DeliveryPolicy × ScheduledOutcome` 全矩阵;
|
||
- 终结工具 schema、空 message、额外字段、重复调用;
|
||
- 终结工具之后的同批工具归约为 Cancelled;
|
||
- 普通文本、所有 `NO_REPLY` 变体都不能产生 `ok`;
|
||
- named/root Agent 解析和工具收窄;
|
||
- Scheduled origin 继承到多层子 Agent,而 completion sink 只在顶层;
|
||
- Scheduled background delegate 在任意嵌套深度自动前台化;
|
||
- Scheduled Agent/子 Agent 不注入 emit_signal、不预留 slot、不创建 agent_session_state;
|
||
- 两个并发 Scheduler 对同一到期 Job 只有一个 claim 成功;
|
||
- At claim 后禁用;
|
||
- 新建/更新过去 At 被拒绝,过期 At 无法直接 enable;
|
||
- disable 不影响活动 Run,remove 在活动 Run/租约存在时返回 conflict;
|
||
- recurring claim 时推进到未来;
|
||
- 长 Run 不阻塞其他到期 Job 的 claim 或已完成 Run 的投递;
|
||
- lease owner 条件提交;
|
||
- 迟到结果不能覆盖 unknown/timeout;
|
||
- status/outcome 非法组合被数据库约束拒绝;
|
||
- Scheduled 恢复在一个事务中完成 JobRun unknown、AgentRun interrupted、delivery 决策和租约清理;
|
||
- delivery claim、瞬态重试、永久失败和尝试上限;
|
||
- ChannelError 到类型化 delivery receipt 的完整映射;
|
||
- `deliver_outbound` 只有收到 Channel 成功回执才返回 Delivered,入队成功不算送达;
|
||
- 回执保留 transient/permanent 分类且不泄露渠道敏感响应;
|
||
- 第一次投递后 target_session_id 固定,切换当前 dialog 不改变重试目标;
|
||
- 本地历史稳定 message ID 去重,重复调用不重复推进 Session metadata;
|
||
- `cron_runs` 摘要/单条详情、limit 边界、Job 所属校验和 read_only 声明。
|
||
|
||
### Migration 测试
|
||
|
||
- 空数据库直接得到 v11;
|
||
- fresh v11 与迁移所得 v11 的 `sqlite_master` schema 等价;
|
||
- `user_version=0` 的历史数据库先规范化再升级,缺少旧 Scheduler 表时直接创建 v11 Scheduler;
|
||
- 真实 v10 fixture 保留所有 Job 和历史 Run;
|
||
- `task/monitor` 列被物理删除;
|
||
- `direct` 全部转为 `always`,运行时无法解析 direct;
|
||
- 非空 model 计数被记录且列被删除;
|
||
- delete_after_run 列被删除;
|
||
- `last_error` 只补入最近历史 Run diagnostic,v11 Job 表不保留该列;
|
||
- 旧 delivery_error 按“执行是否已有结果”分别映射 completed/failed,且 delivery 始终为 failed;
|
||
- 迁移不会生成非法的 status/outcome 组合;
|
||
- 历史 Run 不产生 pending delivery;
|
||
- 旧锁生成 unknown run 并按策略决定 pending/not_requested;
|
||
- 内置维护 prompt 不含 `NO_REPLY`;
|
||
- 损坏 schedule/policy 使迁移原子失败且 `user_version` 不变;
|
||
- `PRAGMA foreign_key_check` 为空;
|
||
- v11 重启迁移幂等;
|
||
- v12+ 数据库被拒绝。
|
||
|
||
### 集成测试
|
||
|
||
- `always + ok` 实际投递;
|
||
- `on_alert + ok` 静默但 JobRun 可查询;
|
||
- `on_alert + alert/failed/refused` 实际投递;
|
||
- `never` 在所有 Outcome 下均不投递;
|
||
- `never` 的完整结果可以通过 `cron_runs` 和管理 API 查询;
|
||
- Provider 普通文本完成被转为协议失败通知;
|
||
- Gateway 在结果提交后、发送前重启,pending 在启动后恢复;
|
||
- Gateway 在发送后、ack 前重启,允许有标识的重复但不丢失;
|
||
- Gateway 在 Agent 执行中退出,恢复为 unknown 且不重跑同一 occurrence;
|
||
- 恢复后关联 AgentRun 为 interrupted,但任何消费者都不能用它覆盖 JobRun unknown;
|
||
- 命名 Agent 被删除后 Job 明确失败且 Gateway 仍能启动;
|
||
- Cron 内 foreground 子 Agent 汇总成功,嵌套 background 参数不会产生 Inbox/agent_session_state;
|
||
- 投递失败后切换当前 dialog,重试仍写入并指向首次固定的 target_session_id;
|
||
- 一个执行到超时上限的 Job 不延迟其他 Job 和 pending delivery。
|
||
|
||
### 必跑验证
|
||
|
||
```bash
|
||
cargo test --lib
|
||
cargo test --test test_scheduler
|
||
cargo test --test test_request_format
|
||
cargo clippy --all-targets --all-features -- -D warnings
|
||
cd webui && npm run check && npm run build
|
||
cargo build
|
||
git diff --check
|
||
```
|
||
|
||
## 21. 验收标准
|
||
|
||
1. 新建和更新任务的 API/工具中没有 `kind`、`monitor` 或 `model` 字段。
|
||
2. 数据库 `scheduled_jobs` 不存在 `job_kind`、`model`、`delete_after_run`。
|
||
3. `DeliveryPolicy` 只有 `always/on_alert/never`。
|
||
4. 所有 Scheduled Run 都必须调用 `complete_scheduled_run`;无调用时 fail-closed。
|
||
5. 全仓库运行时代码不再出现 `NO_REPLY` 结果协议。
|
||
6. Scheduler 只有一条 Agent 执行和一条中央投递路径。
|
||
7. `on_alert + ok` 不投递,其余矩阵行为与本设计一致。
|
||
8. 结果提交后崩溃不会丢失待投递消息。
|
||
9. 执行中崩溃产生 `unknown`,同一 occurrence 不自动重跑。
|
||
10. Scheduled 恢复原子地产生 JobRun unknown 和 AgentRun interrupted,JobRun 始终是投递权威。
|
||
11. Scheduled origin 贯穿所有后代,但 completion sink 只属于顶层;数据库中不产生合成 agent_session_state 或 Inbox event。
|
||
12. Scheduler 的长 Run 不阻塞其他到期 Job 或 pending delivery。
|
||
13. 过期 At 不能直接新建、更新或重新启用。
|
||
14. 每次投递重试复用固定 target_session_id 和 `scheduled:<job_run_id>` message ID,本地历史和 Session metadata 均不重复。
|
||
15. `cron_runs`、管理 API 和 WebUI 都能查询 `never` 等静默任务的结构化结果。
|
||
16. v10 数据库首次启动自动迁移到 v11,失败时完整回滚;fresh/migrated v11 schema 一致,迁移后没有运行时兼容代码。
|