From e521f832c4f9c793b75b7cc9c9f3e581b207ed9c Mon Sep 17 00:00:00 2001 From: caoqianming Date: Wed, 19 Aug 2026 10:55:19 +0800 Subject: [PATCH] feat(node): persist professional software workspaces --- CHANGELOG.md | 2 + DESIGN.md | 8 +- PROGRESS.md | 2 + RUN.md | 10 +- core/software_contracts.py | 46 ++- core/software_jobs.py | 339 +++++++++++++++++- core/software_nodes.py | 11 +- core/storage/models.py | 63 ++++ core/tool_registry.py | 2 + .../20260819_0900_0035_software_workspaces.py | 125 +++++++ ...ansys.mechanical.static_structural.v1.json | 1 + software-contracts/origin.plot.v2.json | 9 +- tests/test_software_contracts.py | 15 +- tests/test_software_job_tools.py | 27 ++ tests/test_software_nodes.py | 26 +- tests/test_software_output_publish.py | 37 +- tests/test_windows_node_source.py | 17 +- tools/software_jobs.py | 53 ++- web/routers/software_nodes.py | 317 +++++++++++++++- web/software_followups.py | 49 ++- web/static/js/media.js | 23 +- web/static/js/software_jobs.js | 16 +- windows-node/README.md | 4 +- .../Zcbot.WindowsNode/AdapterProcessRunner.cs | 31 +- .../Zcbot.WindowsNode/JobInboxStore.cs | 61 +++- .../Zcbot.WindowsNode/JobOutputUploader.cs | 114 +++++- .../Zcbot.WindowsNode/NodeAdapters.cs | 34 +- .../Zcbot.WindowsNode/NodeConnectionLoop.cs | 87 ++++- windows-node/Zcbot.WindowsNode/NodeModels.cs | 16 +- .../Zcbot.WindowsNode/WorkspaceStore.cs | 166 +++++++++ .../adapters/origin.plot@v2/adapter.json | 2 +- .../adapters/origin.plot@v2/worker.py | 34 +- 32 files changed, 1666 insertions(+), 81 deletions(-) create mode 100644 db/migrations/versions/20260819_0900_0035_software_workspaces.py create mode 100644 windows-node/Zcbot.WindowsNode/WorkspaceStore.cs diff --git a/CHANGELOG.md b/CHANGELOG.md index 68cf2c5..38b10b9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,6 +8,8 @@ ## Unreleased +- Origin 专业软件任务现在会在原 Windows Node 上保留可继续加工的工程;后续调整直接打开上一版工程,过程只提供轻量预览,工程和其他大文件仅在明确要求下载或交付时才传回并登记为正式产物。 + - 右侧工作区改为“文件 / 软件作业”双面板,标准页面和 embed 模式保持一致,折叠后仍可查看作业数量与完成提示;文件面板新增“当前对话 / 全部文件”切换、当前目录搜索、排序和修改时间展示。 - 专业软件任务提交后由系统统一发送完成通知和产物,Agent 不再自行等待并重复报告;Origin PNG 改由清理异常文字基线后的矢量页渲染,三维图在指定毫米画布时也会同步缩放模板图层,避免横线及坐标标题、色标和文字裁切。 diff --git a/DESIGN.md b/DESIGN.md index c3bad76..df723be 100644 --- a/DESIGN.md +++ b/DESIGN.md @@ -476,11 +476,15 @@ Node 通过 `Authorization: Bearer` 与 `X-Node-Id` 建立 `/v1/software-nodes/c Recipe 是专业软件的声明式目标状态,不是任意脚本或逐次鼠标命令。`origin.plot@v2` 以新增 `plot.type=recipe + recipe_version=1` 复用已验证的 `layout/panels/series/axis/legend` 组合结构;旧 `multi_panel` 与 Recipe 进入同一个受控执行器,单图和特殊统计图继续保留兼容入口。模型可以组合契约允许的原子图形能力,但不能提供 Python、LabTalk、解释器、模板名或文件路径。其他软件可采用同一通用 Job 外壳和各自的声明式 Recipe,不建设跨软件万能 DSL;API/SDK executor 优先,未来的 UI Automation 或 Computer Use 仅作为 adapter 内部执行后端,不改变云端 Recipe。 -产物修改第一阶段采用不可变重新生成:`software_job_status(job_id)` 对当前 user/task 返回规范化 `editable_request`,`software_job_revise(source_job_id, operation, outputs)` 在服务端复用源 Job 已登记的 inputs,并重新经过当前 capability 契约、artifact 权限和调度校验创建新 Job。它不打开旧 OPJU、不覆盖旧产物、不接受模型重新指定输入,也暂不增加版本树 migration;若以后出现必须保留用户在 GUI 内手工编辑的真实需求,再单独设计工程副本、稳定对象 ID 和增量执行。 +专业软件的连续加工以独立 `software_workspaces` 为稳定边界,Job 只记录一次操作。首个 Job 创建独立 `workspace_id`,两者不共用 ID;服务端保存 user 归属、capability、`home_node_id`、`head_job_id`、容量与留存元数据,Node 物理目录只使用 `%ProgramData%/Zcbot/WindowsNode/workspaces//`,不再按 user 建上层目录。这样用户统计集中在数据库完成,节点只管理实际落到本机的 workspace,也避免多节点各自维护一套用户目录树。Workspace 一旦首次调度即粘在 home node;后续 Job 必须以当前 head 为 source、串行回到同一节点,节点离线时等待而不复制工程或静默改派。 + +Node 对每个 Workspace 只保留 `current` 和 `rollback` 两代。续作开始前把 `current` 原子切换为 `rollback`,Worker 打开该工程并将新状态写入 Job 输出区;成功后输出区提升为新 `current`,失败或进程中断则恢复 `rollback`。因此连续加工不会为每个 Job 永久复制一份工程,物理空间上限约为当前工程加一份回滚工程;Job 账本仍完整保留参数、摘要和 local manifest。`software_job_revise(source_job_id, operation, outputs)` 只允许从 Workspace 当前 head 继续,复用已登记输入并重新经过当前契约校验。当前 Origin adapter 以 OPJU 为 workspace state;ANSYS 静力 capability 仍是无状态一次性执行,契约中的 `workspace: null` 明确保持原发布流程。 科研统计图继续以独立 feature 增量扩展同一契约:`box/histogram` 的系列只绑定原始 Y 列,统计规则由 Origin 固定模板决定;`bubble` 增加正数 `size` 数据角色,由受控 modifier column 驱动符号尺寸;`band` 要求同一输入的 `x/y/lower/upper`,先画上下界并填充到下一曲线,再叠加中心线,Worker 在打开 Origin 前拒绝非有限尺寸、非正尺寸和倒置边界。`stacked_column/stacked_area/stacked_bar` 统一把两条以上 XY 系列复制进内部连续 XYY 工作表,严格校验横坐标相同后建立 plot group,并只执行 Worker 内置的固定累计图层命令;请求不能提供命令、模板或工作表范围。`area/polar/pie` 继续使用固定 Origin 类型 ID。未经过目标 Origin 版本真机验证的统计属性不进入公共 schema,避免暴露看似可配但不能稳定复现的参数。 -第五阶段完成输出上传与发布:Node 只按固定 manifest ID 逐项流式 PUT,并携带 Node、lease、request digest 与内容摘要;云端重新绑定任务身份,不信 Node 提供的路径或媒体类型。文件先进入用户根下隐藏暂存区,固定文件名、单文件/总大小和 SHA-256 全部验证后,把 plot spec、provenance 整理进 `.meta/`,再将完整目录原子移动到 `/origin//`。PNG/SVG/PDF/OPJU 等正式输出登记平台 artifact UUID 和 `software_job_id`,`.meta/` 只落真实文件;成功状态返回 task-relative `output_dir`,Agent 以该目录为起点按需搜索。重复 PUT、complete 和重连均按摘要幂等;部分上传不可见,只有完整集合才能发布。Origin 执行槽与上传确认是两个正交状态:本地已有终态且固定 Worker 已退出时即释放软件执行槽,成功但尚无 `upload-complete.json` 的任务继续后台补传;若云端已经是 succeeded,重复 PUT/complete 必须按数据库持久化 manifest 校验并直接确认,不得按新版本目录规则重新发布旧 Job。 +Workspace capability 的默认完成协议不发布正式产物。Node 只自动上传契约声明的轻量 `preview_outputs`,云端校验后放入 `.zcbot_cache/software-previews/` 并记录 `preview_manifest`;预览可在对话中展示,但没有 artifact UUID,不进入正式文件列表,也不作为后续加工的事实源。完整工程和其他重输出保留在 home node 的 `current` 中,`local_manifest` 只记录固定 output ID、摘要、大小和媒体类型。只有用户明确要求下载、交付或导出时,Agent 调用 `software_job_export`,服务端向 home node 下发所选 output ID,Node 才流式上传,云端复核后移入 task 输出目录并登记正式 Artifact。预览缓存与 Artifact 因此是两条独立生命周期,避免中间截图、低清视频和多轮工程文件污染用户产物。 + +无状态 capability 继续使用原完成即发布协议:Node 按固定 manifest ID 逐项流式 PUT,云端重新绑定路径与媒体类型,校验完整集合后原子发布。两种协议都以 Node、lease、request digest、大小和 SHA-256 做幂等校验;Origin 执行槽在本地终态和 Workspace 提升完成后释放,预览或显式导出的传输失败可在重连时继续,不重复驱动专业软件。 第六阶段增加用户级 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`;终态写入仍由云端账本裁决。 diff --git a/PROGRESS.md b/PROGRESS.md index 3f4e705..7984f1a 100644 --- a/PROGRESS.md +++ b/PROGRESS.md @@ -20,6 +20,8 @@ --- ## 已完成关键能力 +- **08-19 / Unreleased / Windows Node 持久 Workspace 与按需导出**:新增独立 `software_workspaces` 注册表,首个 Job 建立 Workspace、后续 Job 只从当前 head 在同一 home node 串行续作;Node 以 `workspaces//current + rollback` 保存两代状态,Origin adapter 1.0.0 可打开上一版 OPJU 再加工,失败自动回滚且不会按 Job 永久复制工程。Workspace 完成时只上传契约声明的预览缓存,预览不注册 Artifact;用户明确要求下载或交付时才由 `software_job_export` 触发原节点上传并登记正式产物。专项 101 项 unittest 与 .NET build 通过;全量 672 项仍仅 3 个既有数据库集成模块因显式测试库缺少 `users` 表未通过(另跳过 4 项),未连接或迁移生产数据库。 + - **08-18 / Unreleased / Windows Node 能力自动同步**:注册码只建立节点身份,不再要求管理员选择 Origin/ANSYS 能力;Node 从 adapter registry 自动发现全部本机能力,并在首次连接及每次心跳上报,服务端以双方已有共享契约的能力覆盖节点调度列表,管理后台改为只读展示。现有节点无需清身份或重新注册,替换新版 Node 并重启一次即可同步。 - **08-18 / Unreleased / Windows Node 多 runtime 安装修复**:修复统一 BAT 安装器在连续处理 Origin 与 ANSYS 时依赖跨子程序 `PYTHON_EXE/PYTHON_ARGS` 状态、第二个 venv 可能执行成空命令的问题;每个缺失 runtime 现在都直接用已验证的显式 Python 3.12 命令创建,不再共享可变命令状态。 diff --git a/RUN.md b/RUN.md index 75ed172..ac06381 100644 --- a/RUN.md +++ b/RUN.md @@ -1054,7 +1054,7 @@ sudo xfs_quota -x -c "limit -p bhard=10g zcbot_" /opt ### Windows Node 内网 MVP(开发中) -先执行 `.venv/Scripts/python.exe main.py db upgrade head` 创建 `software_node_enrollments`、`software_nodes`、`software_jobs`,为 artifact 增加专业软件来源字段,并加入 Job 完成策略与续跑状态。不要在未确认目标数据库时运行迁移;本机 `.env` 的 `ZCBOT_DB_URL` 可能是生产隧道。 +先执行 `.venv/Scripts/python.exe main.py db upgrade head` 创建或升级 `software_node_enrollments`、`software_nodes`、`software_jobs` 与 `software_workspaces`,为 artifact 增加专业软件来源字段,并加入 Job 完成策略、Workspace head、本机 manifest、预览和按需导出状态。不要在未确认目标数据库时运行迁移;本机 `.env` 的 `ZCBOT_DB_URL` 可能是生产隧道。 0032 旧版曾把 `tasks.working_dir` 二次拼进 user root,已完成 Job 的文件和 artifact 路径可能位于重复的 `workspace/users//...` 子树。升级到 0033 后,用显式迁移库地址先 dry-run;确认所有目标目录无冲突后再加 `--apply`。脚本不加载 `.env`,文件和数据库更新均按 Job manifest 精确处理,目标目录同时存在时会停止,不能覆盖。 @@ -1098,7 +1098,9 @@ cd /d D:\ZcbotNode install-windows-node.bat ``` -若 Python 未加入 PATH,可把绝对路径作为第一个参数,例如 `install-windows-node.bat "C:\Python312\python.exe"`。默认解释器分别为 `%ProgramData%\Zcbot\WindowsNode\runtimes\origin\Scripts\python.exe` 与 `%ProgramData%\Zcbot\WindowsNode\runtimes\ansys\Scripts\python.exe`。如需使用其他受管解释器,使用机器级 `ZCBOT_ADAPTER_ORIGIN_PYTHON` 或 `ZCBOT_ADAPTER_ANSYS_PYTHON`;旧名 `ZCBOT_ORIGIN_PYTHON` 暂时兼容。`node.json`、可恢复任务和 runtime 集中保存在 `%ProgramData%\Zcbot\WindowsNode\`,不会因替换程序目录而丢失。运行时固定依赖见各 adapter 的 `requirements.txt`;任务请求无权选择解释器、脚本或路径。Origin Worker 支持 CSV/XLSX/JSON 输入及 OPJU/PNG/SVG/PDF 输出,成功产物发布到 `origin//`,plot spec 与 provenance 位于其 `.meta/`;上传中断会在重连时幂等续传。 +若 Python 未加入 PATH,可把绝对路径作为第一个参数,例如 `install-windows-node.bat "C:\Python312\python.exe"`。默认解释器分别为 `%ProgramData%\Zcbot\WindowsNode\runtimes\origin\Scripts\python.exe` 与 `%ProgramData%\Zcbot\WindowsNode\runtimes\ansys\Scripts\python.exe`。如需使用其他受管解释器,使用机器级 `ZCBOT_ADAPTER_ORIGIN_PYTHON` 或 `ZCBOT_ADAPTER_ANSYS_PYTHON`;旧名 `ZCBOT_ORIGIN_PYTHON` 暂时兼容。`node.json`、可恢复任务、workspace 和 runtime 集中保存在 `%ProgramData%\Zcbot\WindowsNode\`,不会因替换程序目录而丢失。运行时固定依赖见各 adapter 的 `requirements.txt`;任务请求无权选择解释器、脚本或路径。 + +Workspace 目录固定为 `%ProgramData%\Zcbot\WindowsNode\workspaces\\`,不使用 user ID 上层目录。`current` 是当前可续作工程,`rollback` 只保留上一版用于失败恢复;每次成功续作替换 `current`,不会按 Job 永久堆积工程副本。不要手工移动 workspace 到其他节点:服务端会把后续 Job 固定调度到其 home node,节点离线时任务保持等待。备份或迁移节点时必须把数据库中的 workspace 归属和该目录作为一个整体处理。 ANSYS 能力为 `ansys.mechanical.static_structural@v1`,只面向 Mechanical 2024 R2(revision 242)。adapter 0.2.0 已实现 PyMechanical 0.11.3 真实静力流程;Node 会从 `adapters/` 自动发现并在每次心跳上报全部能力,服务端管理页只读展示双方已有共享契约的能力,不需要管理员勾选,也不需要现有节点清身份或重新注册。默认 probe 仍返回不可用;在目标机器完成真实许可证与基准算例验收前,不得设置 `ZCBOT_ANSYS_242_VALIDATED=1`。 @@ -1125,11 +1127,11 @@ ANSYS 能力为 `ansys.mechanical.static_structural@v1`,只面向 Mechanical 2 重启后 probe 应报告 `health=ready`,再从服务端提交一个同规格小算例验证调度、上传与发布闭环。取消验收若提示 `CANCELLATION_JOB_FINISHED_BEFORE_CANCEL`,说明基准在设定秒数前已经完成,应减小 `--cancel-after` 后使用新的 `--work-root` 重跑,而不是把该次视为通过。 -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。两者均不阻塞提交轮次。请求支持 1–16 个输入、跨输入系列和多个显式输出;`plot.canvas` 可指定毫米画布,轴可指定范围、步长、尺度、刻度角度/字号、标题字号和网格,`legend` 可控制显隐、位置和字号,`series[].style` 可控制颜色、线宽/线型、点型/点大小和透明度。所有排版字段可选,旧请求保持默认样式。成功状态提供 `output_dir`,正式输出的 artifact 带 `software_job_id`,供结果卡和产物详情展示来源。Node 输出上传的逐任务诊断日志位于 `%ProgramData%\Zcbot\WindowsNode\jobs\\logs\node-output-upload.log`;日志包含上传阶段、产物文件名、重试次数和 Windows `HRESULT`,单文件达到 1 MiB 后轮转一份 `.1`,不记录 Node Token 或认证请求头。 +Web 用户登录后,文件栏 Job 中心会聚合本人最近任务。活动任务或完成后的自动分析约 4 秒刷新一次,空闲时降为约 30 秒;卡片优先显示状态、所属对话、输入输出、完成策略和耗时,Job ID、Workspace ID、执行节点及软件版本收在“任务详情”中。停止已派发任务是协作取消,状态先显示“正在停止”,Node 在线时立即接收,断线后在下次连接或心跳时重放。Agent 可调用 `software_capability_list`、`register_artifact`、`software_job_submit`、`software_job_status`、`software_job_revise`、`software_job_export` 和 `software_job_cancel`。Origin 输入必须是 artifact:已有 UUID 可直接提交,普通 task 文件先逐个用相对路径登记;登记不会发布聊天交付卡片。提交工具接收 `inputs`、`operation`、`outputs` 和可选的 `completion_action=report|analyze`。Workspace 任务成功后只自动上传契约声明的 PNG 等轻量预览,预览保存在隐藏缓存且不注册 Artifact;OPJU、SVG、PDF 等完整输出继续留在 Node。只有用户明确要求下载、发送或交付时,才调用 `software_job_export(job_id, output_ids)`,Node 会上传所选输出并在 task 目录登记正式 Artifact。Node 输出上传的逐任务诊断日志位于 `%ProgramData%\Zcbot\WindowsNode\jobs\\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` 数量为 1–4;两个 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。 -声明式绘图使用 `plot.type=recipe`、`recipe_version=1`,第一版有意复用上述 `layout/panels` 结构和限制,作为新的组合入口而不是开放脚本。对已有结果继续调整时,先用 `software_job_status(job_id)` 取得 `editable_request`,在其中形成完整的新 `operation` 与 `outputs`,再调用 `software_job_revise`;服务端自动复用源 Job 的 inputs,重新校验并生成新 Job。旧 Job、OPJU 和图片不会被覆盖,当前也不会保留用户下载后在 Origin GUI 中做的手工修改。 +声明式绘图使用 `plot.type=recipe`、`recipe_version=1`,第一版有意复用上述 `layout/panels` 结构和限制,作为新的组合入口而不是开放脚本。对已有结果继续调整时,先用 `software_job_status(job_id)` 取得 `editable_request`,在其中形成完整的新 `operation` 与 `outputs`,再调用 `software_job_revise`;服务端自动复用源 Job 的 inputs,并要求 source Job 是该 Workspace 当前 head。Origin Worker 打开 Node 上的上一版 OPJU,将新工作表和图添加到工程后保存为新的 `current`;失败时恢复 `rollback`。Job 历史不被覆盖,但 Node 只保留当前和上一版工程文件。 科研统计图使用独立单图 type。`box` 和 `histogram` 的系列填写 `input/y`,直接使用原始观测列并采用 Origin 默认箱线统计与自动分箱;`stacked_column` 至少提供两条同类 `input/x/y` 系列;`bubble` 使用 `input/x/y/size`,每行 size 必须是有限正数;`band` 使用 `input/x/y/lower/upper`,其中 y 是中心线且每行 lower 不得大于 upper。当前 schema 不开放自定义箱线百分位、直方图分箱和分布拟合参数,也不允许把这五类统计图作为 `multi_panel.series[].kind`;在目标 Origin 版本完成真机验证后再增量开放。 diff --git a/core/software_contracts.py b/core/software_contracts.py index d3d0a8e..ec8eddb 100644 --- a/core/software_contracts.py +++ b/core/software_contracts.py @@ -32,6 +32,13 @@ class OutputSpec: required: bool +@dataclass(frozen=True) +class WorkspaceSpec: + state_output: str + preview_outputs: tuple[str, ...] + required_adapter_version: str + + @dataclass(frozen=True) class CapabilityContract: capability: str @@ -45,6 +52,7 @@ class CapabilityContract: features: dict[str, str] summary: dict[str, Any] legacy_runtime: dict[str, Any] | None + workspace: WorkspaceSpec | None def normalize_request(self, request: object) -> tuple[dict[str, Any], str]: errors = sorted( @@ -162,6 +170,7 @@ def _load_contract(path: Path) -> CapabilityContract: required = { "capability", "display_name", "default_enrollment", "output_namespace", "request_schema", "input_policy", "outputs", "feature_path", "features", "summary", "legacy_runtime", + "workspace", } if not isinstance(raw, dict) or set(raw) != required: raise RuntimeError(f"invalid software contract fields: {path.name}") @@ -199,6 +208,33 @@ def _load_contract(path: Path) -> CapabilityContract: outputs[output_id] = OutputSpec(output_id=output_id, **item) if len({item.relative_path for item in outputs.values()}) != len(outputs): raise RuntimeError(f"duplicate output paths: {path.name}") + workspace_raw = raw["workspace"] + workspace: WorkspaceSpec | None = None + if workspace_raw is not None: + if not isinstance(workspace_raw, dict) or set(workspace_raw) != { + "state_output", "preview_outputs", "required_adapter_version" + }: + raise RuntimeError(f"invalid workspace spec: {path.name}") + state_output = workspace_raw["state_output"] + preview_outputs = workspace_raw["preview_outputs"] + required_adapter_version = workspace_raw["required_adapter_version"] + if ( + not isinstance(state_output, str) + or state_output not in outputs + or not outputs[state_output].required + or not isinstance(preview_outputs, list) + or not preview_outputs + or any(not isinstance(item, str) or item not in outputs for item in preview_outputs) + or len(set(preview_outputs)) != len(preview_outputs) + or any(not outputs[item].required for item in preview_outputs) + or state_output in preview_outputs + or not isinstance(required_adapter_version, str) + or not re.fullmatch(r"[0-9]+\.[0-9]+\.[0-9]+", required_adapter_version) + ): + raise RuntimeError(f"unsafe workspace spec: {path.name}") + workspace = WorkspaceSpec( + state_output, tuple(preview_outputs), required_adapter_version + ) Draft202012Validator.check_schema(raw["request_schema"]) return CapabilityContract( capability=capability, @@ -212,6 +248,7 @@ def _load_contract(path: Path) -> CapabilityContract: features=raw["features"], summary=raw["summary"], legacy_runtime=raw["legacy_runtime"], + workspace=workspace, ) @@ -317,7 +354,14 @@ def node_supports_request( feature = contract.feature(request) if isinstance(advertised_features, list) and feature not in advertised_features: return False - return version_at_least(actual_version, contract.required_adapter_version(request)) + required_version = contract.required_adapter_version(request) + if contract.workspace is not None and version_at_least( + contract.workspace.required_adapter_version, required_version + ): + required_version = contract.workspace.required_adapter_version + return version_at_least(actual_version, required_version) + if contract.workspace is not None: + return False # 兼容尚未升级 capability_runtime 的 Node;兼容路径由版本化 contract 声明。 legacy_config = contract.legacy_runtime or {} if not legacy_config: diff --git a/core/software_jobs.py b/core/software_jobs.py index f3545a4..7a3bf32 100644 --- a/core/software_jobs.py +++ b/core/software_jobs.py @@ -17,11 +17,19 @@ from core.software_contracts import ( node_supports_request, ) from core.storage.engine import session_scope -from core.storage.models import Artifact, SoftwareJob, SoftwareNode, Task +from core.storage.models import ( + Artifact, + SoftwareJob, + SoftwareNode, + SoftwareWorkspace, + Task, +) OFFER_SECONDS = 60 MAX_OUTPUT_ARTIFACT_BYTES = 256 * 1024 * 1024 MAX_OUTPUT_TOTAL_BYTES = 512 * 1024 * 1024 +MAX_EXPORT_ARTIFACT_BYTES = 10 * 1024 * 1024 * 1024 +MAX_EXPORT_TOTAL_BYTES = 20 * 1024 * 1024 * 1024 class SoftwareJobError(Exception): @@ -82,6 +90,12 @@ def _job_dict(row: SoftwareJob) -> dict: return { "job_id": str(row.job_id), "task_id": str(row.task_id), + "workspace_id": ( + str(row.workspace_id) if getattr(row, "workspace_id", None) else None + ), + "source_job_id": ( + str(row.source_job_id) if getattr(row, "source_job_id", None) else None + ), "capability": row.capability, "request_digest": row.request_digest, "node_id": str(row.node_id) if row.node_id else None, @@ -92,9 +106,17 @@ def _job_dict(row: SoftwareJob) -> dict: "execution_runtime": _job_execution_runtime(row), "error": row.error, "artifact_manifest": row.artifact_manifest, + "local_manifest": getattr(row, "local_manifest", []), + "preview_manifest": getattr(row, "preview_manifest", []), + "export_status": getattr(row, "export_status", "none"), + "export_outputs": getattr(row, "export_outputs", []), "completion_action": getattr(row, "completion_action", "report"), "followup_status": getattr(row, "followup_status", "none"), - "output_dir": f"{contract.output_namespace}/{row.job_id}", + "output_dir": ( + None + if contract.workspace is not None and getattr(row, "workspace_id", None) + else 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, "terminal_at": row.terminal_at.isoformat() if row.terminal_at else None, @@ -227,6 +249,8 @@ def create_job( capability: str, request: dict, completion_action: str = "report", + workspace_id: UUID | None = None, + source_job_id: UUID | None = None, ) -> tuple[dict, bool]: key = idempotency_key.strip() if not key or len(key) > 200: @@ -244,6 +268,70 @@ def create_job( ).first() if task is None: raise SoftwareJobError("task not found") + existing = session.execute( + select(SoftwareJob).where( + SoftwareJob.user_id == user_id, + SoftwareJob.idempotency_key == key, + ) + ).scalar_one_or_none() + if existing is not None: + if ( + existing.task_id != task_id + or existing.capability != capability + or existing.request_digest != digest + or getattr(existing, "completion_action", "report") != completion_action + or workspace_id is not None and existing.workspace_id != workspace_id + or source_job_id is not None and existing.source_job_id != source_job_id + ): + raise SoftwareJobError( + "idempotency key was already used for a different request" + ) + return _job_dict(existing), False + if (workspace_id is None) != (source_job_id is None): + raise SoftwareJobError( + "workspace_id and source_job_id must be provided together" + ) + workspace: SoftwareWorkspace | None = None + source: SoftwareJob | None = None + if workspace_id is not None: + workspace = session.execute( + select(SoftwareWorkspace).where( + SoftwareWorkspace.workspace_id == workspace_id, + SoftwareWorkspace.user_id == user_id, + ).with_for_update() + ).scalar_one_or_none() + if workspace is None: + raise SoftwareJobError("software workspace not found") + if workspace.capability != capability: + raise SoftwareJobError("software workspace capability does not match") + if workspace.status in {"deleting", "deleted", "lost"}: + raise SoftwareJobError("software workspace is not available") + source = session.execute( + select(SoftwareJob).where( + SoftwareJob.job_id == source_job_id, + SoftwareJob.workspace_id == workspace_id, + SoftwareJob.user_id == user_id, + ) + ).scalar_one_or_none() + if source is None or source.status != "succeeded": + raise SoftwareJobError("source software job is not a successful workspace revision") + if workspace.head_job_id != source.job_id: + raise SoftwareJobError("source software job is not the workspace head") + if source.export_status == "pending": + raise SoftwareJobError( + "software workspace must finish its pending export before revision" + ) + active = session.execute( + select(SoftwareJob.job_id).where( + SoftwareJob.workspace_id == workspace_id, + SoftwareJob.status.in_({ + "queued", "offered", "dispatched", "running", + "disconnected", "cancelling", + }), + ).limit(1) + ).first() + if active is not None: + raise SoftwareJobError("software workspace already has an active job") bindings = contract.input_bindings(normalized) artifact_ids = [UUID(item["artifact_id"]) for item in bindings] artifacts = session.execute( @@ -302,13 +390,32 @@ def create_job( or existing.capability != capability or existing.request_digest != digest or getattr(existing, "completion_action", "report") != completion_action + or workspace_id is not None and existing.workspace_id != workspace_id + or source_job_id is not None and existing.source_job_id != source_job_id ): raise SoftwareJobError("idempotency key was already used for a different request") return _job_dict(existing), False + effective_workspace_id = ( + workspace_id + or (uuid4() if contract.workspace is not None else None) + ) + new_workspace: SoftwareWorkspace | None = None + if workspace is None and effective_workspace_id is not None: + new_workspace = SoftwareWorkspace( + workspace_id=effective_workspace_id, + user_id=user_id, + capability=capability, + status="pending", + size_bytes=0, + state_manifest=[], + retention_until=datetime.now(timezone.utc) + timedelta(days=30), + ) row = SoftwareJob( job_id=uuid4(), user_id=user_id, task_id=task_id, + workspace_id=effective_workspace_id, + source_job_id=source_job_id, idempotency_key=key, capability=capability, request=normalized, @@ -319,11 +426,17 @@ def create_job( metrics={}, error={}, artifact_manifest=[], + local_manifest=[], + preview_manifest=[], + export_status="none", + export_outputs=[], completion_action=completion_action, followup_status="none", ) try: with session.begin_nested(): + if new_workspace is not None: + session.add(new_workspace) session.add(row) session.flush() return _job_dict(row), True @@ -339,6 +452,8 @@ def create_job( or existing.capability != capability or existing.request_digest != digest or getattr(existing, "completion_action", "report") != completion_action + or workspace_id is not None and existing.workspace_id != workspace_id + or source_job_id is not None and existing.source_job_id != source_job_id ): raise SoftwareJobError( "idempotency key was already used for a different request" @@ -389,7 +504,12 @@ def revise_job( ).scalar_one_or_none() if source is None: raise SoftwareJobError("source software job not found") + if source.status != "succeeded" or source.workspace_id is None: + raise SoftwareJobError("source software job has no available workspace state") capability = source.capability + if get_contract(capability).workspace is None: + raise SoftwareJobError("capability does not support workspace continuation") + workspace_id = source.workspace_id completion_action = getattr(source, "completion_action", "report") source_request = source.request inputs = deepcopy(source_request.get("inputs") or []) @@ -402,6 +522,8 @@ def revise_job( idempotency_key=idempotency_key, capability=capability, completion_action=completion_action, + workspace_id=workspace_id, + source_job_id=source_job_id, request={ "schema_version": schema_version, "inputs": inputs, @@ -430,6 +552,122 @@ def request_job_analysis(user_id: UUID, job_id: UUID) -> dict: return _job_dict(job) +def request_job_export(user_id: UUID, job_id: UUID, output_ids: list[str]) -> dict: + """请求原 Node 将当前 Workspace head 的指定本地输出发布为正式 Artifact。""" + requested = list(dict.fromkeys(output_ids)) + if not requested or len(requested) > 16 or any(not isinstance(item, str) for item in requested): + raise SoftwareJobError("export output_ids must contain 1 to 16 output identities") + 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") + contract = get_contract(job.capability) + if job.status != "succeeded" or contract.workspace is None or job.workspace_id is None: + raise SoftwareJobError("job has no exportable workspace state") + workspace = session.execute( + select(SoftwareWorkspace).where( + SoftwareWorkspace.workspace_id == job.workspace_id + ).with_for_update() + ).scalar_one_or_none() + if workspace is None or workspace.head_job_id != job.job_id: + raise SoftwareJobError("only the current workspace head can be exported") + active = session.execute( + select(SoftwareJob.job_id).where( + SoftwareJob.workspace_id == job.workspace_id, + SoftwareJob.job_id != job.job_id, + SoftwareJob.status.in_({ + "queued", "offered", "dispatched", "running", + "disconnected", "cancelling", + }), + ).limit(1) + ).first() + if active is not None: + raise SoftwareJobError("software workspace has an active revision") + local_by_id = { + item.get("artifact_id"): item + for item in job.local_manifest + if isinstance(item, dict) and isinstance(item.get("artifact_id"), str) + } + already_exported = { + item.get("source_artifact_id") + for item in job.artifact_manifest + if isinstance(item, dict) and item.get("artifact_id") + } + for output_id in requested: + if output_id not in local_by_id: + raise SoftwareJobError(f"local output is unavailable: {output_id}") + if not contract.output_spec(output_id).publish: + raise SoftwareJobError(f"output is internal and cannot be exported: {output_id}") + pending = [item for item in requested if item not in already_exported] + if not pending: + job.export_status = "completed" + return _job_dict(job) + if job.export_status == "pending": + if pending == list(job.export_outputs): + return _job_dict(job) + raise SoftwareJobError("software job already has a pending export") + job.export_outputs = pending + job.export_status = "pending" + return _job_dict(job) + + +def pending_node_exports(node_id: UUID) -> list[dict]: + with session_scope() as session: + jobs = session.execute( + select(SoftwareJob).where( + SoftwareJob.node_id == node_id, + SoftwareJob.status == "succeeded", + SoftwareJob.export_status == "pending", + ).order_by(SoftwareJob.updated_at, SoftwareJob.job_id) + ).scalars() + return [ + { + "job_id": str(job.job_id), + "lease_id": str(job.lease_id), + "request_digest": job.request_digest, + "output_ids": list(job.export_outputs), + } + for job in jobs + if job.lease_id is not None + ] + + +def record_job_export( + node_id: UUID, job_id: UUID, output_ids: list[str], published: list[dict] +) -> dict: + 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.node_id != node_id + or job.status != "succeeded" + or job.export_status not in {"pending", "completed"} + ): + raise SoftwareJobError("software export request is stale") + requested = list(job.export_outputs) + if requested != output_ids: + raise SoftwareJobError("software export outputs do not match the request") + by_source = { + item.get("source_artifact_id"): item + for item in job.artifact_manifest + if isinstance(item, dict) + } + for item in published: + by_source[item["source_artifact_id"]] = item + job.artifact_manifest = list(by_source.values()) + job.export_status = "completed" + job.completion_action = "report" + job.followup_status = "pending" + return _job_dict(job) + + def offer_next_job(node_ids: set[UUID]) -> dict | None: """选择最早可执行的 Job–Node 组合,避免跨 capability 队首阻塞。""" if not node_ids: @@ -487,10 +725,23 @@ def offer_next_job(node_ids: set[UUID]) -> dict | None: contract = get_contract(candidate.capability) except SoftwareContractError: continue + workspace = ( + session.get(SoftwareWorkspace, candidate.workspace_id) + if getattr(candidate, "workspace_id", None) is not None + else None + ) + if getattr(candidate, "workspace_id", None) is not None and ( + workspace is None + or workspace.status in {"deleting", "deleted", "lost"} + or (workspace.head_job_id is not None and workspace.home_node_id is None) + ): + continue + home_node_id = workspace.home_node_id if workspace is not None else None node = next( ( item for item in available_nodes if candidate.capability in (item.capabilities or []) + and (home_node_id is None or item.node_id == home_node_id) and node_supports_request( contract, candidate.request, item.runtime or {} ) @@ -519,6 +770,11 @@ def offer_next_job(node_ids: set[UUID]) -> dict | None: job.lease_id = lease_id job.lease_expires_at = expires_at job.status = "offered" + workspace = ( + session.get(SoftwareWorkspace, job.workspace_id) + if getattr(job, "workspace_id", None) is not None + else None + ) return { "node_id": node.node_id, "payload": { @@ -528,6 +784,17 @@ def offer_next_job(node_ids: set[UUID]) -> dict | None: "capability": job.capability, "request_digest": job.request_digest, "request": job.request, + "workspace": ( + { + "workspace_id": str(job.workspace_id), + "source_job_id": ( + str(job.source_job_id) if job.source_job_id else None + ), + "mode": "continue" if job.source_job_id else "new", + } + if workspace is not None + else None + ), "input_transfers": [ { **item, @@ -603,6 +870,11 @@ def get_job_output_context(node_id: UUID, job_id: UUID, lease_id: UUID, digest: "request": job.request, "status": job.status, "artifact_manifest": job.artifact_manifest, + "local_manifest": getattr(job, "local_manifest", []), + "preview_manifest": getattr(job, "preview_manifest", []), + "workspace_id": getattr(job, "workspace_id", None), + "export_status": getattr(job, "export_status", "none"), + "export_outputs": getattr(job, "export_outputs", []), } @@ -755,6 +1027,19 @@ def respond_to_offer(node_id: UUID, *, accepted: bool, payload: dict) -> None: job.status = "dispatched" job.stage = "accepted" job.error = {} + if getattr(job, "workspace_id", None) is not None: + workspace = session.execute( + select(SoftwareWorkspace).where( + SoftwareWorkspace.workspace_id == job.workspace_id + ).with_for_update() + ).scalar_one_or_none() + if workspace is None: + raise SoftwareJobError("software workspace no longer exists") + if workspace.home_node_id not in {None, node_id}: + raise SoftwareJobError("software workspace belongs to another node") + workspace.home_node_id = node_id + if workspace.status == "pending": + workspace.status = "assigned" snapshot = _execution_runtime_snapshot( session.get(SoftwareNode, node_id), job.capability ) @@ -814,10 +1099,19 @@ def record_job_terminal(node_id: UUID, payload: dict) -> None: raise SoftwareJobError("invalid job terminal status") error = payload.get("error") or {} manifest = payload.get("artifact_manifest") or [] + local_manifest = payload.get("local_manifest") or manifest + preview_manifest = payload.get("preview_manifest") or [] if not isinstance(error, dict) or len(json.dumps(error, ensure_ascii=False)) > 64 * 1024: raise SoftwareJobError("job terminal error is invalid") if not isinstance(manifest, list) or len(json.dumps(manifest, ensure_ascii=False)) > 256 * 1024: raise SoftwareJobError("job artifact manifest is invalid") + if ( + not isinstance(local_manifest, list) + or len(json.dumps(local_manifest, ensure_ascii=False)) > 256 * 1024 + or not isinstance(preview_manifest, list) + or len(json.dumps(preview_manifest, ensure_ascii=False)) > 64 * 1024 + ): + raise SoftwareJobError("job local or preview manifest is invalid") now = datetime.now(timezone.utc) with session_scope() as session: job = session.execute( @@ -830,17 +1124,22 @@ def record_job_terminal(node_id: UUID, payload: dict) -> None: job.request, [ { - "artifact_id": item.get("source_artifact_id"), + "artifact_id": ( + item.get("source_artifact_id") or item.get("artifact_id") + ), "filename": item.get("filename"), "media_type": item.get("media_type"), "size_bytes": item.get("size_bytes"), "sha256": item.get("sha256"), } - for item in manifest + for item in local_manifest if isinstance(item, dict) ], ) - if len(expected) != len(manifest) or any( + contract = get_contract(job.capability) + if len(expected) != len(local_manifest): + raise SoftwareJobError("successful job local outputs are incomplete") + if contract.workspace is None and any( not _published_output_is_valid(job.capability, job.job_id, item) for item in manifest ): @@ -856,9 +1155,39 @@ def record_job_terminal(node_id: UUID, payload: dict) -> None: job.progress = 100 if terminal_status == "succeeded" else job.progress job.error = error job.artifact_manifest = manifest + job.local_manifest = local_manifest + job.preview_manifest = preview_manifest job.terminal_at = now if terminal_status == "succeeded": job.followup_status = "pending" + if getattr(job, "workspace_id", None) is not None: + workspace = session.execute( + select(SoftwareWorkspace).where( + SoftwareWorkspace.workspace_id == job.workspace_id + ).with_for_update() + ).scalar_one_or_none() + if workspace is None or workspace.home_node_id != node_id: + raise SoftwareJobError("software workspace ownership is inconsistent") + if job.source_job_id is not None and workspace.head_job_id != job.source_job_id: + raise SoftwareJobError("software workspace head changed during execution") + workspace.head_job_id = job.job_id + workspace.status = "active" + workspace.last_used_at = now + workspace.state_manifest = local_manifest + workspace.size_bytes = int(payload.get("workspace_size_bytes") or 0) + workspace.last_reported_at = now + workspace.retention_until = now + timedelta(days=30) + elif ( + get_contract(job.capability).workspace is not None + and getattr(job, "workspace_id", None) is not None + ): + workspace = session.execute( + select(SoftwareWorkspace).where( + SoftwareWorkspace.workspace_id == job.workspace_id + ).with_for_update() + ).scalar_one_or_none() + if workspace is not None and workspace.head_job_id is None: + workspace.status = "failed" def _is_uuid(value: str) -> bool: diff --git a/core/software_nodes.py b/core/software_nodes.py index 6e8c0ea..d7b40a1 100644 --- a/core/software_nodes.py +++ b/core/software_nodes.py @@ -8,11 +8,11 @@ from hashlib import sha256 from uuid import UUID, uuid4 import bcrypt -from sqlalchemy import select +from sqlalchemy import select, update from core.software_contracts import supported_capabilities from core.storage.engine import session_scope -from core.storage.models import SoftwareNode, SoftwareNodeEnrollment +from core.storage.models import SoftwareNode, SoftwareNodeEnrollment, SoftwareWorkspace MAX_ENROLLMENT_FAILURES = 5 @@ -187,11 +187,16 @@ def set_node_disabled(node_id: UUID, disabled: bool) -> bool: def delete_node(node_id: UUID) -> bool: - """撤销并物理删除节点身份;当前节点表没有任务历史外键。""" + """撤销节点身份,并把无法再定位的持久 Workspace 标为丢失。""" with session_scope() as session: node = session.get(SoftwareNode, node_id) if node is None: return False + session.execute( + update(SoftwareWorkspace) + .where(SoftwareWorkspace.home_node_id == node_id) + .values(status="lost", home_node_id=None) + ) session.delete(node) return True diff --git a/core/storage/models.py b/core/storage/models.py index 9f126dd..46895f1 100644 --- a/core/storage/models.py +++ b/core/storage/models.py @@ -470,6 +470,53 @@ class SoftwareNode(Base): ) +class SoftwareWorkspace(Base): + """专业软件在单一 Windows Node 上持有的可继续加工工程。""" + + __tablename__ = "software_workspaces" + __table_args__ = ( + Index("ix_software_workspaces_user_status", "user_id", "status"), + Index("ix_software_workspaces_node_status", "home_node_id", "status"), + ) + + workspace_id: Mapped[UUID] = mapped_column( + PG_UUID(as_uuid=True), primary_key=True, default=uuid4 + ) + user_id: Mapped[UUID] = mapped_column( + PG_UUID(as_uuid=True), ForeignKey("users.user_id", ondelete="CASCADE"), nullable=False + ) + capability: Mapped[str] = mapped_column(Text, nullable=False) + home_node_id: Mapped[Optional[UUID]] = mapped_column( + PG_UUID(as_uuid=True), + ForeignKey("software_nodes.node_id", ondelete="SET NULL"), + nullable=True, + ) + head_job_id: Mapped[Optional[UUID]] = mapped_column( + PG_UUID(as_uuid=True), + ForeignKey( + "software_jobs.job_id", + name="fk_software_workspaces_head_job_id", + ondelete="SET NULL", + use_alter=True, + ), + nullable=True, + ) + status: Mapped[str] = mapped_column( + Text, nullable=False, default="pending", server_default="pending" + ) + size_bytes: Mapped[int] = mapped_column(BigInteger, nullable=False, default=0, server_default="0") + state_manifest: Mapped[list[Any]] = mapped_column(JSONB, nullable=False, default=list) + last_used_at: Mapped[Optional[datetime]] = mapped_column(DateTime(timezone=True), nullable=True) + last_reported_at: Mapped[Optional[datetime]] = mapped_column(DateTime(timezone=True), nullable=True) + retention_until: Mapped[Optional[datetime]] = mapped_column(DateTime(timezone=True), nullable=True) + created_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), server_default=func.now(), nullable=False + ) + updated_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), server_default=func.now(), onupdate=func.now(), nullable=False + ) + + class SoftwareJob(Base): """专业软件任务账本;请求只保存规范化参数和输入引用。""" @@ -488,6 +535,16 @@ class SoftwareJob(Base): task_id: Mapped[UUID] = mapped_column( PG_UUID(as_uuid=True), ForeignKey("tasks.task_id", ondelete="CASCADE"), nullable=False ) + workspace_id: Mapped[Optional[UUID]] = mapped_column( + PG_UUID(as_uuid=True), + ForeignKey("software_workspaces.workspace_id", ondelete="SET NULL"), + nullable=True, + ) + source_job_id: Mapped[Optional[UUID]] = mapped_column( + PG_UUID(as_uuid=True), + ForeignKey("software_jobs.job_id", ondelete="SET NULL"), + nullable=True, + ) idempotency_key: Mapped[str] = mapped_column(Text, nullable=False) capability: Mapped[str] = mapped_column(Text, nullable=False) request: Mapped[dict[str, Any]] = mapped_column(JSONB, nullable=False) @@ -504,6 +561,12 @@ 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) + local_manifest: Mapped[list[Any]] = mapped_column(JSONB, nullable=False, default=list) + preview_manifest: Mapped[list[Any]] = mapped_column(JSONB, nullable=False, default=list) + export_status: Mapped[str] = mapped_column( + Text, nullable=False, default="none", server_default="none" + ) + export_outputs: Mapped[list[Any]] = mapped_column(JSONB, nullable=False, default=list) # 0034:专业软件完成后的用户侧动作。report 只发布固定结果通知;analyze 会在 # 原 task 空闲后自动续跑 Agent。followup_status 同时充当可恢复的轻量 outbox。 completion_action: Mapped[str] = mapped_column( diff --git a/core/tool_registry.py b/core/tool_registry.py index d68c8b8..6dedf96 100644 --- a/core/tool_registry.py +++ b/core/tool_registry.py @@ -51,6 +51,7 @@ from tools.schedule import ( from tools.software_jobs import ( SoftwareCapabilityListTool, SoftwareJobCancelTool, + SoftwareJobExportTool, SoftwareJobReviseTool, SoftwareJobStatusTool, SoftwareJobSubmitTool, @@ -228,6 +229,7 @@ def build_tools(ctx: ToolContext) -> dict[str, Any]: SoftwareJobReviseTool(ctx.uid, ctx.task_id, **base), SoftwareJobStatusTool(ctx.uid, ctx.task_id, **base), SoftwareJobCancelTool(ctx.uid, ctx.task_id, **base), + SoftwareJobExportTool(ctx.uid, ctx.task_id, **base), ] def _send_email() -> list: diff --git a/db/migrations/versions/20260819_0900_0035_software_workspaces.py b/db/migrations/versions/20260819_0900_0035_software_workspaces.py new file mode 100644 index 0000000..ddde8e4 --- /dev/null +++ b/db/migrations/versions/20260819_0900_0035_software_workspaces.py @@ -0,0 +1,125 @@ +"""Add persistent professional-software workspaces and job lineage. + +Revision ID: 0035 +Revises: 0034 +Create Date: 2026-08-19 +""" +from collections.abc import Sequence + +import sqlalchemy as sa +from alembic import op +from sqlalchemy.dialects import postgresql + +revision: str = "0035" +down_revision: str | None = "0034" +branch_labels: str | Sequence[str] | None = None +depends_on: str | Sequence[str] | None = None + + +def upgrade() -> None: + op.create_table( + "software_workspaces", + sa.Column("workspace_id", postgresql.UUID(as_uuid=True), primary_key=True), + sa.Column( + "user_id", + postgresql.UUID(as_uuid=True), + sa.ForeignKey("users.user_id", ondelete="CASCADE"), + nullable=False, + ), + sa.Column("capability", sa.Text(), nullable=False), + sa.Column( + "home_node_id", + postgresql.UUID(as_uuid=True), + sa.ForeignKey("software_nodes.node_id", ondelete="SET NULL"), + nullable=True, + ), + sa.Column( + "head_job_id", + postgresql.UUID(as_uuid=True), + sa.ForeignKey( + "software_jobs.job_id", + name="fk_software_workspaces_head_job_id", + ondelete="SET NULL", + ), + nullable=True, + ), + sa.Column("status", sa.Text(), nullable=False, server_default="pending"), + sa.Column("size_bytes", sa.BigInteger(), nullable=False, server_default="0"), + sa.Column("state_manifest", postgresql.JSONB(), nullable=False, server_default="[]"), + sa.Column("last_used_at", sa.DateTime(timezone=True), nullable=True), + sa.Column("last_reported_at", sa.DateTime(timezone=True), nullable=True), + sa.Column("retention_until", sa.DateTime(timezone=True), nullable=True), + sa.Column("created_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()), + sa.Column("updated_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()), + ) + op.create_index( + "ix_software_workspaces_user_status", + "software_workspaces", + ["user_id", "status"], + ) + op.create_index( + "ix_software_workspaces_node_status", + "software_workspaces", + ["home_node_id", "status"], + ) + op.add_column( + "software_jobs", + sa.Column("workspace_id", postgresql.UUID(as_uuid=True), nullable=True), + ) + op.add_column( + "software_jobs", + sa.Column("source_job_id", postgresql.UUID(as_uuid=True), nullable=True), + ) + op.add_column( + "software_jobs", + sa.Column("local_manifest", postgresql.JSONB(), nullable=False, server_default="[]"), + ) + op.add_column( + "software_jobs", + sa.Column("preview_manifest", postgresql.JSONB(), nullable=False, server_default="[]"), + ) + op.add_column( + "software_jobs", + sa.Column("export_status", sa.Text(), nullable=False, server_default="none"), + ) + op.add_column( + "software_jobs", + sa.Column("export_outputs", postgresql.JSONB(), nullable=False, server_default="[]"), + ) + op.create_foreign_key( + "fk_software_jobs_workspace_id", + "software_jobs", + "software_workspaces", + ["workspace_id"], + ["workspace_id"], + ondelete="SET NULL", + ) + op.create_foreign_key( + "fk_software_jobs_source_job_id", + "software_jobs", + "software_jobs", + ["source_job_id"], + ["job_id"], + ondelete="SET NULL", + ) + op.create_index("ix_software_jobs_workspace_created", "software_jobs", ["workspace_id", "created_at"]) + + +def downgrade() -> None: + op.drop_index("ix_software_jobs_workspace_created", table_name="software_jobs") + op.drop_constraint("fk_software_jobs_source_job_id", "software_jobs", type_="foreignkey") + op.drop_constraint("fk_software_jobs_workspace_id", "software_jobs", type_="foreignkey") + op.drop_column("software_jobs", "export_outputs") + op.drop_column("software_jobs", "export_status") + op.drop_column("software_jobs", "preview_manifest") + op.drop_column("software_jobs", "local_manifest") + op.drop_column("software_jobs", "source_job_id") + op.drop_column("software_jobs", "workspace_id") + op.drop_constraint( + "fk_software_workspaces_head_job_id", + "software_workspaces", + type_="foreignkey", + ) + op.drop_index("ix_software_workspaces_node_status", table_name="software_workspaces") + op.drop_index("ix_software_workspaces_user_status", table_name="software_workspaces") + op.drop_table("software_workspaces") diff --git a/software-contracts/ansys.mechanical.static_structural.v1.json b/software-contracts/ansys.mechanical.static_structural.v1.json index 9e8f1f5..a60563c 100644 --- a/software-contracts/ansys.mechanical.static_structural.v1.json +++ b/software-contracts/ansys.mechanical.static_structural.v1.json @@ -75,6 +75,7 @@ "title_path": ["operation", "analysis", "title"] }, "legacy_runtime": null, + "workspace": null, "request_schema": { "$schema": "https://json-schema.org/draft/2020-12/schema", "type": "object", diff --git a/software-contracts/origin.plot.v2.json b/software-contracts/origin.plot.v2.json index acc2722..bc1887d 100644 --- a/software-contracts/origin.plot.v2.json +++ b/software-contracts/origin.plot.v2.json @@ -15,14 +15,14 @@ "media_type": "application/x-origin-project", "relative_path": "project.opju", "publish": true, - "required": false + "required": true }, "figure_png": { "filename": "figure.png", "media_type": "image/png", "relative_path": "figure.png", "publish": true, - "required": false + "required": true }, "figure_svg": { "filename": "figure.svg", @@ -87,6 +87,11 @@ "slots_path": ["available_slots"], "assumed_adapter_version": "0.3.0" }, + "workspace": { + "state_output": "project", + "preview_outputs": ["figure_png"], + "required_adapter_version": "1.0.0" + }, "request_schema": { "$schema": "https://json-schema.org/draft/2020-12/schema", "type": "object", diff --git a/tests/test_software_contracts.py b/tests/test_software_contracts.py index 1600d9a..8d8a032 100644 --- a/tests/test_software_contracts.py +++ b/tests/test_software_contracts.py @@ -39,9 +39,12 @@ class SoftwareContractTests(unittest.TestCase): self.assertEqual(normalized, request) self.assertEqual(len(digest), 64) self.assertEqual(contract.output_namespace, "origin") + self.assertEqual(contract.workspace.state_output, "project") + self.assertEqual(contract.workspace.preview_outputs, ("figure_png",)) + self.assertEqual(contract.workspace.required_adapter_version, "1.0.0") self.assertEqual( set(contract.expected_outputs(request)), - {"figure_png", "plot_spec", "provenance"}, + {"project", "figure_png", "plot_spec", "provenance"}, ) self.assertEqual(DEFAULT_CAPABILITIES, ("origin.plot@v2",)) @@ -58,7 +61,7 @@ class SoftwareContractTests(unittest.TestCase): "origin.plot@v2": { "health": "ready", "available_slots": 1, - "adapter_version": "0.5.0", + "adapter_version": "1.0.0", "features": ["line", "heatmap"], } } @@ -85,7 +88,7 @@ class SoftwareContractTests(unittest.TestCase): }) self.assertFalse(node_supports_request(contract, multi_panel, current_runtime)) current_runtime["capability_runtime"]["origin.plot@v2"].update({ - "adapter_version": "0.8.1", + "adapter_version": "1.0.0", }) self.assertTrue(node_supports_request(contract, multi_panel, current_runtime)) @@ -102,7 +105,7 @@ class SoftwareContractTests(unittest.TestCase): }) self.assertFalse(node_supports_request(contract, recipe, current_runtime)) current_runtime["capability_runtime"]["origin.plot@v2"].update({ - "adapter_version": "0.9.0", + "adapter_version": "1.0.0", }) self.assertTrue(node_supports_request(contract, recipe, current_runtime)) @@ -121,7 +124,7 @@ class SoftwareContractTests(unittest.TestCase): self.assertEqual(normalized, bubble) self.assertFalse(node_supports_request(contract, bubble, current_runtime)) current_runtime["capability_runtime"]["origin.plot@v2"].update({ - "adapter_version": "0.7.0", + "adapter_version": "1.0.0", "features": ["bubble"], }) self.assertTrue(node_supports_request(contract, bubble, current_runtime)) @@ -132,7 +135,7 @@ class SoftwareContractTests(unittest.TestCase): self.assertEqual(normalized, area) self.assertFalse(node_supports_request(contract, area, current_runtime)) current_runtime["capability_runtime"]["origin.plot@v2"].update({ - "adapter_version": "0.8.0", + "adapter_version": "1.0.0", "features": ["area"], }) self.assertTrue(node_supports_request(contract, area, current_runtime)) diff --git a/tests/test_software_job_tools.py b/tests/test_software_job_tools.py index 6bdc2e0..1de10b3 100644 --- a/tests/test_software_job_tools.py +++ b/tests/test_software_job_tools.py @@ -11,6 +11,7 @@ from core.software_jobs import revise_job from tools.software_jobs import ( SoftwareCapabilityListTool, SoftwareJobCancelTool, + SoftwareJobExportTool, SoftwareJobReviseTool, SoftwareJobStatusTool, SoftwareJobSubmitTool, @@ -197,6 +198,8 @@ class SoftwareJobToolTests(unittest.TestCase): def test_revise_service_copies_inputs_and_revalidates_as_new_job(self): source_job_id = uuid4() source = SimpleNamespace( + status="succeeded", + workspace_id=uuid4(), capability="origin.plot@v2", request={ "schema_version": 2, @@ -238,6 +241,8 @@ 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["workspace_id"], source.workspace_id) + self.assertEqual(create.call_args.kwargs["source_job_id"], source_job_id) self.assertEqual(create.call_args.kwargs["request"], { "schema_version": 2, "inputs": source.request["inputs"], @@ -262,6 +267,28 @@ class SoftwareJobToolTests(unittest.TestCase): self.assertEqual(result["status"], "cancelled") request_cancel.assert_called_once_with(self.user_id, job_id) + def test_export_requires_explicit_output_identities(self): + job_id = uuid4() + current = {"job_id": str(job_id), "task_id": str(self.task_id)} + exported = { + **current, + "export_status": "pending", + "export_outputs": ["project"], + } + with ( + patch("tools.software_jobs.get_job", return_value=current), + patch( + "tools.software_jobs.request_job_export", + return_value=exported, + ) as request_export, + ): + result = json.loads(SoftwareJobExportTool( + self.user_id, self.task_id + ).execute(str(job_id), ["project"])) + self.assertEqual(result["export_status"], "pending") + self.assertEqual(result["export_delivery"], "automatic") + request_export.assert_called_once_with(self.user_id, job_id, ["project"]) + if __name__ == "__main__": unittest.main() diff --git a/tests/test_software_nodes.py b/tests/test_software_nodes.py index 4a8a089..22dc277 100644 --- a/tests/test_software_nodes.py +++ b/tests/test_software_nodes.py @@ -217,6 +217,29 @@ class SoftwareNodeMigrationTests(unittest.TestCase): self.assertIn("followup_status", rendered) self.assertIn("ix_software_jobs_followup", rendered) + def test_0035_adds_workspace_lineage_preview_and_export_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.20260819_0900_0035_software_workspaces" + ) + with patch.object(migration, "op", operations): + migration.upgrade() + + rendered = "\n".join(statements) + self.assertIn("software_workspaces", rendered) + self.assertIn("workspace_id", rendered) + self.assertIn("source_job_id", rendered) + self.assertIn("preview_manifest", rendered) + self.assertIn("local_manifest", rendered) + self.assertIn("export_status", rendered) + self.assertIn("ix_software_jobs_workspace_created", rendered) + class SoftwareJobProtocolTests(unittest.TestCase): def test_succeeded_replay_uses_persisted_manifest_across_layout_versions(self) -> None: @@ -691,7 +714,7 @@ class SoftwareJobProtocolTests(unittest.TestCase): node.node_id = uuid4() node.capabilities = ["origin.plot@v2"] node.runtime = {"capability_runtime": {"origin.plot@v2": { - "health": "ready", "available_slots": 1, "adapter_version": "0.5.0", + "health": "ready", "available_slots": 1, "adapter_version": "1.0.0", "features": ["line"], }}} jobs = [] @@ -854,6 +877,7 @@ class SoftwareNodeDeleteTests(unittest.TestCase): session.get.return_value = node self.assertTrue(delete_node(uuid4())) + session.execute.assert_called_once() session.delete.assert_called_once_with(node) @patch("core.software_nodes.session_scope") diff --git a/tests/test_software_output_publish.py b/tests/test_software_output_publish.py index 454c527..b4b8f07 100644 --- a/tests/test_software_output_publish.py +++ b/tests/test_software_output_publish.py @@ -7,10 +7,45 @@ from pathlib import Path from unittest.mock import patch from uuid import uuid4 -from web.routers.software_nodes import _publish_software_job_outputs +from web.routers.software_nodes import ( + _publish_software_job_outputs, + _publish_software_job_previews, +) class SoftwareOutputPublishTests(unittest.TestCase): + def test_workspace_preview_uses_hidden_cache_without_artifact_registration(self) -> None: + with tempfile.TemporaryDirectory() as directory: + root = Path(directory) + job_id = uuid4() + staging = root / ".zcbot_software_job_staging" / str(job_id) + staging.mkdir(parents=True) + content = b"preview" + (staging / "figure.png").write_bytes(content) + manifest = [{ + "artifact_id": "figure_png", + "filename": "figure.png", + "media_type": "image/png", + "size_bytes": len(content), + "sha256": hashlib.sha256(content).hexdigest(), + }] + context = { + "capability": "origin.plot@v2", + "user_id": uuid4(), + } + with ( + patch("web.routers.software_nodes.load_user_root", return_value=root), + patch("web.routers.software_nodes.register_published_artifacts") as register, + ): + previews = _publish_software_job_previews(job_id, context, manifest) + self.assertFalse(register.called) + self.assertNotIn("artifact_id", previews[0]) + self.assertEqual(previews[0]["preview_id"], "figure_png") + self.assertEqual( + (root / ".zcbot_cache" / "software-previews" / str(job_id) / "figure.png").read_bytes(), + content, + ) + def test_complete_set_moves_atomically_and_can_be_replayed(self) -> None: with tempfile.TemporaryDirectory() as directory: root = Path(directory) diff --git a/tests/test_windows_node_source.py b/tests/test_windows_node_source.py index 99d6283..9dedb8b 100644 --- a/tests/test_windows_node_source.py +++ b/tests/test_windows_node_source.py @@ -35,6 +35,18 @@ class WindowsNodeSourceTests(unittest.TestCase): ): self.assertIn(marker, source) + def test_workspace_storage_is_global_by_workspace_not_user_directory(self) -> None: + models = (PROJECT / "NodeModels.cs").read_text(encoding="utf-8") + store = (PROJECT / "WorkspaceStore.cs").read_text(encoding="utf-8") + runner = (PROJECT / "AdapterProcessRunner.cs").read_text(encoding="utf-8") + self.assertIn('Path.Combine(root, "workspaces")', models) + self.assertIn('Binding(job).WorkspaceId.ToString("D")', store) + self.assertNotIn("UserId", store) + self.assertIn('Path.Combine(root, "current")', store) + self.assertIn('Path.Combine(root, "rollback")', store) + self.assertIn("workspaceStore.Promote", runner) + self.assertIn("workspaceStore.Restore", runner) + def test_runtime_auto_reports_all_discovered_adapter_capabilities(self) -> None: connection = (PROJECT / "NodeConnectionLoop.cs").read_text(encoding="utf-8") config_store = (PROJECT / "NodeConfigStore.cs").read_text(encoding="utf-8") @@ -353,8 +365,9 @@ class WindowsNodeSourceTests(unittest.TestCase): self.assertIn('DefaultRequestHeaders.Add("X-Lease-Id"', uploader) self.assertIn("SHA256.HashDataAsync", uploader) self.assertIn("upload-complete.json", connection + uploader) - self.assertIn("UploadAsync(refreshed, recoveringOutputs)", connection) - self.assertIn("manifest, allowConflict: true", uploader) + self.assertIn("adapter.PreviewOutputIds", connection) + self.assertIn("adapter.WorkspaceStateFilename", connection) + self.assertIn("WorkspaceSize(jobDirectory), allowConflict: true", uploader) self.assertIn("completeResponse.StatusCode == HttpStatusCode.Conflict", uploader) self.assertIn("OpenOutputWithRetryAsync", uploader) self.assertIn("IsSharingViolation", uploader) diff --git a/tools/software_jobs.py b/tools/software_jobs.py index 5a26524..11daa9d 100644 --- a/tools/software_jobs.py +++ b/tools/software_jobs.py @@ -17,6 +17,7 @@ from core.software_jobs import ( get_job_request, list_jobs, request_job_cancel, + request_job_export, revise_job, ) from core.software_nodes import list_nodes @@ -164,8 +165,9 @@ class SoftwareJobStatusTool(_SoftwareJobTool): name = "software_job_status" description = ( "Check one software job, or list recent jobs in the current task when job_id is omitted. " - "A succeeded job returns output_dir as the authoritative starting directory; " - "inspect or search within that directory to analyze its files." + "A succeeded workspace job returns preview_manifest and local_manifest. Local outputs " + "remain on the Windows node; call software_job_export only after an explicit user " + "request to download or deliver selected files." ) parameters = { "type": "object", @@ -195,8 +197,8 @@ class SoftwareJobReviseTool(_SoftwareJobTool): name = "software_job_revise" description = ( "Create a new professional-software job from a prior job in the current task. " - "Reuse the prior registered inputs, provide a complete replacement operation and " - "outputs, and leave the prior job and artifacts unchanged. Call software_job_status " + "Continue from the prior job's node-local workspace state, reuse its registered inputs, " + "and provide a complete next operation and outputs. Call software_job_status " "first when the prior editable_request is not already in context. A successful revision " "is the terminal action for the current run; completion and artifacts are delivered " "automatically in a later system-managed message." @@ -276,3 +278,46 @@ class SoftwareJobCancelTool(_SoftwareJobTool): return json.dumps(job, ensure_ascii=False) except (ValueError, SoftwareJobError) as exc: return f"[Error] {exc}" + + +class SoftwareJobExportTool(_SoftwareJobTool): + name = "software_job_export" + description = ( + "Export selected node-local outputs from the current successful workspace head. " + "Call this only when the user explicitly asks to download, receive, or deliver files. " + "The transfer is asynchronous and completion is reported automatically." + ) + parameters = { + "type": "object", + "properties": { + "job_id": {"type": "string", "format": "uuid"}, + "output_ids": { + "type": "array", + "minItems": 1, + "maxItems": 16, + "uniqueItems": True, + "items": {"type": "string"}, + "description": "Output identities from software_job_status.local_manifest.", + }, + }, + "required": ["job_id", "output_ids"], + "additionalProperties": False, + } + + def execute(self, job_id: str, output_ids: list[str] | None = None) -> str: + try: + current = get_job(self.user_id, UUID(job_id.strip())) + if current is None or current["task_id"] != str(self.task_id): + return "[Error] software job not found" + job = request_job_export( + self.user_id, + UUID(job_id.strip()), + output_ids or [], + ) + return json.dumps({ + **job, + "export_delivery": "automatic", + "next_action": "end_turn_after_reporting_export_queued", + }, ensure_ascii=False) + except (ValueError, SoftwareJobError) as exc: + return f"[Error] {exc}" diff --git a/web/routers/software_nodes.py b/web/routers/software_nodes.py index 794eddb..bad8a18 100644 --- a/web/routers/software_nodes.py +++ b/web/routers/software_nodes.py @@ -30,6 +30,8 @@ from core.software_contracts import ( from core.software_jobs import ( MAX_OUTPUT_ARTIFACT_BYTES, MAX_OUTPUT_TOTAL_BYTES, + MAX_EXPORT_ARTIFACT_BYTES, + MAX_EXPORT_TOTAL_BYTES, SoftwareJobError, abandon_offer, create_job, @@ -40,6 +42,8 @@ from core.software_jobs import ( mark_node_jobs_disconnected, offer_next_job, pending_node_cancellations, + pending_node_exports, + record_job_export, record_job_terminal, replay_succeeded_outputs, request_job_analysis, @@ -291,6 +295,121 @@ def _publish_software_job_outputs(job_id: UUID, context: dict, manifest: list[di ] +def _publish_software_job_previews( + job_id: UUID, context: dict, manifest: list[dict] +) -> list[dict]: + contract = get_contract(context["capability"]) + if contract.workspace is None: + raise SoftwareJobError("capability does not define workspace previews") + preview_ids = set(contract.workspace.preview_outputs) + selected = [item for item in manifest if item["artifact_id"] in preview_ids] + if len(selected) != len(preview_ids): + raise SoftwareJobError("software job preview manifest is incomplete") + root = load_user_root(context["user_id"]) + staging = safe_join(root, f".zcbot_software_job_staging/{job_id}") + destination = safe_join(root, f".zcbot_cache/software-previews/{job_id}") + _reject_symlink_path(root, staging) + _reject_symlink_path(root, destination) + destination.mkdir(parents=True, exist_ok=True) + previews: list[dict] = [] + for item in selected: + source = staging / item["filename"] + target = destination / item["filename"] + if target.is_file(): + if ( + target.stat().st_size != item["size_bytes"] + or _hash_file(target) != item["sha256"] + ): + raise SoftwareJobError("stored software preview conflicts with upload") + else: + if ( + not source.is_file() + or source.stat().st_size != item["size_bytes"] + or _hash_file(source) != item["sha256"] + ): + raise SoftwareJobError("uploaded software preview is missing or invalid") + os.replace(source, target) + previews.append({ + "preview_id": item["artifact_id"], + "filename": item["filename"], + "media_type": item["media_type"], + "size_bytes": item["size_bytes"], + "sha256": item["sha256"], + "url": f"/v1/software-jobs/{job_id}/previews/{item['artifact_id']}", + }) + if staging.is_dir(): + try: + staging.rmdir() + staging.parent.rmdir() + except OSError: + pass + return previews + + +def _publish_software_job_export( + job_id: UUID, context: dict, manifest: list[dict] +) -> list[dict]: + capability = context["capability"] + contract = get_contract(capability) + root = load_user_root(context["user_id"]) + working_dir = _task_working_dir(root, context["working_dir"]) + staging = safe_join(root, f".zcbot_software_job_staging/{job_id}") + relative_output = Path(contract.output_namespace) / str(job_id) + destination = safe_join(working_dir, relative_output.as_posix()) + _reject_symlink_path(root, staging) + _reject_symlink_path(root, destination) + destination.mkdir(parents=True, exist_ok=True) + for item in manifest: + source = staging / item["filename"] + target = destination / software_job_output_path( + capability, item["artifact_id"] + ) + target.parent.mkdir(parents=True, exist_ok=True) + if target.is_file(): + if target.stat().st_size != item["size_bytes"] or _hash_file(target) != item["sha256"]: + raise SoftwareJobError("published software export conflicts with upload") + else: + if ( + not source.is_file() + or source.stat().st_size != item["size_bytes"] + or _hash_file(source) != item["sha256"] + ): + raise SoftwareJobError("uploaded software export is missing or invalid") + os.replace(source, target) + refs = tuple({ + "path": ( + relative_output / software_job_output_path(capability, item["artifact_id"]) + ).as_posix(), + "label": item["filename"], + "media_type": item["media_type"], + } for item in manifest) + published_refs = register_published_artifacts( + user_id=context["user_id"], + task_id=context["task_id"], + user_root=root, + working_dir=working_dir, + refs=refs, + software_job_id=job_id, + ) + refs_by_path = {item["path"]: item for item in published_refs} + return [ + { + **item, + "source_artifact_id": item["artifact_id"], + "artifact_id": refs_by_path[ + ( + relative_output + / software_job_output_path(capability, item["artifact_id"]) + ).as_posix() + ]["artifact_id"], + "path": ( + relative_output / software_job_output_path(capability, item["artifact_id"]) + ).as_posix(), + } + for item in manifest + ] + + def register_software_node_routes(app, *, require_user, require_admin) -> None: @app.post( "/v1/software-nodes/enroll", @@ -368,22 +487,49 @@ def register_software_node_routes(app, *, require_user, require_admin) -> None: raise HTTPException(400, str(exc)) from exc if artifact_id not in requested_ids: raise HTTPException(400, "unsupported output artifact identity") + workspace_preview = ( + contract.workspace is not None + and artifact_id in contract.workspace.preview_outputs + ) + workspace_export = ( + contract.workspace is not None + and context.get("export_status") == "pending" + and artifact_id in (context.get("export_outputs") or []) + ) + if contract.workspace is not None and not (workspace_preview or workspace_export): + raise HTTPException(400, "only workspace previews may be uploaded automatically") filename = output_spec.filename - if not 1 <= x_content_length <= MAX_OUTPUT_ARTIFACT_BYTES: + max_artifact_bytes = ( + MAX_EXPORT_ARTIFACT_BYTES if workspace_export else MAX_OUTPUT_ARTIFACT_BYTES + ) + if not 1 <= x_content_length <= max_artifact_bytes: raise HTTPException(400, "output artifact size is invalid") if len(x_content_sha256) != 64 or any(c not in "0123456789abcdef" for c in x_content_sha256): raise HTTPException(400, "output artifact digest is invalid") - try: - if succeeded_output_upload_matches( - context, - artifact_id, - size_bytes=x_content_length, - digest=x_content_sha256, - ): - return - except SoftwareJobError as exc: - raise HTTPException(409, str(exc)) from exc + if contract.workspace is None: + try: + if succeeded_output_upload_matches( + context, + artifact_id, + size_bytes=x_content_length, + digest=x_content_sha256, + ): + return + except SoftwareJobError as exc: + raise HTTPException(409, str(exc)) from exc root = load_user_root(context["user_id"]) + if workspace_preview: + cached = safe_join( + root, + f".zcbot_cache/software-previews/{job_id}/{filename}", + ) + if cached.is_file(): + if ( + cached.stat().st_size == x_content_length + and _hash_file(cached) == x_content_sha256 + ): + return + raise HTTPException(409, "stored preview conflicts with uploaded artifact") working_dir = _task_working_dir(root, context["working_dir"]) published = safe_join( working_dir, @@ -410,7 +556,8 @@ def register_software_node_routes(app, *, require_user, require_admin) -> None: return raise HTTPException(409, "uploaded output conflicts with existing staging file") staged_total = sum(item.stat().st_size for item in staging.rglob("*") if item.is_file()) - if staged_total + x_content_length > MAX_OUTPUT_TOTAL_BYTES: + max_total_bytes = MAX_EXPORT_TOTAL_BYTES if workspace_export else MAX_OUTPUT_TOTAL_BYTES + if staged_total + x_content_length > max_total_bytes: raise HTTPException(413, "software job outputs exceed the total size limit") temporary = destination.with_name(destination.name + ".tmp-" + os.urandom(8).hex()) digest = sha256() @@ -419,7 +566,7 @@ def register_software_node_routes(app, *, require_user, require_admin) -> None: with temporary.open("xb") as handle: async for chunk in request.stream(): total += len(chunk) - if total > x_content_length or total > MAX_OUTPUT_ARTIFACT_BYTES: + if total > x_content_length or total > max_artifact_bytes: raise HTTPException(413, "output artifact exceeded declared size") digest.update(chunk) handle.write(chunk) @@ -452,6 +599,43 @@ def register_software_node_routes(app, *, require_user, require_admin) -> None: manifest = validate_output_manifest( context["capability"], context["request"], body.get("artifact_manifest") ) + contract = get_contract(context["capability"]) + if contract.workspace is not None: + if context["status"] == "succeeded": + return { + "status": "succeeded", + "artifact_manifest": context.get("artifact_manifest") or [], + "preview_manifest": context.get("preview_manifest") or [], + } + workspace_size = body.get("workspace_size_bytes") + if ( + not isinstance(workspace_size, int) + or isinstance(workspace_size, bool) + or workspace_size < 1 + or workspace_size > 10 * 1024**4 + ): + raise SoftwareJobError("software workspace size is invalid") + previews = await asyncio.to_thread( + _publish_software_job_previews, job_id, context, manifest + ) + terminal = { + "job_id": str(job_id), + "lease_id": str(lease_id), + "request_digest": x_request_digest, + "status": "succeeded", + "error": {}, + "artifact_manifest": [], + "local_manifest": manifest, + "preview_manifest": previews, + "workspace_size_bytes": workspace_size, + } + await asyncio.to_thread(record_job_terminal, node_id, terminal) + await dispatch_followup(request.app, job_id) + return { + "status": "succeeded", + "artifact_manifest": [], + "preview_manifest": previews, + } replayed = replay_succeeded_outputs(context, manifest) if replayed is not None: return {"status": "succeeded", "artifact_manifest": replayed} @@ -472,6 +656,105 @@ def register_software_node_routes(app, *, require_user, require_admin) -> None: raise HTTPException(409, str(exc)) from exc return {"status": "succeeded", "artifact_manifest": published} + @app.post( + "/v1/software-jobs/{job_id}/exports/complete", + tags=["software-nodes"], + ) + async def complete_software_job_export( + job_id: UUID, + request: Request, + authorization: str | None = Header(default=None), + x_node_id: str = Header(default=""), + x_lease_id: str = Header(default=""), + x_request_digest: str = Header(default=""), + ): + node_id, _, context = await asyncio.to_thread( + _authenticate_output_request, + job_id, authorization, x_node_id, x_lease_id, x_request_digest, + ) + body = await request.json() + submitted = body.get("artifact_manifest") if isinstance(body, dict) else None + requested = context.get("export_outputs") or [] + local_by_id = { + item.get("artifact_id"): item + for item in context.get("local_manifest") or [] + if isinstance(item, dict) + } + if ( + context.get("export_status") not in {"pending", "completed"} + or not isinstance(submitted, list) + or len(submitted) != len(requested) + or {item.get("artifact_id") for item in submitted if isinstance(item, dict)} + != set(requested) + ): + raise HTTPException(409, "software export manifest does not match the request") + identity_fields = ("filename", "media_type", "size_bytes", "sha256") + if any( + not isinstance(item, dict) + or item.get("artifact_id") not in local_by_id + or any( + item.get(field) != local_by_id[item["artifact_id"]].get(field) + for field in identity_fields + ) + for item in submitted + ): + raise HTTPException(409, "software export metadata changed after execution") + if context.get("export_status") == "completed": + return { + "status": "completed", + "artifact_manifest": context.get("artifact_manifest") or [], + } + try: + published = await asyncio.to_thread( + _publish_software_job_export, job_id, context, submitted + ) + job = await asyncio.to_thread( + record_job_export, node_id, job_id, requested, published + ) + await dispatch_followup(request.app, job_id) + return {"status": "completed", "artifact_manifest": job["artifact_manifest"]} + except (SoftwareJobError, KeyError, TypeError) as exc: + raise HTTPException(409, str(exc)) from exc + + @app.get( + "/v1/software-jobs/{job_id}/previews/{preview_id}", + tags=["software-jobs"], + ) + def read_software_job_preview( + job_id: UUID, + preview_id: str, + user_id: UUID = Depends(require_user), # noqa: B008 + ): + job = get_job(user_id, job_id) + if job is None: + raise HTTPException(404, "software job not found") + item = next( + ( + value for value in job.get("preview_manifest") or [] + if isinstance(value, dict) and value.get("preview_id") == preview_id + ), + None, + ) + if item is None: + raise HTTPException(404, "software job preview not found") + root = load_user_root(user_id) + target = safe_join( + root, + f".zcbot_cache/software-previews/{job_id}/{item['filename']}", + ) + if ( + not target.is_file() + or target.stat().st_size != item["size_bytes"] + or _hash_file(target) != item["sha256"] + ): + raise HTTPException(404, "software job preview is unavailable") + return FileResponse( + path=str(target), + filename=item["filename"], + media_type=item["media_type"], + headers={"Cache-Control": "private, max-age=3600"}, + ) + @app.websocket("/v1/software-nodes/connect") async def node_connect(websocket: WebSocket): try: @@ -497,6 +780,10 @@ def register_software_node_routes(app, *, require_user, require_admin) -> None: await node_connections.send_on( node_id, websocket, {"type": "job_cancel", "payload": cancel} ) + for export in await asyncio.to_thread(pending_node_exports, node_id): + await node_connections.send_on( + node_id, websocket, {"type": "job_export", "payload": export} + ) while True: message = await websocket.receive_json() message_type = message.get("type") @@ -560,6 +847,10 @@ def register_software_node_routes(app, *, require_user, require_admin) -> None: await node_connections.send_on( node_id, websocket, {"type": "job_cancel", "payload": cancel} ) + for export in await asyncio.to_thread(pending_node_exports, node_id): + await node_connections.send_on( + node_id, websocket, {"type": "job_export", "payload": export} + ) offer = await asyncio.to_thread( offer_next_job, await node_connections.node_ids() ) diff --git a/web/software_followups.py b/web/software_followups.py index 399482f..034b1a4 100644 --- a/web/software_followups.py +++ b/web/software_followups.py @@ -39,8 +39,43 @@ def _artifact_refs(manifest: list) -> list[dict]: return refs +def _preview_refs(manifest: list) -> list[dict]: + return [ + { + "path": item["filename"], + "label": item["filename"], + "preview_url": item["url"], + "kind": "software_preview", + } + for item in manifest + if isinstance(item, dict) + and item.get("filename") + and item.get("url") + ] + + def _report_text(job: SoftwareJob) -> str: contract = get_contract(job.capability) + if contract.workspace is not None and getattr(job, "workspace_id", None) is not None: + exported = [ + str(item.get("filename")) + for item in job.artifact_manifest + if isinstance(item, dict) and item.get("artifact_id") and item.get("filename") + ] + if exported: + return ( + f"专业软件文件已导出:{contract.display_name}\n\n" + f"- Job ID:`{job.job_id}`\n" + f"- 导出文件:{'、'.join(exported)}\n" + "- Workspace 工程状态仍保留在 Windows Node,可继续加工" + ) + return ( + f"专业软件任务已完成:{contract.display_name}\n\n" + f"- Job ID:`{job.job_id}`\n" + f"- Workspace ID:`{job.workspace_id}`\n" + "- 工程状态:已保存在 Windows Node,可继续加工\n" + "- 交付方式:需要源文件或最终格式时请明确提出导出" + ) output_dir = f"{contract.output_namespace}/{job.job_id}" names = [ str(item.get("filename")) @@ -58,6 +93,14 @@ def _report_text(job: SoftwareJob) -> str: def _analysis_prompt(job: SoftwareJob) -> str: contract = get_contract(job.capability) + if contract.workspace is not None and getattr(job, "workspace_id", None) is not None: + return ( + "[专业软件任务完成事件]\n" + f"任务 {job.job_id}({contract.display_name})已成功完成,工程保存在 " + f"Workspace {job.workspace_id}。请调用 software_job_status 获取输入、" + "本地输出清单和预览信息,结合原始数据向用户报告结果;只有用户明确要求" + "取回文件时才请求导出。" + ) output_dir = f"{contract.output_namespace}/{job.job_id}" return ( "[专业软件任务完成事件]\n" @@ -130,7 +173,11 @@ def claim_followup(job_id: UUID) -> FollowupClaim | None: task_id=task.task_id, idx=next_idx, payload={"role": "assistant", "content": _report_text(job)}, - artifact_refs=_artifact_refs(job.artifact_manifest), + artifact_refs=( + _artifact_refs(job.artifact_manifest) + if job.artifact_manifest + else _preview_refs(job.preview_manifest) + ), kind="software_job_report", )) job.followup_status = "completed" diff --git a/web/static/js/media.js b/web/static/js/media.js index 8e0966f..2c5788a 100644 --- a/web/static/js/media.js +++ b/web/static/js/media.js @@ -191,6 +191,9 @@ export function renderArtifactBarHtml(rels, inlineMode = true, taskId = "", lega const artifactAttr = ref.artifact_id ? ` data-artifact-id="${escapeHtml(String(ref.artifact_id))}"` : ""; + const previewAttr = ref.preview_url + ? ` data-preview-url="${escapeHtml(String(ref.preview_url))}"` + : ""; if (!rel) return ""; const name = String(ref.label || rel.split("/").pop() || rel); const cat = _categorize(rel); @@ -203,7 +206,7 @@ export function renderArtifactBarHtml(rels, inlineMode = true, taskId = "", lega if (inlineMode === true && (cat === "image" || cat === "video")) { // 占位元素;插入 DOM 后 upgradeMediaArtifacts 异步 fetch blob → 填 /