Onyx Craft 定时任务(Scheduled Tasks)架构解析:从 Cron 调度到无头 Agent 执行的完整实现指南
【免费下载链接】danswerOpen Source AI Platform - AI Chat with advanced features that works with every LLM项目地址: https://gitcode.com/GitHub_Trending/da/danswer
本文以 docs/craft/features/scheduled-tasks/overview.md 为核心,结合 Onyx(原 danswer)仓库中该功能的完整源码实现,系统讲解 Craft 定时任务的产品目标、核心设计决策、Celery 调度架构、数据模型、REST API 规范、Cron 调度语义、运行生命周期与测试策略。读完本文,你将掌握"保存 Prompt + 定时执行"这一产品形态在 Onyx 中的落地方案,并能在本地部署、调用 API、阅读源码和编写测试时做到心中有数。
一、产品目标与 V1 范围
Scheduled Tasks(定时任务)是 Onyx Craft 的一个产品面:用户保存一段 Prompt 与一个调度计划,系统按定时器把这段 Prompt 作为 Craft 任务无头(headless)执行。每一次触发都会创建一个全新的 Craft 会话(BuildSession),无头执行完整条 Agent 流程,并把执行过程完整记录下来。
V1 的验收标准(来自 overview.md):
- 用户能在
/craft/v1/tasks创建任务; - 任务能在浏览器关闭后照常触发执行;
- 用户点击任意一次历史运行记录,即可打开那次执行产生的完整会话视图。
V1 明确不做以下能力:实时挂接(live-attach)、事件触发、共享任务、模板、重试策略、预算上限(budget caps)、运行对比(run diffing)、日历视图、连续失败自动禁用、外部任务管理 API。
二、核心设计决策:理解该功能的关键
overview 文档沉淀了该功能最重要的十余条设计决策,它们是理解整个实现的钥匙:
- 一次运行(run)就是一个
BuildSession。scheduled_task_run.session_id是对build_session.id的外键。点击已完成运行的记录,打开的就是现有的会话视图,无需新建 transcript UI。 - 每次触发都创建全新会话,复用现有的
SessionManager.create_session__no_commit路径,因此沙箱预置、工作区初始化、技能物化、AGENTS.md 生成、数据包(packet)日志等流程原样复用。 - 无头执行器复用
send_message的持久化半段:把_stream_cli_agent_response拆成_yield_acp_events(纯 ACP 事件生成器)与_persist_acp_events(BuildStreamingState消费者,写入BuildMessage行)两部分。SSE 端点把两者与 SSE 格式化器组合;执行器则用 drain-to-completion 方式包裹。这样转录(transcript)与交互式完全一致,且没有重复代码。 - 专用
scheduled_tasksCelery worker——只跑执行器。长时执行器run_scheduled_task运行在新的celery_worker_scheduled_tasks进程上(已在 supervisord、dev runner 与 Helm chart 中注册)。因为无头 Agent 触发是长时操作(LLM + 沙箱内工具调用),若与heavy队列(pruning、权限同步、CSV 导出)共用,少量触发就可能饿死 heavy 队列的其他任务。专用 worker 意味着独立的线程池、独立的 HPA/KEDA 扩缩容,以及独立的 Prometheus 端口(9098)。而派发器(dispatcher)与卡死运行清扫器(stuck-run sweeper)是纯 DB 协调工作,运行在 primary 队列——若把派发逻辑也路由到专用池,执行器饱和时会反过来拖住派发。 - 调度存储:
(cron_expression, editor_mode)二元组。三种编辑器模式保存时统一编译为 cron;Cron 表达式按 UTC 求值;editor_mode仅是 UI 提示。next_run_at在每次触发与每次编辑时重算;暂停时置为 NULL。 - 并发策略:周期性触发用
SKIP_IF_RUNNING(上一次运行仍在进行 → 写入skipped行并携带skip_reason,同时仍推进next_run_at);Run Now 用QUEUE_ONE(暂停状态下也能工作,且不触碰next_run_at)。 - 运行以任务作者身份执行:
create_session__no_commit(user_id=task.user_id)——技能、Onyx 搜索、OAuth 授权、审批策略全部走与交互式 UI 相同的用户作用域路径。 - 软删除保留历史:
deleted=true停止派发;运行记录与会话保留,用户仍可从任务运行历史中打开过去的运行。 BuildSession.origin把定时运行排除出侧边栏:新增枚举列(INTERACTIVE | SCHEDULED,默认INTERACTIVE,存量行服务端默认'interactive'),由执行器在创建会话时设置;侧边栏查询过滤origin = INTERACTIVE。定时触发与 Run Now(都走执行器、都带origin=SCHEDULED)一并覆盖。之所以用列而不是对scheduled_task_run做NOT EXISTS反连接,是因为派发器先于执行器写 run 行——基于 join 的过滤会短暂"泄漏"。未来非交互式 origin(eval 运行、自动化)可复用同一接缝。- V1 无重试:失败就是一行记录;用户点击 Run Now 或等待下一次触发。
- 通知复用现有
Notification模型:新增两种类型(SCHEDULED_TASK_FAILED、SCHEDULED_TASK_AWAITING_APPROVAL)。V1 不发邮件、不发 Slack。 - 卡死运行清扫器(每小时):
queued > 15 分钟与running > budget→ 标记为failed (stuck),兜底捕获死亡 worker。
三、调度架构:Beat + 双队列 + 专用 Worker
overview 用一张 ASCII 架构图说明了数据流,结合 tasks.py 与 beat_schedule.py 的源码,可以还原出精确的执行时序:
Beat (30s, per tenant) Celery "scheduled_tasks" queue Primary queue (served by celery_worker_scheduled_tasks) ────────────────────── ──────────────────────────────── dispatch_due_scheduled_tasks run_scheduled_task(run_id) BEGIN; if run.status != 'queued': return SELECT FROM scheduled_task mark 'running' WHERE active AND due session = SessionManager FOR UPDATE SKIP LOCKED; .create_session__no_commit( for each row: user_id=task.user_id) ├─ if prior run in flight run.session_id ← session.id │ → insert skipped row for event in _yield_acp_events( ├─ insert queued run session, task.prompt): ├─ next_run_at = croniter _persist_acp_events([event]) │ .next(now) if budget exceeded → failed └─ enqueue run_scheduled_task ───► if approval required → (run_id, expires=900, awaiting_approval queue=scheduled_tasks) mark succeeded / failed COMMIT; emit Notification if failed Stuck-run sweep (hourly, primary queue) ────────────────────── │ cleanup_stuck_scheduled_runs ▼ queued > 15m → failed (stuck) BuildMessage rows (existing tables, running > budget → failed (timeout) written by shared persist consumer)3.1 派发器(dispatcher):每租户每 30 秒
dispatch_due_scheduled_tasks注册在 beat_schedule.py 中,schedule=timedelta(seconds=30)、队列为 primary、优先级 MEDIUM、expires=60。注释说明 30 秒是规格契约:60 秒对分钟级 cron 太粗,15 秒会让没有到期任务的租户过度占用FOR UPDATE SKIP LOCKED路径。
派发逻辑(见 tasks.py):
- 调用
claim_due_scheduled_tasks,用FOR UPDATE SKIP LOCKED原子认领到期任务,批量上限DISPATCH_BATCH_SIZE = 50。 - 对每个认领的任务,按优先级判定:
- 任务所有者 Craft 被禁用(部署级开关、用户级覆盖或工作区默认)→ 插入
SKIPPED行(skip_reason=owner_craft_disabled),调度保持存活,重新启用后恢复; - 存在进行中的运行(
has_in_flight_run_for_task检查 QUEUED/RUNNING)→ 插入SKIPPED行(skip_reason=prior_in_flight); - 否则插入
QUEUEDrun 行并收集 run_id 待入队。
- 任务所有者 Craft 被禁用(部署级开关、用户级覆盖或工作区默认)→ 插入
- 无论命中哪条分支都推进
next_run_at(advance_next_run_at),否则下一次 tick 会重复认领同一行。 - 全部在同一事务内提交后,再对每个 run_id 发送
SCHEDULED_TASKS_RUN到scheduled_tasks队列,expires=QUEUE_RESIDENCY_SECONDS。如果 commit 与入队之间崩溃,卡死清扫器会在约 15 分钟后回收这些 QUEUED 行。
claim_due_scheduled_tasks的实现细节值得注意(db/scheduled_task.py):查询条件为status='active' AND deleted=false AND next_run_at IS NOT NULL AND next_run_at <= now,按next_run_at升序取批,with_for_update(skip_locked=True)保证并发 tick 不会双触发。
3.2 执行器(executor):专用 worker
run_scheduled_task是薄包装(tasks.py),acks_late=False(worker 崩溃不触发 Celery 重试,V1 无重试),实际逻辑全部在run_scheduled_task_logic(executor.py),拆出来是为了让 Agent 驱动逻辑无需 Celery worker 即可被外部依赖单元测试直接实例化。
执行器状态机(源码确认):
- 幂等守卫:run 状态不是 QUEUED 就直接返回(Celery 可能重投递,或清扫器已先标记失败)。
- 沙箱就绪:
SessionManager.ensure_sandbox_running——没有沙箱则创建,等待并发 provisioner(上限PROVISION_WAIT_SECONDS = 120s),原地唤醒 SLEEPING/TERMINATED/FAILED 的沙箱。等待窗口内仍处于 PROVISIONING 则 SKIP;其他失败标记error_class=sandbox_wake_failed。 - 过渡到 RUNNING 并提交(让 UI 与清扫器可见)。
- 以
origin=SCHEDULED创建全新BuildSession,会话名Scheduled: {task_name},并把 session_id 写回 run 行。 - 通过共享的
yield_sandbox_events生成器驱动 Agent,用persist_sandbox_event持久化每个事件;强制执行预算:默认硬上限SCHEDULED_RUN_HARD_CAP_SECONDS = 60 * 60(60 分钟,见 timeouts.py),软预算为其 40%(24 分钟)。注释特别说明 Celery 线程池会静默忽略soft_time_limit/time_limit,所以预算必须在任务体内实现。 - 命中
RequestPermissionRequest(审批门):标记AWAITING_APPROVAL、发通知、不写终态返回。恢复机制归审批项目所有,在那之前"仅作展示用终态"。 - 正常流结束:
finalize_persist→ 提炼约 120 字符的摘要(_clip_summary,截断时按单词边界断句并加省略号)→ 标记 SUCCEEDED。 - 驱动循环内任何异常:标记 FAILED(带异常类名 + 详情)并发通知。刻意吞掉异常,防止 Celery 自动重试。
摘要生成有两条路径:优先取BuildStreamingState.message_chunks中最近流式文本;若内存中无待刷 chunk(如中途已 flush),则回退扫描持久化消息,倒序查找最后一个agent_message元数据 blob(见_summary_from_session_messages)。
3.3 卡死运行清扫器(stuck-run sweeper)
cleanup_stuck_scheduled_runs每小时运行(beat_schedule.py),判定规则(tasks.py):
- QUEUED 超过
QUEUE_RESIDENCY_SECONDS = 15 * 60(15 分钟)→failed (stuck); - RUNNING 超过
SCHEDULED_RUN_HARD_CAP_SECONDS + TURN_RECLAIM_SLACK_SECONDS(60 分钟 + 15 分钟)→failed (stuck)。
两者与 Celery 的expires构成同一策略的三处执行点(见 timeouts.py 的注释):队列驻留超限的消息被丢弃且行被同阈值回收;正常跑满自身预算的运行一定会先把自己标记 FAILED,清扫器只是兜底。
3.4 专用 worker 进程
apps/scheduled_tasks.py 定义了独立 Celery app,配置来自onyx.background.celery.configs.scheduled_tasks,autodiscover_tasks只加载onyx.background.celery.tasks.scheduled_tasks模块,worker_ready时调用start_metrics_server("scheduled_tasks")(即 Prometheus 端口 9098)。worker 初始化沿用app_base的启动检查(等 Redis、等 DB、等文档索引),并设置 PostgreSQL 应用名POSTGRES_CELERY_WORKER_SCHEDULED_TASKS_APP_NAME。
四、数据模型
overview 给出了核心 ORM 模型的骨架,与 db/enums.py 和 db/scheduled_task.py 的源码完全对应:
class SessionOrigin(str, Enum): # INTERACTIVE / SCHEDULED / SLACK # 加到 BuildSession;默认 INTERACTIVE class ScheduledTaskStatus(str, Enum): # ACTIVE / PAUSED class ScheduledTaskRunStatus(str, Enum): # QUEUED / RUNNING / SUCCEEDED / # FAILED / SKIPPED / AWAITING_APPROVAL class ScheduledTaskTriggerSource(str, Enum): # SCHEDULED / MANUAL_RUN_NOW class ScheduledTask(Base): __tablename__ = "scheduled_task" id, user_id (FK user, CASCADE) name (str), prompt (text) cron_expression (str), editor_mode (str) status (enum, default ACTIVE) next_run_at (DateTime tz, nullable) # dispatcher 唯一读取的字段 deleted (bool, default False) created_at, updated_at runs ← back-populated, cascade all,delete-orphan __table_args__ = ( Index("ix_scheduled_task_dispatch", "status", "deleted", "next_run_at"), Index("ix_scheduled_task_user_created", "user_id", desc("created_at")), ) class ScheduledTaskRun(Base): __tablename__ = "scheduled_task_run" id, task_id (FK scheduled_task, CASCADE) session_id (FK build_session, SET NULL) # 执行器创建会话后回填 status (enum, default QUEUED), trigger_source (enum) skip_reason / error_class / error_detail (nullable) started_at (default now), finished_at (nullable) summary (str, ~120 chars of final agent message) __table_args__ = ( Index("ix_scheduled_task_run_task_started", "task_id", desc("started_at")), Index("ix_scheduled_task_run_status", "status"), )设计要点(源码佐证):
next_run_at是派发器唯一读取的字段:暂停置 NULL(update_scheduled_task中status=PAUSED分支)、编辑时重算、deleted=true从认领查询中排除(claim_due_scheduled_tasks的 WHERE 条件)。session_id可空:在派发器 INSERT 与执行器创建会话之间的短暂窗口内为 NULL。- 没有
attempts计数:V1 无重试。 ScheduledTaskRunStatus.is_terminal()只把 SUCCEEDED/FAILED/SKIPPED 视为终态(enums.py),AWAITING_APPROVAL不在其中。error_class是封闭枚举(task_missing/sandbox_wake_failed/executor_error/timeout/stuck/agent_exception),让仪表盘与排障查询可以按已知词表透视;意外的运行时异常统一记agent_exception,真实异常类名与消息放入error_detail。skip_reason也是封闭枚举:prior_in_flight/owner_craft_disabled。- 派发器的认领查询用
selectinload(ScheduledTask.user)预加载所有者,以便在事务内做 Craft 启用判定。 mark_run_status在写入终态时自动填充finished_at;insert_run对 SKIPPED 行也当场填finished_at,避免 UI 特判。
数据库操作全部集中在 db/scheduled_task.py(按项目 CLAUDE.md 约定,所有查询下沉到 db 层,所有函数首个参数为db_session),包括:任务 CRUD(create_scheduled_task/update_scheduled_task/soft_delete_scheduled_task,所有权与 NOT_FOUND 抛出语义一致)、派发热路径(claim_due_scheduled_tasks/advance_next_run_at)、run CRUD(insert_run/mark_run_status/list_runs_for_task)、卡死扫描(find_stuck_runs),以及会话横幅辅助(get_scheduled_run_context)与审批预授权查询(get_live_scheduled_run_grants)。
五、API 规范
所有端点抛出OnyxError,返回类型化 FastAPI 响应;挂载在/api/build/scheduled-tasks(复用现有/build前缀与require_onyx_craft_enabled门控),作用域限定在认证用户(V1 无管理员视图)。路由实现见 api.py,权限依赖为require_permission(Permission.BASIC_ACCESS)。
| 方法 | 路径 | 说明 |
|---|---|---|
| GET | /scheduled-tasks | 任务列表(id、名称、人类可读调度、状态、next_run_at、最近一次运行摘要) |
| POST | /scheduled-tasks | 创建。编辑器输入编译为 cron,可选run_immediately |
| GET | /scheduled-tasks/{id} | 任务详情 + 未来 3 次触发时间(UI 预览) |
| PATCH | /scheduled-tasks/{id} | 部分编辑。调度变化时重算next_run_at;暂停→NULL;恢复→重算。进行中的运行不受影响 |
| DELETE | /scheduled-tasks/{id} | 软删除(幂等,返回 204) |
| POST | /scheduled-tasks/{id}/run-now | 插入manual_run_nowrun 并入队执行器;暂停时可用;不触碰next_run_at |
| GET | /scheduled-tasks/{id}/runs | 分页运行历史,50/页,cursor=参数,最新在前 |
| GET | /build/sessions/{id}/scheduled-run-context | 若会话来自定时运行,返回可选的任务名 + id + 定时started_at;否则 404。供会话视图横幅使用 |
V1 没有 live-attach 端点。
请求体校验细节(源码确认):
ScheduledTaskCreate:name1–200 字符、prompt非空、editor_mode为 Literal 枚举、editor_payload是受类型约束的联合体、status默认 ACTIVE、run_immediately默认 false,另有pre_approved_app_ids与pre_approved_mcp_server_ids(预授权目标,配合审批门使用)。模型extra="forbid",未知字段直接 422。- 一个巧妙的校验管道:
_dispatch_editor_payload作为model_validator(mode="before"),把原始editor_payloaddict 按editor_mode路由到对应的类型化模型(IntervalPayload/DailyWeeklyPayload/AdvancedPayload),不匹配的形状在 FastAPI 层表现为 PydanticValidationError→ 422。 - PATCH 用
_editor_pair_consistency强制editor_mode与editor_payload必须成对出现。 _validate_app_ids拒绝未知外部应用 id;_validate_mcp_server_ids拒绝该用户在 Craft 中不可用的 MCP server id(已授权的存量 grant 可保留)。- runs 列表的
cursor是 ISO-8601 时间(上一页最后一行started_at),next_cursor为空表示已到末页;游标分页与ix_scheduled_task_run_task_started (task_id, started_at DESC)索引精确对齐。 run_immediately与run-now都遵循"先落库再入队"原则:run 行持久化提交后才_enqueue_executor,避免 worker 跑到自己的输入前面。入队时附带tenant_id,由TenantAwareTask在任务体内设置租户 contextvar,否则会在默认 schema 上执行。
六、调度语义:Cron 编译与求值
调度相关的纯函数集中在 schedule.py,它自述为 overview 文档中 cron 语义的单一事实来源:DB 存规范 5 字段 cron 字符串 +editor_mode(UI 提示);三种编辑器模式保存时统一编译为同一 cron 形式;compute_next_run_at返回 UTC 时间,比较也在 UTC 进行。
6.1 三种编辑器模式
- 间隔(Interval,N 分钟/小时/天):
*/N * * * */0 */N * * */M H */N * *(日节奏时 M:H 是编辑器要求的具体时刻)。 - 每日/每周(Daily/Weekly):
M H * * <weekdays>(如0 9 * * 1,3,5)。 - 高级(Advanced):用户手写,用
croniter校验。
compute_next_run_at=croniter(cron, after).get_next(datetime),存 UTC、UTC 比较。
6.2 源码级校验规则
Pydantic 模型在 HTTP 边界完成了全部校验,因此compile_to_cron是纯变换(源码注释明确:没有dict[str, Any]查找、没有临时字符串解析):
_HH_MM_RE = ^([01]?\d|2[0-3]):([0-5]\d)$:24 小时制HH:MM格式,与 UI 的<input type="time">对齐。IntervalPayload:every >= 1;unit == "days"时time_of_day必填(model_validator 强制)。DailyWeeklyPayload:weekdays是 0–6 的整数列表(0=周日,cron 约定);空列表 = 每天,等价于*。AdvancedPayload:必须恰好 5 个字段且croniter.is_valid通过。compile_to_cron的边界兜底:分钟间隔every > 59折叠为0 * * * *(*/60 * * * *非法);小时间隔every > 23折叠为0 0 * * *——与 UI 依赖的历史行为一致。- 读侧(
compute_next_run_at/next_n_fires/human_readable)也有_validate_cron防御,把非法表达式、非 5 字段、无未来触发(如不可能表达式)统一转成OnyxError(INVALID_INPUT)。 human_readable用cron-descriptor的ExpressionDescriptor渲染人类可读描述,失败时回退原始 cron 字符串。
6.3 时间语义要点
- 暂停中途触发:进行中的运行照常完成,
next_run_at→ NULL;恢复时从恢复时刻起重算(不回火)。 - 过期任务:只触发一次并向前推进(标准 cron 行为)。
- 没有未来触发的 cron:在 POST/PATCH 时直接拒绝。
- 调度变更重算规则(
update_scheduled_task):cron 变更且任务为(或变为)ACTIVE → 从 now 重算;变为 PAUSED → NULL;变为 ACTIVE → 重算。各预授权目标字段遵循普通 patch 语义(提供的替换该类目标,省略的不动)。
七、运行生命周期
queued ─► running ─► succeeded │ └─► failed (executor crash | budget | ACP error) └─► awaiting_approval ─► (approvals project resumes) ─► running ─► ... dispatcher also writes: skipped (prior run still in flight; next_run_at advances)补充说明(源码确认):
AWAITING_APPROVAL的恢复机制归审批项目所有;在该项目落地前,仅作展示用终态。- 多个定时运行与交互式
send_message路径可以并发执行在同一沙箱上——没有序列化租约。执行器仍会为每个 BuildSession 获取 prompt slot 锁(session_manager.prompt_slot),这是与交互式路径对称的防御性做法:当前每个定时运行都有全新 BuildSession、锁不会竞争,但未来若让定时与交互共享 BuildSession,此锁自动提供保护。若锁未获取成功(并发回合在飞),run 被标记 FAILED(error_class=agent_exception)。 - 每次触发都会
stamp_turn_deadline(软预算 = min(软预算, budget),硬上限 = budget);离开时clear_turn_deadline(slot 未丢失时),避免后续回合继承过期的定时截止时间。 - 超时与完成的竞态处理很精细:先检查
PromptResponse/Error再检查预算——截止时间恰好在 Agent 完成瞬间触发时,不能把成功误标为超时;预算检查先于持久化,防止失控 Agent 无限增长转录。 stop_reason == "cancelled"的PromptResponse被视为失败而非成功(opencode 中止 Agent,如MessageAbortedError)。
八、UI 设计
- 列表页(
/craft/v1/tasks):名称、调度、状态、last_run(相对时间 + 成功/失败图标)、next_run、行操作(Run now、Pause/Resume、Edit、Delete)。空状态展示产品文档中的三个 starter prompts。 - 编辑器(
/new、/:id/edit):单一表单 + 三 Tab 分段控件选择调度模式;Prompt 输入框复用 Craft 聊天输入(产品需求 10);"未来 3 次运行"预览由前端用cron-parser计算。 - 详情页(
/:id):头部(名称、调度、状态开关、next_run、操作按钮)+ 分页运行历史表格。行点击 → 打开会话视图(succeeded/failed);queued/running/awaiting_approval/skipped不可点击、显示 tooltip。 - 会话视图横幅:当 scheduled-run-context 端点返回结果时,在转录上方渲染"This session was started by scheduled task X at Y. ← Back to task."。聊天输入保持可用——用户可以在定时运行的会话上发送后续消息。
- 通知:两个新的铃铛条目——"Task X failed"、"Task X needs approval",深链到 run 行。通知实现在执行器的
_notify中(标题Scheduled task "{task_name}" failed/needs approval,additional_data携带task_id/run_id),且刻意 best-effort——通知失败绝不掩盖底层 run 状态更新。 - 移动端:列表/详情可容忍;编辑器明确要求桌面端(产品需求 8)。
九、测试策略
详见 tests.md,覆盖计划、既有测试套件布局与手工冒烟清单。
E2E 冒烟:web/tests/e2e/scheduled-tasks.spec.ts——单一 speccreate, run, and verify a run row exists:登录标准 worker 用户 → 访问/craft/v1/tasks(若重定向到/app则 Craft 功能开关关闭,spec 软跳过)→ 用唯一名称 + promptsay hi+ 每 5 分钟间隔创建任务 → 点save-and-run-now→ 断言详情页 active 状态芯片 → 最多等 60 秒看到data-run-status="succeeded"或="failed"的 run 行(两者皆可,证明 dispatcher → executor → 运行历史链路端到端可达)。前端选择器已锁定:new-task-button、task-name-input、task-prompt-input、interval-every、save-and-run-now、task-status-active、data-run-status。
运行方式:
npx playwright test scheduled-tasks需要完整本地 Onyx 栈:web、API、Postgres、Redis,以及专用celery_worker_scheduled_tasksworker——缺了它 run 永远到不了终态,测试在第 6 步超时。
刻意不自动化(依赖 code review + 手工清单):FOR UPDATE SKIP LOCKED派发并发、卡死清扫器、审批路径、侧边栏过滤、HTTP API 的用户所有权边界。
手工冒烟清单(合并前逐项过一遍):
- 每 2 分钟任务 + Onyx 搜索 prompt:离开 6 分钟,回来确认三次运行都有完整会话与合理
summary; - 周一/三/五 9 AM:验证
next_run_at正确、UTC 9 AM 强制 tick 能触发; - 中途暂停:进行中运行完成、无新触发;恢复 → 下次触发从
now()向前计算(不回火); - 审批边界:run 停在
AWAITING_APPROVAL,通知铃铛有条目,同一沙箱上的交互式 Craft 仍可用; - 中途杀 worker:run 处于 RUNNING 时停掉
celery_worker_scheduled_tasks,一小时内清扫器把它转为FAILED、error_class="stuck"。
十、部署与运行要求
- 调度功能由 Celery Beat 驱动(30 秒派发 + 每小时清扫),需要 Beat 进程在跑;
- 必须部署专用 worker 进程
celery_worker_scheduled_tasks(supervisord / dev runner / Helm chart 均已注册),执行器只在其上运行; - 多租户模式下派发与执行都依赖租户 contextvar(
tenant_id随任务入队),确保任务落到正确 schema; - V1 无重试,Celery 消息带
expires=,死消费者不会无限堆积 backlog(每个入队点都传 expires,符合项目 CLAUDE.md 约定)。
十一、深入阅读的源码地图
- 功能总设计:docs/craft/features/scheduled-tasks/overview.md
- 测试计划:docs/craft/features/scheduled-tasks/tests.md
- 数据库操作层:backend/onyx/db/scheduled_task.py
- 调度纯函数(cron 编译/校验/求值):backend/onyx/server/features/build/scheduled_tasks/schedule.py
- REST 路由:backend/onyx/server/features/build/scheduled_tasks/api.py
- 无头执行器:backend/onyx/server/features/build/scheduled_tasks/executor.py
- Celery 任务(派发/执行/清扫):backend/onyx/background/celery/tasks/scheduled_tasks/tasks.py
- Beat 调度注册:backend/onyx/background/celery/tasks/beat_schedule.py
- 专用 worker app:backend/onyx/background/celery/apps/scheduled_tasks.py
- 枚举定义(状态机/错误词表):backend/onyx/db/enums.py
- 超时注册表(预算与清扫阈值):backend/onyx/server/features/build/timeouts.py
- E2E 冒烟测试:
web/tests/e2e/scheduled-tasks.spec.ts
【免费下载链接】danswerOpen Source AI Platform - AI Chat with advanced features that works with every LLM项目地址: https://gitcode.com/GitHub_Trending/da/danswer
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考