Compare commits

...

4 Commits

34 changed files with 1297 additions and 146 deletions

View File

@ -8,7 +8,11 @@
## Unreleased
- 改进 Origin 多面板图排版:由 Origin 统一排列图层,共享横轴时仅在底行显示横轴标题和刻度标签,图例可自动避让数据;同时消除中文标题和坐标轴文字在导出图片中的异常横线。
- 专业软件任务完成后可按用户意图仅报告产物或自动读取图表和数据并给出分析Job 中心重新整理了状态、输入输出、耗时和执行详情,支持查看结果、复制 Job ID、深入分析及重新分析。
- 对话步骤进度改为按每轮任务保存完整计划;长任务、刷新或网络重连后可恢复当前步骤,不再因历史分页或首个实时事件错过而出现进度消失、串到上一轮或无法完成。正常完成后进度面板自动收起,等待确认、停止或异常时仍可查看停留步骤。
- 改进 Origin 多面板图排版:由 Origin 统一排列图层,共享横轴时仅在底行显示横轴标题和刻度标签,旧任务中的固定角落图例也会自动避让数据;同时消除 SVG 和 PDF 文本对象的异常横线。
- 修复 Origin 对数坐标图可能从 `1E-10` 开始、导致有效数据挤在图形右侧的问题;多面板未单独填写纵轴名称时,也会优先使用面板标题或数据列名,不再显示笼统的 `Y`

View File

