Danswer Craft 定时任务实时运行查看:基于 opencode SSE 代理的 Live Session 架构设计
Danswer Craft 定时任务实时运行查看基于 opencode SSE 代理的 Live Session 架构设计【免费下载链接】danswerOpen Source AI Platform - AI Chat with advanced features that works with every LLM项目地址: https://gitcode.com/GitHub_Trending/da/danswer导读本文基于仓库中的设计文档 docs/craft/fix-live-scheduled-task-runs.md完整解读 DanswerOnyxCraft 模块中定时任务Scheduled Tasks实时运行查看的功能设计。该设计解决的核心痛点是此前用户必须等待定时任务进入终态SUCCEEDED / FAILED后才能点击运行记录打开会话本文描述的方案让用户可以在任务仍在执行中就打开 Craft 会话视图以 SSE 方式实时观察 Agent 进展并在运行结束后在同一会话中继续追问。读者读完本文将掌握实时查看的总体架构api-server 代理 opencode/event流、会话视图的运行态上下文设计、运行历史可点击规则的变化、前后端实现步骤与完整测试用例清单以及仓库中对应的源码印证。一、背景V1 定时任务的只能事后查看限制Craft 定时任务Scheduled Tasks的核心目标是用户保存一条 prompt 调度规则系统按 cron 定时以无头headless方式驱动 Agent 执行每次触发都创建一个全新的 Craft 会话并记录结果。V1 版本的明确定位是schedule-only、无 live-attach用户只能等运行结束后点击历史记录打开已完成的会话见 docs/craft/features/scheduled-tasks/overview.md。V1 的设计约束带来了三个与实时查看直接相关的现状运行历史表格每 5 秒轮询最新一页已加载的旧页保持稳定见 RunHistoryTable.tsx 中的RUN_HISTORY_REFRESH_INTERVAL_MS 5000。前端会阻止RUNNING与AWAITING_APPROVAL状态的行被点击即使后端已经在该运行上关联了session_id。可点击判定逻辑在 utils.ts 的getNonClickableReason中实现QUEUED行提示排队中、SKIPPED行提示已跳过、其余状态在无session_id时提示无会话。会话视图只在打开时一次性加载已保存的消息既不挂接到执行器的实时事件流也不感知定时来源的会话是否仍在被后台执行器驱动。也就是说后端在运行完成之前其实已经创建了会话见下文执行器时序但产品与前端都不允许用户在运行中打开它——这正是本设计文档要修复的产品契约。二、核心目标What要做什么设计文档将本次变更收敛为四条明确的产品行为允许用户在定时任务运行尚未结束时打开该运行无需等待终态。复用现有 Craft 会话视图作为运行查看器不新建独立的 run-detail 界面。仅在定时运行仍在执行期间阻止后续消息。运行一旦结束无论SUCCEEDED还是FAILED用户就可以在同一个 Craft 会话中正常继续对话。保持 QUEUED 与 SKIPPED 运行不可打开——因为它们尚未或永远不会创建 Craft 会话。关键事实执行器在完成前就已创建会话设计文档Important Notes第 3 条指出定时任务执行器scheduled-task executor已经会先创建BuildSession、写入初始用户 prompt、把ScheduledTaskRun.session_id关联到该会话并在驱动 Agent 循环之前提交事务。也就是说在运行完成之前就存在一个活的会话。这一点可在源码中得到印证执行器 executor.py 中_drive_agent通过session_manager.create_session(user_idtask_user_id, originSessionOrigin.SCHEDULED, namefScheduled: {task_name})创建会话随后create_message写入用户 promptturn 0再mark_run_status(..., statusRUNNING, session_idsession_id)将会话 ID 关联到运行行最后db_session.commit()——此时事务已提交运行处于RUNNING且已带session_id。会话视图 banner 使用的上下文接口 session/api.py 中ScheduledRunContextResponse已返回run_id、task_id、task_name、status、started_at、finished_at六个字段GET /sessions/{session_id}/scheduled-run-context返回 200 表示渲染 banner 并应用定时运行状态404 表示这是交互式会话正常行为。因此运行中打开会话在数据层面早已具备条件缺的只是前端点击规则、实时事件通道与输入禁用逻辑。三、总体架构Key Decisions关键决策设计文档的核心决策是采用 api-server 的 SSE 端点直接代理 opencode-serve 的/event流作为实时查看架构。给定一个定时运行的session_idapi-server 需要解析出该会话所属的沙箱 Podsandbox pod使用现有鉴权路径打开该 Pod 的/event流按会话 ID 过滤事件将匹配的 ACP 事件以 SSE 方式流式转发到浏览器。为什么选择直接代理 opencode/event文档给出了四条理由opencode-serve 本身就已经在发射实时事件流。定时运行与交互式路径的区别只在于由 Celery 驱动 prompt而不是底层实时数据源发生了变化。直接代理避免了重新发明轮子。单一事实来源实时进度的唯一权威来源就是沙箱 Pod 上的 opencode/event流api-server 只是以另一个查看者的身份挂接并按会话过滤。避免引入 Redis 作为新的实时传输层。api-server 已经可以直连沙箱 Pod只有在该网络假设api-server 到 Pod 可达预期会消失时才值得引入 Redis——否则它就是多余的额外基础设施。定时任务 worker 保持职责单一专注运行任务并持久化 durable 消息无需向另一个消息总线发布实时事件。数据库始终是恢复路径recovery path页面加载或 SSE 重连时会话视图先从已持久化的消息中水合hydrate再恢复直连 SSE 订阅以获取新事件直到运行进入终态。其他被否决的方案Other Options Considered方案优点被否决原因在 api-server 内复用PodEventBus每个 api-server 副本只需持有每个 Pod 一条上游/event连接本地扇出fan-out多查看者时更高效若并发查看成为常态是好的优化但首版直接代理更简单Redis 支撑的 SSE解耦 api-server 可达性与沙箱 Pod 网络适合跨集群 / serverless / 边缘路由的 api-server 部署只要 api-server 到 Pod 的可达性保持稳定Redis 就显得过重浏览器直连 opencode 隧道实现薄、简单/event是 Pod 级而非会话级作用域必须在到达浏览器前由服务端过滤最终又绕回 api-server 代理直接轮询数据库最容易构建延迟高、依赖执行器冲刷部分进度的频率且用户观看期间会造成重复读压力从仓库现状看会话事件订阅的基础能力已经存在SessionManager.subscribe_to_existing_session_events被 session/api.py 中已有的GET /sessions/{session_id}/scheduled-run-eventsSSE 端点使用该端点会等待BuildSession.opencode_session_id就绪通过SSE_KEEPALIVE心跳 POLL_INTERVAL_SECONDS轮询保持流存活然后订阅事件并逐块转发期间持续检查运行状态是否仍为RUNNING一旦进入非运行态即结束流。这正是设计文档第 4 条实现项实时定时会话事件路径在仓库中的雏形。四、实现方案Implementation设计文档给出了 7 条实现步骤下面逐条展开并结合仓库源码说明落点。1. 扩展会话视图使用的定时运行上下文Craft 会话视图需要的最小运行状态应新增关键字段运行状态run status让 UI 能区分RUNNING/AWAITING_APPROVAL与SUCCEEDED/FAILEDfinished_at用于展示也用于停止实时订阅run_id用于订阅或使能invalidate某一次确切的运行。对照仓库现状ScheduledRunContextResponse 已包含run_id、status、finished_at等字段前端类型定义 interfaces.ts 也已同步说明该上下文模型已就位本次改动主要是让会话视图消费这些字段来控制 banner、输入框与实时订阅的生命周期。2. 更新运行历史的可点击规则设计目标是让RUNNING、FAILED、SUCCEEDED、AWAITING_APPROVAL四种状态在拥有session_id时均可打开QUEUED行保持阻塞执行器尚未创建会话SKIPPED行保持阻塞永远不会创建会话。当前 utils.ts 的getNonClickableReason已经覆盖了有会话则可点、无会话给出原因的框架且 RunHistoryTable.tsx 的行点击处理为if (!getNonClickableReason(row, tReason) row.session_id) { router.push(buildSessionPath(row.session_id)); }——即点击跳转到会话视图。本次改动只需将RUNNING/AWAITING_APPROVAL从整体不可点中移出统一纳入有session_id即可点的判定同时保留不可点行上的锁定图标与 Tooltip 提示NonClickableCell组件逻辑不变。3. 更新定时运行 banner 与 Craft 聊天面板运行中显示定时运行 banner该会话由定时任务 X 于 Y 启动 ← 返回任务并禁用普通聊天输入框——因为 Agent 仍由后台执行器驱动此时发送消息会产生竞争运行进入终态后重新启用输入框用户可以在同一会话中继续追问追问的回复以正常 Craft 后续消息follow-up方式流式返回。前端工具函数 utils.ts 已提供isScheduledRunInFlight(status)与isScheduledRunContextInFlight(context)辅助判断分别基于RUNNING/AWAITING_APPROVAL状态可直接用于控制输入框禁用与订阅启停。4. 新增实时定时会话事件路径api-server SSE这是本次变更的核心新增端点。流程为校验会话所有权沿用现有会话与定时任务所有权检查路径解析沙箱 Pod打开 opencode-serve 的/event按会话 ID 过滤事件通过 SSE 将匹配事件流式转发到浏览器前端将事件合并进现有会话 store同时保持滚动行为运行进入终态后停止实时订阅。仓库中的 session/api.py 已经实现了几乎一致的形态get_session_scheduled_run_events先校验上下文存在且状态为RUNNING生成器内部先等待opencode_session_id就绪期间发心跳随后用subscribe_to_existing_session_events逐块 yield每块后expire_all()并复查运行状态非RUNNING即结束错误统一以_format_stream_error序列化为 SSEmessage事件返回响应头包含Cache-Control: no-cache, no-transform、Connection: keep-alive、X-Accel-Buffering: no媒体类型为text/event-stream。5. 保证执行器的持久化进度足够用于恢复运行期间要持续提交已持久化的工具进度tool progress、计划plans和最终确定的消息块finalized message chunks使页面刷新与 SSE 重连都能从 durable 状态水合。对照执行器源码 executor.pyBuildStreamingState(turn_index0)维护流式状态persist_sandbox_event逐事件持久化后紧跟db_session.commit()见循环内session_manager.persist_sandbox_event(session_id, state, sandbox_event); db_session.commit()终态前还会finalize_persist冲刷未落盘的 chunks——这与数据库是恢复路径的设计完全一致。此外执行器还实现了约 120 字符的摘要机制_clip_summary/_summary_from_state/_summary_from_session_messages运行结束或进入AWAITING_APPROVAL时写入summary字段供运行历史表格展示。6. 保持既有边界不变定时来源会话不进 Craft 侧边栏由BuildSession.origin SCHEDULED保证见 docs/craft/features/scheduled-tasks/overview.md 的SessionOrigin设计侧边栏查询过滤origin INTERACTIVE所有权检查继续沿用现有会话与定时任务的所有权路径所有后端错误统一使用OnyxError仓库中 session/api.py 与 scheduled_tasks/api.py 均遵循此约定例如get_session_scheduled_run_context在无上下文时抛出OnyxError(OnyxErrorCode.NOT_FOUND, ...)事件端点在状态不符时抛出OnyxError(OnyxErrorCode.CONFLICT, Scheduled run is not running)。7. 更新产品文档设计文档明确要求修改定时任务产品文档中运行中不可打开 / 等待完成的旧表述改为描述实时查看 运行结束后可继续追问的新产品契约使产品契约与实际行为保持一致。五、测试用例Test Cases设计文档给出 8 条验收测试可作为实现完成度的核对清单运行中可打开验证一个已关联会话的、仍在运行的定时运行可在其完成前从任务详情页的运行历史中打开。排队与跳过仍不可打开验证无会话的 QUEUED 运行与 SKIPPED 运行保持不可点击且有清晰的禁用提示锁定图标 Tooltip对应NonClickableCell组件。实时视图行为验证实时运行视图显示定时运行 banner、运行期间禁用普通聊天输入、无需刷新页面即可收到实时进度。订阅收尾验证运行进入终态后实时订阅停止/收敛且普通聊天输入恢复可用可发送追问消息。刷新恢复验证实时订阅结束后刷新页面加载并使用 Postgres 中已保存的消息而非 live opencode 事件流。追问流式返回验证运行结束后可发送新消息响应以正常 Craft 后续消息方式流式返回。侧边栏隔离验证定时来源会话仍不出现在正常 Craft 侧边栏历史中。静态检查与测试切片运行聚焦的前端类型检查frontend type check以及相关的定时任务后端/API 测试切片。仓库中的相关测试基础设施包括 test_executor.py、test_schedule.py 以及外部依赖单元测试 test_scheduled_task_executor.py执行器逻辑拆分为可导入、可测试的run_scheduled_task_logic无需 Celery worker 即可实例化前端侧则有 scheduledTaskRuns.test.ts 与 schedule.test.tsx 等。六、与现有定时任务体系的衔接本设计不是孤立的它建立在 V1 定时任务的既有体系之上理解以下几点有助于把握改动范围运行即会话ScheduledTaskRun.session_id是到build_session.id的外键SET NULL点击任意运行打开的就是现有会话视图无新 UIoverview.md 的数据模型一节。运行状态机queued → running → succeeded / failed / awaiting_approvaldispatcher 还会写skipped前一次运行仍在飞时。AWAITING_APPROVAL的恢复机制由 approvals 项目负责在它上线前该状态为展示而终态terminal-for-display。执行时序dispatcher 先写运行行再入队执行器消费队列后确保沙箱运行ensure_sandbox_running等待窗口由PROVISION_WAIT_SECONDS控制、置RUNNING、创建会话、关联session_id、提交——因此实时查看的窗口从session_id落库那一刻起就存在。超时与预算执行器有 30 分钟单调预算SCHEDULED_RUN_HARD_CAP_SECONDS与软预算SCHEDULED_RUN_SOFT_BUDGET_SECONDS超时以error_classTIMEOUT标记失败并通知用户运行期间persist_sandbox_event会逐条落库为实时查看的重连恢复提供数据基础。七、总结Live Scheduled Task Run Viewing是一次克制的架构决策不新增 Redis 传输、不新建查看器 UI、不改变执行器职责而是复用 opencode-serve 已有的/event实时流由 api-server 作为会话级过滤代理通过 SSE 转发给浏览器数据库继续充当刷新与重连的恢复路径。配合会话视图的 banner、运行态输入禁用/恢复与运行历史点击规则的调整最终达成运行中可观看、终态后可追问、侧边栏仍隔离的产品契约。对于希望二次开发或深度使用该能力的读者建议从 fix-live-scheduled-task-runs.md设计文档、session/api.pySSE 端点与上下文接口、executor.py执行器时序以及 utils.ts前端状态判定四个入口入手阅读。【免费下载链接】danswerOpen Source AI Platform - AI Chat with advanced features that works with every LLM项目地址: https://gitcode.com/GitHub_Trending/da/danswer创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考