修复外部生成任务连接竞争

并列轮询业务任务、心跳和截止期限,避免单连接相互等待

在截止期限写失败态前释放业务与心跳 Future

新增连接竞争回归测试

明确队列告警展示投影、历史快照与同版本发布约束
This commit is contained in:
2026-07-16 06:13:10 +00:00
parent 830df44381
commit ad44aaf896
5 changed files with 123 additions and 29 deletions
@@ -85,6 +85,7 @@
- 背景:角色形象、图标 spritesheet 和 UI 素材提取会先得到带纯色背景的原图,再执行抠图或拆分;图片修改会先得到模型对齐尺寸的原始输出,角色动作会先得到绿幕预览视频,再抽帧和抠图。此前部分原始产物只登记到 OSS,或者要等后处理成功后才进入项目资源,用户无法在失败后找回已经生成成功的内容。
- 决策:凡一次资产生成任务产生多个具有独立复用价值的产物,后端必须把上游已返回的中间产物写入 OSS、`asset_object`、项目资源和账号素材库,再执行抠图、抽帧或拆分;未指定素材文件夹时进入默认“项目”文件夹。角色形象、图标 spritesheet 和 UI 素材提取在透明背景处理正常成功时同时保留纯色背景原图与透明后处理结果;透明背景处理最终失败时只保留已经持久化的 provider 原图,并按下一条降级规则收口。普通图片和图片修改的纯尺寸变换不属于独立产物:provider 回图保留在内存,变换成功只上传变换结果,变换失败只上传 provider 原图,整个流程只写一次 OSS 并只创建一个素材,不能制造重复“原始输出”。`nanobanana2` 使用标量清晰度档位和独立比例,保留 provider 输出尺寸,不按 `WIDTHxHEIGHT` 解析。角色动作把绿幕预览视频作为一个可复用素材保存,逐帧源图继续留在同一任务 OSS 路径,不把 32 至 48 帧逐张灌入素材库。去背景、音频等没有独立上游中间产物的任务不制造重复副本。
- 画布、成本与降级:有项目上下文的图片多产物继续由同一次 `canvasCompletion` 写入权威画布快照,正常成功时生成器 `generatedLayerId` 锚定主后处理结果。角色形象、图标 spritesheet 或 UI 素材提取已经保存 provider 原图、但透明背景处理最终失败时,任务以 `completed + warning` 收口,原图作为唯一主图完成画布占位;不写入不存在的透明处理图,图标和 UI 也不继续拆分。透明处理成功后的图标和 UI 图集自动拆分仍是 best-effort;识别或切片持久化失败继续完成整张透明图集,并在 inline、队列轮询和刷新后任务列表中提示非阻断 warning,不得借用失败错误字段。provider 原图或角色动作预览视频承载该任务的模型生成成本,抠图、逐帧处理、透明图集和切片等后处理派生产物的 `generation_cost_mud_points = 0`,避免把生图成本误显示成抠图成本;所有中间产物沿用所属任务的真实 `asset_kind`,角色原图仍为 `character`、图标和 UI 图集原图仍为 `icon-spritesheet`、角色动作预览仍为 `character-animation`,不得再写新的“原图类型”。后台素材查询按任务分页,最终产物作为父行并显示任务总成本,每个中间产物作为可展开的独立子行显示阶段生成器和阶段成本。扣费确认边界保持为 provider 成功,OSS、尺寸恢复和画布回填不延长退款保护。
- 2026-07-16 告警契约补充:inline / external v1 继续返回结构化原始诊断;queue 有意把通用 `warning``sliceWarning` 归一为展示就绪字符串,通用 `warning.reason` 原样保留,`sliceWarning.reason` 由 worker 添加“图集已生成,但自动拆分未完成:”前缀,摘要与 BFF 原样投影,Web 直接展示。历史值保留写入时快照,不按新格式回填或推断;该内部字符串契约通过 API/worker 与 Web 同一维护窗口、同版本发布收口,不增加混部兼容层。
- 影响范围:`server-rs/crates/api-server/src/editor_project.rs``character_animation_assets.rs`、外部生成任务摘要、图片画布完成快照、账号素材库和前端生成提示。
- 验证方式:覆盖中间产物登记先于后处理、默认素材文件夹、图集拆分降级、inline / queue warning 和主结果锚定的定向测试,并运行 `cargo check -p api-server --manifest-path server-rs/Cargo.toml``npm run check:spacetime-schema`、前端定向测试、`npm run check:encoding``git diff --check`
- 关联文档:`docs/【编辑器】生成类面板Lovart统一改造方案-2026-06-17.md``docs/technical/【前端架构】图片画布编辑器MVP接入方案-2026-06-11.md`
@@ -44,7 +44,7 @@
- `GET /api/runtime/external-generation/jobs?limit=20&includeAcknowledgedTerminal=false`:当前账号正式生成任务列表,用于 `我的` 页签任务列表和完成 / 失败提示。返回每个任务的 job id、kind、source、可展示 label、状态、进度、错误、可选 `warning``priceMudPoints``refundLedgerId``notificationAcknowledgedAt` 和时间戳。默认不返回已确认的终态任务;需要拆分活跃和完成列表时可追加 `statuses=running,queued``statuses=completed,failed`BFF 仍只返回当前账号任务。
- 任务被 claim 后默认处于 `generating`,BFF 显示“正在生成”;真实进入 BgFilter、逐帧抠图或手动去背景时切换为 `processing`BFF 显示“正在处理”。旧任务 `phase=None``generating` 兼容,前端不得按耗时或 job kind 推断阶段。
- `POST /api/runtime/external-generation/jobs/acknowledge`:生成完成 / 失败提示展示后由前端后台调用,BFF 只传当前账号 job ids,后端只确认属于当前账号且已终态的任务。
- `GET /api/runtime/external-generation/jobs/{jobId}`:单 job 状态,用于生成页轮询某次动作。返回 `operationId`(即任务 ID)、`status``phaseLabel``phaseDetail``progress``error``updatedAtMicros`,以及可选的 `warning` 完整原因字符串。生成页轮询只依赖状态、阶段、进度、错误和警告;`jobKind`、source 和完整时间信息继续由任务列表接口或业务快照提供。`attempt` / `maxAttempts` 属于 worker 调度事实,不向该前端契约暴露;若未来需要面向用户展示,必须单独完成产品、契约和摘要投影设计。
- `GET /api/runtime/external-generation/jobs/{jobId}`:单 job 状态,用于生成页轮询某次动作。返回 `operationId`(即任务 ID)、`status``phaseLabel``phaseDetail``progress``error``updatedAtMicros`,以及可选、可直接展示`warning` 完整文案。生成页轮询只依赖状态、阶段、进度、错误和警告;`jobKind`、source 和完整时间信息继续由任务列表接口或业务快照提供。`attempt` / `maxAttempts` 属于 worker 调度事实,不向该前端契约暴露;若未来需要面向用户展示,必须单独完成产品、契约和摘要投影设计。
BFF 只做鉴权、授权裁剪、字段脱敏和契约映射;worker 调度、lease、执行和计费事实仍以 `external_generation_job` 为准,用户可见任务列表、单任务状态、执行阶段和通知确认的正式读取事实源为 `external_generation_job_summary`,业务结果仍以玩法 session / work profile 为准。生成页 / 进度页只展示当前玩法业务进度;用户可见任务列表放在 `我的` 页签,必要时再用单 job 状态补充排障信息,并继续按原玩法 session/detail 接口收敛到 ready 或 failed。队列接口不替代玩法恢复接口,也不把 private `request_payload_json` 原样传给前端。终态提示的弹出与否以后端 `notification_acknowledged_at` 为准;前端在提示展示后后台调用 acknowledge 接口,关闭按钮只负责收起本地弹窗,不能只靠本地 dismiss 永久吞掉任务。
@@ -196,7 +196,7 @@ controller 配置:
角色形象、图标 spritesheet 和 UI 素材提取在 provider 原图已经持久化后,如果透明背景处理最终失败,只用原图完成 `canvasCompletion`,不创建或回填透明处理图,图标和 UI 也不继续拆分,任务保持 `completed`。这个 source-only 降级只包住透明背景处理的最终失败;phase 上报、provider 原图持久化、透明处理图持久化或画布写回失败仍按任务错误传播。
inline 成功响应使用结构化 `warning.code/reason`queue worker 只把同一结构写入有界的 `result_payload_json.warning`,任务摘要提取其中`reason``warning_message`,单 job 状态和刷新后的任务列表 BFF 再以 `warning: string` 返回完整原因。图标 / UI 的透明图已经成功、只有自动拆分失败时,继续保留现有 `sliceWarning` 兼容契约;通用 `warning` 优先,只有不存在通用 `warning` 时,worker 和前端才把 `sliceWarning` 归一为自动拆分告警
inline 与 external v1 成功响应继续使用结构化 `warning.code/reason`图标 / UI 的透明图已经成功、只有自动拆分失败时,继续返回结构化 `sliceWarning.code/reason`,其中 `sliceWarning.reason` 保留原始诊断。queue worker 把两类告警归一为有界的 `result_payload_json.warning`:通用 `warning` 优先并原样保留完整 `reason`;只有不存在通用 `warning` 时,才给 `sliceWarning.reason` 添加“图集已生成,但自动拆分未完成:”前缀。任务摘要将该展示就绪`reason` 原样提取`warning_message`,单 job 状态和刷新后的任务列表 BFF 再以 `warning: string` 返回;Web 必须直接展示,不再补前缀或按 code 推断类型。历史任务保留写入时的 `reason` 快照,摘要 backfill 不按当前格式重新解释或补写前缀。该字符串语义是 worker / BFF / Web 的内部同版本契约,三者必须协调发布,不承诺滚动混部或旧 Web 缓存下的跨版本字符串兼容
## 验收
@@ -282,14 +282,14 @@ npm run check:server-rs-ddd
- 源码:`server-rs/crates/spacetime-module/src/external_generation.rs`
- 用途:外部生成 worker 的内部持久任务队列;`GENARRATIVE_EXTERNAL_GENERATION_MODE=queue` 时,`api-server` HTTP 角色只入队,`external-generation-worker` 角色通过 claim lease 领取、续租、执行,并用 `lease_token` 栅栏回写阶段、完成 / 失败。队列行继续保存 worker 执行、计费与滚动发布兼容所需字段,末尾可选 `phase` 只取 `generating / processing`claim 写 `generating`,真实进入抠图处理时由受 `job_id + worker_id + lease_token` 保护的 procedure 写 `processing`。phase procedure 以结构化结果区分 `LeaseFencingRejected``OtherRejected``LeaseFencingRejected` 立即终止,`OtherRejected` 以及 SDK 的 `Procedure` / `Runtime` 错误不重试,只有 `Build` / `ConnectDropped` / `Timeout` 在同一个 job attempt 内重试一次。该重试只重新上报 phase,不把任务写回 `pending`,也不重新调用 provider;编辑器 job 入队固定 `max_attempts=1`,第二次传输失败后任务进入 `failed`,不会回到 `pending` 或从 provider 生成起点重跑。用户可见任务列表、价格、状态、阶段、未确认终态数量和通知确认时间的正式读取事实源已经迁到 `external_generation_job_summary`;BFF 不得再为列表 / 详情 / acknowledge 读取该大表。拼图 `compile_puzzle_draft` 的前置 `compile_puzzle_agent_draft``generate_puzzle_images``generate_puzzle_ui_background` 的业务写回也在对应 SpacetimeDB transaction 内校验 `job_id + worker_id + lease_token`、job kind、owner 和 source entity,避免过期 worker 写 session / work profile;图片画布编辑器的 `editor_image_generation``editor_image_edit``editor_background_removal``editor_icon_spritesheet_generation``editor_ui_design_asset_extraction``editor_character_animation_generation``editor_video_generation``editor_sound_effect_generation``editor_background_music_generation` 复用同一队列表,worker 成功后经 `api-server` facade 写入 `editor_project_resource` / `editor_asset` / `editor_canvas.layers_json`,前端只通过 BFF job 状态轮询和项目快照读取恢复完成态。`GENARRATIVE_EXTERNAL_GENERATION_MODE=inline` 时不创建该队列行,三个 external generation guard 字段必须同时为空才允许 api-server 受控同步写回,半空 guard 仍会拒绝。worker 成功写回业务事实后才能 complete job;业务失败态写回成功后才能 fail job,失败态未写回时保留租约等待后续重领。
- 载荷约束:本次先对 `source_module = editor-canvas``request_payload_json` / `result_payload_json` 实施有限大小合法 JSON、任意层级禁止 `data:` / `blob:` 的双层门禁,只保存 worker 执行必需的普通参数和已登记媒体引用;其它玩法在完成各自参考图资源化之前不由本次门禁静默改变既有请求契约。该主表只供 worker claim / 执行和受控维护读取;正式用户任务列表、单任务状态、队列概览与 acknowledge 不得再返回或解析这两个 payload。
- 非阻断告警:角色形象、图标图集和 UI 素材提取已保存 provider 原图、但透明背景处理最终失败时,以原图唯一主图完成任务;透明图和切片不写入画布。这个 source-only 降级只包住透明背景处理的最终失败,phase 上报、provider 原图持久化、透明处理图持久化或画布写回失败仍按任务错误传播。图标 / UI 透明图集成功但自动拆分降级时仍保留透明图集和既有 `sliceWarning` 兼容契约;通用 `warning``sliceWarning` 互斥。两类成功降级都以既有 `completed` 状态收口,不新增状态值:source-only 的 inline 响应使用结构化 `warning.code/reason`,仅拆分失败的 inline 响应继续使用既有 `sliceWarning.code/reason`queue worker 才把两者归一为有界的 `result_payload_json.warning`,且通用 `warning` 优先不保存图片、切片列表或媒体 URL。
- 非阻断告警:角色形象、图标图集和 UI 素材提取已保存 provider 原图、但透明背景处理最终失败时,以原图唯一主图完成任务;透明图和切片不写入画布。这个 source-only 降级只包住透明背景处理的最终失败,phase 上报、provider 原图持久化、透明处理图持久化或画布写回失败仍按任务错误传播。图标 / UI 透明图集成功但自动拆分降级时仍保留透明图集;通用 `warning``sliceWarning` 互斥。两类成功降级都以既有 `completed` 状态收口,不新增状态值:source-only 的 inline / external v1 响应使用结构化 `warning.code/reason`,仅拆分失败的 inline / external v1 响应继续使用既有 `sliceWarning.code/reason`,其 `reason` 保留原始诊断queue worker 才把两者归一为有界的 `result_payload_json.warning`,且通用 `warning` 优先并原样保留完整 `reason`,只有 `sliceWarning.reason` 由 worker 添加“图集已生成,但自动拆分未完成:”前缀。队列结果不保存图片、切片列表或媒体 URL。
### `external_generation_job_summary`
- Rust 结构体:`ExternalGenerationJobSummary`
- 源码:`server-rs/crates/spacetime-module/src/external_generation.rs`
- 用途:外部生成正式任务列表的轻量投影,按 `job_id` 保存 owner、来源、状态、可选 `phase`、价格、有界错误摘要、通知确认时间、各阶段时间和入队时提取的 `request_prompt`,不包含 request/result payload、worker lease 或 dedupe 内部字段。错误摘要统一拒绝内联媒体并限制为 2048 字符;列表在单次 owner 扫描中同时计数并只保留请求 limit 的固定大小 top-N,不得先收集全量历史再截断。enqueue、claim、renew、phase update、complete、fail 事务同步投影;acknowledge 只更新该轻量表并写审计事件,后续主任务同步必须保留已有确认时间,禁止为了写确认时间加载 / 重写大 payload 行。BFF 的列表、状态和确认只调用 summary procedure`running + processing` 映射为“正在处理”,其它 running(含旧行 `phase=None`)映射为“正在生成”。历史终态任务由迁移操作员的游标分批 maintenance procedure 在压缩 payload 时同步回填摘要,正式列表不得为兼容旧数据回扫完整主表。
- 非阻断告警:摘要字段 `warning_message` 由完成任务的轻量 `result_payload_json.warning.reason` 提取,complete 和历史 backfill 共用同一构建路径;单 job 状态和任务列表 BFF 以 `warning: string` 返回该完整原因,不再返回结构化 code。错误与告警摘要都不复制内联媒体并限制为 2048 字符。`phase``warning_message` 分别表示当前执行阶段和成功降级提示,不得混用。
- 非阻断告警:摘要字段 `warning_message` 是展示投影,由完成任务的轻量 `result_payload_json.warning.reason` 原样提取,不等同于公开 inline / external v1 的原始结构化诊断字段。complete 和历史 backfill 共用同一构建路径;历史任务按其结果载荷中已写入的 `reason` 快照投影,不为格式升级重写或补前缀。单 job 状态和任务列表 BFF 以 `warning: string` 返回该可直接展示的完整文案,不再返回结构化 code,Web 不得再次补前缀或按字符串推断告警类型。错误与告警摘要都不复制内联媒体并限制为 2048 字符。`phase``warning_message` 分别表示当前执行阶段和成功降级提示,不得混用worker / BFF / Web 必须同版本协调发布,不保证滚动混部或旧 Web 缓存下的字符串语义兼容
- 正式读取 procedure 为 `get_external_generation_job_summary_and_return``list_external_generation_job_summaries_and_return``acknowledge_external_generation_job_summaries_and_return`。历史维护 procedure 为 `compact_external_generation_job_payloads_and_return``backfill_external_generation_job_summaries_and_return`,仅 migration operator 可调用;运维入口统一使用 `npm run spacetime:external-generation:maintain -- ...`,默认 dry-run、单批最多 25 条。B-tree cursor 选择阶段最多反序列化 `limit + 1` 行,apply 再按主键逐条读取选中行;怀疑存在单行异常巨型 JSON 时必须先使用 `--limit 1`。payload 压缩额外固定使用 `source_module = editor-canvas` 的复合 cursor 索引,不得静默改写其它玩法历史任务。
### `external_generation_job_event`
@@ -76,7 +76,7 @@ lease 过期后不代表任务一定再次执行:claim transaction 只有在 `
自 2026-07-11 起,`Genarrative-Full-Build-And-Deploy` 的每日 04:00 timer 默认以 `DEPLOY_TARGET=development``STDB_API_ROLLOUT_MODE=normal` 对仅供开发使用的 dev 服务器执行 Stdb → API → Web 完整发布,不进入人工 rollout gate。三个下游 Build 都由 Full Job 显式传 `PUBLISH_AFTER_BUILD=false`,不得依赖下游 Job 默认值或提前各自发布;统一 Build 完成后仍由 Full Job 按固定顺序发布。人工维护窗口才选择 `pause-after-stdb`,且必须配置 `STDB_API_ROLLOUT_APPROVERS`。上文“定时构建缺少审批人时失败”的旧口径不再作为当前 dev 定时发布行为。
Full Job 通过 `EXIT_MAINTENANCE_MODE_AFTER_COMPLETION` 明确选择完整发布成功后是否退出维护,默认勾选以保持历史行为。Full 对 Stdb Publish 和 API Deploy 两个下游阶段都固定传 `KEEP_MAINTENANCE_MODE=true`,让 maintenance marker 持续覆盖 Stdb → API → Web 整段发布;Web Deploy 成功后才进入独立 `Exit Maintenance` 阶段。取消勾选时跳过最终退出阶段,便于内网验收完成后人工恢复公网。`Genarrative-Api-Deploy` 也单独暴露 `KEEP_MAINTENANCE_MODE` 参数,并转换为随发布包脚本的 `--keep-maintenance-mode`;失败路径仍按既有 current 切换边界保留或退出维护,不受成功态选项覆盖。
Full Job 通过 `EXIT_MAINTENANCE_MODE_AFTER_COMPLETION` 明确选择完整发布成功后是否退出维护,默认勾选以保持历史行为。Full 对 Stdb Publish 和 API Deploy 两个下游阶段都固定传 `KEEP_MAINTENANCE_MODE=true`,让 maintenance marker 持续覆盖 Stdb → API → Web 整段发布;Web Deploy 成功后才进入独立 `Exit Maintenance` 阶段。取消勾选时跳过最终退出阶段,便于内网验收完成后人工恢复公网。`Genarrative-Api-Deploy` 也单独暴露 `KEEP_MAINTENANCE_MODE` 参数,并转换为随发布包脚本的 `--keep-maintenance-mode`;失败路径仍按既有 current 切换边界保留或退出维护,不受成功态选项覆盖。外部生成 queue 的 `warning` 由 API/worker 固化为可直接展示的完整文案,Web 不再补前缀,因此 API/worker 与 Web 必须在同一维护窗口按同一版本协调发布;分开运行 Job 时先保持维护态完成 API/worker,再发布 Web,二者完成后才能恢复公网,不得在公网可用期间只滚动其中一侧。
需要验证“更新 API 不停 worker”和“worker 是否持续消费队列”时,优先使用隔离容器 smoke:`npm run container:worker-smoke -- smoke`。该脚本生成 gitignored 的 `deploy/container/worker-smoke/api-server.env`,启动独立 compose project 与独立 SpacetimeDB,发布当前 `spacetime-module` 后写入 `worker_smoke_unsupported` 测试 job;预期 worker claim 后执行 unsupported 失败分支,再执行 API-only recreate 并确认 worker 容器 ID 不变,最后再次入队验证 API 更新后队列仍可消费。`external_generation_job` 是 private table,脚本通过 worker 日志确认 job_id 被消费,不用 CLI SQL 查询私表。该 smoke 不读取 `.env.local`,也不依赖真实 VectorEngine / OSS 密钥;真实生图链路联调再在本地私有 env 中补齐 provider 配置。worker-smoke 默认把本机 `spacetime` CLI 打成轻量 SpacetimeDB 镜像,避免本机首次 smoke 依赖官方大镜像下载。若容器内 Cargo 拉取 crates.io 依赖不稳定,可用 `npm run container:worker-smoke -- smoke --local-binary` 让容器内 Cargo 复用本机 Cargo 缓存构建当前二进制,再打入 Debian bookworm smoke runtime 临时镜像;可用 `GENARRATIVE_WORKER_SMOKE_LOCAL_BASE_IMAGE` 覆盖运行时基础镜像;若隔离端口或库数据需要重建,追加 `--force`。完成 queue 链路验证时,还要用队列概览 BFF 和单 job 状态接口确认 job 从 queued/running 收敛,并用对应玩法 session/detail 接口确认业务状态同步完成。
@@ -8,10 +8,7 @@ use spacetime_client::{
ExternalGenerationJobFailRecordInput, ExternalGenerationJobRecord,
ExternalGenerationJobRenewLeaseRecordInput, ExternalGenerationQueueWakeSubscription,
};
use tokio::{
task::JoinSet,
time::{Instant, sleep},
};
use tokio::{task::JoinSet, time::sleep};
use tracing::{error, info, warn};
const MAX_EDITOR_GENERATION_WARNING_CHARS: usize = 2_048;
@@ -331,32 +328,69 @@ async fn process_external_generation_job(
billing_price_mud_points,
process_external_generation_job_once(state.clone(), worker_id.clone(), job.clone()),
);
let heartbeat =
maintain_external_generation_job_lease(&state, &worker_id, &job, lease, heartbeat_interval);
match await_external_generation_job_execution(work, heartbeat, job_timeout).await {
ExternalGenerationJobExecutionOutcome::Finished(result) => result,
ExternalGenerationJobExecutionOutcome::TimedOut => {
// work 与 heartbeat future 已在调度函数返回前被 drop,它们持有的
// SpacetimeDB 连接租约也已释放;此时才写失败态,避免单连接池自等待。
let message = external_generation_worker_timeout_message(&job, job_timeout);
warn!(
job_id = %job.job_id,
job_kind = %job.job_kind,
timeout_seconds = job_timeout.as_secs(),
"external generation worker 任务超过执行预算,停止当前尝试并释放 worker 槽位"
);
fail_job(&state, &worker_id, &job, message.clone()).await?;
Err(message)
}
ExternalGenerationJobExecutionOutcome::LeaseRenewalFailed(error) => Err(error),
}
}
#[derive(Debug, PartialEq, Eq)]
enum ExternalGenerationJobExecutionOutcome {
Finished(Result<(), String>),
TimedOut,
LeaseRenewalFailed(String),
}
async fn await_external_generation_job_execution<W, H>(
work: W,
heartbeat: H,
job_timeout: Duration,
) -> ExternalGenerationJobExecutionOutcome
where
W: Future<Output = Result<(), String>>,
H: Future<Output = Result<(), String>>,
{
tokio::pin!(work);
let heartbeat = sleep(heartbeat_interval);
tokio::pin!(heartbeat);
let job_deadline = sleep(job_timeout);
tokio::pin!(job_deadline);
tokio::select! {
biased;
result = &mut work => ExternalGenerationJobExecutionOutcome::Finished(result),
_ = &mut job_deadline => ExternalGenerationJobExecutionOutcome::TimedOut,
result = &mut heartbeat => match result {
Ok(()) => unreachable!("external generation heartbeat monitor should not finish"),
Err(error) => ExternalGenerationJobExecutionOutcome::LeaseRenewalFailed(error),
},
}
}
async fn maintain_external_generation_job_lease(
state: &AppState,
worker_id: &str,
job: &ExternalGenerationJobRecord,
lease: Duration,
heartbeat_interval: Duration,
) -> Result<(), String> {
loop {
tokio::select! {
biased;
result = &mut work => return result,
_ = &mut job_deadline => {
let message = external_generation_worker_timeout_message(&job, job_timeout);
warn!(
job_id = %job.job_id,
job_kind = %job.job_kind,
timeout_seconds = job_timeout.as_secs(),
"external generation worker 任务超过执行预算,停止当前尝试并释放 worker 槽位"
);
fail_job(&state, &worker_id, &job, message.clone()).await?;
return Err(message);
}
_ = &mut heartbeat => {
renew_job_lease(&state, &worker_id, &job, lease).await?;
heartbeat.as_mut().reset(Instant::now() + heartbeat_interval);
}
}
sleep(heartbeat_interval).await;
renew_job_lease(state, worker_id, job, lease).await?;
}
}
@@ -1373,6 +1407,65 @@ mod tests {
assert_eq!(message, "VECTOR_ENGINE_API_KEY 未配置");
}
#[tokio::test]
async fn worker_heartbeat_wait_does_not_stop_work_progress() {
let connection = std::sync::Arc::new(tokio::sync::Semaphore::new(1));
let work_started = std::sync::Arc::new(tokio::sync::Notify::new());
let work_connection = connection.clone();
let work_started_signal = work_started.clone();
let work = async move {
let permit = work_connection
.acquire_owned()
.await
.expect("work should acquire the only connection");
work_started_signal.notify_one();
tokio::time::sleep(Duration::from_millis(25)).await;
drop(permit);
Ok(())
};
let heartbeat_connection = connection.clone();
let heartbeat = async move {
work_started.notified().await;
let _permit = heartbeat_connection
.acquire_owned()
.await
.expect("heartbeat should acquire after work releases the connection");
std::future::pending::<Result<(), String>>().await
};
let outcome =
await_external_generation_job_execution(work, heartbeat, Duration::from_secs(1)).await;
assert_eq!(
outcome,
ExternalGenerationJobExecutionOutcome::Finished(Ok(()))
);
}
#[tokio::test]
async fn worker_deadline_drops_work_connection_before_failure_writeback() {
let connection = std::sync::Arc::new(tokio::sync::Semaphore::new(1));
let work_connection = connection.clone();
let work = async move {
let _permit = work_connection
.acquire_owned()
.await
.expect("work should acquire the only connection");
std::future::pending::<Result<(), String>>().await
};
let heartbeat = std::future::pending::<Result<(), String>>();
let outcome =
await_external_generation_job_execution(work, heartbeat, Duration::from_millis(25))
.await;
assert_eq!(outcome, ExternalGenerationJobExecutionOutcome::TimedOut);
assert!(
connection.try_acquire().is_ok(),
"deadline 返回前应先 drop work 并释放连接"
);
}
#[test]
fn editor_generation_result_payload_keeps_only_lightweight_slice_warning() {
let job = external_generation_job_record_fixture(Some("lease-1"));