@ -199,7 +199,7 @@ Admin GET /v1/admin/*(require_admin;overview + usage/models|users + storage/
Export GET /v1/tasks/{id}/export(docx)
```
**SSE 事件**:`run_start / llm_start / text{delta} / reasoning{delta}(thinking 模型推理流,前端灰色折叠卡)/ tool_call / tool_result(预览,完整走 DB)/ llm_end / model_switch / warn{msg}(熔断·重复拦截·折叠失败等运行时提醒)/ context_fold{phase,...}(§8.8 Phase 2 折叠 start/done)/ cancelled / error / done`。fan-out:每订阅独立 queue;迟到订阅立收 done。事件不持久化(messages 走 PG)。
**SSE 事件**:`run_start / llm_start / text{delta} / reasoning{delta}(thinking 模型推理流,前端灰色折叠卡)/ progress_snapshot{run_id,steps,waiting}(当前 user message 即 run 边界,从 messages 投影恢复)/ tool_call / tool_result(预览,完整走 DB)/ llm_end / model_switch / warn{msg}(熔断·重复拦截·折叠失败等运行时提醒)/ context_fold{phase,...}(§8.8 Phase 2 折叠 start/done)/ cancelled / error / done`。`task_progress` 每次提交完整步骤快照,前端整体替换;旧 `set_plan/update_step` 仅在历史投影时兼容。进度 dock 是运行状态而非历史消息:活跃时紧凑显示,正常完成后隐藏,`ask_user` 等待确认/取消/异常时折叠保留。fan-out:每订阅独立 queue;迟到订阅先从 PG 恢复当前 run 最新进度,终态迟到订阅立收 done。普通直播事件不持久化(messages 走 PG)。
**版本化**:`/v1` minor 半年兼容,major 6 个月 deprecation。**CORS**:本地 `*`,部署收紧。
### 7.3 认证
@ -234,7 +234,7 @@ scheduled_jobs(§8.5) channel_bindings(§8.7,判别列+JSONB)
- working_dir 存相对 ROOT posix 串,读写统一过 `core/paths.py`;入口 `validate_task_name` 拒空/`/\NUL`/`.` 起头。
- `auto_title_pending`(0023)只是一轮 UI 命名闸,不是 task 状态机;旧创建入口/存量行恒 false快速入口首发后消费人工改名优先清闸。
- **0004 简化**:runs 表只写不读、run_id 单活 run 下全冗余 → 合并 `run_status/run_error` 入 tasks。**0006**:`tasks.model_profile` 为 source-of-truth(PATCH 切、下条 send 生效);usage_events 重建 v2 多态形态,统计 source-of-truth;tasks 三列保留作粗概览。run_status 终态:ok 收回 idle,error(出错)与 cancelled(用户停止)是持久终态 —— 前端 `renderPersistedRunTerminal` 据此在每次重渲后补持久卡(扛过收尾 loadMessages 整屏重建),刷新/切任务仍在;下次起新 run(post_message 写 running)覆盖清掉。
- **0004 简化**:runs 表只写不读、独立 run 实体在单活形态下冗余 → 合并 `run_status/run_error` 入 tasks;需要前端关联本轮时直接复用该轮 user message UUID 为 `run_id`,不恢复 runs 表。**0006**:`tasks.model_profile` 为 source-of-truth(PATCH 切、下条 send 生效);usage_events 重建 v2 多态形态,统计 source-of-truth;tasks 三列保留作粗概览。run_status 终态:ok 收回 idle,error(出错)与 cancelled(用户停止)是持久终态 —— 前端 `renderPersistedRunTerminal` 据此在每次重渲后补持久卡(扛过收尾 loadMessages 整屏重建),刷新/切任务仍在;下次起新 run(post_message 写 running)覆盖清掉。
- **0029 消息序号**:`tasks.next_message_idx` 在 task 行锁下统一分配 `messages.idx`Web、agent 与渠道追加不再各自维护序号或依赖冲突重试;分配时仍与 `max(idx)` 校准,允许蓝绿发布窗口内旧实例继续写入。清空消息与计数器在同一事务归零。
- **No-subtask**:同 user 下前缀互含即拒(归一 posix 后 Python 端比对);同 working_dir 允许。
- **文件面板先备料**:user_root 与普通目录都可显式新建直接子目录;创建成功后前端进入该目录,用户可先上传/选入资料,再把顶层空目录选作新对话 working_dir。目录 leaf 复用 `validate_task_name`,不允许借 UI 创建点目录或路径式名称。
@ -484,6 +484,8 @@ Recipe 是专业软件的声明式目标状态,不是任意脚本或逐次鼠
第六阶段增加用户级 Job 中心与 Agent typed tools。`software_capability_list` 只暴露固定能力及当前在线空闲节点数,`software_job_submit/status/cancel` 在构造时绑定当前 user/task模型不能跨用户或跨对话指定归属。`software_job_submit` 的唯一入口为通用 `inputs[] + operation + outputs[]`;输入和输出使用任务内稳定 key具体 selector、type、format 和 options 由 capability 校验。`register_artifact` 是普通文件获得输入身份的唯一入口,并明确返回 UUID。右下角 Job 中心按用户聚合各对话任务,活动期短轮询、空闲期降频;终态变化通知用户,成功任务可回到原对话发起分析。取消采用协作协议:未派发任务直接终止,已派发任务先进入 `cancelling`,云端通过 WebSocket 发送并在心跳时重放 `job_cancel`Node 杀死固定 Worker 进程树后回报 `cancelled`;终态写入仍由云端账本裁决。
Job 完成后的对话闭环复用 `software_jobs` 账本,不增加通用事件表。提交时用 `completion_action=report|analyze` 区分“只报告产物”和“报告并分析”,未明确分析时默认 `report`;成功终态将 `followup_status` 置为 `pending`。分发器只在原 task 空闲时领取:`report` 直接写入带 artifact refs 的固定 assistant 消息,不产生模型费用;`analyze` 写入前端隐藏的内部完成事件,并复用 task 单活锁、持久消息和 SSE run 自动续跑 Agent。`pending/running/completed/failed` 状态使服务重启、对话繁忙和重复完成回执均不会造成重复回复Job 中心的手动深入分析通过同一入口重新排入分析,而不是模拟点击或只预填输入框。
后续仍需实现 Token 轮换;不得以任意命令或脚本接口临时代替。当前 Job 中心采用轮询而非用户事件推送,单活与 offer 选择仍只覆盖单 Web 进程;生产启用多实例前必须增加 Redis/PG fencing 或固定路由到单一控制面实例。
---

View File

@ -2,7 +2,7 @@
> 配合 `DESIGN.md`。本文件只记 phase 状态、决策偏差、文件量、下一步。每条 1-2 句:做了啥 + 关键判断;细节查 `git log` / `git diff` / `DESIGN §7.9`
最后更新:2026-08-17(Origin 多面板自动布局与导出基线修复完成,未发版)
最后更新:2026-08-17(专业软件 Job 完成自动报告/分析与结果中心优化完成,未发版)
---
@ -20,9 +20,11 @@
---
## 已完成关键能力
- **08-17 / Unreleased / Software Job 完成闭环与结果中心优化**`software_job_submit` 新增默认 `report`、可选 `analyze` 的完成策略0034 在原 Job 账本加入可恢复 follow-up 状态;成功任务可直接向原对话报告产物,或在 task 空闲后复用单活锁和 SSE 自动续跑分析,内部完成事件不冒充用户消息,手动深入/重新分析走同一幂等入口。Job 中心改为结果优先布局,显示短 Job ID、输入输出、完成方式和耗时折叠展示完整 ID/输出目录/执行版本,并支持复制、查看结果及状态化分析按钮。相关专项 unittest、Python/JavaScript 语法、Ruff 致命规则和 diff 检查通过;未连接生产 DB、未执行 migration。
### 2026-08-17
- **08-17 / Unreleased / Origin 多面板自动布局与导出基线修复**adapter 提升至 0.9.2,多面板图改用 Origin `layarrange` 根据行列、页边距和间距统一排列图层,并读取 Origin 计算后的实际图层位置放置标题与图例;共享 X 轴时仅底行保留标题和刻度标签,图例新增 `auto` 位置以调用 Origin 智能避让。每次任务显式设置 `@U=1`,消除 Origin 默认打印基线在中文标题、轴标题和图例上形成的异常横线,且不修改用户机器的持久配置
- **08-17 / Unreleased / Origin 多面板自动布局与导出基线修复**adapter 0.9.2 完成 `layarrange` 图层排列与共享 X 轴收敛,但真机复核发现旧请求仍固定右上图例,且 `@U=1` 不能消除 Origin 2024 导出的文本基线。adapter 0.9.3 因此在多面板中将 `auto` 及历史四角图例统一交给 Origin 智能避让,并从 SVG、PDF 中精确清理 Origin 额外生成的独立基线图元2×2、四系列、误差棒、中文标题的真机回归已确认两种矢量格式横线消失、字形完整。PNG 的 Origin 原生导出与剪贴板渲染均仍会绘制该基线,暂不采用会破坏汉字笔画的像素擦除方案
- **08-17 / Unreleased / Origin 对数坐标自动缩放修复**adapter 提升至 0.9.1,单图、多面板及右 Y 轴统一在 Origin 自动计算范围前应用坐标轴尺度,避免线性范围中的零值切换为对数轴后被强制展开到 `1E-10`;显式范围仍在自动缩放后覆盖,保持请求权威。多面板缺少 `y_axis` 时改用面板标题或数据列名回退,并在契约中提示对数轴提供正数范围及纵轴标题。相关 88 项 unittest、Python 编译、契约 JSON 和 diff 检查通过,未连接或写入生产数据库。

10
RUN.md
View File

@ -363,7 +363,7 @@ $env:ZCBOT_EVAL_TOKEN = "<dedicated-eval-user-jwt>"
| `GET/POST /v1/admin/external-system-definitions` | 管理员列出或新增可信外部系统目录 | admin |
| `PUT/DELETE /v1/admin/external-system-definitions/{id}` | 管理员编辑、停用或删除目录项;已有用户连接时拒绝删除 | admin |
| `GET /v1/tasks/{id}/messages` | LiteLLM payload 透传;另带 `artifact_refs`(助手产物)与 `attachment_refs`(用户附件)。两者均以 `null` 表示旧消息、`[]` 表示新消息明确为空、非空数组表示 task-relative 结构化引用 | 必填 |
| `POST /v1/tasks/{id}/messages` | `{content, attachments?:[{path,kind,label?}], image_model?=""}` 发消息;`attachments` 路径由后端按当前 task working_dir 校验并规范化,允许纯附件消息;旧客户端省略该字段继续兼容。返 `{events_url}`**`run_status` 是 running/cancelling → 409**UI 应 disable send 直到 SSE `done` | 必填 |
| `POST /v1/tasks/{id}/messages` | `{content, attachments?:[{path,kind,label?}], image_model?=""}` 发消息;`attachments` 路径由后端按当前 task working_dir 校验并规范化,允许纯附件消息;旧客户端省略该字段继续兼容。返 `{events_url,run_id}`,其中 `run_id` 是本轮 user message UUID**`run_status` 是 running/cancelling → 409**UI 应 disable send 直到 SSE `done` | 必填 |
| `GET /v1/tasks/{id}/events` | SSE 流(`event: <type>` + `data: <json>`);订阅 task 当前活动 | 必填 |
| `POST /v1/tasks/{id}/cancel` | 协作式 cancel;`run_status != running` → 409;LLM 走 streaming,chunk 间 poll cancel — 延迟 100ms 级,基本秒退 | 必填 |
| `GET /v1/procs` | 当前用户全部后台进程(bg proc,§8.12;shell/run_python `background=true` 启动);纯文件系统读取,前端运行条 5s 轮询用 | 必填 |
@ -384,7 +384,7 @@ $env:ZCBOT_EVAL_TOKEN = "<dedicated-eval-user-jwt>"
| `GET /v1/models` | 列 chat LLM 模型清单(扫 `config/models/*.yaml`),前端顶栏切换 / 新建对话框下拉用 | 必填 |
| `GET /v1/image_models` | 列图像生成 variant 清单(扫 `config/media/doubao.yaml` image 段),前端"生图"下拉用;yaml 无 image variant → 空列表 → UI 隐藏下拉 | 必填 |
**SSE 事件**(每帧 `event: <type>` + `data: <JSON>`):`run_start{}` → `llm_start{}``text{delta}` / `tool_call{name,args,args_preview}` / `tool_result{name,preview,truncated}``llm_end{prompt_tokens,completion_tokens}``done{}`;cancel 走 `cancelled{}` 后随 `done{}` 收流;异常走 `error{msg}`。30s 无 event 服务端发 `: ping` 心跳。nginx 反代记得关 buffering(响应头已带 `X-Accel-Buffering: no` 默认起效)。
**SSE 事件**(每帧 `event: <type>` + `data: <JSON>`):建连时若当前 run 已发布计划,先补 `progress_snapshot{run_id,steps,waiting}``run_start{}``llm_start{}``text{delta}` / `tool_call{name,args,args_preview}` / `tool_result{name,preview,truncated}``llm_end{prompt_tokens,completion_tokens}``done{}`;cancel 走 `cancelled{}` 后随 `done{}` 收流;异常走 `error{msg}`。`task_progress` 新协议每次携带完整 `steps`,客户端整体替换;消息分页响应也附加 `progress_snapshot`,刷新不依赖当前 30 条窗口。`waiting=true` 表示本轮已调用 `ask_user` 等待确认;正常完成回看时隐藏进度,等待/取消/异常则折叠保留。30s 无 event 服务端发 `: ping` 心跳。nginx 反代记得关 buffering(响应头已带 `X-Accel-Buffering: no` 默认起效)。
**SSE 客户端注意**:浏览器原生 `EventSource` 不支持自定义 header,无法塞 Bearer token。要么 `fetch + ReadableStream` 自解 SSE 帧(dev.html 走的就是这条),要么后端日后加 `?token=...` query(目前不支持,避免 token 进 access log)。
@ -1054,7 +1054,7 @@ sudo xfs_quota -x -c "limit -p bhard=10g zcbot_<user_uuid>" /opt
### Windows Node 内网 MVP开发中
先执行 `.venv/Scripts/python.exe main.py db upgrade head` 创建 `software_node_enrollments`、`software_nodes`、`software_jobs`为 artifact 增加专业软件来源字段。不要在未确认目标数据库时运行迁移;本机 `.env``ZCBOT_DB_URL` 可能是生产隧道。
先执行 `.venv/Scripts/python.exe main.py db upgrade head` 创建 `software_node_enrollments`、`software_nodes`、`software_jobs`,为 artifact 增加专业软件来源字段,并加入 Job 完成策略与续跑状态。不要在未确认目标数据库时运行迁移;本机 `.env``ZCBOT_DB_URL` 可能是生产隧道。
0032 旧版曾把 `tasks.working_dir` 二次拼进 user root已完成 Job 的文件和 artifact 路径可能位于重复的 `workspace/users/<uid>/...` 子树。升级到 0033 后,用显式迁移库地址先 dry-run确认所有目标目录无冲突后再加 `--apply`。脚本不加载 `.env`,文件和数据库更新均按 Job manifest 精确处理,目标目录同时存在时会停止,不能覆盖。
@ -1070,7 +1070,7 @@ $env:ZCBOT_MIGRATION_DB_URL="postgresql+psycopg://user:pass@host:5432/zcbot"
- Node `POST /v1/software-nodes/enroll` 注册并一次性取得 `node_id`、`node_token`
- Node 携带 `Authorization: Bearer <node_token>``X-Node-Id` 连接 `WS /v1/software-nodes/connect`
- 用户或 Agent 通过 `POST /v1/tasks/{task_id}/software-jobs` 提交专业软件任务;
- 用户通过 `GET /v1/software-jobs` 查看本人跨对话任务,可用 `task_id`、`active_only` 和 `limit` 筛选,`POST /v1/software-jobs/{job_id}/cancel` 请求停止;
- 用户通过 `GET /v1/software-jobs` 查看本人跨对话任务,可用 `task_id`、`active_only` 和 `limit` 筛选,`POST /v1/software-jobs/{job_id}/cancel` 请求停止`POST /v1/software-jobs/{job_id}/analyze` 对成功结果发起或重新发起分析
- 管理员 `GET /v1/admin/software-nodes` 查看节点,`PATCH /v1/admin/software-nodes/{node_id}` 启停节点,`DELETE /v1/admin/software-nodes/{node_id}` 永久删除节点身份。
Node API 只能绑定受控内网地址并由安全组限制来源 IP。当前 HTTP/WS 链路不加密;跨安全域、公网或不可信终端接入前,必须先升级 HTTPS/WSS。多 Web 实例部署时Node API 暂时固定路由到单一实例,直至 Connection Manager 增加跨实例 fencing。
@ -1100,7 +1100,7 @@ install-windows-node.bat
若 Python 未加入 PATH可把绝对路径作为第一个参数例如 `install-windows-node.bat "C:\Python312\python.exe"`。默认解释器为 `%ProgramData%\Zcbot\WindowsNode\runtimes\origin\Scripts\python.exe`。如需使用其他受管解释器,优先设置机器级 `ZCBOT_ADAPTER_ORIGIN_PYTHON`;旧名 `ZCBOT_ORIGIN_PYTHON` 暂时兼容。`node.json`、可恢复任务和 runtime 集中保存在 `%ProgramData%\Zcbot\WindowsNode\`,不会因替换程序目录而丢失。运行时固定依赖见发布目录的 `adapters/origin.plot@v2/requirements.txt`;任务请求无权选择解释器、脚本或路径。当前 Worker 支持 CSV/XLSX/JSON 输入及 OPJU/PNG/SVG/PDF 输出。成功产物由 Node 流式上传,全部校验通过后发布到任务工作目录 `origin/<job_id>/`plot spec 与 provenance 位于其 `.meta/`;上传中断会在重连时幂等续传。
Web 用户登录后,文件栏 Job 中心会聚合本人最近任务。活动任务约 4 秒刷新一次,空闲时降为约 30 秒停止已派发任务是协作取消状态先显示“正在停止”Node 在线时立即接收断线后在下次连接或心跳时重放。Agent 可调用 `software_capability_list`、`register_artifact`、`software_job_submit`、`software_job_status`、`software_job_revise` 和 `software_job_cancel`。Origin 输入必须是 artifact已有 UUID 可直接提交,普通 task 文件先逐个用相对路径登记;登记不会发布聊天交付卡片。提交工具接收 `inputs`、`operation`、`outputs`,支持 116 个输入、跨输入系列和多个显式输出;`plot.canvas` 可指定毫米画布,轴可指定范围、步长、尺度、刻度角度/字号、标题字号和网格,`legend` 可控制显隐、位置和字号,`series[].style` 可控制颜色、线宽/线型、点型/点大小和透明度。所有排版字段可选,旧请求保持默认样式。任务只创建固定 v2 schema 的持久任务,不会阻塞当前对话等待完成。成功状态提供 `output_dir`Agent可在该目录内搜索并分析正式输出的 artifact 带 `software_job_id`供结果卡和产物详情展示来源。Node 输出上传的逐任务诊断日志位于 `%ProgramData%\Zcbot\WindowsNode\jobs\<job-id>\logs\node-output-upload.log`;日志包含上传阶段、产物文件名、重试次数和 Windows `HRESULT`,单文件达到 1 MiB 后轮转一份 `.1`,不记录 Node Token 或认证请求头。
Web 用户登录后,文件栏 Job 中心会聚合本人最近任务。活动任务或完成后的自动分析约 4 秒刷新一次,空闲时降为约 30 秒;卡片优先显示状态、所属对话、输入输出、完成策略和耗时Job ID、输出目录、执行节点及软件版本收在“任务详情”中。成功任务可查看结果目录、复制完整 Job ID、深入分析或重新分析。停止已派发任务是协作取消状态先显示“正在停止”Node 在线时立即接收断线后在下次连接或心跳时重放。Agent 可调用 `software_capability_list`、`register_artifact`、`software_job_submit`、`software_job_status`、`software_job_revise` 和 `software_job_cancel`。Origin 输入必须是 artifact已有 UUID 可直接提交,普通 task 文件先逐个用相对路径登记;登记不会发布聊天交付卡片。提交工具接收 `inputs`、`operation`、`outputs` 和可选的 `completion_action=report|analyze`:只要求生成或保存时用默认 `report`任务成功后平台向原对话写入固定产物报告;明确要求解释、总结或结论时用 `analyze`,平台在原 task 空闲后自动续跑 Agent。两者均不阻塞提交轮次。请求支持 116 个输入、跨输入系列和多个显式输出;`plot.canvas` 可指定毫米画布,轴可指定范围、步长、尺度、刻度角度/字号、标题字号和网格,`legend` 可控制显隐、位置和字号,`series[].style` 可控制颜色、线宽/线型、点型/点大小和透明度。所有排版字段可选,旧请求保持默认样式。成功状态提供 `output_dir`,正式输出的 artifact 带 `software_job_id`供结果卡和产物详情展示来源。Node 输出上传的逐任务诊断日志位于 `%ProgramData%\Zcbot\WindowsNode\jobs\<job-id>\logs\node-output-upload.log`;日志包含上传阶段、产物文件名、重试次数和 Windows `HRESULT`,单文件达到 1 MiB 后轮转一份 `.1`,不记录 Node Token 或认证请求头。
二维复合图使用 `plot.type=multi_panel`。`layout` 支持 1×1、1×2、2×1 和 2×2`panels` 数量为 14两个 panel 必须选择横排或竖排,三个和四个 panel 使用 2×2。每个 panel 的 `series[]` 必须显式给出 `kind``line/scatter/line_scatter/column/bar`),可用 `y_axis=right` 绑定右 Y 轴,并用 `x_error`、`y_error` 指定对称误差列;顶层 `x_axis/y_axis/legend` 作为全部 panel 的默认值panel 内同名设置可单独覆盖,`right_y_axis` 仍按 panel 设置。一个请求最多仍为 16 条系列,且每个 panel 至少有一条左轴系列。`share_x/share_y` 会统一各主图层自动缩放后的范围。复合图不接受 Origin 模板名也不支持把等高线、3D 曲面、三元图或热图混入 panel。

View File

@ -92,6 +92,8 @@ def _job_dict(row: SoftwareJob) -> dict:
"execution_runtime": _job_execution_runtime(row),
"error": row.error,
"artifact_manifest": row.artifact_manifest,
"completion_action": getattr(row, "completion_action", "report"),
"followup_status": getattr(row, "followup_status", "none"),
"output_dir": f"{contract.output_namespace}/{row.job_id}",
"created_at": row.created_at.isoformat() if row.created_at else None,
"started_at": row.started_at.isoformat() if row.started_at else None,
@ -224,10 +226,13 @@ def create_job(
idempotency_key: str,
capability: str,
request: dict,
completion_action: str = "report",
) -> tuple[dict, bool]:
key = idempotency_key.strip()
if not key or len(key) > 200:
raise SoftwareJobError("idempotency_key must contain 1 to 200 characters")
if completion_action not in {"report", "analyze"}:
raise SoftwareJobError("completion_action must be report or analyze")
try:
contract = get_contract(capability)
normalized, digest = contract.normalize_request(request)
@ -296,6 +301,7 @@ def create_job(
existing.task_id != task_id
or existing.capability != capability
or existing.request_digest != digest
or getattr(existing, "completion_action", "report") != completion_action
):
raise SoftwareJobError("idempotency key was already used for a different request")
return _job_dict(existing), False
@ -313,6 +319,8 @@ def create_job(
metrics={},
error={},
artifact_manifest=[],
completion_action=completion_action,
followup_status="none",
)
try:
with session.begin_nested():
@ -330,6 +338,7 @@ def create_job(
existing.task_id != task_id
or existing.capability != capability
or existing.request_digest != digest
or getattr(existing, "completion_action", "report") != completion_action
):
raise SoftwareJobError(
"idempotency key was already used for a different request"
@ -381,6 +390,7 @@ def revise_job(
if source is None:
raise SoftwareJobError("source software job not found")
capability = source.capability
completion_action = getattr(source, "completion_action", "report")
source_request = source.request
inputs = deepcopy(source_request.get("inputs") or [])
schema_version = get_contract(capability).request_schema[
@ -391,6 +401,7 @@ def revise_job(
task_id,
idempotency_key=idempotency_key,
capability=capability,
completion_action=completion_action,
request={
"schema_version": schema_version,
"inputs": inputs,
@ -400,6 +411,25 @@ def revise_job(
)
def request_job_analysis(user_id: UUID, job_id: UUID) -> dict:
"""为已成功任务排入一次分析续跑pending/running 请求保持幂等。"""
with session_scope() as session:
job = session.execute(
select(SoftwareJob).where(
SoftwareJob.job_id == job_id,
SoftwareJob.user_id == user_id,
).with_for_update()
).scalar_one_or_none()
if job is None:
raise SoftwareJobError("job not found")
if job.status != "succeeded":
raise SoftwareJobError("only a succeeded job can be analyzed")
job.completion_action = "analyze"
if job.followup_status not in {"pending", "running"}:
job.followup_status = "pending"
return _job_dict(job)
def offer_next_job(node_ids: set[UUID]) -> dict | None:
"""选择最早可执行的 JobNode 组合,避免跨 capability 队首阻塞。"""
if not node_ids:
@ -827,6 +857,8 @@ def record_job_terminal(node_id: UUID, payload: dict) -> None:
job.error = error
job.artifact_manifest = manifest
job.terminal_at = now
if terminal_status == "succeeded":
job.followup_status = "pending"
def _is_uuid(value: str) -> bool:

View File

@ -478,6 +478,7 @@ class SoftwareJob(Base):
UniqueConstraint("user_id", "idempotency_key", name="uq_software_jobs_user_idempotency"),
Index("ix_software_jobs_status_created", "status", "created_at"),
Index("ix_software_jobs_node_status", "node_id", "status"),
Index("ix_software_jobs_followup", "followup_status", "terminal_at"),
)
job_id: Mapped[UUID] = mapped_column(PG_UUID(as_uuid=True), primary_key=True, default=uuid4)
@ -503,6 +504,14 @@ class SoftwareJob(Base):
metrics: Mapped[dict[str, Any]] = mapped_column(JSONB, nullable=False, default=dict)
error: Mapped[dict[str, Any]] = mapped_column(JSONB, nullable=False, default=dict)
artifact_manifest: Mapped[list[Any]] = mapped_column(JSONB, nullable=False, default=list)
# 0034:专业软件完成后的用户侧动作。report 只发布固定结果通知analyze 会在
# 原 task 空闲后自动续跑 Agent。followup_status 同时充当可恢复的轻量 outbox。
completion_action: Mapped[str] = mapped_column(
Text, nullable=False, default="report", server_default="report"
)
followup_status: Mapped[str] = mapped_column(
Text, nullable=False, default="none", server_default="none"
)
started_at: Mapped[Optional[datetime]] = mapped_column(DateTime(timezone=True), nullable=True)
terminal_at: Mapped[Optional[datetime]] = mapped_column(DateTime(timezone=True), nullable=True)
created_at: Mapped[datetime] = mapped_column(

View File

@ -0,0 +1,41 @@
"""Add durable completion actions for professional software jobs.
Revision ID: 0034
Revises: 0033
Create Date: 2026-08-17
"""
from collections.abc import Sequence
import sqlalchemy as sa
from alembic import op
revision: str = "0034"
down_revision: str | None = "0033"
branch_labels: str | Sequence[str] | None = None
depends_on: str | Sequence[str] | None = None
def upgrade() -> None:
op.add_column(
"software_jobs",
sa.Column(
"completion_action", sa.Text(), nullable=False, server_default="report"
),
)
op.add_column(
"software_jobs",
sa.Column(
"followup_status", sa.Text(), nullable=False, server_default="none"
),
)
op.create_index(
"ix_software_jobs_followup",
"software_jobs",
["followup_status", "terminal_at"],
)
def downgrade() -> None:
op.drop_index("ix_software_jobs_followup", table_name="software_jobs")
op.drop_column("software_jobs", "followup_status")
op.drop_column("software_jobs", "completion_action")

View File

@ -199,7 +199,7 @@
"enabled": {"type": "boolean"},
"position": {
"enum": ["auto", "top_left", "top_right", "bottom_left", "bottom_right"],
"description": "图例位置;auto 由 Origin 根据当前数据和图层空间智能定位。"
"description": "图例位置;单图支持固定四角,多面板由 Origin 自动避开数据(包括兼容旧请求中的四角值)。"
},
"font_size": {"type": "number", "minimum": 6, "maximum": 72}
}
@ -525,7 +525,7 @@
"enabled": {"type": "boolean"},
"position": {
"enum": ["auto", "top_left", "top_right", "bottom_left", "bottom_right"],
"description": "图例位置;auto 由 Origin 根据当前面板数据智能定位。"
"description": "图例位置;单图支持固定四角,多面板由 Origin 自动避开数据(包括兼容旧请求中的四角值)。"
},
"font_size": {"type": "number", "minimum": 6, "maximum": 72}
}

View File

@ -4,6 +4,7 @@ import test from "node:test";
import {
applyProgressAction,
enforceMonotonicProgress,
progressPanelMode,
progressActionsFromToolCalls,
} from "../web/static/js/progress.js";
@ -51,6 +52,31 @@ test("tool calls can apply progress updates on top of previous task progress", (
]);
});
test("a full snapshot replaces the prior run plan instead of merging it", () => {
const previous = [
{ id: "old", title: "上一轮计划", status: "completed" },
];
const updated = applyProgressAction(previous, {
steps: [
{ id: "s1", title: "分析需求", status: "completed" },
{ id: "s2", title: "执行修改", status: "in_progress" },
],
});
assert.deepEqual(updated, [
{ id: "s1", title: "分析需求", status: "completed" },
{ id: "s2", title: "执行修改", status: "in_progress" },
]);
});
test("progress panel is live-only after successful completion", () => {
assert.equal(progressPanelMode({ hasLiveRun: true, runStatus: "running" }), "active");
assert.equal(progressPanelMode({ runStatus: "idle" }), "hidden");
assert.equal(progressPanelMode({ runStatus: "idle", waiting: true }), "waiting");
assert.equal(progressPanelMode({ runStatus: "cancelled" }), "cancelled");
assert.equal(progressPanelMode({ runStatus: "error" }), "error");
});
test("a completed step force-completes earlier dangling steps (monotonic heal)", () => {
const steps = [
{ id: "s1", title: "摄取素材", status: "in_progress" },

View File

@ -4,6 +4,7 @@ import importlib.util
import json
import tempfile
import unittest
import zlib
from pathlib import Path
from unittest.mock import patch
@ -921,10 +922,11 @@ class OriginWorkerUnitTests(unittest.TestCase):
"Strength",
{"font_size": 12},
)
self.assertEqual(layer.legend.values["attach"], 1)
self.assertEqual(layer.legend.values["attach"], 0)
self.assertEqual(layer.legend.values["background"], 0)
self.assertEqual(layer.legend.values["left"], 655)
self.assertEqual(layer.legend.values["top"], 307)
self.assertEqual(layer.legend.values["smartpos"], 1)
self.assertNotIn("left", layer.legend.values)
self.assertNotIn("top", layer.legend.values)
self.assertEqual(layer.title.values["attach"], 1)
self.assertEqual(layer.title.values["background"], 0)
self.assertEqual(layer.title.values["fsize"], 12)
@ -943,6 +945,42 @@ class OriginWorkerUnitTests(unittest.TestCase):
worker._configure_origin_session(origin)
self.assertEqual(origin.variables, [("@U", 1)])
def test_svg_normalizer_removes_only_text_baseline_path(self) -> None:
source = (
b'<svg viewBox="0,0 100,100" width="100" height="100">'
b'<text x="10" y="30" font-size="20">Title</text>\n'
b'<path d="M 10,14.4 L 70,14.4" '
b'style="stroke: black;stroke-width: 2;stroke-linecap:butt" '
b'fill="none"/>'
b'<path d="M 0,90 L 100,90" '
b'style="stroke: black;stroke-width: 2" fill="none"/>'
b'</svg>'
)
normalized, baselines = worker._strip_origin_svg_text_baselines(source)
self.assertEqual(len(baselines), 1)
self.assertNotIn(b'M 10,14.4 L 70,14.4', normalized)
self.assertIn(b'M 0,90 L 100,90', normalized)
def test_pdf_normalizer_preserves_stream_size_and_removes_baseline(self) -> None:
decoded = b"4 w\n0 J\n10 20 m\n90 20 l\nS\nQ\n"
compressed = zlib.compress(decoded)
original = b"%PDF-1.2\nstream\n" + compressed + b"\nendstream\n%%EOF"
with tempfile.TemporaryDirectory() as temp_dir:
path = Path(temp_dir) / "figure.pdf"
path.write_bytes(original)
removed = worker._strip_origin_pdf_text_baselines(path)
content = path.read_bytes()
start = content.index(b"stream\n") + len(b"stream\n")
end = content.index(b"\nendstream", start)
normalized = zlib.decompress(content[start:end])
self.assertEqual(removed, 1)
self.assertNotIn(b"10 20 m", normalized)
self.assertEqual(len(content), len(original))
def test_origin_arranges_panel_layers_and_reports_resulting_geometry(self) -> None:
class Layer:
def __init__(self, geometry):

View File

@ -43,6 +43,7 @@ class ClaimRunTests(unittest.TestCase):
self.assertEqual(claim.metadata, {"marker": "ok"})
added = session.add.call_args.args[0]
self.assertEqual(claim.run_id, added.message_id)
self.assertEqual(added.task_id, tid)
self.assertEqual(added.idx, 7)
self.assertEqual(

View File

@ -0,0 +1,124 @@
from __future__ import annotations
import unittest
from contextlib import contextmanager
from types import SimpleNamespace
from unittest.mock import patch
from uuid import uuid4
from web.software_followups import claim_followup
class _Result:
def __init__(self, value=None):
self.value = value
def scalar_one_or_none(self):
return self.value
class _Session:
def __init__(self, results):
self.results = list(results)
self.added = []
self.statements = []
def execute(self, statement):
self.statements.append(statement)
return _Result(self.results.pop(0) if self.results else None)
def add(self, row):
self.added.append(row)
def _job(action: str):
job_id = uuid4()
task_id = uuid4()
return SimpleNamespace(
job_id=job_id,
task_id=task_id,
user_id=uuid4(),
capability="origin.plot@v2",
status="succeeded",
followup_status="pending",
completion_action=action,
artifact_manifest=[
{
"artifact_id": str(uuid4()),
"filename": "figure.png",
"path": f"origin/{job_id}/figure.png",
},
{
"artifact_id": None,
"filename": "plot-spec.json",
"path": f"origin/{job_id}/.meta/plot-spec.json",
},
],
)
class SoftwareFollowupTests(unittest.TestCase):
def test_report_claim_persists_fixed_assistant_message_with_artifacts(self):
job = _job("report")
task = SimpleNamespace(task_id=job.task_id, run_status="idle")
session = _Session([job.task_id, task, job])
@contextmanager
def scope():
yield session
with (
patch("web.software_followups.session_scope", scope),
patch("web.software_followups.allocate_message_idx", return_value=7),
):
claim = claim_followup(job.job_id)
self.assertEqual(claim.action, "report")
self.assertEqual(job.followup_status, "completed")
self.assertEqual(len(session.added), 1)
message = session.added[0]
self.assertEqual(message.payload["role"], "assistant")
self.assertIn(str(job.job_id), message.payload["content"])
self.assertEqual(len(message.artifact_refs), 1)
self.assertEqual(message.artifact_refs[0]["label"], "figure.png")
def test_analyze_claim_persists_internal_turn_and_locks_task(self):
job = _job("analyze")
task = SimpleNamespace(task_id=job.task_id, run_status="idle")
session = _Session([job.task_id, task, job, None])
@contextmanager
def scope():
yield session
with (
patch("web.software_followups.session_scope", scope),
patch("web.software_followups.allocate_message_idx", return_value=8),
):
claim = claim_followup(job.job_id)
self.assertEqual(claim.action, "analyze")
self.assertEqual(job.followup_status, "running")
self.assertEqual(session.added[0].payload["role"], "user")
self.assertIn("software_job_status", session.added[0].payload["content"])
self.assertGreaterEqual(len(session.statements), 4)
def test_busy_task_leaves_followup_pending(self):
job = _job("analyze")
task = SimpleNamespace(task_id=job.task_id, run_status="running")
session = _Session([job.task_id, task])
@contextmanager
def scope():
yield session
with patch("web.software_followups.session_scope", scope):
claim = claim_followup(job.job_id)
self.assertIsNone(claim)
self.assertEqual(job.followup_status, "pending")
self.assertFalse(session.added)
if __name__ == "__main__":
unittest.main()

View File

@ -61,6 +61,7 @@ class SoftwareJobToolTests(unittest.TestCase):
self.assertTrue(result["created"])
self.assertEqual(create.call_args.args[:2], (self.user_id, self.task_id))
self.assertEqual(create.call_args.kwargs["capability"], "origin.plot@v2")
self.assertEqual(create.call_args.kwargs["completion_action"], "report")
self.assertEqual(
create.call_args.kwargs["request"],
{
@ -96,6 +97,25 @@ class SoftwareJobToolTests(unittest.TestCase):
self.assertIn("outputs", required)
self.assertNotIn("input_path", SoftwareJobSubmitTool.parameters["properties"])
self.assertNotIn("request", SoftwareJobSubmitTool.parameters["properties"])
self.assertEqual(
SoftwareJobSubmitTool.parameters["properties"]["completion_action"]["enum"],
["report", "analyze"],
)
def test_submit_can_request_automatic_analysis(self):
artifact_id = uuid4()
with patch(
"tools.software_jobs.create_job",
return_value=({"job_id": str(uuid4())}, True),
) as create:
SoftwareJobSubmitTool(self.user_id, self.task_id).execute(
"origin.plot@v2",
inputs=[{"key": "sample", "artifact_id": str(artifact_id)}],
operation={"plot": {"type": "line", "series": []}},
outputs=[{"key": "figure_png", "type": "figure", "format": "png"}],
completion_action="analyze",
)
self.assertEqual(create.call_args.kwargs["completion_action"], "analyze")
def test_status_and_cancel_reject_cross_task_job(self):
foreign = {"job_id": str(uuid4()), "task_id": str(uuid4())}
@ -201,6 +221,7 @@ class SoftwareJobToolTests(unittest.TestCase):
)
self.assertEqual(result, (created, True))
self.assertEqual(create.call_args.kwargs["capability"], "origin.plot@v2")
self.assertEqual(create.call_args.kwargs["completion_action"], "report")
self.assertEqual(create.call_args.kwargs["request"], {
"schema_version": 2,
"inputs": source.request["inputs"],

View File

@ -155,6 +155,25 @@ class SoftwareNodeMigrationTests(unittest.TestCase):
self.assertIn("ix_artifacts_software_job_id", rendered)
self.assertIn("jsonb_array_elements", rendered)
def test_0034_adds_software_job_followup_state(self) -> None:
statements: list[str] = []
def capture(sql, *multiparams, **params):
statements.append(str(sql.compile(dialect=postgresql.dialect())))
engine = create_mock_engine("postgresql+psycopg://", capture)
operations = Operations(MigrationContext.configure(engine.connect()))
migration = importlib.import_module(
"db.migrations.versions.20260817_1000_0034_software_job_followups"
)
with patch.object(migration, "op", operations):
migration.upgrade()
rendered = "\n".join(statements)
self.assertIn("completion_action", rendered)
self.assertIn("followup_status", rendered)
self.assertIn("ix_software_jobs_followup", rendered)
class SoftwareJobProtocolTests(unittest.TestCase):
def test_succeeded_replay_uses_persisted_manifest_across_layout_versions(self) -> None:
@ -738,6 +757,41 @@ class SoftwareJobProtocolTests(unittest.TestCase):
})
self.assertEqual(job.status, "failed")
@patch("core.software_jobs._published_output_is_valid", return_value=True)
@patch("core.software_jobs.validate_output_manifest")
@patch("core.software_jobs.session_scope")
def test_successful_terminal_queues_completion_followup(
self, session_scope, validate_manifest, _published
) -> None:
session = session_scope.return_value.__enter__.return_value
node_id = uuid4()
lease_id = uuid4()
digest = "d" * 64
manifest = [{"source_artifact_id": "figure_png"}]
validate_manifest.return_value = manifest
job = type("Job", (), {})()
job.job_id = uuid4()
job.node_id = node_id
job.lease_id = lease_id
job.request_digest = digest
job.status = "running"
job.capability = "origin.plot@v2"
job.request = {}
job.progress = 80
session.execute.return_value.scalar_one_or_none.return_value = job
record_job_terminal(node_id, {
"job_id": str(uuid4()),
"lease_id": str(lease_id),
"request_digest": digest,
"status": "succeeded",
"error": {},
"artifact_manifest": manifest,
})
self.assertEqual(job.status, "succeeded")
self.assertEqual(job.followup_status, "pending")
@patch("core.software_jobs.session_scope")
def test_disconnect_does_not_requeue_active_jobs(self, session_scope) -> None:
session = session_scope.return_value.__enter__.return_value

View File

@ -43,16 +43,19 @@ class StaticVendorTests(unittest.TestCase):
self.assertIn("before_job_id", source)
self.assertIn("onListScroll", source)
self.assertIn("/cancel`", source)
self.assertIn("分析结果", source)
self.assertIn("refreshCurrentTaskFiles(job.task_id)", source)
self.assertIn("深入分析", source)
self.assertIn("查看结果", source)
self.assertIn("Job ID", source)
self.assertIn("refreshCurrentTaskAfterSoftwareJob(job.task_id)", source)
self.assertIn("job.execution_runtime", source)
self.assertIn("执行版本未记录", source)
self.assertIn("版本:${escapeHtml(jobRuntimeText(job))}", source)
self.assertIn("执行版本</dt><dd>${escapeHtml(jobRuntimeText(job))}", source)
self.assertIn(
'document.addEventListener("software-job-submitted", refreshSoftwareJobs)',
source,
)
self.assertIn("export function refreshCurrentTaskFiles(taskId)", chat_source)
self.assertIn("export async function openSoftwareJobResults", chat_source)
self.assertIn(
'document.dispatchEvent(new Event("software-job-submitted"))',
chat_source,

View File

@ -0,0 +1,65 @@
from __future__ import annotations
import json
import unittest
from web.task_progress import progress_waiting_for_user, project_progress_payloads
def _call(args: dict) -> dict:
return {
"role": "assistant",
"tool_calls": [{
"function": {
"name": "task_progress",
"arguments": json.dumps(args, ensure_ascii=False),
},
}],
}
class TaskProgressProjectionTests(unittest.TestCase):
def test_latest_full_snapshot_replaces_prior_snapshot(self) -> None:
steps, seen = project_progress_payloads([
_call({"steps": [
{"id": "old", "title": "旧计划", "status": "in_progress"},
]}),
_call({"steps": [
{"id": "s1", "title": "分析", "status": "completed"},
{"id": "s2", "title": "实现", "status": "in_progress"},
]}),
])
self.assertTrue(seen)
self.assertEqual([step["id"] for step in steps], ["s1", "s2"])
def test_legacy_updates_remain_replayable(self) -> None:
steps, seen = project_progress_payloads([
_call({"action": "set_plan", "steps": [
{"id": "s1", "title": "分析", "status": "in_progress"},
{"id": "s2", "title": "实现", "status": "pending"},
]}),
_call({"action": "update_step", "step": {
"id": "s2", "status": "completed",
}}),
])
self.assertTrue(seen)
self.assertEqual([step["status"] for step in steps], ["completed", "completed"])
def test_ask_user_marks_the_run_as_waiting(self) -> None:
payloads = [
_call({"steps": [
{"id": "s1", "title": "确认方案", "status": "in_progress"},
]}),
{
"role": "assistant",
"tool_calls": [{"function": {"name": "ask_user", "arguments": "{}"}}],
},
]
self.assertTrue(progress_waiting_for_user(payloads))
if __name__ == "__main__":
unittest.main()

View File

@ -5,17 +5,18 @@ from tools.task_progress import TaskProgressTool
class TaskProgressToolTests(unittest.TestCase):
def test_schema_exposes_set_update_and_clear_actions(self) -> None:
def test_schema_requires_complete_steps_snapshot(self) -> None:
schema = TaskProgressTool().schema
fn = schema["function"]
self.assertEqual(fn["name"], "task_progress")
action_enum = fn["parameters"]["properties"]["action"]["enum"]
self.assertEqual(action_enum, ["set_plan", "update_step", "clear"])
params = fn["parameters"]
self.assertEqual(params["required"], ["steps"])
self.assertNotIn("action", params["properties"])
self.assertNotIn("step", params["properties"])
def test_execute_returns_short_ui_only_result(self) -> None:
out = TaskProgressTool().execute(
action="set_plan",
steps=[
{"id": "s1", "title": "理解需求", "status": "completed"},
{"id": "s2", "title": "实现功能", "status": "in_progress"},
@ -23,19 +24,9 @@ class TaskProgressToolTests(unittest.TestCase):
)
data = json.loads(out)
self.assertEqual(data, {"ok": True, "action": "set_plan", "step_count": 2})
self.assertEqual(data, {"ok": True, "step_count": 2})
self.assertLess(len(out), 80)
def test_execute_normalizes_update_step_without_echoing_title(self) -> None:
out = TaskProgressTool().execute(
action="update_step",
step={"id": "s2", "title": "实现功能", "status": "completed"},
)
data = json.loads(out)
self.assertEqual(data, {"ok": True, "action": "update_step", "step_id": "s2"})
self.assertNotIn("实现功能", out)
if __name__ == "__main__":
unittest.main()

View File

@ -80,6 +80,15 @@ class SoftwareJobSubmitTool(_SoftwareJobTool):
"inputs": _contract_property_schema("inputs"),
"operation": _contract_property_schema("operation"),
"outputs": _contract_property_schema("outputs"),
"completion_action": {
"type": "string",
"enum": ["report", "analyze"],
"default": "report",
"description": (
"Use report to notify with output files only. Use analyze when the "
"user also asked for interpretation, conclusions, or data analysis."
),
},
"idempotency_key": {
"type": "string",
"description": (
@ -106,6 +115,7 @@ class SoftwareJobSubmitTool(_SoftwareJobTool):
inputs: list[dict] | None = None,
operation: dict | None = None,
outputs: list[dict] | None = None,
completion_action: str = "report",
idempotency_key: str = "",
) -> str:
try:
@ -134,6 +144,7 @@ class SoftwareJobSubmitTool(_SoftwareJobTool):
self.task_id,
idempotency_key=idempotency_key.strip() or str(uuid4()),
capability=capability,
completion_action=completion_action,
request=normalized_request,
)
return json.dumps({**job, "created": created}, ensure_ascii=False)

View File

@ -7,7 +7,7 @@ arguments for Web rendering and is compacted out of older LLM context.
from __future__ import annotations
import json
from typing import Any
from typing import Any, ClassVar
from .base import Tool
@ -15,24 +15,23 @@ from .base import Tool
class TaskProgressTool(Tool):
name = "task_progress"
description = (
"Publish or update a concise user-visible progress checklist for the current task. "
"Use only for meaningful multi-step work: set the plan once, update a step when it "
"starts or completes, and when all work is done mark the final step completed (do NOT "
"clear). Use clear only when the plan is no longer relevant. This is a UI progress "
"signal, not a work product."
"Publish the complete current user-visible progress checklist for this run. "
"Use only for meaningful multi-step work. Every call must include the full checklist, "
"including unchanged steps; never send a partial step patch. Keep stable step ids while "
"revising the plan, allow at most one in_progress step, and mark every step completed "
"before a successful final answer. This is a UI progress signal, not a work product."
)
parameters = {
parameters: ClassVar[dict[str, Any]] = {
"type": "object",
"additionalProperties": False,
"properties": {
"action": {
"explanation": {
"type": "string",
"enum": ["set_plan", "update_step", "clear"],
"description": "set_plan replaces the checklist; update_step changes one step; clear removes it.",
"description": "Optional short reason when the plan materially changes.",
},
"steps": {
"type": "array",
"description": "Required for set_plan. Keep to 3-7 user-meaningful steps.",
"description": "The complete current checklist. Keep to 3-7 user-meaningful steps.",
"items": {
"type": "object",
"additionalProperties": False,
@ -47,32 +46,14 @@ class TaskProgressTool(Tool):
"required": ["id", "title", "status"],
},
},
"step": {
"type": "object",
"description": "Required for update_step.",
"additionalProperties": False,
"properties": {
"id": {"type": "string"},
"title": {"type": "string"},
"status": {
"type": "string",
"enum": ["pending", "in_progress", "completed"],
},
},
"required": ["id", "status"],
},
},
"required": ["action"],
"required": ["steps"],
}
def execute(self, **kwargs: Any) -> str:
action = str(kwargs.get("action") or "")
out: dict[str, Any] = {"ok": True, "action": action}
if action == "set_plan":
steps = kwargs.get("steps")
out["step_count"] = len(steps) if isinstance(steps, list) else 0
elif action == "update_step":
step = kwargs.get("step")
if isinstance(step, dict) and step.get("id"):
out["step_id"] = str(step["id"])
steps = kwargs.get("steps")
out: dict[str, Any] = {
"ok": True,
"step_count": len(steps) if isinstance(steps, list) else 0,
}
return json.dumps(out, ensure_ascii=False, separators=(",", ":"))

View File

@ -49,7 +49,6 @@ from .background import (
from .broker import broker
from .routers.asr import register_asr_routes
from .routers.authroutes import register_auth_routes
from .routers.software_nodes import register_software_node_routes
from .routers.external_systems import register_external_system_routes
from .routers.files import register_file_routes
from .routers.kb import register_kb_routes
@ -58,9 +57,14 @@ from .routers.misc import register_misc_routes
from .routers.models import register_model_routes
from .routers.schedules import register_schedule_routes
from .routers.skills_memory import register_skill_memory_routes
from .routers.software_nodes import register_software_node_routes
from .routers.tasks import register_task_routes
from .routers.wechat import register_wechat_routes
from .scheduler_runner import start_scheduler
from .software_followups import (
requeue_stale_analysis_followups,
start_software_followup_dispatcher,
)
from .static_files import NoCacheStaticFiles
from .wechat_runner import start_wechat_inbound
@ -113,6 +117,9 @@ def create_app() -> FastAPI:
# 启动钩子 + 后台协程群(本体见 web/background.py 等;None=该项未启用)
reap_stale_runs()
requeued_followups = requeue_stale_analysis_followups()
if requeued_followups:
print(f"[startup] requeued {requeued_followups} software followup(s)")
disk_scanner_task = start_disk_scanner(_cfg)
stats_logger_task = start_stats_logger(app, run_max_workers)
toolfail_task = start_toolfail_scanner()
@ -120,12 +127,14 @@ def create_app() -> FastAPI:
wechat_task, wechat_stop = start_wechat_inbound(app)
sandbox_reaper_task = init_sandbox(app, _cfg)
proc_sweeper_task = start_proc_sweeper(_cfg)
software_followup_task = start_software_followup_dispatcher(app)
try:
yield
finally:
# 先拒新 run + drain in-flight(细节见 background.drain_inflight)
app.state.draining.set()
await cancel_and_wait(software_followup_task)
await drain_inflight(app, drain_timeout, cancel_grace)
await cancel_and_wait(disk_scanner_task)

View File

@ -42,6 +42,7 @@ from ..run_lifecycle import (
schedule_claimed_run,
)
from ..schemas import MessageRequest, OptimizePromptRequest
from ..task_progress import latest_run_progress
from ..userfiles import load_user_root
@ -101,6 +102,7 @@ def register_message_routes(app, *, require_user) -> None:
Message.tokens_out, Message.model_profile, Message.created_at,
Message.artifact_refs,
Message.attachment_refs,
Message.kind,
)
if limit is None:
# 旧行为:升序全量
@ -137,9 +139,11 @@ def register_message_routes(app, *, require_user) -> None:
.where(Message.task_id == tid, Message.idx > last_idx)
.limit(1)
).first() is not None
progress_snapshot = latest_run_progress(s, tid)
return {
"has_more": has_more,
"has_more_after": has_more_after,
"progress_snapshot": progress_snapshot,
"messages": [
{
"idx": r.idx,
@ -150,6 +154,7 @@ def register_message_routes(app, *, require_user) -> None:
"created_at": iso(r.created_at),
"artifact_refs": r.artifact_refs,
"attachment_refs": r.attachment_refs,
"kind": r.kind,
}
for r in rows
]
@ -174,6 +179,7 @@ def register_message_routes(app, *, require_user) -> None:
.where(
Message.task_id == tid,
Message.payload["role"].astext == "user",
Message.kind.is_distinct_from("software_job_followup"),
)
.order_by(Message.idx)
).all()
@ -320,7 +326,10 @@ def register_message_routes(app, *, require_user) -> None:
))
app.state.aux_tasks.add(title_task)
title_task.add_done_callback(app.state.aux_tasks.discard)
return {"events_url": f"/v1/tasks/{tid}/events"}
return {
"events_url": f"/v1/tasks/{tid}/events",
"run_id": str(claim.run_id),
}
@app.post("/v1/tasks/{task_id}/cancel", status_code=202, tags=["tasks"])
def cancel_task(
@ -638,6 +647,12 @@ def register_message_routes(app, *, require_user) -> None:
return
q = broker.subscribe(tid)
try:
# Durable catch-up for events emitted before this subscriber existed.
# Subscribing first keeps updates emitted during the DB read queued.
with session_scope() as s:
snapshot = latest_run_progress(s, tid)
if snapshot is not None:
yield sse_event("progress_snapshot", snapshot)
while True:
try:
ev = await asyncio.wait_for(q.get(), timeout=30.0)

View File

@ -22,7 +22,11 @@ from fastapi.responses import FileResponse
from core.artifact_lifecycle import register_published_artifacts
from core.paths import from_db_path
from core.software_contracts import SoftwareContractError, default_capabilities, get_contract
from core.software_contracts import (
SoftwareContractError,
default_capabilities,
get_contract,
)
from core.software_jobs import (
MAX_OUTPUT_ARTIFACT_BYTES,
MAX_OUTPUT_TOTAL_BYTES,
@ -38,6 +42,7 @@ from core.software_jobs import (
pending_node_cancellations,
record_job_terminal,
replay_succeeded_outputs,
request_job_analysis,
request_job_cancel,
respond_to_offer,
software_job_output_path,
@ -62,6 +67,7 @@ from web.schemas import (
SoftwareNodeDisableRequest,
SoftwareNodeEnrollRequest,
)
from web.software_followups import dispatch_followup
from web.userfiles import load_user_root, safe_join
@ -461,6 +467,7 @@ def register_software_node_routes(app, *, require_user, require_admin) -> None:
"artifact_manifest": published,
}
await asyncio.to_thread(record_job_terminal, node_id, terminal)
await dispatch_followup(request.app, job_id)
except (SoftwareJobError, KeyError, TypeError) as exc:
raise HTTPException(409, str(exc)) from exc
return {"status": "succeeded", "artifact_manifest": published}
@ -654,6 +661,20 @@ def register_software_node_routes(app, *, require_user, require_admin) -> None:
)
return job
@app.post("/v1/software-jobs/{job_id}/analyze", tags=["software-jobs"])
async def analyze_software_job(
job_id: UUID,
request: Request,
user_id: UUID = Depends(require_user), # noqa: B008
):
try:
await asyncio.to_thread(request_job_analysis, user_id, job_id)
except SoftwareJobError as exc:
detail = str(exc)
raise HTTPException(404 if detail == "job not found" else 409, detail) from exc
await dispatch_followup(request.app, job_id)
return get_job(user_id, job_id)
@app.get("/v1/software-jobs/{job_id}", tags=["software-jobs"])
def read_software_job(
job_id: UUID,

View File

@ -8,7 +8,7 @@ from __future__ import annotations
import asyncio
from dataclasses import dataclass, field
from typing import Any, Callable, Optional
from uuid import UUID
from uuid import UUID, uuid4
from sqlalchemy import select, update
@ -44,6 +44,7 @@ PrepareClaim = Callable[[Any, Task], tuple[dict[str, Any], dict[str, Any]]]
class RunClaim:
"""抢占成功后返回给调用方的领域元数据。"""
run_id: UUID
metadata: dict[str, Any] = field(default_factory=dict)
@ -79,7 +80,9 @@ def claim_run_with_message(
# task 行锁串行化同一 task 的 idx 分配,无需额外序列表或 advisory lock。
next_idx = allocate_message_idx(s, task_id, locked_task=task)
run_id = uuid4()
s.add(Message(
message_id=run_id,
task_id=task_id,
idx=int(next_idx),
payload={"role": "user", "content": user_message},
@ -92,7 +95,7 @@ def claim_run_with_message(
**extra_values,
}
s.execute(update(Task).where(Task.task_id == task_id).values(**values))
return RunClaim(metadata=dict(metadata))
return RunClaim(run_id=run_id, metadata=dict(metadata))
def schedule_claimed_run(

View File

@ -107,7 +107,6 @@ def run_agent_bg(
})
except Exception as e:
err = f"{type(e).__name__}: {e}"
broker.emit(task_id, {"type": "error", "msg": err})
mp = ""
try:
with session_scope() as s:
@ -120,7 +119,10 @@ def run_agent_bg(
)
)
except Exception:
pass # 已 emit error 给前端,DB 写失败不放大噪声
pass # DB 写失败不阻断后续可见 error 事件
# error 是 SSE 终止事件,客户端收到后会立即回读 task meta先落 DB 终态再 emit
# 避免它抢读到旧 running导致持久错误卡和中断进度面板偶发不显示。
broker.emit(task_id, {"type": "error", "msg": err})
# 留痕 + 告警(0.58.21,反 2a1bc25d 教训:Zai 余额不足连挂 3 次续跑无人知):
# run_error 列只留最后一次,usage_events(kind=run_error)才是聚合面板/巡检
# 邮件的完整数据源;余额/认证类 provider 级错误另走即时邮件(6h 签名冷却)。

View File

@ -1,7 +1,7 @@
"""/v1 请求体 Pydantic 模型(从 app.py 析出,2026-07-23 拆分)。"""
from __future__ import annotations
from typing import Optional
from typing import Literal, Optional
from uuid import UUID
from pydantic import BaseModel, Field
@ -141,3 +141,4 @@ class SoftwareJobCreateRequest(BaseModel):
idempotency_key: str
capability: str = Field(default_factory=lambda: default_capabilities()[0])
request: dict = Field(default_factory=dict)
completion_action: Literal["report", "analyze"] = "report"

239
web/software_followups.py Normal file
View File

@ -0,0 +1,239 @@
"""Software Job 完成后的固定报告与 Agent 自动续跑。"""
from __future__ import annotations
import asyncio
from dataclasses import dataclass
from uuid import UUID
from sqlalchemy import select, update
from core.software_contracts import get_contract
from core.storage import session_scope
from core.storage.message_index import allocate_message_idx
from core.storage.models import Message, SoftwareJob, Task
from .common import INSTANCE
from .run_lifecycle import RunScheduleError, schedule_claimed_run
@dataclass(frozen=True)
class FollowupClaim:
job_id: UUID
task_id: UUID
user_id: UUID
action: str
prompt: str = ""
def _artifact_refs(manifest: list) -> list[dict]:
refs: list[dict] = []
for item in manifest:
if not isinstance(item, dict) or not item.get("artifact_id") or not item.get("path"):
continue
refs.append({
"path": item["path"],
"label": item.get("filename") or item["path"].rsplit("/", 1)[-1],
"artifact_id": item["artifact_id"],
"version": 2,
})
return refs
def _report_text(job: SoftwareJob) -> str:
contract = get_contract(job.capability)
output_dir = f"{contract.output_namespace}/{job.job_id}"
names = [
str(item.get("filename"))
for item in job.artifact_manifest
if isinstance(item, dict) and item.get("artifact_id") and item.get("filename")
]
files = "".join(names) if names else "无可发布文件"
return (
f"专业软件任务已完成:{contract.display_name}\n\n"
f"- Job ID`{job.job_id}`\n"
f"- 输出目录:`{output_dir}`\n"
f"- 结果文件:{files}"
)
def _analysis_prompt(job: SoftwareJob) -> str:
contract = get_contract(job.capability)
output_dir = f"{contract.output_namespace}/{job.job_id}"
return (
"[专业软件任务完成事件]\n"
f"任务 {job.job_id}{contract.display_name})已成功完成,"
f"输出目录为 {output_dir}。请调用 software_job_status 获取完整产物清单,"
"读取输出图表和输入数据,向用户报告结果并分析主要趋势、结论及必要的限制。"
)
def pending_followup_ids(limit: int = 20) -> list[UUID]:
with session_scope() as session:
return list(session.execute(
select(SoftwareJob.job_id)
.where(
SoftwareJob.status == "succeeded",
SoftwareJob.followup_status == "pending",
)
.order_by(SoftwareJob.terminal_at, SoftwareJob.job_id)
.limit(limit)
).scalars())
def requeue_stale_analysis_followups() -> int:
"""启动 reaper 已确认 run 随进程丢失时,把自动分析退回可重试状态。"""
with session_scope() as session:
job_ids = list(session.execute(
select(SoftwareJob.job_id)
.join(Task, Task.task_id == SoftwareJob.task_id)
.where(
SoftwareJob.followup_status == "running",
Task.run_status == "error",
Task.run_error == "server restarted before run finished",
)
).scalars())
if job_ids:
session.execute(
update(SoftwareJob)
.where(SoftwareJob.job_id.in_(job_ids))
.values(followup_status="pending")
)
return len(job_ids)
def claim_followup(job_id: UUID) -> FollowupClaim | None:
"""task 空闲时原子领取回调report 在事务内直接落固定 assistant 消息。"""
with session_scope() as session:
candidate = session.execute(
select(SoftwareJob.task_id).where(SoftwareJob.job_id == job_id)
).scalar_one_or_none()
if candidate is None:
return None
task = session.execute(
select(Task).where(Task.task_id == candidate).with_for_update()
).scalar_one_or_none()
if task is None or task.run_status in {"running", "cancelling"}:
return None
job = session.execute(
select(SoftwareJob).where(SoftwareJob.job_id == job_id).with_for_update()
).scalar_one_or_none()
if (
job is None
or job.status != "succeeded"
or job.followup_status != "pending"
):
return None
next_idx = allocate_message_idx(session, task.task_id, locked_task=task)
if job.completion_action == "report":
session.add(Message(
task_id=task.task_id,
idx=next_idx,
payload={"role": "assistant", "content": _report_text(job)},
artifact_refs=_artifact_refs(job.artifact_manifest),
kind="software_job_report",
))
job.followup_status = "completed"
return FollowupClaim(
job_id=job.job_id,
task_id=task.task_id,
user_id=job.user_id,
action="report",
)
prompt = _analysis_prompt(job)
session.add(Message(
task_id=task.task_id,
idx=next_idx,
payload={"role": "user", "content": prompt},
kind="software_job_followup",
))
session.execute(update(Task).where(Task.task_id == task.task_id).values(
run_status="running",
run_error=None,
run_owner=INSTANCE or None,
))
job.followup_status = "running"
return FollowupClaim(
job_id=job.job_id,
task_id=task.task_id,
user_id=job.user_id,
action="analyze",
prompt=prompt,
)
def mark_followup_failed(job_id: UUID) -> None:
with session_scope() as session:
session.execute(
update(SoftwareJob)
.where(
SoftwareJob.job_id == job_id,
SoftwareJob.followup_status.in_({"pending", "running"}),
)
.values(followup_status="failed")
)
def finish_analysis_followup(job_id: UUID) -> None:
with session_scope() as session:
job = session.execute(
select(SoftwareJob).where(SoftwareJob.job_id == job_id).with_for_update()
).scalar_one_or_none()
if job is None or job.followup_status != "running":
return
task_status = session.execute(
select(Task.run_status).where(Task.task_id == job.task_id)
).scalar_one_or_none()
job.followup_status = (
"failed" if task_status in {"error", "cancelled"} else "completed"
)
def _track(app, task: asyncio.Task) -> None:
app.state.aux_tasks.add(task)
task.add_done_callback(app.state.aux_tasks.discard)
async def dispatch_followup(app, job_id: UUID) -> bool:
claim = await asyncio.to_thread(claim_followup, job_id)
if claim is None:
return False
if claim.action == "report":
return True
try:
run_task = schedule_claimed_run(
app,
claim.task_id,
claim.user_id,
claim.prompt,
scheduled=False,
)
except RunScheduleError:
await asyncio.to_thread(mark_followup_failed, claim.job_id)
return False
async def finish() -> None:
try:
await run_task
finally:
await asyncio.to_thread(finish_analysis_followup, claim.job_id)
_track(app, asyncio.create_task(finish()))
return True
async def dispatch_pending_followups(app) -> None:
for job_id in await asyncio.to_thread(pending_followup_ids):
if app.state.draining.is_set():
return
await dispatch_followup(app, job_id)
def start_software_followup_dispatcher(app) -> asyncio.Task:
async def loop() -> None:
while True:
await dispatch_pending_followups(app)
await asyncio.sleep(3)
return asyncio.create_task(loop(), name="software-followups")

View File

@ -1096,7 +1096,7 @@
#software-job-backdrop.open { opacity: 1; visibility: visible; }
#software-job-panel {
position: fixed; z-index: 130; top: 0; right: 0; bottom: 0;
width: min(420px, 92vw); display: flex; flex-direction: column; background: #fff;
width: min(480px, 94vw); display: flex; flex-direction: column; background: #fff;
box-shadow: -8px 0 28px rgba(0,0,0,.16); transform: translateX(105%);
visibility: hidden; transition: transform .2s ease, visibility .2s ease;
}
@ -1111,6 +1111,9 @@
.sj-title { display: flex; align-items: center; justify-content: space-between; gap: 10px; }
.sj-title-main { min-width: 0; display: flex; align-items: center; gap: 7px; }
.sj-title-main strong { overflow: hidden; text-overflow: ellipsis; white-space: nowrap; }
.sj-context { display: flex; align-items: center; justify-content: space-between; gap: 8px; margin-top: 5px; color: #555; font-size: 11px; }
.sj-context span { overflow: hidden; text-overflow: ellipsis; white-space: nowrap; }
.sj-context code { flex-shrink: 0; color: var(--muted); font-size: 10px; }
.sj-status { flex-shrink: 0; display: inline-flex; align-items: center; gap: 4px; padding: 2px 7px; border-radius: 999px; color: #6b7280; background: #f3f4f6; font-size: 11px; white-space: nowrap; }
.sj-status svg { width: 12px; height: 12px; }
.sj-card.running .sj-status, .sj-card.dispatched .sj-status, .sj-card.offered .sj-status, .sj-card.queued .sj-status { color: #9b4316; background: #fff3e7; }
@ -1125,11 +1128,28 @@
.sj-sub svg, .sj-meta svg { width: 12px; height: 12px; flex-shrink: 0; vertical-align: -2px; }
.sj-progress { height: 4px; margin-top: 7px; border-radius: 4px; background: var(--panel-muted); overflow: hidden; }
.sj-progress i { display: block; height: 100%; background: var(--accent); transition: width .25s ease; }
.sj-actions { display: flex; justify-content: flex-end; gap: 6px; margin-top: 9px; }
.sj-facts { display: grid; grid-template-columns: minmax(0, 1fr) minmax(0, 1fr); gap: 4px 12px; margin-top: 7px; padding: 7px 8px; border-radius: 6px; background: #f8f8f8; color: #62666d; font-size: 11px; }
.sj-facts span { min-width: 0; overflow-wrap: anywhere; }
.sj-details { margin-top: 6px; font-size: 11px; }
.sj-details > summary { width: fit-content; color: var(--muted); cursor: pointer; user-select: none; }
.sj-details dl { margin: 7px 0 0; padding: 7px 8px; border: 1px solid var(--border-soft); border-radius: 6px; background: #fcfcfc; }
.sj-details dl > div { display: grid; grid-template-columns: 72px minmax(0, 1fr); gap: 8px; padding: 3px 0; }
.sj-details dt { color: var(--muted); }
.sj-details dd { min-width: 0; margin: 0; overflow-wrap: anywhere; }
.sj-details code { font-size: 10px; }
.sj-copy { margin-left: 6px; padding: 1px 5px; border: 0; border-radius: 4px; background: transparent; color: var(--muted); cursor: pointer; font-size: 10px; }
.sj-copy:hover { color: var(--accent); background: var(--accent-soft); }
.sj-copy svg { width: 10px; height: 10px; vertical-align: -2px; }
.sj-actions { display: flex; flex-wrap: wrap; justify-content: flex-end; gap: 6px; margin-top: 9px; }
.sj-actions button { display: inline-flex; align-items: center; gap: 4px; }
.sj-actions button:disabled { cursor: default; opacity: .62; }
.sj-actions svg { width: 12px; height: 12px; }
.sj-empty { padding: 22px; text-align: center; color: var(--muted); }
.sj-loading { padding: 14px; text-align: center; color: var(--muted); font-size: 11px; }
@media (max-width: 520px) {
.sj-facts { grid-template-columns: 1fr; }
.sj-actions button { flex: 1 1 auto; justify-content: center; }
}
/* media tool 摘要 banner(model / size / cost / elapsed,折叠态也可见) */
.tool-banner {
display: inline-flex; flex-wrap: wrap; gap: 6px;

View File

@ -34,7 +34,12 @@ import {
import { loadFiles, scheduleFilesRefresh, uploadFiles, formatUploadProgress } from "./files.js";
import { toolActivityLabel, _workingDirName, extractMediaBanner, extractArtifactRels, renderArtifactBarHtml, renderArtifactChipContent, upgradeMediaArtifacts, ARTIFACT_PRODUCING_TOOLS, _flushMediaArtifactCache } from "./media.js";
import { parseUserAttachments } from "./attachments.js";
import { applyProgressAction, cloneProgressSteps, progressActionsFromToolCalls } from "./progress.js";
import {
applyProgressAction,
cloneProgressSteps,
progressActionsFromToolCalls,
progressPanelMode,
} from "./progress.js";
import { refreshProcs, decorateBgprocCard, hasRunningProc, killTaskProcs } from "./procs.js";
export async function loadModels() {
@ -1023,6 +1028,7 @@ function alignedEarlierLimit(firstIdx) {
async function loadMessages({ render = true } = {}) {
const data = await api("GET", `/v1/tasks/${state.taskId}/messages?limit=${MSG_PAGE}`);
state.loadedMessages = data.messages || [];
state.taskProgressSnapshot = data.progress_snapshot || null;
state.msgHasMore = !!data.has_more;
state.msgHasMoreNewer = !!data.has_more_after; // 尾部窗口通常为 false
state.msgLoadingEarlier = false;
@ -1432,7 +1438,7 @@ function renderLiveRunIfVisible() {
// card 已持有全部文字段/工具卡 DOM(切走再切回只需重新挂载,不重渲);
// 新建的重连 card 由 createLiveAssistantCard 自行渲染已累积文字。
const card = run.card || createLiveAssistantCard(run);
renderTaskProgressDock(run.progressSteps || []);
renderTaskProgressDock(run.progressSteps || [], currentProgressPanelMode(run.taskId));
if (card.parentElement !== wrap) wrap.appendChild(card);
wrap.scrollTop = wrap.scrollHeight;
setActionMode(run.cancelling ? "cancelling" : "streaming");
@ -1454,7 +1460,7 @@ function ensureRunningTaskSubscribed(taskId, url, seed = {}) {
curSeg: null,
cancelling: seed.run_status === "cancelling",
workingDir: seed.working_dir || "",
progressSteps: cloneProgressSteps(state.taskProgressByTask.get(taskId)),
progressSteps: [],
};
state.liveRuns.set(taskId, run);
state.streaming = true;
@ -1469,19 +1475,35 @@ function setRunHint(run, text) {
if (state.taskId === run.taskId) $("chat-hint").textContent = text;
}
// 进度只在对话区顶部的单一 dock 里渲染(codex 式钉顶面板),不再内联进每条消息卡。
// 进行中:展开实时显示 pending/in_progress/completed;全部完成:折叠成一行摘要,点开看清单。
function renderTaskProgressDock(steps) {
function currentProgressPanelMode(taskId) {
const run = getLiveRun(taskId);
const meta = state.taskId === taskId ? state.taskMeta : null;
const snapshot = state.taskId === taskId ? state.taskProgressSnapshot : null;
return progressPanelMode({
hasLiveRun: !!run,
cancelling: !!(run && run.cancelling),
runStatus: (meta && meta.run_status) || "",
waiting: !!(snapshot && snapshot.waiting),
});
}
// 进度是当前 run 的状态面板,不是历史消息:活跃时紧凑展示;正常完成后隐藏;
// 等待确认/取消/异常时折叠保留。用户手动展开后,同一状态内更新不强制折回。
function renderTaskProgressDock(steps, mode = "active") {
const dock = $("task-progress-dock");
if (!dock) return;
if (!Array.isArray(steps) || !steps.length) {
if (mode === "hidden" || !Array.isArray(steps) || !steps.length) {
dock.innerHTML = "";
delete dock.dataset.progressMode;
dock.classList.remove("show");
return;
}
const previousMode = dock.dataset.progressMode || "";
const wasOpen = previousMode === mode && !!dock.querySelector("details.task-progress")?.open;
const total = steps.length;
const done = steps.filter(s => s.status === "completed").length;
const allDone = done === total;
const current = steps.find(s => s.status === "in_progress") || steps.find(s => s.status === "pending");
const mark = (status) => status === "completed" ? "✓" : (status === "in_progress" ? "…" : "");
const rows = steps.map((s) => `
<div class="tp-step ${escapeHtml(s.status)}">
@ -1489,18 +1511,24 @@ function renderTaskProgressDock(steps) {
<span class="tp-text">${escapeHtml(s.title)}</span>
</div>
`).join("");
const summary = allDone
? `<summary class="tp-summary tp-done">✓ 全部完成 · ${done}/${total} 步</summary>`
: `<summary class="tp-summary">进度 · ${done}/${total} 步</summary>`;
const openAttr = allDone ? "" : " open"; // 全完成默认折叠,其余展开
let label = current ? `正在执行:${current.title}` : "执行进度";
if (allDone) label = "✓ 全部完成";
else if (mode === "waiting") label = "等待你确认";
else if (mode === "cancelled") label = "已停止";
else if (mode === "error") label = "执行中断";
else if (mode === "cancelling") label = "正在停止";
const summaryClass = allDone ? "tp-summary tp-done" : "tp-summary";
const summary = `<summary class="${summaryClass}">${escapeHtml(label)} · ${done}/${total} 步</summary>`;
const openAttr = wasOpen ? " open" : "";
dock.innerHTML = `<details class="task-progress"${openAttr}>${summary}<div class="tp-list">${rows}</div></details>`;
dock.dataset.progressMode = mode;
dock.classList.add("show");
}
function setTaskProgress(taskId, steps) {
const normalized = cloneProgressSteps(steps);
if (taskId) state.taskProgressByTask.set(taskId, normalized);
if (state.taskId === taskId) renderTaskProgressDock(normalized);
if (state.taskId === taskId) renderTaskProgressDock(normalized, currentProgressPanelMode(taskId));
}
const COPY_ICON = `<svg viewBox="0 0 24 24" aria-hidden="true"><rect x="8" y="8" width="11" height="11" rx="2"></rect><path d="M16 8V6a2 2 0 0 0-2-2H6a2 2 0 0 0-2 2v8a2 2 0 0 0 2 2h2"></path></svg>`;
@ -1653,6 +1681,7 @@ function renderMessages(msgs, { stickBottom = true } = {}) {
const m = msgs[mi];
const p = m.payload || {};
const role = p.role || "?";
if (m.kind === "software_job_followup") continue;
if (role === "system") continue; // 不显示 system
if (role === "assistant" && m.model_profile && m.model_profile !== lastAsstModel) {
const dn = (state.models.find(x => x.profile === m.model_profile) || {}).display_name || m.model_profile;
@ -1798,7 +1827,9 @@ function renderMessages(msgs, { stickBottom = true } = {}) {
}
if (stickBottom) wrap.scrollTop = wrap.scrollHeight;
setTaskProgress(state.taskId, currentProgressSteps);
const snapshotSteps = state.taskProgressSnapshot && Array.isArray(state.taskProgressSnapshot.steps)
? state.taskProgressSnapshot.steps : null;
setTaskProgress(state.taskId, snapshotSteps === null ? currentProgressSteps : snapshotSteps);
upgradeMediaArtifacts(wrap);
renderPersistedRunTerminal(); // 上次 run error/cancelled 终态 → 末尾补持久卡(所有重渲路径统一走这)
renderLiveRunIfVisible();
@ -1877,6 +1908,34 @@ export function refreshCurrentTaskFiles(taskId) {
if (taskId && state.taskId === taskId) scheduleFilesRefresh();
}
// Software Job 终态可能新增一条固定报告,也可能启动自动分析 run。刷新当前 task
// 并在后者场景主动接上 SSE用户无需切走再切回来才能看到回复。
export async function refreshCurrentTaskAfterSoftwareJob(taskId) {
if (!taskId || state.taskId !== taskId) return;
try {
const meta = await api("GET", "/v1/tasks/" + taskId);
if (state.taskId !== taskId) return;
state.taskMeta = meta;
renderChatMeta();
if (meta.run_status === "running" || meta.run_status === "cancelling") {
ensureRunningTaskSubscribed(taskId, `/v1/tasks/${taskId}/events`, meta);
} else {
await loadMessages();
refreshOutline();
}
scheduleFilesRefresh();
} catch (_) { /* Job 中心下一轮轮询或切换 task 时会恢复 */ }
}
export async function openSoftwareJobResults(taskId, outputDir) {
if (!taskId) return;
await selectTask(taskId);
const workingDir = state.taskMeta?.working_dir || "";
const base = workingDir.split("/").filter(Boolean).pop() || "";
state.filesPath = [base, outputDir].filter(Boolean).join("/");
await loadFiles();
}
function chatAction() {
if (isCurrentTaskStreaming()) { cancelCurrentTask(); return; }
if (hasRunningProc(state.taskId)) { killTaskProcs(state.taskId); return; }
@ -2741,8 +2800,10 @@ async function sendMessage(overrideText) {
cancelling: false,
workingDir: state.taskMeta && state.taskMeta.working_dir,
autoTitleEligible: !attachmentOnly,
progressSteps: cloneProgressSteps(state.taskProgressByTask.get(taskId)),
runId: r.run_id || "",
progressSteps: [],
};
setTaskProgress(taskId, []);
// 预建的空占位 .body 即首个文字段(首字到达前显示「思考中」)
run.curSeg = { el: asstCard.querySelector(".body"), acc: "", pending: false };
setRunPhase(run, "llm"); // POST 返回即起跳:覆盖 build_agent + 首轮 TTFT 的空窗
@ -3035,6 +3096,12 @@ function handleSseEvent(ev, asstCard, ctx) {
if (nearBottom) stream.scrollTop = stream.scrollHeight;
});
}
} else if (t === "progress_snapshot") {
const snapshot = ev.data || {};
if (ctx.runId && snapshot.run_id && ctx.runId !== snapshot.run_id) return;
if (!ctx.runId && snapshot.run_id) ctx.runId = snapshot.run_id;
ctx.progressSteps = cloneProgressSteps(snapshot.steps);
setTaskProgress(ctx.taskId, ctx.progressSteps);
} else if (t === "tool_call") {
const fn = (ev.data && ev.data.name) || "?";
const args = (ev.data && ev.data.args) || "";

View File

@ -9,6 +9,21 @@ export function normalizeProgressStatus(status) {
return ["pending", "in_progress", "completed"].includes(status) ? status : "pending";
}
export function progressPanelMode({
hasLiveRun = false,
cancelling = false,
runStatus = "",
waiting = false,
} = {}) {
if (hasLiveRun) return cancelling ? "cancelling" : "active";
if (runStatus === "running") return "active";
if (runStatus === "cancelling") return "cancelling";
if (runStatus === "cancelled") return "cancelled";
if (runStatus === "error") return "error";
if (waiting) return "waiting";
return "hidden";
}
export function normalizeProgressStep(step) {
if (!step || typeof step !== "object") return null;
const id = String(step.id || "").trim();
@ -39,7 +54,9 @@ export function applyProgressAction(progress, args) {
if (!args || typeof args !== "object") return current;
const action = args.action || "";
if (action === "clear") return [];
if (action === "set_plan") {
// Current protocol sends a complete snapshot on every call. Legacy set_plan
// has the same shape, so both formats intentionally converge here.
if (Array.isArray(args.steps)) {
const planned = Array.isArray(args.steps) ? args.steps.map(normalizeProgressStep).filter(Boolean) : [];
return enforceMonotonicProgress(planned);
}

View File

@ -3,7 +3,11 @@ import { api } from "./api.js";
import { state } from "./state.js";
import { $ } from "./dom.js";
import { escapeHtml } from "./format.js";
import { refreshCurrentTaskFiles, selectTask } from "./chat.js";
import {
openSoftwareJobResults,
refreshCurrentTaskAfterSoftwareJob,
selectTask,
} from "./chat.js";
import { dialogConfirm, message } from "./dialog.js";
const ACTIVE = new Set(["queued", "offered", "dispatched", "running", "disconnected", "cancelling"]);
@ -38,6 +42,7 @@ const icons = {
clock: '<svg viewBox="0 0 24 24" fill="none" stroke="currentColor" stroke-width="2" stroke-linecap="round" stroke-linejoin="round" aria-hidden="true"><circle cx="12" cy="12" r="9"/><path d="M12 7v5l3 2"/></svg>',
failed: '<svg viewBox="0 0 24 24" fill="none" stroke="currentColor" stroke-width="2" stroke-linecap="round" stroke-linejoin="round" aria-hidden="true"><circle cx="12" cy="12" r="9"/><path d="M12 8v5m0 3h.01"/></svg>',
open: '<svg viewBox="0 0 24 24" fill="none" stroke="currentColor" stroke-width="2" stroke-linecap="round" stroke-linejoin="round" aria-hidden="true"><path d="M14 5h5v5m0-5-8 8"/><path d="M19 13v5a1 1 0 0 1-1 1H6a1 1 0 0 1-1-1V6a1 1 0 0 1 1-1h5"/></svg>',
copy: '<svg viewBox="0 0 24 24" fill="none" stroke="currentColor" stroke-width="2" stroke-linecap="round" stroke-linejoin="round" aria-hidden="true"><rect x="8" y="8" width="11" height="11" rx="2"/><path d="M16 8V6a2 2 0 0 0-2-2H6a2 2 0 0 0-2 2v8a2 2 0 0 0 2 2h2"/></svg>',
pending: '<svg viewBox="0 0 24 24" fill="none" stroke="currentColor" stroke-width="2" stroke-linecap="round" stroke-linejoin="round" aria-hidden="true"><circle cx="12" cy="12" r="9"/><path d="M12 7v5l3 2"/></svg>',
running: '<svg viewBox="0 0 24 24" fill="none" stroke="currentColor" stroke-width="2" stroke-linecap="round" stroke-linejoin="round" aria-hidden="true"><path d="M21 12a9 9 0 1 1-3-6.7"/></svg>',
succeeded: '<svg viewBox="0 0 24 24" fill="none" stroke="currentColor" stroke-width="2" stroke-linecap="round" stroke-linejoin="round" aria-hidden="true"><circle cx="12" cy="12" r="9"/><path d="m8 12 3 3 5-6"/></svg>',
@ -81,8 +86,26 @@ export async function refreshSoftwareJobs() {
const next = data.results || [];
next.forEach((job) => {
const previous = known.get(job.job_id);
if (previous && ACTIVE.has(previous) && TERMINAL.has(job.status)) notifyTerminal(job);
known.set(job.job_id, job.status);
if (previous && ACTIVE.has(previous.status) && TERMINAL.has(job.status)) {
notifyTerminal(job);
} else if (
previous
&& previous.followup_status !== job.followup_status
&& ["running", "completed", "failed"].includes(job.followup_status)
) {
void refreshCurrentTaskAfterSoftwareJob(job.task_id);
}
if (
job.followup_status === "running"
&& state.taskId === job.task_id
&& !state.liveRuns.has(job.task_id)
) {
void refreshCurrentTaskAfterSoftwareJob(job.task_id);
}
known.set(job.job_id, {
status: job.status,
followup_status: job.followup_status,
});
});
const previousJobs = jobs;
const hadLoadedMore = previousJobs.length > PAGE_SIZE;
@ -113,7 +136,9 @@ async function loadMore() {
function schedule() {
if (timer) clearTimeout(timer);
const delay = jobs.some((job) => ACTIVE.has(job.status)) ? POLL_ACTIVE_MS : POLL_IDLE_MS;
const delay = jobs.some((job) => (
ACTIVE.has(job.status) || ["pending", "running"].includes(job.followup_status)
)) ? POLL_ACTIVE_MS : POLL_IDLE_MS;
timer = setTimeout(refreshSoftwareJobs, document.hidden ? Math.max(delay, 30000) : delay);
}
@ -178,28 +203,86 @@ function jobRuntimeText(job) {
: "执行版本未记录";
}
function fmtDate(value) {
return value ? new Date(value).toLocaleString("zh-CN", { hour12: false }) : "—";
}
function elapsedText(job) {
const start = job.started_at || job.created_at;
const end = job.terminal_at;
if (!start || !end) return "";
const seconds = Math.max(0, Math.round((new Date(end) - new Date(start)) / 1000));
if (seconds < 60) return `${seconds}`;
const minutes = Math.floor(seconds / 60);
return seconds % 60 ? `${minutes}${seconds % 60}` : `${minutes} 分钟`;
}
function followupLabel(job) {
if (job.completion_action === "analyze") {
if (job.followup_status === "pending") return "等待自动分析";
if (job.followup_status === "running") return "正在自动分析";
if (job.followup_status === "completed") return "已报告并分析";
if (job.followup_status === "failed") return "自动分析失败";
return "报告并分析";
}
return job.followup_status === "completed" ? "已报告结果" : "仅报告结果";
}
function jobActions(job, active) {
const task = escapeHtml(job.task_id);
if (active) return `
<button class="small" data-job-action="open" data-task-id="${task}">${icons.open}打开对话</button>
${job.status !== "cancelling" ? `<button class="small danger" data-job-action="cancel">${icons.cancel}停止</button>` : ""}`;
if (job.status !== "succeeded") {
return `<button class="small" data-job-action="open" data-task-id="${task}">${icons.open}${job.status === "failed" ? "查看错误" : "打开对话"}</button>`;
}
const view = `<button class="small primary" data-job-action="results" data-task-id="${task}">${icons.open}查看结果</button>`;
if (["pending", "running"].includes(job.followup_status) && job.completion_action === "analyze") {
const label = job.followup_status === "pending" ? "等待分析" : "分析中";
return `${view}<button class="small" disabled>${icons.analyze}${label}</button>`;
}
if (job.completion_action === "analyze" && job.followup_status === "completed") {
return `${view}<button class="small" data-job-action="open" data-task-id="${task}">${icons.analyze}查看分析</button><button class="small" data-job-action="analyze" data-task-id="${task}">重新分析</button>`;
}
const label = job.followup_status === "failed" ? "重试分析" : "深入分析";
return `${view}<button class="small" data-job-action="analyze" data-task-id="${task}">${icons.analyze}${label}</button>`;
}
function jobCard(job) {
const summary = job.request_summary || {};
const inputs = Array.isArray(job.input) ? job.input : [];
const inputNames = inputs.map((item) => item.filename).filter(Boolean).join("、");
const allInputNames = inputs.map((item) => item.filename).filter(Boolean);
const inputNames = allInputNames.slice(0, 2).join("、")
+ (allInputNames.length > 2 ? `${allInputNames.length} 个文件` : "");
const active = ACTIVE.has(job.status);
const progress = Math.max(0, Math.min(100, Number(job.progress || 0)));
const detail = stageLabel[job.stage] || statusLabel[job.status] || "正在处理";
const error = job.error && (job.error.detail || job.error.code);
const opened = job.created_at;
const openedText = opened ? new Date(opened).toLocaleString("zh-CN", { hour12: false }) : "";
const artifacts = Array.isArray(job.artifact_manifest)
? job.artifact_manifest.filter((item) => item && item.artifact_id) : [];
const formats = (summary.formats || []).filter(Boolean).join("、");
const elapsed = elapsedText(job);
const shortId = String(job.job_id || "").slice(0, 8);
return `<article class="sj-card ${escapeHtml(job.status)}" data-job-id="${escapeHtml(job.job_id)}">
<div class="sj-title"><div class="sj-title-main"><strong>${escapeHtml(summary.display_name || job.capability)}</strong></div>
<span class="sj-status">${statusIcon(job.status)}${escapeHtml(statusLabel[job.status] || job.status)}</span></div>
<div class="sj-context"><span>${escapeHtml(job.task_name || "未命名对话")}</span><code>Job ${escapeHtml(shortId)}</code></div>
<div class="sj-sub">${icons.activity}<span>${escapeHtml(detail)}${error ? ` · ${escapeHtml(error)}` : ""}</span></div>
${active ? `<div class="sj-progress"><i style="width:${progress}%"></i></div>` : ""}
<div class="sj-meta">${icons.clock}${escapeHtml(job.task_name || "未命名对话")}${inputNames ? ` · ${escapeHtml(inputNames)}` : ""}${openedText ? ` · 开启于 ${escapeHtml(openedText)}` : ""}</div>
<div class="sj-meta">版本${escapeHtml(jobRuntimeText(job))}</div>
<div class="sj-actions">
<button class="small" data-job-action="open" data-task-id="${escapeHtml(job.task_id)}">${icons.open}打开对话</button>
${job.status === "succeeded" ? `<button class="small primary" data-job-action="analyze" data-task-id="${escapeHtml(job.task_id)}">${icons.analyze}分析结果</button>` : ""}
${active && job.status !== "cancelling" ? `<button class="small danger" data-job-action="cancel">${icons.cancel}停止</button>` : ""}
<div class="sj-facts">
${inputNames ? `<span>输入:${escapeHtml(inputNames)}</span>` : ""}
${job.status === "succeeded" ? `<span>输出:${artifacts.length} 个文件${formats ? ` · ${escapeHtml(formats)}` : ""}</span>` : ""}
<span>完成方式${escapeHtml(followupLabel(job))}</span>
<span>${active ? `创建于 ${escapeHtml(fmtDate(job.created_at))}` : `${elapsed ? `耗时 ${escapeHtml(elapsed)} · ` : ""}${job.terminal_at ? `完成于 ${escapeHtml(fmtDate(job.terminal_at))}` : `创建于 ${escapeHtml(fmtDate(job.created_at))}`}`}</span>
</div>
<details class="sj-details"><summary>任务详情</summary><dl>
<div><dt>Job ID</dt><dd><code>${escapeHtml(job.job_id)}</code><button type="button" class="sj-copy" data-job-action="copy">${icons.copy}</button></dd></div>
<div><dt>Capability</dt><dd><code>${escapeHtml(job.capability)}</code></dd></div>
<div><dt>输出目录</dt><dd><code>${escapeHtml(job.output_dir || "")}</code></dd></div>
<div><dt>执行版本</dt><dd>${escapeHtml(jobRuntimeText(job))}</dd></div>
<div><dt>时间</dt><dd> ${escapeHtml(fmtDate(job.created_at))}<br> ${escapeHtml(fmtDate(job.started_at))}<br> ${escapeHtml(fmtDate(job.terminal_at))}</dd></div>
</dl></details>
<div class="sj-actions">${jobActions(job, active)}</div>
</article>`;
}
@ -208,6 +291,13 @@ async function handleAction(event, button) {
const card = button.closest("[data-job-id]");
const jobId = card.dataset.jobId;
const action = button.dataset.jobAction;
if (action === "copy") {
try {
await navigator.clipboard.writeText(jobId);
message("Job ID 已复制", "success");
} catch (_) { message("复制失败,请手动复制", "error"); }
return;
}
if (action === "cancel") {
if (!await dialogConfirm({
title: "停止专业软件任务",
@ -221,17 +311,28 @@ async function handleAction(event, button) {
return;
}
const taskId = button.dataset.taskId;
if (action === "results") {
await openSoftwareJobResults(taskId, jobs.find((item) => item.job_id === jobId)?.output_dir || "");
closeDrawer();
return;
}
if (action === "analyze") {
button.disabled = true;
try {
await api("POST", `/v1/software-jobs/${jobId}/analyze`);
if (taskId) await selectTask(taskId);
await refreshCurrentTaskAfterSoftwareJob(taskId);
message("已排队分析结果", "success");
closeDrawer();
refreshSoftwareJobs();
} catch (error) {
button.disabled = false;
message(error.message || "启动分析失败", "error");
}
return;
}
if (taskId) await selectTask(taskId);
closeDrawer();
if (action === "analyze") {
setTimeout(() => {
const input = $("chat-input");
if (!input) return;
input.value = `请分析专业软件任务 ${jobId} 的结果,结合输出图表和输入数据总结主要结论。`;
input.focus();
input.dispatchEvent(new Event("input", { bubbles: true }));
}, 400);
}
}
function notifyTerminal(job) {
@ -239,5 +340,5 @@ function notifyTerminal(job) {
const label = ok ? "已完成" : (job.status === "cancelled" ? "已取消" : "失败");
const summary = job.request_summary || {};
message(`${summary.display_name || "专业软件任务"}${label}`, ok ? "success" : "error", 6000);
refreshCurrentTaskFiles(job.task_id);
void refreshCurrentTaskAfterSoftwareJob(job.task_id);
}

View File

@ -55,6 +55,7 @@ export const state = {
streaming: false, // 兼容旧判断:任一 task 是否在流式中
liveRuns: new Map(), // task_id -> 当前浏览器会话内运行中的回复卡/累计文本
taskProgressByTask: new Map(), // task_id -> 历史消息重放后的当前进度步骤
taskProgressSnapshot: null, // 当前 task 最近一轮的分页外完整进度快照
// 消息分页(尾部窗口 + 向上滚动加载更早):切 task 默认只拉最近一批,
// 顶部 sentinel 进视口自动往前补。loadedMessages 是当前已加载的升序窗口,
// renderMessages 对它做纯函数渲染(时序累积逻辑无需改)。

141
web/task_progress.py Normal file
View File

@ -0,0 +1,141 @@
"""Project the latest run-scoped progress snapshot from append-only messages."""
from __future__ import annotations
import json
from collections.abc import Iterable
from typing import Any
from uuid import UUID
from sqlalchemy import select
from core.storage.models import Message
_VALID_STATUSES = {"pending", "in_progress", "completed"}
def _normalize_step(value: Any) -> dict[str, str] | None:
if not isinstance(value, dict):
return None
step_id = str(value.get("id") or "").strip()
title = str(value.get("title") or "").strip()
status = str(value.get("status") or "pending")
if not step_id or not title:
return None
if status not in _VALID_STATUSES:
status = "pending"
return {"id": step_id, "title": title, "status": status}
def _heal_monotonic(steps: list[dict[str, str]]) -> list[dict[str, str]]:
last_completed = max(
(i for i, step in enumerate(steps) if step["status"] == "completed"),
default=-1,
)
return [
{**step, "status": "completed"} if i < last_completed else dict(step)
for i, step in enumerate(steps)
]
def apply_progress_args(
current: list[dict[str, str]], args: Any,
) -> list[dict[str, str]]:
"""Apply current full snapshots plus legacy set/update/clear calls."""
if not isinstance(args, dict):
return [dict(step) for step in current]
action = args.get("action") or ""
if action == "clear":
return []
if isinstance(args.get("steps"), list):
normalized: list[dict[str, str]] = []
for raw in args["steps"]:
step = _normalize_step(raw)
if step is not None:
normalized.append(step)
return _heal_monotonic(normalized)
if action != "update_step" or not isinstance(args.get("step"), dict):
return [dict(step) for step in current]
raw = args["step"]
step_id = str(raw.get("id") or "").strip()
if not step_id:
return [dict(step) for step in current]
next_steps: list[dict[str, str]] = []
found = False
for step in current:
if step["id"] != step_id:
next_steps.append(dict(step))
continue
found = True
status = str(raw.get("status") or step["status"])
next_steps.append({
"id": step_id,
"title": str(raw.get("title") or step["title"]).strip(),
"status": status if status in _VALID_STATUSES else "pending",
})
if not found:
normalized = _normalize_step(raw)
if normalized is not None:
next_steps.append(normalized)
return _heal_monotonic(next_steps)
def project_progress_payloads(
payloads: Iterable[dict[str, Any]],
) -> tuple[list[dict[str, str]], bool]:
steps: list[dict[str, str]] = []
seen = False
for payload in payloads:
if not isinstance(payload, dict) or payload.get("role") != "assistant":
continue
for call in payload.get("tool_calls") or []:
function = call.get("function") if isinstance(call, dict) else None
if not isinstance(function, dict) or function.get("name") != "task_progress":
continue
raw_args = function.get("arguments") or "{}"
try:
args = json.loads(raw_args) if isinstance(raw_args, str) else raw_args
except (TypeError, json.JSONDecodeError):
args = {}
steps = apply_progress_args(steps, args)
seen = True
return steps, seen
def progress_waiting_for_user(payloads: Iterable[dict[str, Any]]) -> bool:
for payload in payloads:
if not isinstance(payload, dict) or payload.get("role") != "assistant":
continue
for call in payload.get("tool_calls") or []:
function = call.get("function") if isinstance(call, dict) else None
if isinstance(function, dict) and function.get("name") == "ask_user":
return True
return False
def latest_run_progress(session, task_id: UUID) -> dict[str, Any] | None:
"""Return the latest user message id and that run's projected plan."""
run = session.execute(
select(Message.message_id, Message.idx)
.where(
Message.task_id == task_id,
Message.payload["role"].astext == "user",
)
.order_by(Message.idx.desc())
.limit(1)
).first()
if run is None:
return None
payloads = session.execute(
select(Message.payload)
.where(Message.task_id == task_id, Message.idx > run.idx)
.order_by(Message.idx)
).scalars().all()
steps, seen = project_progress_payloads(payloads)
if not seen:
return None
return {
"run_id": str(run.message_id),
"steps": steps,
"waiting": progress_waiting_for_user(payloads),
}

View File

@ -1,6 +1,6 @@
{
"capability": "origin.plot@v2",
"adapter_version": "0.9.2",
"adapter_version": "0.9.3",
"runtime": "python",
"runtime_id": "origin",
"entrypoint": "worker.py",

View File

@ -14,6 +14,7 @@ import math
import os
import re
import sys
import zlib
from datetime import datetime, timezone
from importlib.metadata import PackageNotFoundError, version
from itertools import pairwise
@ -74,7 +75,7 @@ PANEL_LAYOUT = {
(2, 1): {"left": 13, "right": 13, "top": 14, "bottom": 12, "xgap": 0, "ygap": 10},
(2, 2): {"left": 9, "right": 4, "top": 14, "bottom": 12, "xgap": 11, "ygap": 10},
}
ADAPTER_VERSION = "0.9.2"
ADAPTER_VERSION = "0.9.3"
def _server_executable(command: str) -> Path:
@ -422,11 +423,117 @@ def _apply_canvas(graph: Any, canvas: Any) -> None:
def _configure_origin_session(op: Any) -> None:
"""Normalize export behavior that otherwise depends on machine defaults."""
# Origin's default @U=0 prints a baseline under axis/data labels. It is
# especially conspicuous with CJK fonts and can look like an underline.
# Keep labels flush with their frames. Origin 2024 still emits a separate
# baseline drawing primitive; exported files are normalized below.
op.set_lt_var("@U", 1)
_SVG_TEXT_BASELINE_RE = re.compile(
rb'(?P<text><text\b(?P<attrs>[^>]*)>.*?</text>\s*)'
rb'(?P<path><path\s+d="M\s*(?P<x1>-?[\d.]+),(?P<y1>-?[\d.]+)\s+'
rb'L\s*(?P<x2>-?[\d.]+),(?P<y2>-?[\d.]+)"\s+'
rb'style="stroke:\s*black;stroke-width:\s*(?P<width>[\d.]+);[^\"]*"'
rb'\s+fill="none"[^>]*/>)',
re.DOTALL,
)
def _strip_origin_svg_text_baselines(
content: bytes,
) -> tuple[bytes, list[tuple[float, float, float, float, float]]]:
"""Remove Origin 2024's separate printing-baseline paths after text nodes."""
baselines: list[tuple[float, float, float, float, float]] = []
def replace(match: re.Match[bytes]) -> bytes:
attrs = match.group("attrs")
x1, y1, x2, y2, width = (
float(match.group(name)) for name in ("x1", "y1", "x2", "y2", "width")
)
font_match = re.search(rb'font-size="([\d.]+)"', attrs)
if font_match is None:
return match.group(0)
font_size = float(font_match.group(1))
transform = re.search(
rb'transform="matrix\([^\"]+\s(-?[\d.]+)\s(-?[\d.]+)\)"', attrs
)
text_x = re.search(rb'\bx="(-?[\d.]+)"', attrs)
text_y = re.search(rb'\by="(-?[\d.]+)"', attrs)
horizontal = abs(y1 - y2) < 0.01
vertical = abs(x1 - x2) < 0.01
matches_text_frame = False
if horizontal and text_x and text_y:
expected_x = float(text_x.group(1))
expected_y = float(text_y.group(1)) - 0.78 * font_size
matches_text_frame = (
abs(x1 - expected_x) <= max(2, font_size * 0.08)
and abs(y1 - expected_y) <= max(3, font_size * 0.12)
)
elif vertical and transform:
translate_x, translate_y = map(float, transform.groups())
expected_x = translate_x - 0.78 * font_size
matches_text_frame = (
abs(x1 - expected_x) <= max(3, font_size * 0.12)
and abs(y1 - translate_y) <= max(3, font_size * 0.08)
)
if not matches_text_frame or max(abs(x2 - x1), abs(y2 - y1)) < 1.5 * font_size:
return match.group(0)
if not 0.04 <= width / font_size <= 0.14:
return match.group(0)
baselines.append((x1, y1, x2, y2, width))
return match.group("text")
return _SVG_TEXT_BASELINE_RE.sub(replace, content), baselines
_PDF_BASELINE_RE = re.compile(
rb'(?m)(?P<width>[\d.]+) w\s*\n0 J\s*\n'
rb'(?P<x1>-?[\d.]+) (?P<y1>-?[\d.]+) m\s*\n'
rb'(?P<x2>-?[\d.]+) (?P<y2>-?[\d.]+) l\s*\nS\s*\nQ'
)
def _strip_origin_pdf_text_baselines(path: Path) -> int:
"""Remove isolated long black text-baseline strokes without rewriting xrefs."""
content = bytearray(path.read_bytes())
removed = 0
for stream_match in reversed(list(re.finditer(rb'stream\r?\n', content))):
start = stream_match.end()
end = content.find(b"endstream", start)
if end < 0:
continue
compressed = bytes(content[start:end]).rstrip(b"\r\n")
try:
decoded = zlib.decompress(compressed)
except zlib.error:
continue
def replace(match: re.Match[bytes]) -> bytes:
nonlocal removed
width, x1, y1, x2, y2 = (
float(match.group(name))
for name in ("width", "x1", "y1", "x2", "y2")
)
length = max(abs(x2 - x1), abs(y2 - y1))
if min(abs(x2 - x1), abs(y2 - y1)) > 0.02 or length < 10 * width:
return match.group(0)
removed += 1
replacement = b"Q"
return replacement + b" " * (len(match.group(0)) - len(replacement))
normalized = _PDF_BASELINE_RE.sub(replace, decoded)
if normalized == decoded:
continue
recompressed = zlib.compress(normalized, 9)
if len(recompressed) > len(compressed):
raise RuntimeError("PDF_BASELINE_NORMALIZATION_FAILED")
content[start : start + len(compressed)] = recompressed + b"\n" * (
len(compressed) - len(recompressed)
)
if removed:
path.write_bytes(content)
return removed
def _arrange_panel_layers(
layers: list[Any], grid: tuple[int, int]
) -> list[tuple[float, float, float, float]]:
@ -549,18 +656,11 @@ def _apply_legend(
label.set_int("smartpos", 1)
elif "position" in legend:
if attach_to_layer:
x_fraction, y_fraction = {
"top_left": (0.03, 0.04),
"top_right": (0.66, 0.04),
"bottom_left": (0.03, 0.68),
"bottom_right": (0.66, 0.68),
}[legend["position"]]
if vertical_offset:
y_fraction += 0.22 if legend["position"].startswith("top") else -0.22
left, top = _panel_page_pixel(graph, geometry, x_fraction, y_fraction)
label.set_int("attach", 1)
label.set_int("left", round(left))
label.set_int("top", round(top))
# Fixed page coordinates disable Origin's collision avoidance.
# Delegate multi-panel placement even for historical requests
# that specified a corner before automatic layout existed.
label.set_int("attach", 0)
label.set_int("smartpos", 1)
else:
left, top = LEGEND_POSITIONS[legend["position"]]
if vertical_offset and legend["position"].startswith("bottom"):
@ -1198,20 +1298,29 @@ def run(job_dir: Path) -> list[dict[str, Any]]:
canvas_width_mm = (plot_spec.get("canvas") or {}).get("width_mm", 160)
pixel_width = round(dpi * canvas_width_mm / 25.4)
for extension in ("png", "svg", "pdf"):
if extension in formats:
target = output / f"figure.{extension}"
exported = Path(
graph.save_fig(
str(target),
type=extension,
width=pixel_width if extension == "png" else 0,
ratio=100 if extension in {"svg", "pdf"} else 0,
)
).resolve()
if exported != target.resolve() or not target.is_file():
raise RuntimeError(f"{extension.upper()}_EXPORT_FAILED")
_validate_artifact(target, extension)
artifacts.append(_manifest(target, media[extension]))
if extension not in formats:
continue
target = output / f"figure.{extension}"
_configure_origin_session(op)
exported = Path(
graph.save_fig(
str(target),
type=extension,
width=pixel_width if extension == "png" else 0,
ratio=0 if extension == "png" else 100,
)
).resolve()
if exported != target.resolve() or not target.is_file():
raise RuntimeError(f"{extension.upper()}_EXPORT_FAILED")
if extension == "svg":
normalized_svg, _ = _strip_origin_svg_text_baselines(
target.read_bytes()
)
target.write_bytes(normalized_svg)
elif extension == "pdf":
_strip_origin_pdf_text_baselines(target)
_validate_artifact(target, extension)
artifacts.append(_manifest(target, media[extension]))
try:
originpro_version = version("originpro")
except PackageNotFoundError: