docs(schedule): rewrite the module walkthrough for the finished slice

Bring SCHEDULE_DESIGN_zh.md to the same reference-implementation
walkthrough shape as FEEDBACK_DESIGN_zh.md: it now assumes the spec has
been read, drops the mid-migration 'current state' warnings (the legacy
code is deleted), and covers what the old version predated -- the
commands chapter (UNSET three-state updates, ContextChange, why the
clock- and callback-driven writes stay plain methods), the adapters and
composition-root chapter (CAS field ownership, the corrupt-row error
split, the two anti-corruption layers and the inbound completion
listener), the dispatch journey with the three-driver macro diagram the
spec's §5.3 points at, and the test-layering table for the current
suites.

The spec's §6 marks the schedule slice done and drops the three
completed to-dos; the docs index gains the schedule walkthrough next to
its sibling.
This commit is contained in:
rayhpeng 2026-07-29 22:34:43 +08:00
parent d254e7fc63
commit 563245d96f
3 changed files with 297 additions and 417 deletions

View File

@ -268,7 +268,7 @@ sequenceDiagram
| 模块 | 状态 |
|---|---|
| **Feedback** | ✅ 参考实现(`domain/feedback/` + `app/adapters/feedback/` |
| **Scheduling** | ✅ 参考实现(`domain/schedule/` + `app/adapters/schedule/`);旧代码已停用待删 |
| **Scheduling** | ✅ 参考实现(`domain/schedule/` + `app/adapters/schedule/`);旧代码已删 |
| Run / ThreadMeta / RunEvent / Channel 等 | 旧模式,待迁移——阅读时勿以其为模板,新模块一律照 §2 |
待办,按优先级:
@ -276,9 +276,6 @@ sequenceDiagram
| # | 待办 |
|---|---|
| 1 | 补防腐层对真实上游组件的契约测试(防腐层的门本身没被测过) |
| 2 | 删除 schedule 旧代码(`app/scheduler/service.py``gateway/routers/scheduled_tasks.py`、旧 `sql.py` 仓储部分) |
| 3 | feedback 依赖注入对齐 `Annotated[Service, Depends(...)]` 别名形态 |
| 4 | schedule 写用例按 §3.1 command 化(含部分更新的 `Unset` sentinel |
## 7. 引用出处

View File

@ -9,6 +9,7 @@ This directory contains detailed documentation for the DeerFlow backend.
| [ARCHITECTURE.md](ARCHITECTURE.md) | System architecture overview |
| [HEXAGONAL_ARCHITECTURE_zh.md](HEXAGONAL_ARCHITECTURE_zh.md) | 六边形Ports & Adapters分层规范标准结构AWS 三文件夹 + domain 七件套、Commands/Events 设计、规则清单与执法、调用关系 |
| [FEEDBACK_DESIGN_zh.md](FEEDBACK_DESIGN_zh.md) | 用户反馈模块设计:首个完成的六边形切片,聚合/端口/适配器逐层走读与二次开发指引 |
| [SCHEDULE_DESIGN_zh.md](SCHEDULE_DESIGN_zh.md) | 定时任务模块设计:第二个六边形切片——两聚合两状态机、三驱动源、四道并发防线的完整走读 |
| [API.md](API.md) | Complete API reference |
| [AUTH_DESIGN.md](AUTH_DESIGN.md) | User authentication, CSRF, platform-trust (IM / Internal Auth), and per-user isolation |
| [SSO.md](SSO.md) | OIDC / SSO single sign-on |
@ -41,7 +42,7 @@ This directory contains detailed documentation for the DeerFlow backend.
2. **Configuring the system?** See [CONFIGURATION.md](CONFIGURATION.md)
3. **Understanding the architecture?** Read [ARCHITECTURE.md](ARCHITECTURE.md)
4. **Building integrations?** Check [API.md](API.md) for API reference
5. **Wondering why the layers are split the way they are?** Read [HEXAGONAL_ARCHITECTURE_zh.md](HEXAGONAL_ARCHITECTURE_zh.md) for the rules, then [FEEDBACK_DESIGN_zh.md](FEEDBACK_DESIGN_zh.md) for a worked example
5. **Wondering why the layers are split the way they are?** Read [HEXAGONAL_ARCHITECTURE_zh.md](HEXAGONAL_ARCHITECTURE_zh.md) for the rules, then [FEEDBACK_DESIGN_zh.md](FEEDBACK_DESIGN_zh.md) for the minimal worked example and [SCHEDULE_DESIGN_zh.md](SCHEDULE_DESIGN_zh.md) for the full-complexity one
## Document Organization
@ -51,6 +52,7 @@ docs/
├── ARCHITECTURE.md # System architecture
├── HEXAGONAL_ARCHITECTURE_zh.md # Hexagonal layering rules (zh)
├── FEEDBACK_DESIGN_zh.md # Feedback module design (zh) — first hexagonal slice
├── SCHEDULE_DESIGN_zh.md # Schedule module design (zh) — second hexagonal slice
├── API.md # API reference
├── AUTH_DESIGN.md # User authentication and isolation design
├── CONFIGURATION.md # Configuration guide

View File

@ -1,120 +1,60 @@
# 定时任务Schedule模块设计
> 面向想理解或扩展定时任务模块的人。读完你将能回答:一条定时规则从创建到执行经历了什么、它的每条业务规则住在哪个文件、以及你要改它时该动哪里
> [六边形设计规范](HEXAGONAL_ARCHITECTURE_zh.md) 的**参考实现走读**。本文假定你已读过规范术语端口、适配器、聚合、command、防腐层与通用规则不再重复解释只讲三件事——schedule 特有的产品决定、每条规范规则落在哪个文件、以及 schedule 特有的陷阱。读完你将能回答:一条定时规则从创建到执行经历了什么、它的每条并发防线由谁守卫、想加一种调度类型该改哪几个文件
>
> 配套文档:[`HEXAGONAL_ARCHITECTURE_zh.md`](HEXAGONAL_ARCHITECTURE_zh.md)(本文遵循的架构分层)、`backend/AGENTS.md`(编码规约与调度相关的运行时约定)。
>
> **本文覆盖整个内圈**`domain/schedule/`:模型、端口、应用服务)。适配器与入口仍在旧位置,见 §2。
> 姊妹文档:[FEEDBACK_DESIGN_zh.md](FEEDBACK_DESIGN_zh.md)——首个切片、规范的最小实例。schedule 是同样的结构,复杂度高一个量级:两个聚合、两个状态机、三种驱动源、四道并发防线。
---
## 1. 这个模块做什么
一句话:
让用户注册一条「到点、或按 cron 用这段 prompt 起一次 agent run」的规则由后台轮询器按时把它派发进**既有的** Gateway run 生命周期。
> 让用户注册一条「到点、或按 cron 用这段 prompt 起一次 agent run」的规则由后台轮询器按时把它派发进**既有的** Gateway run 生命周期。
产品层面的四条决定,解释了后面所有设计:
一条最简单的使用路径:
| 决定 | 后果 |
|---|---|
| **调度器只决定 when不建第二套执行栈**`backend/AGENTS.md` 的硬约束) | 派发最终只做一件事:在正确的时刻调用一次现有的 run 启动入口,然后记账。全部复杂度在"正确的时刻"与"记账"的并发与崩溃语义上 |
| **规则与执行是两回事** | 两个聚合、两张表:`ScheduledTask`(规则,长期存在)与 `ScheduledRun`(一次触发的历史记录)。互相只有 `task_id` 字符串引用,各自独立事务 |
| **重叠就跳过**MVP 固定 `overlap_policy="skip"` | 一个任务同时最多一条活跃执行;到点时上一次还没结束,这一次被跳过并留终态墓碑 |
| **定时运行不可交互** | 派发走 `non_interactive=True` 的内部启动路径lead-agent 工具集不含 `ask_clarification` |
对应的 HTTP 面(`app/gateway/routers/schedule/`
```
用户在 /workspace/scheduled-tasks 建一条任务
"每天 09:00Asia/Shanghai帮我总结昨天的 GitHub issue"
后台轮询器每 5 秒扫一次,发现它到点了
起一个 agent run新 thread、非交互模式
run 结束后回写执行记录,并算出下一次的时间
GET /api/scheduled-tasks 列表
POST /api/scheduled-tasks 创建
GET /api/scheduled-tasks/{id} 详情
PATCH /api/scheduled-tasks/{id} 部分更新
POST /api/scheduled-tasks/{id}/pause 暂停
POST /api/scheduled-tasks/{id}/resume 恢复
POST /api/scheduled-tasks/{id}/trigger 立即执行409/502/200 三种码)
DELETE /api/scheduled-tasks/{id} 删除
GET /api/scheduled-tasks/{id}/runs 执行历史
GET /api/threads/{tid}/scheduled-tasks 某 thread 绑定的任务
```
**一条硬约束**(写在 `backend/AGENTS.md` 里):调度器只决定 **when**,不得引入第二套执行栈。它最终只做一件事——在正确的时刻调用一次现有的 run 启动入口,然后记账。模块的全部复杂度都在"正确的时刻"和"记账"的并发与崩溃语义上,而不在执行本身
但 HTTP 只是三种驱动源之一——轮询器与运行完成回调走同样的用例,见 §8
两个核心概念,对应两张表、两个聚合:
| 概念 | 是什么 | 聚合 | 表 |
|---|---|---|---|
| **任务**task | 用户注册的**规则**,长期存在 | `ScheduledTask` | `scheduled_tasks` |
| **执行**run | 规则的**一次触发**,只读历史 | `ScheduledRun` | `scheduled_task_runs` |
---
## 2. 当前状态:生产已切换,旧代码待删
**这一节请先读,否则你会在代码库里迷路。**
生产路径**已经走新架构**。旧代码还留在仓库里,但已经没有任何东西装配它,下一个提交负责删除:
```mermaid
flowchart LR
subgraph IN["✅ 内圈 · harness"]
direction TB
SV["service.py<br/>ScheduleService用例编排"]
PO["ports.py<br/>4 个 Protocol + 2 个 DTO"]
M["model/<br/>ScheduleSpec · ScheduledTask · ScheduledRun"]
SV --> PO
SV --> M
end
subgraph PRIM["✅ 主适配器 · app"]
R["gateway/routers/schedule/<br/>router · models"]
PL["scheduler/poller.py<br/>轮询时钟"]
RC["adapters/schedule/run_completion.py<br/>运行完成回调"]
end
subgraph SEC["✅ 从适配器 · app"]
AD["adapters/schedule/<br/>两个仓储 · run_launcher · thread_lookup"]
end
subgraph CR["✅ 组合根"]
CO["composition.py<br/>build_domain_services()"]
end
subgraph DEAD["🗑️ 已停用 · 待删除"]
S["app/scheduler/service.py"]
OR["gateway/routers/scheduled_tasks.py"]
P["persistence/scheduled_task*/sql.py"]
C["deerflow/scheduler/"]
end
R --> SV
PL --> SV
AD -.->|实现| PO
CO -->|装配| SV
style IN fill:#eef6ff
style DEAD fill:#f5f5f5
```
含义很具体:
- **`domain/schedule/` 是唯一的真相声明处**,运行中的定时任务已经走它。改一条业务规则只改这一处。
- `DEAD` 里的文件仍能编译、仍有测试,但**不再被任何装配路径引用**——`app.py` 启动的是 `SchedulePoller`,注册的是 `routers/schedule/`,完成回调走 `composition.py` 装的那个。留着只是为了让删除单独成为一个可审查的提交。
- 旧的重复规则因此不再需要"两边都改"。如果你在 `DEAD` 里发现和内圈不一致的逻辑,以内圈为准。
- 新增业务规则请**只写在内圈**,然后在旧位置调用它——不要再往 router / 旧 service 里加新的判断。
**还差什么**:四个端口都没有真实实现。适配器落地后,`app/scheduler/service.py``deerflow/scheduler/` 整包、两个 `sql.py` 都会消失router 瘦身成协议转换。
## 3. 内圈全景
## 2. 内圈全景
```
domain/schedule/
├── service.py ScheduleService —— 用例编排input port
├── ports.py 4 个 Protocol + LaunchedRun / RunOutcomeoutput ports
└── model/ 领域模型;是包而非单文件,因为这里有两个聚合加值对象
├── errors.py 9 个领域错误,零依赖
├── enums.py 6 个枚举,零依赖
├── spec.py ScheduleSpec值对象· SchedulePolicy值对象
├── task.py ScheduledTask聚合根· TERMINAL_TASK_STATUSES
└── run.py ScheduledRun聚合· ACTIVE/TERMINAL_RUN_STATUSES
domain/schedule/ 对应规范 §2 的 domain 七件套
├── __init__.py 上下文公开 API聚合 + 命令 + 错误 + 服务(端口刻意不导出)
├── exceptions.py 10 个领域错误,一个基类 ScheduleError → 七件套 exceptions/
├── commands.py 6 个写用例 command + UNSET + ContextChange → 七件套 commands/
├── ports.py 4 个 Protocol + LaunchedRun / RunOutcome → 七件套 ports/
├── service.py ScheduleService写方法即 handler → 七件套 command_handlers/
└── model/ 两个聚合 + 值对象;是包而非单文件 → 七件套 model/
├── enums.py 6 个枚举,零依赖
├── spec.py ScheduleSpec值对象· SchedulePolicy值对象
├── task.py ScheduledTask聚合根· TERMINAL_TASK_STATUSES
└── run.py ScheduledRun聚合· ACTIVE/TERMINAL_RUN_STATUSES
```
三层的纪律各不相同:
`model/` 是**包**而非单文件,因为这里有两个聚合加两个值对象——对比 feedback 的单文件 `model.py`。模块规模决定形态,不必强行对称。`model/__init__` 把拆分对调用方隐藏;**子模块只 import 兄弟模块、永不回头 import 包本身**(防止部分初始化循环,`model/__init__` docstring 有说明)。
| 文件 | 职责 | 纪律 |
|---|---|---|
| `model/` | 业务事实与不变量 | 零依赖;不知道存储和 HTTP 存在;不变量在构造期校验 |
| `ports.py` | 领域声明的接口 | 技术中立——签名里不出现 SQL、表名、HTTP 状态码、运行时类型 |
| `service.py` | 用例编排 | 内圈唯一调用 output port 的地方;`user_id` 显式传参;自身不含业务规则 |
一个关键理解:**service 调用 port 是合法的**——port 是领域自己声明、自己拥有的接口,调用自己的抽象不构成对外圈的依赖。运行时注入的实现来自外圈,但 service 只见 Protocol 类型。
**纪律**:这一层零基础设施依赖——没有 SQL、没有 HTTP、没有配置读取、不看时钟`now` 一律由调用方显式传入。CI 的 `tests/test_harness_domain_purity.py` 会 AST 扫描整个 `domain/` 目录执法。唯一的第三方依赖是 `croniter`,它是确定性纯计算库(无 IO、无全局状态与标准库 `zoneinfo` 同性质——日历计算本身就是定时任务的领域知识。
纯度纪律与 feedback 相同(规范 §4 的 AST 测试执法),多一条豁免说明:唯一的第三方依赖 `croniter` 是确定性纯计算库(无 IO、无全局状态与标准库 `zoneinfo` 同性质——日历计算本身就是定时任务的领域知识。
```mermaid
classDiagram
@ -130,27 +70,25 @@ classDiagram
+status_after_failure(trigger)
+status_after_skip()
+status_after_completion(outcome)
+with_schedule(...) ScheduledTask
+with_context(...) ScheduledTask
+with_schedule(...) / with_context(...)
+paused() / resumed()
}
class ScheduleSpec {
<<值对象 · frozen>>
+ScheduleType schedule_type
+str timezone
+str|None cron
+datetime|None run_at
+next_after(now) datetime|None
+ensure_launchable(now, policy)
}
class SchedulePolicy {
<<值对象 · frozen>>
+int min_once_delay_seconds
+min_once_delay_seconds
+max_concurrent_runs
+lease_seconds
}
class ScheduledRun {
<<聚合 · frozen>>
+RunStatus status
+TriggerKind trigger
+queued(...)$
+skipped_tombstone(...)$
+is_active bool
@ -159,125 +97,38 @@ classDiagram
ScheduledTask *-- ScheduleSpec : 持有
ScheduledTask ..> SchedulePolicy : 方法入参
ScheduleSpec ..> SchedulePolicy : 方法入参
note for ScheduledRun "与 ScheduledTask 只有 task_id 字符串引用\n各自独立的一致性边界"
```
三个关系值得注意:
- **`SchedulePolicy` 是"传入"而非"持有"**。它承载运营可调阈值,由组合根从 `config.scheduler` 构造§7.3)。聚合若持有它,同一个任务对象在不同部署配置下语义就不同了。领域默认值刻意宽松(不施加约束),真正的业务阈值只能由外圈注入——`test_composition.py` 钉住"生产拿到的不是领域默认值"。
- **`ScheduledTask` 持有解析后的 `ScheduleSpec`,不是原始 dict**。dict ↔ 值对象的映射归适配器(规范 §2.1)。
- **两个聚合之间只有 `task_id` 字符串引用,没有对象引用。** 一次派发要写两张表,且现状就不在同一个事务里——它们是各自独立的一致性边界。
- **`SchedulePolicy` 是"传入"而非"持有"。** 它承载运营可调阈值(当前只有 `min_once_delay_seconds`),由组合根从 `config.scheduler` 构造。聚合若持有它,同一个任务对象在不同部署配置下语义就不同了。它的领域默认值是 `0`(不施加约束),真正的业务阈值只能由外圈注入。
- **`ScheduledTask` 持有的是解析后的 `ScheduleSpec`,不是原始 dict。** dict ↔ 值对象的映射属于适配器层。
## 3. 聚合与值对象
---
### 3.1 `ScheduleSpec`:调度在什么时候
## 4. `ScheduleSpec`:调度在什么时候
把三个存储字段(`schedule_type` / `schedule_spec` JSON / `timezone`)解析成一个校验过的值对象。两种类型:`cron`5 段表达式,周期性)与 `once``run_at`,一次性)。
它把三个存储字段(`schedule_type` / `schedule_spec` JSON / `timezone`)解析成一个校验过的值对象
**创建即一致**:全部校验在 `__post_init__`(时区必须是合法 IANA 名cron 恰好 5 段、空白折叠once 必须带 `run_at`naive 的 `run_at` 按**任务自己的时区**本地化)。非法输入在任何 IO 之前抛 `InvalidScheduleError`
### 4.1 两种调度类型
**两个算时间的方法,别用反**——这是本模块最容易写错的一处:
| 类型 | 依据字段 | 语义 |
|---|---|---|
| `cron` | `cron`5 段) | 周期性,永远有下一次 |
| `once` | `run_at` | 一次性,用完即止 |
- `next_after(now)`——算下一次是什么时候。cron 在**任务时区**里求值后转 UTC夏令时切换被正确吸收纽约的 `0 9 * * *` 在 EST 是 14:00 UTC、EDT 是 13:00 UTC用户看到的始终是"每天早上九点"once 只在还没过期时返回。派发后重排用它。
- `ensure_launchable(now, policy)`——同样的计算,外加**只在用户提交时才成立**的约束once 必须在未来、且至少提前 `min_once_delay_seconds`。cron 从不受这个下限约束。创建/更新用它。
### 4.2 创建即一致
用反的后果:拿 `ensure_launchable` 做派发后重排,会让一个正常执行完的 cron 任务被"提前量不足"拒绝。
所有校验与规范化都在 `__post_init__` 里,**不在工厂方法里**。原因很实际frozen dataclass 仍然可以逐字段直接构造,把规则放在 `cron_schedule()` / `once_at()` 里等于留了一条绕过通道
**它不知道自己怎么被存储**`ScheduleSpec` 上没有序列化方法。`from_primitives()` 是入向解析的唯一领域面;出向的「本格式用哪两个键」各归各的适配器——存储侧 `_spec_to_column`SQL 适配器)、线上侧 `_spec_to_wire`api model两者今天恰好相似但主适配器不许 import 从适配器所以刻意不合并§7.1
构造期做四件事:
### 3.2 `ScheduledTask`:规则本体与状态机
1. 时区必须是合法 IANA 名(先做,因为下面本地化 `run_at` 要用它)
2. `cron` 必须存在且恰好 5 段,空白折叠后写回
3. `once` 必须带 `run_at`
4. naive 的 `run_at` 按**任务自己的时区**本地化——`2026-08-01T09:00``Asia/Shanghai` 意为"上海时间早上九点",不是 UTC 九点
字段分组:身份(`task_id` / `user_id`)、内容(`title` / `prompt`)、调度(`schedule`)、执行上下文(`context_mode` / `thread_id` / `assistant_id`)、状态(`status` / `overlap_policy`)、调度游标(`next_run_at`——**认领的唯一依据**:为空或在未来 ⇒ 不会被派发)、回执(`last_*` / `run_count`,仅供展示)。
非法输入在**任何 IO 发生之前**就抛 `InvalidScheduleError`
执行上下文两种:`fresh_thread_per_run`(默认,每次派发新 thread`reuse_thread`(同一 threadagent 能看到历史;构造期强制必须带 `thread_id`)。`resolve_execution_thread()` **不是幂等的**——fresh 模式每次调用生成新 UUID派发流程必须只调一次、结果三处复用。
### 4.3 两个算时间的方法,别用反
**时间戳的真相源**(对照 FEEDBACK_DESIGN §3.4`created_at` 归聚合构造时刻——`create` 工厂从显式收到的 `now` 打戳,适配器原样持久化、绝不自己读时钟(契约用例双实现钉住)。`updated_at` 归**存储写路径**——CAS 方法(`record_launch` / `record_completion`)更新行时手里根本没有聚合,只有适配器能盖"最后写入"戳;聚合上的转换方法因此刻意不动它。两个戳都不是规则输入;规则看的时钟永远是显式 `now=` 参数。
这是本模块最容易写错的一处:
```mermaid
flowchart TD
subgraph CU["创建 / 更新路径(用户提交)"]
A1["ensure_launchable(now, policy)"] --> A2{"ONCE?"}
A2 -->|是| A3{"没有未来的一次?"}
A3 -->|是| A4["InvalidScheduleError<br/>必须是未来时间"]
A3 -->|否| A5{"距今 < min_once_delay?"}
A5 -->|是| A6["InvalidScheduleError<br/>至少提前 N 秒"]
A5 -->|否| A7["返回时间"]
A2 -->|否| A8["CRON 不受提前量约束"]
end
subgraph DP["派发路径(重排下一次)"]
B1["next_after(now)"] --> B2["纯计算,不做任何提交期校验"]
end
style A4 fill:#ffe0e0
style A6 fill:#ffe0e0
```
- `next_after(now)` —— 算下一次是什么时候。cron 在**任务时区**里算,返回 UTConce 只在还没过期时返回。
- `ensure_launchable(now, policy)` —— 同样的计算,外加**只在用户提交时才成立**的约束once 必须在未来,且至少提前 `min_once_delay_seconds`(默认配置 60 秒防止用户建一个立刻就要跑的任务。cron 从不受这个下限约束。
**用反的后果**:拿 `ensure_launchable` 做派发后重排,会让一个正常执行完的 cron 任务被"提前量不足"拒绝。
### 4.4 时区是真的按时区算
cron 表达式在任务声明的时区里求值,然后转成 UTC 存储。这意味着夏令时切换会被正确吸收:
```
America/New_York 的 "0 9 * * *"
2026-03-07EST, UTC-5→ 14:00 UTC
2026-03-09EDT, UTC-4→ 13:00 UTC
```
用户看到的始终是"每天早上九点"。
### 4.5 它不知道自己怎么被存储
`ScheduleSpec` 上**没有**序列化方法。数据库里那个 `schedule_spec` JSON 列(以及 HTTP 请求/响应里的同名字段)与值对象之间的双向映射,属于适配器层:
```
{"cron": "0 9 * * *"} ←──→ ScheduleSpec(CRON, "Asia/Shanghai", cron="0 9 * * *")
JSON 列 / HTTP 字段) (值对象)
ScheduleSpec.from_primitives()(领域,一份)
+ 每个适配器自己那几行「本格式用哪两个键」
```
这不是洁癖,是一条可执行的判据:**一旦领域方法的签名里出现 `Mapping[str, Any]`,就说明领域在处理持久化/传输格式了**。解析这件事天然可以切成两半——结构校验键在不在值是不是字符串属于边界值校验cron 是不是 5 段、时区认不认识、`run_at` 有没有)属于 `__post_init__`。切开之后领域完全不需要看见 dict签名全部强类型。
同样的形状在 Feedback 上下文里也成立:`Feedback` 聚合对 ORM 行一无所知,转换全在 `app/adapters/feedback/feedback_repository.py`
---
## 5. `ScheduledTask`:规则本体
聚合根,持有全部不变量。
### 5.1 字段分组
| 组 | 字段 | 说明 |
|---|---|---|
| 身份 | `task_id` `user_id` | 一切读写按 `user_id` 隔离 |
| 内容 | `title` `prompt` | |
| 调度 | `schedule``ScheduleSpec` | |
| 执行上下文 | `context_mode` `thread_id` `assistant_id` | 见 5.2 |
| 状态 | `status` `overlap_policy` | 见 5.3、5.5 |
| 调度游标 | `next_run_at` | **认领的唯一依据**:为空或在未来 ⇒ 不会被派发 |
| 回执 | `last_run_at` `last_run_id` `last_thread_id` `last_error` `run_count` | 仅供展示 |
### 5.2 执行上下文:每次新会话,还是复用同一个
| `context_mode` | 语义 |
|---|---|
| `fresh_thread_per_run`(默认) | 每次派发新建一个 thread各次执行互不干扰 |
| `reuse_thread` | 所有执行都落在同一个 thread 里agent 能看到历史 |
不变量:`reuse_thread` **必须**带 `thread_id`,构造期强制。
`resolve_execution_thread()` 回答"这次派发用哪个 thread"。注意它**不是幂等的**——`fresh_thread_per_run` 每次调用都生成新 UUID。调用方必须每次派发只调一次把结果存进局部变量供 run 记录、启动调用、返回结果三处复用。
### 5.3 状态机
状态机:
```mermaid
stateDiagram-v2
@ -302,15 +153,13 @@ stateDiagram-v2
三件反直觉的事,理解了它们就理解了这个状态机:
**① `running` 不表示 agent 在执行。** 轮询器认领任务的那一刻就写 `running`,此时 run 还没创建。真正的含义是"某个轮询进程持有这一轮的调度所有权(租约)"。`ensure_mutable()` 因此在这个状态拒绝编辑——正在被派发的任务改不得
**① `running` 不表示 agent 在执行。** 轮询器认领任务的那一刻就写 `running`,此时 run 还没创建。真正的含义是"某个轮询进程持有这一轮的调度所有权(租约)"。`ensure_mutable()` 因此在这个状态拒绝编辑。
**② cron 任务几乎从不停在 `running`。** 派发一完成立刻回 `enabled`。只有 `once` 会停在 `running`完成回调——因为在启动那一刻宣布 `completed` 会在 run 失败或进程崩溃时永久说谎。
**② cron 任务几乎从不停在 `running`。** 派发一完成立刻回 `enabled`。只有 `once` 会停在 `running` 等完成回调——在启动那一刻宣布 `completed` 会在 run 失败或进程崩溃时永久说谎。
**③ 终态可以被重新武装。** 把一个 `completed` / `failed` / `cancelled` 任务的调度改到未来时间,状态被强制拉回 `enabled``TERMINAL_TASK_STATUSES`)。不这么做的话,接口返回 200、`next_run_at` 有值、但**永远不触发**——静默死亡。
**③ 终态可以被重新武装。** 把 `completed` / `failed` / `cancelled` 任务的调度改到未来,状态被强制拉回 `enabled``TERMINAL_TASK_STATUSES`——一个常量服务重新武装与 `protect_terminal` 两条规则,因为"终态"正是"可重新武装"的定义域)。不这么做,接口返回 200、`next_run_at` 有值、但**永远不触发**——静默死亡。
### 5.4 四条状态推导规则
派发流程的每个出口都要回答"任务接下来是什么状态"。这四个方法就是答案,**判定顺序即语义**(现状是 `if/elif/else`,写成并列的 `if` 会静默改变行为):
**四条状态推导规则**——派发流程的每个出口都要回答"任务接下来是什么状态"**判定顺序即语义**`if/elif/else` 写成并列 `if` 会静默改变行为):
```mermaid
flowchart TD
@ -341,62 +190,62 @@ flowchart TD
style F1 fill:#ffd9d9
```
两个**判定顺序陷阱**(红色节点),都有专门的测试锁住:
两个**判定顺序陷阱**(红色节点),都有专门的测试锁住:`status_after_launch` 先判 ONCE暂停的 once 任务被手动触发 → `RUNNING` 而非 `PAUSED``status_after_failure` 先判 MANUALonce 任务手动触发失败 → 保持原状态——失败的手动触发不能吃掉任务本来的调度未来)。另外三处意图:`status_after_skip` 没有 trigger 参数跳过只发生在自动路径上手动重叠是直接拒绝、不留记录once 被跳过是 `FAILED` 而非 `COMPLETED`(唯一的机会丢了,说"完成"等于谎称执行过);`INTERRUPTED` 映射 `CANCELLED` 而非 `FAILED`(用户主动取消不是执行失败)。
- **`status_after_launch` 先判 ONCE**:一个 `paused` 的 once 任务被手动触发,结果是 `RUNNING` 而不是 `PAUSED`
- **`status_after_failure` 先判 MANUAL**:一个 once 任务被手动触发且失败,保持原状态而不是 `FAILED`——失败的手动触发不能吃掉这个任务本来的调度未来。
### 3.3 `ScheduledRun`:一次执行的记账
另外三处设计意图:
- `status_after_skip` **没有 trigger 参数**。跳过只发生在自动调度路径上——手动触发遇到重叠是直接拒绝、不留记录的,所以那个分支根本不存在,加参数会暗示一个虚构的可能性。
- once 被跳过是 `FAILED` 而非 `COMPLETED`:唯一的那次机会丢了,说"完成"等于谎称执行过。
- 完成回调里 `INTERRUPTED` 映射到 `CANCELLED` 而非 `FAILED`:用户主动取消、或同 thread 被新 run 抢占,都不是执行失败。
### 5.5 重叠策略
`overlap_policy` 目前固定为 `"skip"`:一个任务同时最多有一个活跃执行,到点时若上一次还没结束,这一次就被跳过。
聚合上用 `skips_on_overlap` 属性封装这个判断,字符串比较只存在于这一处——将来加第二种策略(比如 `queue`)只改这里。之所以不做成枚举:目前只有一个取值,单值枚举是噪音。
---
## 6. `ScheduledRun`:一次执行的记账
字段几乎是 `scheduled_task_runs` 表的镜像,业务逻辑很少。它存在的核心理由是**两个具名工厂**
字段几乎是 `scheduled_task_runs` 表的镜像。它存在的核心理由是**两个具名工厂**
| 工厂 | 产出状态 | 用在哪 |
|---|---|---|
| `ScheduledRun.queued(...)` | `QUEUED`(活跃) | 正常派发 |
| `ScheduledRun.skipped_tombstone(...)` | `SKIPPED`**直接终态** | 因重叠被跳过 |
**为什么墓碑必须是第二个工厂,而不是"先 queued 再改状态"**——这是全模块最容易写错的一处:
**为什么墓碑必须是第二个工厂,而不是"先 queued 再改状态"**:数据库上的部分唯一索引 `uq_scheduled_task_run_active`(谓词 `status IN ('queued','running')`)保证一个任务最多一条活跃执行。跳过发生时,上一次执行**还占着**那个槽位——墓碑若先建成 `queued`,它自己就会撞上索引。`skipped` 落在谓词之外,永不冲突。两个具名工厂是用类型系统把这条规则焊死,让错误实现写不出来。
数据库上有一个部分唯一索引 `uq_scheduled_task_run_active`,谓词是 `status IN ('queued','running')`,保证一个任务最多一条活跃执行。跳过发生时,上一次执行**还占着**那个槽位。如果墓碑先建成 `queued`,它自己就会撞上这个索引。`skipped` 落在谓词之外,永不冲突
`ACTIVE_RUN_STATUSES` 常量必须与索引谓词逐字一致——它是"跳过"判断的快路径依据,索引是并发下的最终仲裁者,两者漂移会让它们对不上(一致性由独立测试钉住,域测试保持零依赖)。
用两个具名工厂而不是一个可变状态的构造器,就是用类型系统把这条规则焊死,让错误实现写不出来。
## 4. Commands写用例的具名载体
`ACTIVE_RUN_STATUSES` 这个常量必须与上述索引谓词保持逐字一致——它是"跳过"判断的快路径依据,而索引是并发下的最终仲裁者,两者漂移会让它们对不上。
`commands.py` 是规范 §3.1 的落地,六个 HTTP 驱动的写用例:
执行状态共六个:`queued → running → success | failed | interrupted`,外加旁路的 `skipped`
| Command | Handler | 请求模型 |
|---|---|---|
| `CreateScheduledTask` | `ScheduleService.create_scheduled_task` | `CreateScheduledTaskRequest` |
| `UpdateScheduledTask` | `ScheduleService.update_scheduled_task` | `UpdateScheduledTaskRequest` |
| `PauseTask` | `ScheduleService.pause_task` | 无 body端点内直接构造 |
| `ResumeTask` | `ScheduleService.resume_task` | 同上 |
| `DeleteTask` | `ScheduleService.delete_task` | 同上 |
| `TriggerTask` | `ScheduleService.trigger_task` | 同上 |
---
命名链三拼写互为变体(规范 §2.1grep 任何一个就能找到用例全部三层。command 是哑数据(`context_mode="not-a-mode"` 能构造出来——校验归聚合,错误归因顺序归 handler 的构造顺序);身份不进请求模型(`user_id` 服务端解析后经 `to_command` 参数注入,测试钉住两个请求模型上都没有该字段)。
## 7. 端口:领域对外的四个依赖
schedule 在 feedback 之上多出三个设计点:
`ports.py` 声明领域需要外界做什么,由外圈实现。签名一律技术中立——出现 SQL、表名、HTTP 状态码就是越界了。
**① `UNSET` sentinel 表达三态**(规范 §3.1 明文)。`UpdateScheduledTask` 的四个可选字段默认 `UNSET`"未提供"),与 `None` 严格区分。wire 保持历史契约(显式 `null` = 未提供),`to_command``None``UNSET` 翻译的唯一发生处——三态在 command 层无歧义,两态在 wire 层不破坏兼容
| 端口 | 回答什么问题 |
|---|---|
| `ScheduledTaskRepository` | 规则存在哪、怎么按用户隔离、怎么原子地认领到期任务 |
| `ScheduledRunRepository` | 执行记录存在哪、谁来仲裁"一个任务只能有一条活跃执行" |
| `RunLauncher` | 怎么真正启动一次 agent run |
| `ThreadLookup` | 这个 thread 存在吗、这个用户能用吗 |
**② `ContextChange` 是字段联动的对象化先例。** `context_mode``thread_id` 总是一起变(`with_context` 同时接收两者,切到 fresh 模式就意味着清空绑定)。打包之后,"解绑 thread"与"没传 thread"再无歧义——`thread_id``None` 有含义,但它只在 `ContextChange` 内部出现。
外加两个 DTO`LaunchedRun`(启动成功后拿到的 run 身份)、`RunOutcome`(一次执行到达终态的领域表述)
**③ `now` 不是 command 字段。** 它是服务端时钟读数、规则输入不是客户端意图的一部分——handler 签名保留显式 `now=` 参数(`create_scheduled_task(cmd, *, now)`),与规范 §4 的时钟纪律一脉相承。
### 7.1 三条约定,比签名更重要
**时钟与回调驱动的写用例刻意不 command 化**`run_once` / `dispatch_task` / `handle_run_completion` / `reconcile_on_startup`。这些 driver 没有 wire 形状可翻译(规范 §5.3)——轮询器递给 service 的是已认领的聚合加时钟读数,完成回调递的是已经领域化的 `RunOutcome`。command 化只是把领域词汇再包一层领域词汇。
**① 越权一律表现为"不存在"。** 别人的任务在读取时返回 `None` / `False` / 从列表里消失,而不是抛权限错误——调用方不能借此判断"这个 id 到底存不存在"。
`UpdateScheduledTask` 的组装有一处独特编排schedule 与 context 是复合字段,客户端可以只提交一半(只改时区、只改 thread。router 先读一次当前任务、把它传给 `to_command(task_id, user_id, current)`,由请求模型补出**完整值对象**——api model 保持零 IOservice 收到的永远是整个 `ScheduleSpec` / `ContextChange` 而非松散字段的补丁。
## 5. 端口:领域对外的四个依赖
| 端口 | 回答什么问题 | 实现形态 |
|---|---|---|
| `ScheduledTaskRepository` | 规则存在哪、怎么按用户隔离、怎么原子地认领到期任务 | 自有持久化 |
| `ScheduledRunRepository` | 执行记录存在哪、谁仲裁"一个任务只能有一条活跃执行" | 自有持久化 |
| `RunLauncher` | 怎么真正启动一次 agent run | 防腐层 |
| `ThreadLookup` | 这个 thread 存在吗、这个用户能用吗 | 防腐层 |
外加两个 DTO`LaunchedRun`(启动成功拿到的 run 身份)、`RunOutcome`(一次执行到达终态的领域表述)。
### 5.1 三条约定,比签名更重要
**① 越权一律表现为"不存在"**(规范 §4别人的任务读取返回 `None` / `False` / 从列表消失,不抛权限错误。
**② `RunLauncher` 只允许两种异常逃逸。**
@ -405,40 +254,34 @@ flowchart TD
其他任何失败 -> LaunchFailedError
```
这条约定是**整个重构的支点**。旧代码的编排层`isinstance(exc, HTTPException) and exc.status_code == 409` 嗅探"线程忙"于是业务逻辑依赖了 Web 框架。翻译交给适配器后,领域只认自己的两个错误。两者必须区分,因为结果不同:自动调度遇到线程忙是一次跳过的机会,真正的失败要记为失败。
这条约定是**整个重构的支点**。旧编排层靠 `isinstance(exc, HTTPException) and exc.status_code == 409` 嗅探"线程忙",业务逻辑因此依赖了 Web 框架。翻译交给适配器后,领域只认自己的两个错误。必须区分,因为结果不同:自动调度遇到线程忙是一次跳过的机会,真正的失败要记为失败。
**③ `RunOutcome` 挡住运行时类型。** 旧完成回调直接吃 `deerflow.runtime.RunRecord`——一个基础设施类型,纯度测试会拦下它。现在由外圈先转换成 `RunOutcome`,顺带承担了旧代码内联做的过滤:一次不带定时任务元数据、或还没到终态的 run压根产生不出 `RunOutcome`service 也就不会被调用
**③ `RunOutcome` 挡住运行时类型。** 旧完成回调直接吃 `deerflow.runtime.RunRecord`——纯度测试会拦下它。现在由入站适配器先过滤再翻译§7.2service 因此不需要守卫子句
### 7.2 两个刻意的缺席
### 5.2 两个刻意的缺席
**没有 `Clock` 端口。** `now` 一律由调用方显式传参`run_once(now=...)``dispatch_task(now=...)`,领域从不读时钟,测试天然确定。再加一个 `Clock` 只会制造"到底该用参数还是 `self._clock.now()`"的第二个真相源。
**没有 `Clock` 端口。** `now` 一律显式传参,领域从不读时钟,测试天然确定。再加 `Clock` 只会制造"该用参数还是 `self._clock.now()`"的第二个真相源。
**`claim_due` 不接收"谁在认领"。** 那是**进程身份**,不是规则。适配器可以记一个(诊断用),但没有任何代码读回来——决定认领能否被接管的只有过期时间。相比之下 `lease_seconds`(多久算过期)留在了 `SchedulePolicy` 里,因为它直接决定崩溃后多快能恢复,是领域关心的策略。
**`claim_due` 不接收"谁在认领"。** 那是**进程身份**,不是规则——适配器自己生成 `lease_owner`(诊断用),没有任何代码读回来;决定认领能否被接管的只有过期时间。相比之下 `lease_seconds` 留在 `SchedulePolicy` 里,因为它决定崩溃后多快恢复,是领域关心的策略。
### 7.3 原子性归实现,不归契约
### 5.3 原子性归实现,不归契约
`claim_due` docstring 明说:单线程语义(选哪些行、写什么状态)是契约的一部分,**原子性不是**。一个内存实现可以满足全部单线程规则却毫无并发保证。
`claim_due` 的单线程语义(选哪些行、写什么状态)是契约的一部分,**原子性不是**——内存实现可以满足全部单线程规则却毫无并发保证。**契约测试全绿不代表可以跑多个调度器实例**;并发由真数据库上的 `test_schedule_dispatch_race.py` 单独负责§9
这条边界很重要——`tests/schedule_fakes.py` 的模块 docstring 和 `test_schedule_fakes.py` 都重复了它:**契约测试全绿不代表可以跑多个调度器实例**。并发由真数据库上的 `test_scheduled_task_dispatch_race.py` 单独负责。
## 6. 应用服务:用例编排
---
`ScheduleService` 是这个上下文的 input port。三种主适配器HTTP router、轮询器、完成回调调它把返回值翻译成各自的协议。它自身不含业务规则——每个判断都委托给聚合。
## 8. 应用服务:用例编排
| 分类 | 方法 | 输入 |
|---|---|---|
| 读 | `list_tasks` · `list_tasks_by_thread` · `get_task` · `list_task_runs` | 普通参数(查询不 command 化) |
| 写 | `create_scheduled_task` · `update_scheduled_task` · `pause_task` · `resume_task` · `delete_task` | command§4 |
| 派发 | `trigger_task`command + `now=` · `run_once` · `dispatch_task` | 后两者是时钟驱动,领域词汇直入 |
| 生命周期 | `handle_run_completion` · `reconcile_on_startup` | 回调/启动驱动 |
`ScheduleService` 是这个上下文的 input port。主适配器HTTP router、轮询器、完成回调调它然后把返回值翻译成自己的协议。它**自身不含业务规则**——每个判断都委托给聚合。
### 6.1 `dispatch_task`:四条出口
### 8.1 用例清单
| 分类 | 方法 |
|---|---|
| 读 | `list_tasks` · `list_tasks_by_thread` · `get_task` · `list_task_runs` |
| 写 | `create_task` · `update_task` · `pause_task` · `resume_task` · `delete_task` |
| 派发 | `trigger_task`(手动) · `run_once`(轮询一轮) · `dispatch_task`(单次派发) |
| 生命周期 | `handle_run_completion` · `reconcile_on_startup` |
### 8.2 `dispatch_task`:四条出口
整个模块的风险中心。它把一个到期任务变成一次执行,有且只有四种结局:
整个模块的风险中心。把一个到期任务变成一次执行,有且只有四种结局:
```mermaid
flowchart TD
@ -447,7 +290,7 @@ flowchart TD
C -->|是, 手动| D["CONFLICT<br/>不留任何记录"]
C -->|是, 自动| E["SKIPPED<br/>写终态墓碑"]
C -->|否| F["创建 queued 记录"]
F -->|"活跃槽位被抢"| G{"trigger?"}
F -->|"活跃槽位被抢<br/>ActiveRunConflictError"| G{"trigger?"}
G -->|手动| D
G -->|自动| E
F -->|插入成功| H["launcher.launch(...)"]
@ -462,11 +305,13 @@ flowchart TD
style K fill:#e2f7e2
```
**快路径与槽位仲裁必须产生相同结果。** `has_active` 是非原子快路径——两个并发派发可以都通过它。真正的仲裁者是仓储:第二条活跃记录会被拒绝,抛 `ActiveRunConflictError`。调用方**不能分辨自己被哪一种机制拦下**,否则重试行为就会分叉。`test_schedule_service.py` 里那条断言逐字段比较两条路径的 `DispatchResult`,就是钉死这一点。
**快路径与槽位仲裁必须产生相同结果。** `has_active` 是非原子快路径——两个并发派发可以都通过它。真正的仲裁者是仓储:第二条活跃记录被拒绝,抛 `ActiveRunConflictError`,收敛到与快路径**字节一致**的结局(手动 → CONFLICT 且不留 run 行;自动 → SKIPPED 墓碑)。调用方不能分辨自己被哪种机制拦下,否则重试行为分叉——`test_schedule_service.py` 逐字段比较两条路径的 `DispatchResult` 钉死这一点。
**只有 launch 被 try 包住。** 端口契约保证它只逃逸两种异常,所以后续记账写入的失败是真故障,会如实冒泡。旧代码把 launch 和记账一起包在 `except Exception` 里,结果是一次**已经成功启动**的执行会因记账失败被标记成 failed。
**只有 launch 被 try 包住。** 端口契约保证它只逃逸两种异常,所以后续记账写入的失败是真故障如实冒泡。旧代码把 launch 和记账一起包在 `except Exception` 里,结果是一次**已经成功启动**的执行会因记账失败被标记成 failed。
### 8.3 `run_once`:预算是全局的
**成功分支的两笔回写都带 `protect_terminal=True`**——一个快速失败的 run 能在这两笔写落库之前就到达完成回调CAS 保证回调写下的终态不被启动路径的迟到写覆盖§7.1)。
### 6.2 `run_once`:预算是全局的
```python
active = await runs.count_active() # 跨所有任务
@ -475,189 +320,225 @@ if budget <= 0: return []
claimed = await tasks.claim_due(now=..., lease_seconds=..., limit=budget)
```
`max_concurrent_runs` 限制的是**同时活跃的执行总数**,不是每轮的批量大小。长时间运行的任务会跨轮次累积,所以每一轮只能认领进剩余的额度。把它当成"每轮最多认领 N 个"会让长任务把系统压垮
`max_concurrent_runs` 限制**同时活跃的执行总数**,不是每轮批量。长任务跨轮次累积,每轮只能认领进剩余额度——当成"每轮最多 N 个"会让长任务压垮系统
### 8.4 更新任务:为什么 `context` 是打包的
### 6.3 完成回调与启动清扫
```python
await service.update_task(
task_id, user_id=..., now=...,
title=None, # None = 不改
prompt=None,
schedule=None,
context=ContextChange(ContextMode.REUSE_THREAD, "thread-1"),
)
`handle_run_completion(outcome, now)` 先写执行记录终态再问聚合这个结局意味着什么、只写那个答案once 推向终态cron 保持不变。**但 `last_error` 无条件写入**——cron 保住了调度,仍要报告上次出了什么问题。任务在 run 飞行途中被删掉不是错误,静默返回。
`reconcile_on_startup(error)` 清扫崩溃留下的两种残骸:永远不会结束的活跃执行记录、和停在 `running` 等一个已死回调的 once 任务(后者不被过期认领覆盖——已启动的任务释放了租约,认领查询永远看不见它)。它**不吞异常**——部分清扫失败要不要阻塞启动,是调用方的策略。
## 7. 适配器与组合根
四个从适配器 + 一个入站适配器住在 `app/adapters/schedule/`,一个端口一个文件;两形态判据见规范 §2。
### 7.1 两个自有持久化仓储
`SqlScheduledTaskRepository` / `SqlScheduledRunRepository`——表归本上下文所有、自己写 SQL。除 feedback 同款的 `_to_domain` / `_apply` / `_tz_aware` 外,有四处 schedule 特有:
**① 认领是原子的**`claim_due``FOR UPDATE SKIP LOCKED` 挑出到期行、盖租约、写 `running`——迁移自旧仓储的语句原文,它是模块的并发契约而非风格选择。两个认领分支(`enabled` 且到期、`running` 但租约过期)加上 `cancel_stuck_once_tasks` 只清租约为 NULL 的行,共同构成崩溃恢复语义。
**② CAS 分字段所有权,刻意不走整聚合映射**(规范 §2.1 ③w 的明文例外)。派发路径与完成回调并发写同一行 `scheduled_tasks``record_launch` 拥有调度侧字段(`next_run_at` / `last_run_*` / `run_count` / 租约),`record_completion` 只拥有裁决(终态 status + `last_error`)。任何一方表达成"读-改-写整个聚合"都会重放过期快照、回滚对方的写入——后果是 cron 任务带着已流逝的 `next_run_at` 卡在 `running`,且没有任何恢复路径能到达它。`protect_terminal` 旗标关闭镜像方向的竞态(回调先写终态、启动路径的写后到)。**不要把这两个方法"简化"成 `save(task)`**。
**③ `IntegrityError → ActiveRunConflictError` 的翻译点**在 `add()` 的 commit 上——只有活跃状态的插入可能撞 `uq_scheduled_task_run_active`;终态行(墓碑)在谓词之外,那里的 IntegrityError 是真故障、原样冒泡。
**④ 坏行读取抛 `CorruptStoredScheduleError`**,不是 `InvalidScheduleError`。后者是"客户端提交的 schedule 非法"router 映射 422存储损坏是服务端故障专用错误刻意不进 router 的映射表、落"未分类 → 500"分支(否则 PATCH——修复坏行的唯一 HTTP 路径——自己也会 422坏行变得不可修复。单行读冒泡列表读跳过并记日志。`tests/test_schedule_corrupt_rows.py` 在真 sqlite 上钉住。
**索引双定义**`uq_scheduled_task_run_active` 同时定义在 ORM `__table_args__` 和迁移 0007——空库 bootstrap 走 `create_all` 不执行迁移,改索引必须两处同步。
### 7.2 两个防腐层与一个入站适配器
`GatewayRunLauncher`(防腐层):把 Gateway 启动路径的两种"线程忙"信号run manager 的 `ConflictError`、路由层的 `HTTPException(409)`)都翻译成 `ThreadBusyError`,其余翻译成 `LaunchFailedError`。启动 callable 由组合根注入(生产的那个绑定着 FastAPI app。TODO 写的是触发条件run 上下文发布 DTO 契约时替换类体,端口不动。
`ThreadStoreThreadLookup`(防腐层):把宽接口 `ThreadMetaStore` 收窄成一个问题。`require_existing=True` 是承重参数——store 的默认把"行不存在"当可访问(对还没写过的 thread 合理),对"把任务绑定到它"是错的。存在与有权合并成单个 bool防止调用方探测自己看不见的 thread 是否存在。
`ScheduleRunCompletionListener`**入站适配器**——包里唯一驱动领域而非被领域调用的文件,方向由 docstring 首行声明):进程里每个 run 结束都会到达完成回调,大多数与本上下文无关。它做旧钩子内联做的过滤——不带定时任务元数据、或没到终态的 run 压根产生不出 `RunOutcome`service 因此不需要守卫子句(规范 §3.2 轻量形态的活样本:一个业务事实一个消费者,点对点回调 + 入站适配器,第二个订阅方出现时才升格为事件)。运行时四种终态映射到领域三种:`timeout``error` 对定时任务是同一件事FAILED`interrupted` 刻意不是CANCELLED
### 7.3 组合根
`app/composition.py` 三个纯函数:
- `build_domain_services(session_factory, run_store, thread_store, launch_run, scheduler_config)` —— feedback 与 schedule 两个服务的唯一装配点。memory 后端(`session_factory is None`)两者皆 `None`,路由 503——定时任务重启即蒸发比被拒绝接受更糟。
- `build_run_completion_hook(schedule_service)` —— 入站半边的装配service 为 `None` 时返回 `None`,运行时干脆不装钩子,而不是装一个永远拒绝的。
- `build_schedule_policy(scheduler_config)` —— 配置到领域的完整翻译面:领域声明需要哪些阈值,这里点名每个来自哪个配置键。`poll_interval_seconds` 刻意缺席——多久看一眼是轮询器的事,不是任何任务受制于的规则。
`deps.py::langgraph_runtime` 启动时调用;`ScheduleServiceDep = Annotated[ScheduleService, Depends(get_schedule_service)]` 是路由拿服务的别名形态feedback 的 `FeedbackServiceDep` 同款)。
## 8. 一次派发的旅程
宏观上,三种驱动源共存于同一个六边形(规范 §5.3 所指的图):
![三驱动源与派发闭环](assets/hexagonal_dispatch_relation.png)
HTTP 入口(用户管理规则、手动触发)、时钟入口(轮询器按 `poll_interval_seconds` 醒来、回调入口run 运行时报告终态)都走同一个 `ScheduleService`——第三条是**闭环**:派发启动的 run 结束后经完成回调流回本上下文,写下裁决。
一次自动派发的完整时序:
```mermaid
sequenceDiagram
participant PL as SchedulePoller<br/>(时钟入口)
participant S as ScheduleService
participant TR as SqlScheduledTaskRepository
participant RR as SqlScheduledRunRepository
participant L as GatewayRunLauncher<br/>(防腐层)
participant GW as Gateway run 生命周期
participant RC as ScheduleRunCompletionListener<br/>(入站适配器)
PL->>S: run_once(now)
S->>RR: count_active() — 全局预算
S->>TR: claim_due(now, lease, budget)
Note over TR: FOR UPDATE SKIP LOCKED<br/>盖租约、写 running
S->>S: dispatch_task(task, now, SCHEDULED)
S->>RR: has_active? (快路径)
S->>RR: add(queued 记录)
Note over RR: uq_scheduled_task_run_active<br/>是原子仲裁者
S->>L: launch(thread, prompt, metadata)
L->>GW: launch_scheduled_thread_run(...)
GW-->>L: run_id / thread_id
S->>RR: update_status(RUNNING, protect_terminal)
S->>TR: record_launch(..., protect_terminal)
Note over S,TR: 两笔回写分字段所有权§7.1 ②)
GW--)RC: run 到达终态(每个 run 都会)
RC->>RC: _to_outcome — 过滤 + 翻译
RC->>S: handle_run_completion(RunOutcome, now)
S->>RR: update_status(终态)
S->>TR: record_completion(status_after_completion, error)
```
`None` 表示"没传",不需要哨兵值。这能成立的唯一原因是:所有可选字段里,只有 `thread_id``None` 本身有含义(解绑),而它和 `context_mode` 本来就一起变化——`with_context` 同时接收两者,切换到 fresh 模式就意味着清空绑定。打包成 `ContextChange` 之后歧义消失,其余字段就能用最朴素的 `None`
四道并发防线在这条时序上各守一段:
这也正好对上 HTTP 层router 的 PATCH 本来就是 `exclude_none` 语义。
| # | 防线 | 守什么 |
|---|---|---|
| 1 | 部分唯一索引 `uq_scheduled_task_run_active` | 两个并发派发(双击、重试、手动撞轮询)最多一个拿到活跃槽位 |
| 2 | `claim_due``FOR UPDATE SKIP LOCKED` + 租约 | 认领本身原子;崩溃的认领者过期后可被接管 |
| 3 | 双向 `protect_terminal` CAS | 快速失败 run 的回调 与 启动路径的迟到写,谁先谁后都不丢裁决 |
| 4 | 启动清扫 `reconcile_on_startup` | 进程死亡留下的活跃记录与卡住的 once 任务 |
### 8.5 完成回调与启动清扫
## 9. 测试分层
`handle_run_completion(outcome, now)` 先写执行记录的终态,再看任务:`once` 推向终态成功→completed / 中断→cancelled / 失败→failed`cron` 保持状态不变。**但 `last_error` 无条件写入**——cron 任务保住了调度,仍然要报告上次出了什么问题。任务在 run 飞行途中被删掉不算错误,静默返回。
分层与架构一一对应,失败定位因此清晰:
`reconcile_on_startup(error)` 依次跑两个清扫并返回修复计数。它**不吞异常**——部分清扫失败要不要阻塞启动,是调用方的策略,不是领域的。
| 测试文件 | 层 | 规模 | 红了说明 |
|---|---|---|---|
| `test_schedule_domain.py` | 聚合与值对象 | 103 | 业务规则错(四张状态真值表逐格覆盖) |
| `test_schedule_service.py` | 用例编排(四 fake 端口) | 44 | 编排顺序、四出口、快路径/仲裁一致性错 |
| `test_schedule_fakes.py` | 端口契约 × 2 实现 | 85 | 存储实现与端口语义漂移 |
| `test_schedule_dispatch_race.py` | 真 sqlite 并发 | 4 | 索引仲裁、TOCTOU 收敛坏了——契约测试拿不住这个 |
| `test_schedule_corrupt_rows.py` | 适配器坏行语义(真 sqlite | 5 | 坏行错误词汇或跳过语义坏了 |
| `test_schedule_router.py` | 主适配器(真 service + fake 端口) | 51 | 协议转换、to_command、错误映射错 |
| `test_schedule_response_models.py` | api model 出向 | 14 | 泄漏关闭/线上兼容断言坏了 |
| `test_schedule_poller.py` / `test_schedule_run_completion.py` / `test_schedule_run_launcher.py` / `test_schedule_thread_lookup.py` | 时钟入口与三个适配器 | 10/21/16/6 | 对应适配器的翻译或韧性语义坏了 |
| `test_composition.py` | 组合根 | 16 | 装配规则错memory → None、policy 逐键映射、hook None 分支) |
---
契约套件的形状与 feedback 相同fake 住 `tests/schedule_fakes.py`,参数化 fixture 双实现各跑一遍),多一条边界声明:**契约不拥有原子性**——真数据库的 race 测试是唯一能证明"两个派发者并发是安全的"的地方。`test_schedule_service.py` 是规范 §1「唯一检验」的活样本完整生命周期创建→认领→派发→重叠跳过→完成→暂停→删除在零 IO 的 fake 上端到端跑通。
## 9. 二次开发指引
## 10. 二次开发指引
### 9.1 给任务加一个字段
### 10.1 给任务加一个字段
例:加一个 `notify_on_failure: bool`
例:加 `notify_on_failure: bool`
1. `model/task.py``ScheduledTask` 加字段(带默认值)
2. 若有约束,写进 `__post_init__`
3. 若用户要能设置它:`service.py``create_task` / `update_task` 加参数
4. `persistence/scheduled_tasks/model.py` 的 ORM 行加列
5. 新增一个 alembic revision`cd backend && make migrate-rev MSG="..."`),用 `_helpers.py` 的幂等 helper
6. router 的请求/响应模型加字段
1. `model/task.py``ScheduledTask` 加字段(带默认值);若有约束,写进 `__post_init__`
2. `test_schedule_domain.py` 先写红的测试TDD 是本仓库硬要求)
3. 用户要能设置它才改 command`commands.py``CreateScheduledTask` 加字段、`UpdateScheduledTask``UNSET` 默认字段;`service.py` 的 handler 跟着传
4. ORM 层——`persistence/scheduled_tasks/model.py` 加列 + alembic revision`make migrate-rev`,用 `_helpers.py` 幂等 helper
5. `scheduled_task_repository.py``_to_domain` / `_apply` 各加一行
6. api model请求模型加字段、`to_command` 跟着传;响应模型**默认不加**,除非前端真的需要
7. 前端 `frontend/src/core/scheduled-tasks/types.ts` 同步
8. 域测试补一条;若走了第 3 步service 测试也补一条
端口通常不用动——`add` / `save` 交换的是整个聚合,多一个字段不改变签名。
端口通常不用动——`add` / `save` 交换整个聚合。**例外**:想让 `record_launch` / `record_completion` 写它,必须先回答它归哪条 CAS 所有§7.1 ②)
### 9.2 加一种调度类型
### 10.2 加一种调度类型
例:加 `interval`(每 N 分钟)。
1. `model/enums.py``ScheduleType` 加成员
2. `model/spec.py`:加承载参数的字段(如 `interval_seconds`)、在 `__post_init__` 加校验、在 `next_after` 加一个分支
3. `ScheduleSpec.from_primitives`(见 §4.5)加一个分支;两个适配器各自的 `_spec_to_column` / `_spec_to_wire` 也各加一个
4. **逐个检查 `task.py` 里四个 `status_after_*`**——它们目前都在问"是不是 ONCE",新类型会落进 else 分支。确认那是你要的语义大概率是interval 与 cron 同属周期性)
5. 域测试:新类型在 §5.4 四张表里各补一行
2. `model/spec.py`:加承载参数的字段`__post_init__` 加校验、`next_after` 加分支、`from_primitives`分支
3. 两个出向映射各加一个分支:适配器的 `_spec_to_column`、api model 的 `_spec_to_wire`
4. **逐个检查 `task.py` 四个 `status_after_*`**——它们都在问"是不是 ONCE",新类型落进 else 分支确认那是你要的语义interval 与 cron 同属周期性,大概率是
5. 域测试在四张真值表里各补一行
不需要改数据库——`schedule_spec` 是 JSON 列。
### 9.3 改一条状态推导规则
### 10.3 加一个用例
只改 `model/task.py` 对应的那个方法,**顺便改 `app/scheduler/service.py` 里的旧副本**(见 §2`domain/schedule/service.py` 一般不用动——它只调用聚合,不复制规则。
1. 端口够用吗?够就不要加方法;不够则先在 `ports.py` 写清语义docstring 是契约测试的依据),契约套件加用例——**两套实现都跑通并断言返回值**
2. 写用例在 `commands.py` 加 command命名链三拼写对齐HTTP 驱动才 command 化,时钟/回调驱动不)
3. `service.py` 加 handler`test_schedule_service.py` 用 fake 测。**不要在 service 里写规则**——出现 `if task.status is ...` 就说明那条判断属于聚合
4. router 加端点,只做协议转换(有 body 则请求模型加 `to_command`
域测试里对应的真值表用例必须同步更新——那张表就是规则的规格说明。
### 10.4 加一种重叠策略
### 9.4 加一个用例
例:加 `queue`(排队而非跳过)。改动面最大,因为触及数据库不变量:
例:加"克隆一个任务"。
1. `model/task.py``skips_on_overlap` 拆成策略判断(字符串比较只在这一处)
2. `uq_scheduled_task_run_active` **必须**改成条件化谓词(`... AND overlap_policy='skip'`)——**ORM `__table_args__` 与迁移两处都要改**(空库 bootstrap 走 `create_all` 不执行迁移)
3. 跳过路径的墓碑逻辑相应分叉
4. 动手前读 `backend/AGENTS.md` 关于这个索引的整段说明
1. `service.py` 加方法:读出聚合 → 用聚合方法或 `replace` 造出新的 → 经 `add` 持久化
2. **不要在 service 里写规则**。如果发现自己在写 `if task.status is ...`,那条判断属于聚合
3. 只有当现有端口回答不了你的问题时才加端口方法——先问"这是新的存储能力,还是我把编排写复杂了"
4. service 测试补一组(全 fake零 IO
5. router 加端点,把领域错误映射成 HTTP 码
### 10.5 不要做的事
### 9.5 加一种重叠策略
规范 §4 全部适用schedule 语境下最常犯的五条:
例:加 `queue`(排队而非跳过)。
- **不要把 `record_launch` / `record_completion` 改成 `save(task)`。** CAS 分字段所有权是规范 §2.1 ③w 的明文例外,读-改-写会重新引入它们要关的竞态。
- **不要在 service 里写业务判断。** 状态推导、重叠语义、re-arm 规则都有归属。
- **不要让 `launch` 逃逸第三种异常。** 领域只认两个;线程忙被记成失败而不是跳过。
- **不要在领域层读时钟或配置。** `now` 显式传参,阈值经 `SchedulePolicy` 注入。
- **不要为定时执行另建运行栈。** 必须复用现有 run 生命周期。
这是改动面最大的一种,因为它触及数据库不变量:
1. `model/task.py``skips_on_overlap` 拆成策略判断
2. `scheduled_task_runs` 的部分唯一索引**必须**改成条件化的(`... AND overlap_policy = 'skip'`)——现在它是纯状态谓词,会把排队的执行也挡掉。索引同时定义在 ORM `__table_args__` 和迁移文件里,**两处都要改**(空库 bootstrap 走 `create_all`,不执行迁移)
3. 跳过路径的墓碑逻辑要相应分叉
动手前先读 `backend/AGENTS.md` 里关于这个索引的整段说明。
### 9.6 不要做的事
- **不要在 router 或旧 service 里新增业务判断**——新规则一律进聚合
- **不要在 `domain/schedule/service.py` 里写规则**——它只编排;出现 `if task.status is ...` 就说明放错地方了
- **不要在领域层读配置或时钟**——阈值通过 `SchedulePolicy` 注入,`now` 显式传参CI 的纯度测试会拦截基础设施导入
- **不要让端口签名沾上技术词汇**——`Mapping[str, Any]`、HTTP 状态码、表名出现在 `ports.py` 里都是信号
- **不要为定时执行另建一套运行栈**——必须复用现有的 run 生命周期
---
## 10. 常见陷阱速查
## 11. 常见陷阱速查
| 陷阱 | 后果 |
|---|---|
| 用 `ensure_launchable` 做派发后重排 | 正常执行完的 cron 任务被"提前量不足"拒绝 |
| 把 `status_after_*``if/elif` 写成并列 `if` | 两处判定顺序失效,静默改变行为 |
| 多次调用 `resolve_execution_thread()` | 每次拿到不同 thread记录与实际执行对不上 |
| 多次调用 `resolve_execution_thread()` | 每次拿到不同 thread记录与实际执行对不上 |
| 墓碑先建 `queued` 再改 `skipped` | 撞上唯一索引,跳过流程直接报错 |
| 改了 `ACTIVE_RUN_STATUSES` 没改索引谓词 | 快路径与数据库仲裁者判断不一致 |
| CAS 写成读-改-写整聚合 | 完成回调重放过期快照cron 任务永久卡死在 `running` |
| 终态任务改了调度但没重新武装 | 接口返回 200任务永不触发 |
| 只改领域模型,忘了 `app/scheduler/service.py` | 迁移完成前,生产行为不变 |
| 适配器让 `launch` 逃逸出第三种异常 | 领域收到不认识的错误;线程忙被记成失败而不是跳过 |
| 快路径与槽位仲裁产生不同结果 | 调用方能分辨被哪种机制拦下,重试行为分叉 |
| 把契约测试全绿当成可以多实例 | fake 没有任何原子性;并发由真数据库的 dispatch_race 测试负责 |
---
## 11. 术语表
| 术语 | 含义 |
|---|---|
| 任务task | 用户注册的定时**规则**,长期存在 |
| 执行run | 规则的**一次**触发,只读历史 |
| 派发dispatch | 把一个到期任务变成一次执行的动作 |
| 触发方式trigger | `scheduled`(轮询器)或 `manual`(用户点"立即执行" |
| 认领claim | 轮询器取得某任务这一轮调度所有权 |
| 租约lease | 认领时盖的带过期时间的戳,用于崩溃恢复 |
| 重叠overlap | 到点时上一次执行还没结束 |
| 墓碑tombstone | 被跳过的那次执行留下的终态记录 |
| 重新武装re-arm | 把终态任务拉回 `enabled` 使其可再次被认领 |
| 端口port | 领域声明、外圈实现的技术中立接口 |
| input port | 用例接口,被入口调用——这里就是 `ScheduleService` |
| output port | 领域对外的依赖,被适配器实现——这里是四个 Protocol |
| 聚合aggregate | 一致性边界,不变量在构造期成立 |
| 活跃槽位active slot | 一个任务同时最多持有一条 `queued`/`running` 执行 |
---
| 适配器让 `launch` 逃逸第三种异常 | 线程忙被记成失败而不是跳过 |
| 坏行错误进了 router 的 422 映射 | 存储故障被报成客户端错误PATCH 修复路径自身 422 |
| 把契约测试全绿当成可以多实例 | fake 没有原子性;并发由真数据库的 race 测试负责 |
| 显式继承 Protocol 拼错方法名 | 静默继承 `...` 体返回 `None``isinstance` 抓不到——契约用例必须断言返回值 |
## 12. 代码索引
**内圈(已迁移,本文覆盖)**
**内圈**
| 文件 | 内容 |
| 关注点 | 文件 |
|---|---|
| [`domain/schedule/service.py`](../packages/harness/deerflow/domain/schedule/service.py) | `ScheduleService` · `DispatchResult` · `ContextChange` |
| [`domain/schedule/ports.py`](../packages/harness/deerflow/domain/schedule/ports.py) | 4 个 Protocol · `LaunchedRun` · `RunOutcome` |
| [`domain/schedule/model/enums.py`](../packages/harness/deerflow/domain/schedule/model/enums.py) | `TaskStatus` `RunStatus` `ScheduleType` `ContextMode` `TriggerKind` `DispatchOutcome` |
| [`domain/schedule/model/errors.py`](../packages/harness/deerflow/domain/schedule/model/errors.py) | 9 个领域错误 |
| [`domain/schedule/model/spec.py`](../packages/harness/deerflow/domain/schedule/model/spec.py) | `ScheduleSpec` `SchedulePolicy` |
| [`domain/schedule/model/task.py`](../packages/harness/deerflow/domain/schedule/model/task.py) | `ScheduledTask` `TERMINAL_TASK_STATUSES` |
| [`domain/schedule/model/run.py`](../packages/harness/deerflow/domain/schedule/model/run.py) | `ScheduledRun` `ACTIVE_RUN_STATUSES` `TERMINAL_RUN_STATUSES` |
| [`tests/test_schedule_domain.py`](../tests/test_schedule_domain.py) | 域测试,全同步零 IO四张真值表逐格覆盖 |
| [`tests/test_schedule_service.py`](../tests/test_schedule_service.py) | 用例测试;完整生命周期跑在 fake 上,是迁移的验收标准 |
| [`tests/schedule_fakes.py`](../tests/schedule_fakes.py) | 四个端口的内存实现 |
| [`tests/test_schedule_fakes.py`](../tests/test_schedule_fakes.py) | 契约测试31 用例 × 内存 fake + 真 sqlite 两套实现 |
| 上下文公开 API | `packages/harness/deerflow/domain/schedule/__init__.py` |
| 领域错误(一族一基类) | `packages/harness/deerflow/domain/schedule/exceptions.py` |
| 命令 + `UNSET` + `ContextChange` | `packages/harness/deerflow/domain/schedule/commands.py` |
| 端口契约(语义写在 docstring 里) | `packages/harness/deerflow/domain/schedule/ports.py` |
| 用例编排(写方法即 command handler | `packages/harness/deerflow/domain/schedule/service.py` |
| 值对象 | `packages/harness/deerflow/domain/schedule/model/spec.py` |
| 任务聚合与四条状态推导 | `packages/harness/deerflow/domain/schedule/model/task.py` |
| 执行聚合与两个具名工厂 | `packages/harness/deerflow/domain/schedule/model/run.py` |
**主适配器(入口)**
**入口(三驱动源)**
| 文件 | 内容 |
| 关注点 | 文件 |
|---|---|
| [`gateway/routers/schedule/router.py`](../app/gateway/routers/schedule/router.py) | 10 个 HTTP 端点,只做协议转换 + 领域错误→状态码 |
| [`gateway/routers/schedule/models.py`](../app/gateway/routers/schedule/models.py) | 请求/响应模型;响应是白名单,不是 ORM 转储;`schedule_spec` 的进出转换是模型自己的方法 |
| [`app/scheduler/poller.py`](../app/scheduler/poller.py) | 轮询时钟 + 启动恢复 |
| [`adapters/schedule/run_completion.py`](../app/adapters/schedule/run_completion.py) | 运行完成回调;过滤掉非本上下文的运行,其余转成 `RunOutcome` 并调用用例。与上面两个同为入站,只是留在上下文包内 |
| HTTP 端点 + 错误映射单表 | `backend/app/gateway/routers/schedule/router.py` |
| api model`to_command` / `from_domain` / `_spec_to_wire` | `backend/app/gateway/routers/schedule/models.py` |
| 时钟入口 + 启动恢复 | `backend/app/scheduler/poller.py` |
| 回调入口(过滤 + 翻译成 `RunOutcome` | `backend/app/adapters/schedule/run_completion.py` |
**从适配器**
**从适配器与组合根**
| 文件 | 内容 |
| 关注点 | 文件 |
|---|---|
| [`adapters/schedule/scheduled_task_repository.py`](../app/adapters/schedule/scheduled_task_repository.py) | 自有持久化;`FOR UPDATE SKIP LOCKED``protect_terminal` CAS |
| [`adapters/schedule/scheduled_run_repository.py`](../app/adapters/schedule/scheduled_run_repository.py) | 自有持久化;`IntegrityError → ActiveRunConflictError` 的翻译点 |
| [`adapters/schedule/run_launcher.py`](../app/adapters/schedule/run_launcher.py) | 防腐层;`ConflictError` / `HTTPException(409)``ThreadBusyError` |
| [`adapters/schedule/thread_lookup.py`](../app/adapters/schedule/thread_lookup.py) | 防腐层;`check_access(require_existing=True)` |
| 任务仓储claim / CAS / 坏行翻译) | `backend/app/adapters/schedule/scheduled_task_repository.py` |
| 执行仓储(`IntegrityError → ActiveRunConflictError` | `backend/app/adapters/schedule/scheduled_run_repository.py` |
| 防腐层:启动 run | `backend/app/adapters/schedule/run_launcher.py` |
| 防腐层thread 归属 | `backend/app/adapters/schedule/thread_lookup.py` |
| 组合根三函数 | `backend/app/composition.py` |
| 运营阈值来源 | `packages/harness/deerflow/config/scheduler_config.py` |
| ORM 行 + 唯一索引(双定义之一) | `packages/harness/deerflow/persistence/scheduled_task{s,_runs}/model.py`、迁移 `0003` / `0007` |
| 零 IO fake | `backend/tests/schedule_fakes.py` |
| 测试分层 | 见 §9 |
**组合根**
| 文件 | 内容 |
|---|---|
| [`app/composition.py`](../app/composition.py) | `build_domain_services()` · `build_run_completion_hook()` · `build_schedule_policy()` |
| [`config/scheduler_config.py`](../packages/harness/deerflow/config/scheduler_config.py) | `enabled` `poll_interval_seconds` `lease_seconds` `max_concurrent_runs` `min_once_delay_seconds` |
**已停用,待删除(见 §2**
| 文件 | 内容 |
|---|---|
| [`app/gateway/routers/scheduled_tasks.py`](../app/gateway/routers/scheduled_tasks.py) | 旧端点,已不注册 |
| [`app/scheduler/service.py`](../app/scheduler/service.py) | 旧轮询 + 派发编排,已不装配 |
| [`deerflow/scheduler/schedules.py`](../packages/harness/deerflow/scheduler/schedules.py) | 时区 / cron 计算,规则已进 `ScheduleSpec` |
| [`persistence/scheduled_task*/sql.py`](../packages/harness/deerflow/persistence/scheduled_tasks/) | 旧仓储(返回裸 dict。**ORM 行与唯一索引定义仍在用,不要一起删** |
**前端**
`frontend/src/core/scheduled-tasks/`类型、API、cron 解析、预设配方)与 `frontend/src/app/workspace/scheduled-tasks/page.tsx`
**前端**`frontend/src/core/scheduled-tasks/`类型、API、cron 解析、预设配方)与 `frontend/src/app/workspace/scheduled-tasks/page.tsx`