收口图片生成任务超时预算

- 将四类 VectorEngine 图片任务纳入长任务预算并预留终态写回窗口
- 通过请求上下文传递绝对截止时间并约束发送重试与图片下载
- 允许显式下调单次请求超时并补齐预算回归测试
- 同步外部生成 Worker 运维文档与项目记忆
This commit is contained in:
2026-07-20 20:56:17 +08:00
parent 06c12f40b6
commit f748dc72c3
19 changed files with 846 additions and 96 deletions
@@ -4340,3 +4340,7 @@
- Vite 全量边界补充:`src/games/**`、`src/data/**`、`src/prompts/**`、旧顶层 App / Playground、旧路由和 `services/ai.ts` 必须由 pre-transform 门禁直接拒绝;所有同源 `/generated-*` 裸读在 dev 与生产统一为空 `404`,历史对象只经现役签名读取接口兼容,不允许 SPA fallback 伪装成资产成功响应。
- 前端混合根目录补充:`src/components`、`src/hooks`、`src/persistence`、`src/routing`、`src/services` 的根级文件实行现役白名单,Vite 与 ESLint 使用同一口径阻断旧 RPG / 玩法根文件;子目录仍按现役目录和退役目录分别管理,新增公共根文件必须显式登记。
- 影响范围:`PlatformEntryActiveFlowShell`、`PlatformActiveProfileView`、编辑器 / 项目搜索、平台 profile clients、`module-runtime`、`platform-llm`、Cargo workspace / resolve graph、Vite / ESLint / Rust 产物门禁及旧业务退役方案。
## 2026-07-20 VectorEngine 图片任务预算收口到 worker deadline
- 决策:`editor_image_generation`、`editor_image_edit`、`editor_icon_spritesheet_generation` 和 `editor_ui_design_asset_extraction` 使用默认 `1800s` long job 预算。worker 从同一起点计算绝对 job deadline,并向 provider 提前保留 `min(60s, job 预算 / 2)` 作为审计、OSS 和终态写回窗口。deadline 只经进程内 `RequestContext` 传递;VectorEngine 单 attempt 取配置 timeout 与剩余预算的较小值,退避加下一次 attempt 无法落在同一 deadline 内时停止重试,参考图和响应图片下载也受同一 deadline 限制。普通 HTTP / `inline` 保持无 deadline 行为;`VECTOR_ENGINE_IMAGE_REQUEST_TIMEOUT_MS` 默认仍为 `1000000`,配置加载层允许显式值更低。lease 续租 / fencing、迟到写回仲裁、attempt 耗尽和原子退款语义不变。
@@ -3240,3 +3240,9 @@
- 原因:现役模块的静态 import 会让 Vite、TypeScript 和打包器递归解析依赖;watch ignore 只停止监听,Tailwind source 只控制 class 扫描,tree-shaking 也发生在模块已经加载之后。经 barrel 只取一个公共函数尤其容易把同文件的旧导出一起带回图中。
- 处理:把仍在用的公共账号 / 钱包 / 设置能力迁到明确的现役 client 与 presentation model;Vite `pre` transform 对退役模块真实路径直接失败,ESLint 在现役源上增加 restricted imports。每次恢复公共 UI 后用 `tsc --listFilesOnly` 和全新浏览器 context 复核,不能用已有 HMR 会话判绿。
- 关联:`vite.config.ts`、`.eslintrc.cjs`、`src/services/platform-entry/`、`docs/technical/【架构下线】旧创作模板业务退役方案-2026-07-17.md`。
## VectorEngine 请求超时不能脱离 worker 绝对预算(2026-07-20)
- 现象:VectorEngine 单次请求超时大于 worker job 执行预算时,worker 已停止续租,provider 才超时或开始重试;最终 lease 过期、任务失败并退款,上游却可能继续消耗资源或迟到成功。
- 原因:单 attempt timeout、重试退避、图片下载与 worker / lease 分别使用独立的相对计时,没有共享同一绝对 deadline;只抬高 worker timeout 或单独压低 provider timeout 都无法保证留出终态写回窗口。
- 处理:实际调用 VectorEngine 的四类图片 job 使用 `1800s` long 预算;从 job 开始的同一起点派生 provider deadline,常规提前 `60s`、短预算提前一半。每次 attempt、退避、下一次 attempt 和图片下载都必须在该 deadline 内;普通 HTTP / `inline` 不伪造 worker deadline。修复时不改动 lease fencing、迟到写回仲裁和原子退款语义。
@@ -124,7 +124,11 @@ worker 配置:
- `GENARRATIVE_EXTERNAL_GENERATION_WORKER_POLL_INTERVAL_MS`:空队列轮询间隔。
- `GENARRATIVE_EXTERNAL_GENERATION_WORKER_LEASE_SECONDS`:任务 lease 时长,默认 `600`;worker 会按约三分之一 lease、最长 30 秒的间隔续租。该值应覆盖一次心跳网络抖动窗口,不需要大于完整外部生成链路耗时。
- `GENARRATIVE_EXTERNAL_GENERATION_WORKER_JOB_TIMEOUT_SECONDS`:普通外部生成 job 的执行预算,默认 `900`。超过预算后当前 worker 停止续租并释放 worker 槽位,但不取消已启动的业务 future,也不主动写入失败 / 重试状态;在途执行交由 lease fencing 仲裁:写回在租约有效期内到达则照常完成,否则被拒绝,租约过期后任务可被重新认领,attempt 耗尽时由认领事务原子标记失败并结算退款。
- `GENARRATIVE_EXTERNAL_GENERATION_WORKER_LONG_JOB_TIMEOUT_SECONDS`:视频、角色动作等长耗时 job 的执行预算,默认 `1800`。
- `GENARRATIVE_EXTERNAL_GENERATION_WORKER_LONG_JOB_TIMEOUT_SECONDS`:VectorEngine 图片生成 / 编辑、图标 spritesheet 生成、UI 素材提取以及角色动作、视频等长耗时 job 的执行预算,默认 `1800`。其中四类 VectorEngine 图片 job 固定为 `editor_image_generation`、`editor_image_edit`、`editor_icon_spritesheet_generation` 和 `editor_ui_design_asset_extraction`;手动去背景等不直接调用 VectorEngine 的 job 继续使用普通预算。
worker 在单次 job 开始执行时从同一个单调时钟起点计算绝对 `job deadline` 和更早的 `provider deadline`:常规情况下为终态审计、OSS 持久化及 `complete/fail` 回写保留 `60` 秒;当整个 job 预算小于 `120` 秒时,保留其一半,避免 provider 预算被全部吃掉。该 deadline 只通过进程内 `RequestContext` 传给 VectorEngine 图片调用,不写入 HTTP DTO、队列 payload 或 SpacetimeDB;普通 HTTP / `inline` 上下文没有 deadline,保持原有行为。
VectorEngine 每次发送的实际 timeout 取 `min(VECTOR_ENGINE_IMAGE_REQUEST_TIMEOUT_MS, provider 剩余预算)`。遇到可重试传输错误或 408 / 429 / 5xx 时,只有当剩余预算还容得下本次退避和下一次 attempt 才继续;否则立即停止重试并返回当前 provider 错误,deadline 耗尽时返回 timeout。同一绝对 deadline 同时覆盖参考图下载、provider 请求 / 响应和响应图片 URL 下载,不允许请求已返回后的图片下载越过 provider 预算。`VECTOR_ENGINE_IMAGE_REQUEST_TIMEOUT_MS` 默认仍为 `1000000`;配置加载层允许显式值低于默认值,不再在读取环境变量时强制抬升。
controller 配置:
@@ -136,7 +140,7 @@ controller 配置:
- `GENARRATIVE_EXTERNAL_GENERATION_CONTROLLER_SERVICE_TEMPLATE`:systemd worker 模板,默认 `genarrative-external-generation-worker@{}.service`。
- `GENARRATIVE_EXTERNAL_GENERATION_CONTROLLER_DRY_RUN`:只记录决策不执行 systemctl,默认 `false`。
动态缩扩容方式:生产默认由 `deploy/systemd/genarrative-external-generation-controller.service` 启动 `GENARRATIVE_PROCESS_ROLE=external-generation-controller`,controller 读取 `get_external_generation_queue_stats_and_return` 后对 `genarrative-external-generation-worker@N.service` 执行精确 `systemctl start/stop`;无需改变 HTTP 进程数。controller 只操作 `@1..@MAX` 中的缺口或最高编号多余实例,保留 `@1` 作为保底 worker。缩容或发布重启 worker 时,进程收到 SIGINT/SIGTERM 后会停止 claim 新任务并等待当前任务完成;若进程被硬杀、机器断电或超过 systemd `TimeoutStopSec`,未完成任务会在 lease 过期后被其它 worker 重新领取。若 worker 内业务 future 长时间无返回,执行预算到期后 worker 会停止续租并释放槽位,在途 future 继续运行至租约仲裁窗口;有效租约内的写回仍可完成,租约过期后才会由其它 worker 重新认领,避免客户端取消与服务端写回竞态。容器链路已有独立 `external-generation-worker` compose service;扩 worker 必须扩这个 worker service,不能只扩 `api-server` HTTP service。
动态缩扩容方式:生产默认由 `deploy/systemd/genarrative-external-generation-controller.service` 启动 `GENARRATIVE_PROCESS_ROLE=external-generation-controller`,controller 读取 `get_external_generation_queue_stats_and_return` 后对 `genarrative-external-generation-worker@N.service` 执行精确 `systemctl start/stop`;无需改变 HTTP 进程数。controller 只操作 `@1..@MAX` 中的缺口或最高编号多余实例,保留 `@1` 作为保底 worker。缩容或发布重启 worker 时,进程收到 SIGINT/SIGTERM 后会停止 claim 新任务并等待当前任务完成;若进程被硬杀、机器断电或超过 systemd `TimeoutStopSec`,未完成任务会在 lease 过期后被其它 worker 重新领取。VectorEngine 图片链路会先于整个 job 执行预算停止 provider 发送 / 重试,以便 worker 在有效 lease 内完成终态写回;若其他业务 future 仍长时间无返回,执行预算到期后 worker 会停止续租并释放槽位,在途 future 继续运行至租约仲裁窗口;有效租约内的写回仍可完成,租约过期后才会由其它 worker 重新认领,避免客户端取消与服务端写回竞态。本次预算收口不改变 lease 续租 / fencing、迟到写回仲裁、attempt 耗尽收口和原子退款语义。容器链路已有独立 `external-generation-worker` compose service;扩 worker 必须扩这个 worker service,不能只扩 `api-server` HTTP service。
## 已接入的拼图纵切
@@ -63,7 +63,7 @@ HTTP 角色的 `GENARRATIVE_SPACETIME_POOL_SIZE` 只表示 procedure / reducer
生产拆分角色时,`external-generation-worker` 和 `external-generation-controller` 的专属 env 示例会把 `GENARRATIVE_SPACETIME_POOL_SIZE` 覆盖为 `1`;非 HTTP 角色不创建 API 缓存读连接,只保留 `external_generation_job` 队列窄订阅作为响应式唤醒信号,实际抢占和扩缩容判断仍走 SpacetimeDB procedure。worker / controller 不执行模型定价 seed,启动时先调用受 runtime writer 鉴权的 queue-stats procedure 做只读预检,身份不匹配时 fail-fast;当前正式 systemd unit 通过共同加载 `/etc/genarrative/api-server.env` 继承同一 `GENARRATIVE_SPACETIME_TOKEN`,专属角色 env 示例不重复配置该 token。`GENARRATIVE_EXTERNAL_GENERATION_WORKER_POLL_INTERVAL_MS` 与 controller poll interval 只作为订阅失效、漏事件和 lease 过期这类时间条件的兜底,不作为正常领取任务的主路径。
生产 worker 默认 `GENARRATIVE_EXTERNAL_GENERATION_WORKER_LEASE_SECONDS=600`,只覆盖 worker 心跳抖动和短暂断连窗口,不再把 lease 当成完整任务时长;默认 `GENARRATIVE_EXTERNAL_GENERATION_WORKER_JOB_TIMEOUT_SECONDS=900`,角色动画 / 视频类长任务使用 `GENARRATIVE_EXTERNAL_GENERATION_WORKER_LONG_JOB_TIMEOUT_SECONDS=1800`。worker 在单次尝试超过执行预算后会停止续租并释放 worker 槽位,但不会取消已启动的业务 future 或主动写入失败 / 重试状态;在途执行由 lease fencing 仲裁,有效租约内写回仍可完成,租约过期后任务才可重新领取,attempt 耗尽时由认领事务标记失败并结算退款。生产部署和 provision 脚本会给 `/etc/genarrative/api-server.env` 与 `/etc/genarrative/external-generation-worker.env` 补齐这些变量;已有自定义值不覆盖,只会把历史旧默认 `3600` 迁移为 `600`。
生产 worker 默认 `GENARRATIVE_EXTERNAL_GENERATION_WORKER_LEASE_SECONDS=600`,只覆盖 worker 心跳抖动和短暂断连窗口,不再把 lease 当成完整任务时长;默认 `GENARRATIVE_EXTERNAL_GENERATION_WORKER_JOB_TIMEOUT_SECONDS=900`。`editor_image_generation`、`editor_image_edit`、`editor_icon_spritesheet_generation`、`editor_ui_design_asset_extraction` 四类 VectorEngine 图片任务与角色动画 / 视频类长任务使用 `GENARRATIVE_EXTERNAL_GENERATION_WORKER_LONG_JOB_TIMEOUT_SECONDS=1800`,手动去背景、音效和背景音乐继续使用普通预算。worker 在单次尝试超过执行预算后会停止续租并释放 worker 槽位,但不会取消已启动的业务 future 或主动写入失败 / 重试状态;在途执行由 lease fencing 仲裁,有效租约内写回仍可完成,租约过期后任务才可重新领取,attempt 耗尽时由认领事务标记失败并结算退款。生产部署和 provision 脚本会给 `/etc/genarrative/api-server.env` 与 `/etc/genarrative/external-generation-worker.env` 补齐这些变量;已有自定义值不覆盖,只会把历史旧默认 `3600` 迁移为 `600`。
lease 过期后不代表任务一定再次执行:claim transaction 只有在 `attempt < max_attempts` 时才会递增 attempt 并返回 worker;如果过期的是最终 attempt,则直接把 job 收口为 `failed`、清理 lease,并按入队冻结价格为当前 attempt 原子退款或写 cancellation intent。该终态任务不会再次进入 provider executor,迟到 consume 会被 settlement intent 拒绝。
@@ -153,9 +153,9 @@ spacetime sql <database> "SELECT * FROM runtime_setting LIMIT 1" --server http:/
本地 `spacetime` CLI / standalone 版本必须和 `server-rs/Cargo.toml` 里锁定的 `spacetimedb` 版本一致;当前统一版本为 `2.6.1`。若版本错配,procedure 返回值可能在宿主侧触发 `Failed to BSATN deserialize procedure return value`,api-server 最终表现为现役 settings、editor project 或 profile procedure 超时。排障时先运行 `spacetime --version`,再对照 `server-rs/Cargo.toml` 的 `spacetimedb = "..."`;遇到版本不匹配时直接执行 `spacetime version install <version> && spacetime version use <version>`,或在目标就是最新版本时执行 `spacetime version upgrade`,升级后重启 `npm run dev:spacetime` 再重试。当前 `scripts/dev.mjs` 会在启动和复用本地 SpacetimeDB 前写入并校验 `dev-spacetime-tool-version`。2.6.1 修复了 procedure context 中调用者 `Identity` / `ConnectionId` 始终为空的回归,依赖 `ctx.sender` 鉴权时必须同时确认宿主已升级。
本地 `.env`、`.env.local` 或 `.env.secrets.local` 修改后必须重启 `api-server` 才会生效;若已经通过 `npm run dev` 启动完整联调,可在该终端输入 `rs api-server`。排查图片编辑器 VectorEngine 生成链路时,确认 `VECTOR_ENGINE_BASE_URL`、`VECTOR_ENGINE_API_KEY` 和 `VECTOR_ENGINE_IMAGE_REQUEST_TIMEOUT_MS` 只在本地或服务器密钥文件中配置,不能写入 Git。VectorEngine `gpt-image-2` 图片协议、URL / base64 响应解析、远端图片下载和 provider 侧结构化日志在 `server-rs/crates/platform-image`;`api-server` 只做编辑器请求编排、OSS / asset 持久化、计费和失败审计落库。`platform-image` 会在 JSON 生成和 multipart 编辑请求发送前归一显式像素尺寸;若请求发送失败,先按同一 `request_id` 查看 provider 日志与 `external_api_call_failure.metadata_json.errorSource`,当前 multipart `/v1/images/edits` 单独强制 HTTP/1.1。
本地 `.env`、`.env.local` 或 `.env.secrets.local` 修改后必须重启 `api-server` 才会生效;若已经通过 `npm run dev` 启动完整联调,可在该终端输入 `rs api-server`。排查图片编辑器 VectorEngine 生成链路时,确认 `VECTOR_ENGINE_BASE_URL`、`VECTOR_ENGINE_API_KEY` 和 `VECTOR_ENGINE_IMAGE_REQUEST_TIMEOUT_MS` 只在本地或服务器密钥文件中配置,不能写入 Git。`VECTOR_ENGINE_IMAGE_REQUEST_TIMEOUT_MS` 是单次 attempt 的配置上限,默认 `1000000`;配置加载层允许显式值低于该默认值,不再在读取环境变量时强制抬高。VectorEngine `gpt-image-2` 图片协议、URL / base64 响应解析、远端图片下载和 provider 侧结构化日志在 `server-rs/crates/platform-image`;`api-server` 只做编辑器请求编排、OSS / asset 持久化、计费和失败审计落库。`platform-image` 会在 JSON 生成和 multipart 编辑请求发送前归一显式像素尺寸;若请求发送失败,先按同一 `request_id` 查看 provider 日志与 `external_api_call_failure.metadata_json.errorSource`,当前 multipart `/v1/images/edits` 单独强制 HTTP/1.1。
VectorEngine 图片生成 / 编辑在 `request_send` 阶段出现 `timeout`、`connect`、libcurl 35 SSL connect reset、libcurl 56 receive error / `unexpected eof while reading`、recv failure 等临时传输错误,或在 `upstream_status` 阶段收到 408 / 429 / 5xx(例如 Nginx HTML `502 Bad Gateway`)时,`platform-image` 会对同一请求最多发送 5 次;multipart 图片编辑每次重试都会重新构造 form,避免复用已消费的 body。日志中 `VectorEngine 图片请求发送失败,准备重试` 或 `VectorEngine 图片上游状态可重试,准备重试` 表示本次失败已进入下一次尝试;最终仍失败时才会写入 `external_api_call_failure` 并返回 504 / 502。排查生产失败时应同时统计 retry 前的尝试日志和最终 audit,避免把一次用户请求内的多次发送误判成多个用户请求。
VectorEngine 图片生成 / 编辑在 `request_send` 阶段出现 `timeout`、`connect`、libcurl 35 SSL connect reset、libcurl 56 receive error / `unexpected eof while reading`、recv failure 等临时传输错误,或在 `upstream_status` 阶段收到 408 / 429 / 5xx(例如 Nginx HTML `502 Bad Gateway`)时,`platform-image` 会对同一请求最多发送 5 次;multipart 图片编辑每次重试都会重新构造 form,避免复用已消费的 body。worker 从 job 开始的同一时钟起点计算绝对 deadline,常规保留最后 `60` 秒给审计、OSS 和终态写回;job 预算小于 `120` 秒时保留一半。VectorEngine 单次 attempt timeout 取配置值和剩余 provider 预算的较小值;退避后已没有下一次 attempt 的预算时立即停止重试。该 deadline 覆盖参考图、provider 请求 / 响应和响应图片下载的整次 provider future,但只在 worker 进程内通过 `RequestContext` 传递;普通 HTTP / `inline` 没有该 deadline,继续保持原有 timeout 和重试行为。日志中 `VectorEngine 图片请求发送失败,准备重试` 或 `VectorEngine 图片上游状态可重试,准备重试` 表示本次失败确有预算进入下一次尝试;预算耗尽或最终仍失败时才会写入 `external_api_call_failure` 并返回 504 / 502。排查生产失败时应同时统计 retry 前的尝试日志和最终 audit,避免把一次用户请求内的多次发送误判成多个用户请求。这项收口不修改 lease 续租 / fencing、迟到写回仲裁、attempt 耗尽与原子退款语义。
图片编辑器生成属于持久队列长任务:提交接口返回 job 后,前端通过 `/api/runtime/external-generation/jobs/{jobId}` 与编辑器项目资源状态收敛。生产排查小程序或 WebView `Failed to fetch` 时,若 Nginx access log 为 `499`、`upstream_status=-`,先按提交请求的 `request_id`、job id、worker 日志和 `external_api_call_failure` 对齐真实任务,不把客户端断开直接判定为 provider 失败。
+5 -10
View File
@@ -1090,9 +1090,8 @@ impl AppConfig {
if let Some(vector_engine_image_request_timeout_ms) =
read_first_positive_u64_env(&["VECTOR_ENGINE_IMAGE_REQUEST_TIMEOUT_MS"])
{
// 中文注释:VectorEngine image-2 实测可能超过 500 秒;旧环境文件中常见的 180 秒值不能再提前截断真实生图。
config.vector_engine_image_request_timeout_ms = vector_engine_image_request_timeout_ms
.max(DEFAULT_VECTOR_ENGINE_IMAGE_REQUEST_TIMEOUT_MS);
// 单次 attempt 上限允许按环境收短;worker 调用还会受整次任务的绝对 deadline 约束。
config.vector_engine_image_request_timeout_ms = vector_engine_image_request_timeout_ms;
}
if let Some(vector_engine_audio_request_timeout_ms) =
@@ -1530,9 +1529,8 @@ mod tests {
DEFAULT_EDITOR_BGFILTER_REQUEST_TIMEOUT_MS,
DEFAULT_EXTERNAL_GENERATION_WORKER_JOB_TIMEOUT_SECONDS,
DEFAULT_EXTERNAL_GENERATION_WORKER_LEASE_SECONDS,
DEFAULT_EXTERNAL_GENERATION_WORKER_LONG_JOB_TIMEOUT_SECONDS,
DEFAULT_VECTOR_ENGINE_IMAGE_REQUEST_TIMEOUT_MS, ExternalGenerationMode, LlmProvider,
ProcessRole, parse_bool, parse_external_generation_mode, parse_process_role,
DEFAULT_EXTERNAL_GENERATION_WORKER_LONG_JOB_TIMEOUT_SECONDS, ExternalGenerationMode,
LlmProvider, ProcessRole, parse_bool, parse_external_generation_mode, parse_process_role,
};
use std::{
fs,
@@ -1773,10 +1771,7 @@ mod tests {
config.vector_engine_base_url,
"https://vector.internal.example"
);
assert_eq!(
config.vector_engine_image_request_timeout_ms,
DEFAULT_VECTOR_ENGINE_IMAGE_REQUEST_TIMEOUT_MS
);
assert_eq!(config.vector_engine_image_request_timeout_ms, 210_000);
assert_eq!(
config.hyper3d_base_url,
"https://model.internal.example/api/v2"
@@ -1049,6 +1049,7 @@ mod tests {
base_url: "https://vector.example".to_string(),
api_key: "secret".to_string(),
request_timeout_ms: 180_000,
request_deadline: None,
external_api_audit_state: None,
external_api_audit_user_id: None,
external_api_audit_profile_id: None,
@@ -1,4 +1,9 @@
use std::{future::Future, io, pin::Pin, time::Duration};
use std::{
future::Future,
io,
pin::Pin,
time::{Duration, Instant},
};
use axum::Json;
use serde_json::{Value, json};
@@ -15,6 +20,8 @@ use tokio::{
use tracing::{error, info, warn};
const MAX_EDITOR_GENERATION_WARNING_CHARS: usize = 2_048;
// provider 必须先结束,给失败审计、计费结算和队列终态写回保留有效 lease 内的收尾窗口。
const EXTERNAL_GENERATION_WORKER_TERMINAL_WRITE_RESERVE: Duration = Duration::from_secs(60);
const EDITOR_GENERATION_SLICE_WARNING_PREFIX: &str = "图集已生成,但自动拆分未完成:";
const EDITOR_GENERATION_WARNING_REDACTED_MESSAGE: &str =
"自动拆分未完成(告警详情含内联媒体引用,已省略)";
@@ -315,6 +322,8 @@ async fn process_external_generation_job(
) -> Result<(), String> {
let heartbeat_interval = external_generation_worker_heartbeat_interval(lease);
let job_timeout = external_generation_worker_job_timeout(&state.config, job.job_kind.as_str());
let (job_deadline, provider_deadline) =
external_generation_worker_deadlines(Instant::now(), job_timeout);
let billing_price_mud_points = match external_generation_billing_price_mud_points(&job) {
Ok(price_mud_points) => price_mud_points,
Err(message) => {
@@ -338,14 +347,19 @@ async fn process_external_generation_job(
job.job_id.clone(),
billing_claim_attempt,
billing_price_mud_points,
process_external_generation_job_once(state.clone(), worker_id.clone(), job.clone()),
process_external_generation_job_once(
state.clone(),
worker_id.clone(),
job.clone(),
provider_deadline,
),
));
let heartbeat =
maintain_external_generation_job_lease(&state, &worker_id, &job, lease, heartbeat_interval);
match await_external_generation_job_execution(
join_external_generation_work(&mut work_handle),
heartbeat,
job_timeout,
job_deadline,
)
.await
{
@@ -446,7 +460,7 @@ enum ExternalGenerationJobExecutionOutcome {
async fn await_external_generation_job_execution<W, H>(
work: W,
heartbeat: H,
job_timeout: Duration,
job_deadline: Instant,
) -> ExternalGenerationJobExecutionOutcome
where
W: Future<Output = Result<(), String>>,
@@ -454,13 +468,13 @@ where
{
tokio::pin!(work);
tokio::pin!(heartbeat);
let job_deadline = sleep(job_timeout);
tokio::pin!(job_deadline);
let job_deadline_sleep = tokio::time::sleep_until(job_deadline.into());
tokio::pin!(job_deadline_sleep);
tokio::select! {
biased;
result = &mut work => ExternalGenerationJobExecutionOutcome::Finished(result),
_ = &mut job_deadline => ExternalGenerationJobExecutionOutcome::TimedOut,
_ = &mut job_deadline_sleep => ExternalGenerationJobExecutionOutcome::TimedOut,
result = &mut heartbeat => match result {
Ok(()) => unreachable!("external generation heartbeat monitor should not finish"),
Err(error) => ExternalGenerationJobExecutionOutcome::LeaseRenewalFailed(error),
@@ -485,6 +499,7 @@ async fn process_external_generation_job_once(
state: AppState,
worker_id: String,
job: ExternalGenerationJobRecord,
provider_deadline: Instant,
) -> Result<(), String> {
match job.job_kind.as_str() {
#[cfg(any())]
@@ -752,7 +767,7 @@ async fn process_external_generation_job_once(
return Err(message);
}
};
let request_context = worker_request_context(&job);
let request_context = worker_request_context(&job, provider_deadline);
match generate_editor_image_for_owner(
&state,
&request_context,
@@ -788,7 +803,7 @@ async fn process_external_generation_job_once(
return Err(message);
}
};
let request_context = worker_request_context(&job);
let request_context = worker_request_context(&job, provider_deadline);
match edit_editor_image_for_owner(
&state,
&request_context,
@@ -819,7 +834,7 @@ async fn process_external_generation_job_once(
}
};
payload.task_id.get_or_insert_with(|| job.job_id.clone());
let request_context = worker_request_context(&job);
let request_context = worker_request_context(&job, provider_deadline);
match remove_editor_image_background_for_owner(
&state,
&request_context,
@@ -849,7 +864,7 @@ async fn process_external_generation_job_once(
return Err(message);
}
};
let request_context = worker_request_context(&job);
let request_context = worker_request_context(&job, provider_deadline);
match generate_editor_icon_spritesheet_for_owner(
&state,
&request_context,
@@ -885,7 +900,7 @@ async fn process_external_generation_job_once(
return Err(message);
}
};
let request_context = worker_request_context(&job);
let request_context = worker_request_context(&job, provider_deadline);
match extract_editor_ui_design_assets_for_owner(
&state,
&request_context,
@@ -922,7 +937,7 @@ async fn process_external_generation_job_once(
return Err(message);
}
};
let request_context = worker_request_context(&job);
let request_context = worker_request_context(&job, provider_deadline);
match generate_editor_character_animation_for_owner(
state.clone(),
request_context,
@@ -954,7 +969,7 @@ async fn process_external_generation_job_once(
return Err(message);
}
};
let request_context = worker_request_context(&job);
let request_context = worker_request_context(&job, provider_deadline);
match generate_editor_video_for_owner(
state.clone(),
request_context,
@@ -985,7 +1000,7 @@ async fn process_external_generation_job_once(
return Err(message);
}
};
let request_context = worker_request_context(&job);
let request_context = worker_request_context(&job, provider_deadline);
match generate_editor_sound_effect_for_owner(
state.clone(),
request_context,
@@ -1016,7 +1031,7 @@ async fn process_external_generation_job_once(
return Err(message);
}
};
let request_context = worker_request_context(&job);
let request_context = worker_request_context(&job, provider_deadline);
match generate_editor_background_music_for_owner(
state.clone(),
request_context,
@@ -1112,13 +1127,17 @@ async fn fail_queue_job_after_worker_error(
Ok(())
}
fn worker_request_context(job: &ExternalGenerationJobRecord) -> RequestContext {
fn worker_request_context(
job: &ExternalGenerationJobRecord,
provider_deadline: Instant,
) -> RequestContext {
RequestContext::new(
format!("external-generation-worker-{}", job.job_id),
format!("external-generation-worker {}", job.job_kind),
std::time::Duration::ZERO,
false,
)
.with_external_call_deadline(provider_deadline)
}
fn editor_generation_worker_caller(
@@ -1409,11 +1428,29 @@ fn external_generation_worker_heartbeat_interval(lease: Duration) -> Duration {
Duration::from_millis(heartbeat_millis)
}
fn external_generation_worker_deadlines(
started_at: Instant,
job_timeout: Duration,
) -> (Instant, Instant) {
let terminal_write_reserve =
EXTERNAL_GENERATION_WORKER_TERMINAL_WRITE_RESERVE.min(job_timeout / 2);
let Some(job_deadline) = started_at.checked_add(job_timeout) else {
return (started_at, started_at);
};
let provider_deadline = started_at
.checked_add(job_timeout.saturating_sub(terminal_write_reserve))
.unwrap_or(started_at);
(job_deadline, provider_deadline)
}
fn external_generation_worker_job_timeout(config: &AppConfig, job_kind: &str) -> Duration {
match job_kind {
EDITOR_CHARACTER_ANIMATION_GENERATION_JOB_KIND | EDITOR_VIDEO_GENERATION_JOB_KIND => {
config.external_generation_worker_long_job_timeout
}
EDITOR_IMAGE_GENERATION_JOB_KIND
| EDITOR_IMAGE_EDIT_JOB_KIND
| EDITOR_ICON_SPRITESHEET_GENERATION_JOB_KIND
| EDITOR_UI_DESIGN_ASSET_EXTRACTION_JOB_KIND
| EDITOR_CHARACTER_ANIMATION_GENERATION_JOB_KIND
| EDITOR_VIDEO_GENERATION_JOB_KIND => config.external_generation_worker_long_job_timeout,
_ => config.external_generation_worker_job_timeout,
}
}
@@ -1556,8 +1593,12 @@ mod tests {
std::future::pending::<Result<(), String>>().await
};
let outcome =
await_external_generation_job_execution(work, heartbeat, Duration::from_secs(1)).await;
let outcome = await_external_generation_job_execution(
work,
heartbeat,
Instant::now() + Duration::from_secs(1),
)
.await;
assert_eq!(
outcome,
@@ -1578,9 +1619,12 @@ mod tests {
};
let heartbeat = std::future::pending::<Result<(), String>>();
let outcome =
await_external_generation_job_execution(work, heartbeat, Duration::from_millis(25))
.await;
let outcome = await_external_generation_job_execution(
work,
heartbeat,
Instant::now() + Duration::from_millis(25),
)
.await;
assert_eq!(outcome, ExternalGenerationJobExecutionOutcome::TimedOut);
assert!(
@@ -1884,7 +1928,7 @@ mod tests {
}
#[test]
fn worker_job_timeout_uses_long_budget_for_video_jobs() {
fn worker_job_timeout_uses_long_budget_for_image_and_video_jobs() {
let config = AppConfig {
external_generation_worker_job_timeout: Duration::from_secs(60),
external_generation_worker_long_job_timeout: Duration::from_secs(600),
@@ -1893,7 +1937,25 @@ mod tests {
assert_eq!(
external_generation_worker_job_timeout(&config, EDITOR_IMAGE_GENERATION_JOB_KIND),
Duration::from_secs(60)
Duration::from_secs(600)
);
assert_eq!(
external_generation_worker_job_timeout(&config, EDITOR_IMAGE_EDIT_JOB_KIND),
Duration::from_secs(600)
);
assert_eq!(
external_generation_worker_job_timeout(
&config,
EDITOR_ICON_SPRITESHEET_GENERATION_JOB_KIND
),
Duration::from_secs(600)
);
assert_eq!(
external_generation_worker_job_timeout(
&config,
EDITOR_UI_DESIGN_ASSET_EXTRACTION_JOB_KIND
),
Duration::from_secs(600)
);
assert_eq!(
external_generation_worker_job_timeout(
@@ -1906,6 +1968,61 @@ mod tests {
external_generation_worker_job_timeout(&config, EDITOR_VIDEO_GENERATION_JOB_KIND),
Duration::from_secs(600)
);
assert_eq!(
external_generation_worker_job_timeout(&config, EDITOR_BACKGROUND_REMOVAL_JOB_KIND),
Duration::from_secs(60)
);
assert_eq!(
external_generation_worker_job_timeout(
&config,
EDITOR_SOUND_EFFECT_GENERATION_JOB_KIND
),
Duration::from_secs(60)
);
}
#[test]
fn worker_provider_deadline_reserves_terminal_write_budget() {
let started_at = Instant::now();
let (job_deadline, provider_deadline) =
external_generation_worker_deadlines(started_at, Duration::from_secs(1_800));
assert_eq!(
job_deadline.duration_since(started_at),
Duration::from_secs(1_800)
);
assert_eq!(
provider_deadline.duration_since(started_at),
Duration::from_secs(1_740)
);
assert_eq!(
job_deadline.duration_since(provider_deadline),
Duration::from_secs(60)
);
let (short_job_deadline, short_provider_deadline) =
external_generation_worker_deadlines(started_at, Duration::from_secs(30));
assert_eq!(
short_job_deadline.duration_since(started_at),
Duration::from_secs(30)
);
assert_eq!(
short_provider_deadline.duration_since(started_at),
Duration::from_secs(15)
);
}
#[test]
fn worker_request_context_carries_provider_deadline() {
let job = external_generation_job_record_fixture(Some("lease-1"));
let provider_deadline = Instant::now() + Duration::from_secs(30);
let request_context = worker_request_context(&job, provider_deadline);
assert_eq!(
request_context.external_call_deadline(),
Some(provider_deadline)
);
}
#[test]
@@ -3005,6 +3005,7 @@ mod tests {
base_url: base_url.clone(),
api_key: api_key.clone(),
request_timeout_ms: 180_000,
request_deadline: None,
};
let http_client = platform_image::build_vector_engine_image_http_client(&settings)
.expect("构建 HTTP 客户端");
@@ -3324,6 +3325,7 @@ mod tests {
base_url: base_url.clone(),
api_key: api_key.clone(),
request_timeout_ms: 180_000,
request_deadline: None,
};
let http_client = platform_image::build_vector_engine_image_http_client(&settings)
.expect("构建 HTTP 客户端");
@@ -13,6 +13,7 @@ use platform_image::{
vector_engine_images_generation_url,
};
use serde_json::{Value, json};
use std::time::Instant;
use time::OffsetDateTime;
use crate::{
@@ -39,6 +40,7 @@ pub(crate) struct OpenAiImageSettings {
pub base_url: String,
pub api_key: String,
pub request_timeout_ms: u64,
pub request_deadline: Option<Instant>,
pub external_api_audit_state: Option<AppState>,
pub external_api_audit_user_id: Option<String>,
pub external_api_audit_profile_id: Option<String>,
@@ -52,6 +54,7 @@ impl std::fmt::Debug for OpenAiImageSettings {
.field("base_url", &self.base_url)
.field("api_key", &"<redacted>")
.field("request_timeout_ms", &self.request_timeout_ms)
.field("request_deadline_enabled", &self.request_deadline.is_some())
.field(
"external_api_audit_enabled",
&self.external_api_audit_state.is_some(),
@@ -107,6 +110,7 @@ pub(crate) fn require_openai_image_settings(
base_url: base_url.to_string(),
api_key: api_key.to_string(),
request_timeout_ms: state.config.vector_engine_image_request_timeout_ms.max(1),
request_deadline: None,
external_api_audit_state: Some(state.clone()),
external_api_audit_user_id: None,
external_api_audit_profile_id: None,
@@ -406,6 +410,7 @@ impl OpenAiImageSettings {
self.external_api_audit_user_id = user_id;
self.external_api_audit_profile_id = profile_id;
self.external_api_audit_request_id = Some(request_context.request_id().to_string());
self.request_deadline = request_context.external_call_deadline();
self
}
@@ -414,6 +419,7 @@ impl OpenAiImageSettings {
base_url: self.base_url.clone(),
api_key: self.api_key.clone(),
request_timeout_ms: self.request_timeout_ms.max(1),
request_deadline: self.request_deadline,
}
}
}
@@ -573,6 +579,35 @@ mod tests {
use super::*;
use base64::{Engine as _, engine::general_purpose::STANDARD as BASE64_STANDARD};
#[test]
fn external_api_audit_context_forwards_external_call_deadline() {
let deadline = Instant::now() + std::time::Duration::from_secs(30);
let request_context = RequestContext::new(
"request-1".to_string(),
"external-generation-worker editor_image_edit".to_string(),
std::time::Duration::ZERO,
false,
)
.with_external_call_deadline(deadline);
let settings = OpenAiImageSettings {
base_url: "https://vector.example".to_string(),
api_key: "test-key".to_string(),
request_timeout_ms: 1_000_000,
request_deadline: None,
external_api_audit_state: None,
external_api_audit_user_id: None,
external_api_audit_profile_id: None,
external_api_audit_request_id: None,
}
.with_external_api_audit_context(&request_context, None, None);
assert_eq!(settings.request_deadline, Some(deadline));
assert_eq!(
settings.provider_settings().request_deadline,
Some(deadline)
);
}
#[test]
fn gpt_image_2_generation_request_uses_create_model_without_reference_images() {
let body = build_openai_image_request_body(
@@ -597,6 +632,7 @@ mod tests {
base_url: "https://vector.example".to_string(),
api_key: "test-key".to_string(),
request_timeout_ms: 1_000_000,
request_deadline: None,
external_api_audit_state: None,
external_api_audit_user_id: None,
external_api_audit_profile_id: None,
@@ -606,6 +642,7 @@ mod tests {
base_url: "https://vector.example/v1".to_string(),
api_key: "test-key".to_string(),
request_timeout_ms: 1_000_000,
request_deadline: None,
external_api_audit_state: None,
external_api_audit_user_id: None,
external_api_audit_profile_id: None,
@@ -628,6 +665,7 @@ mod tests {
base_url: "https://vector.example".to_string(),
api_key: "test-key".to_string(),
request_timeout_ms: 1_000_000,
request_deadline: None,
external_api_audit_state: None,
external_api_audit_user_id: None,
external_api_audit_profile_id: None,
@@ -637,6 +675,7 @@ mod tests {
base_url: "https://vector.example/v1".to_string(),
api_key: "test-key".to_string(),
request_timeout_ms: 1_000_000,
request_deadline: None,
external_api_audit_state: None,
external_api_audit_user_id: None,
external_api_audit_profile_id: None,
@@ -659,6 +698,7 @@ mod tests {
base_url: "https://vector.example".to_string(),
api_key: "test-key".to_string(),
request_timeout_ms: 1_000_000,
request_deadline: None,
external_api_audit_state: None,
external_api_audit_user_id: None,
external_api_audit_profile_id: None,
@@ -30,6 +30,7 @@ pub(crate) struct PuzzleVectorEngineSettings {
pub(crate) base_url: String,
pub(crate) api_key: String,
pub(crate) request_timeout_ms: u64,
pub(crate) request_deadline: Option<std::time::Instant>,
pub(crate) external_api_audit_state: Option<AppState>,
pub(crate) external_api_audit_user_id: Option<String>,
pub(crate) external_api_audit_profile_id: Option<String>,
@@ -101,6 +102,7 @@ impl PuzzleVectorEngineSettings {
base_url: self.base_url.clone(),
api_key: self.api_key.clone(),
request_timeout_ms: self.request_timeout_ms,
request_deadline: self.request_deadline,
external_api_audit_state: self.external_api_audit_state.clone(),
external_api_audit_user_id: self.external_api_audit_user_id.clone(),
external_api_audit_profile_id: self.external_api_audit_profile_id.clone(),
@@ -117,6 +119,7 @@ impl PuzzleVectorEngineSettings {
self.external_api_audit_user_id = user_id;
self.external_api_audit_profile_id = profile_id;
self.external_api_audit_request_id = Some(request_context.request_id().to_string());
self.request_deadline = request_context.external_call_deadline();
self
}
}
@@ -193,6 +196,7 @@ pub(crate) fn require_puzzle_vector_engine_settings(
base_url: base_url.to_string(),
api_key: api_key.to_string(),
request_timeout_ms: state.vector_engine_image_request_timeout_ms().max(1),
request_deadline: None,
external_api_audit_state: Some(state.root_state().clone()),
external_api_audit_user_id: None,
external_api_audit_profile_id: None,
@@ -18,6 +18,7 @@ pub struct RequestContext {
operation: String,
request_started_at: Instant,
wants_envelope: bool,
external_call_deadline: Option<Instant>,
}
impl RequestContext {
@@ -34,9 +35,15 @@ impl RequestContext {
.checked_sub(elapsed_seed)
.unwrap_or_else(Instant::now),
wants_envelope,
external_call_deadline: None,
}
}
pub fn with_external_call_deadline(mut self, deadline: Instant) -> Self {
self.external_call_deadline = Some(deadline);
self
}
pub fn request_id(&self) -> &str {
&self.request_id
}
@@ -49,6 +56,10 @@ impl RequestContext {
self.wants_envelope
}
pub fn external_call_deadline(&self) -> Option<Instant> {
self.external_call_deadline
}
pub fn elapsed(&self) -> u64 {
self.request_started_at
.elapsed()
@@ -111,3 +122,20 @@ fn wants_api_envelope<B>(request: &HttpRequest<B>) -> bool {
.map(str::to_lowercase)
.is_some_and(|value| matches!(value.as_str(), "1" | "true" | "v1" | "envelope"))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn request_context_has_no_external_call_deadline_by_default() {
let context = RequestContext::new(
"request-1".to_string(),
"GET /healthz".to_string(),
Duration::ZERO,
false,
);
assert_eq!(context.external_call_deadline(), None);
}
}
@@ -0,0 +1,179 @@
use std::time::{Duration, Instant};
use super::{
audit::build_failure_audit, constants::VECTOR_ENGINE_PROVIDER, error::PlatformImageError,
};
pub(crate) fn effective_request_timeout_ms(
configured_timeout_ms: u64,
request_deadline: Option<Instant>,
) -> Option<u64> {
effective_request_timeout_ms_at(configured_timeout_ms, request_deadline, Instant::now())
}
fn effective_request_timeout_ms_at(
configured_timeout_ms: u64,
request_deadline: Option<Instant>,
now: Instant,
) -> Option<u64> {
let configured_timeout_ms = configured_timeout_ms.max(1);
let Some(request_deadline) = request_deadline else {
return Some(configured_timeout_ms);
};
let remaining = request_deadline.checked_duration_since(now)?;
let remaining_ms = u64::try_from(remaining.as_millis()).unwrap_or(u64::MAX);
if remaining_ms == 0 {
return None;
}
Some(configured_timeout_ms.min(remaining_ms))
}
pub(crate) fn retry_delay_fits_request_deadline(
request_deadline: Option<Instant>,
delay_ms: u64,
) -> bool {
retry_delay_fits_request_deadline_at(request_deadline, delay_ms, Instant::now())
}
fn retry_delay_fits_request_deadline_at(
request_deadline: Option<Instant>,
delay_ms: u64,
now: Instant,
) -> bool {
let Some(request_deadline) = request_deadline else {
return true;
};
let Some(remaining) = request_deadline.checked_duration_since(now) else {
return false;
};
let required = Duration::from_millis(delay_ms).saturating_add(Duration::from_millis(1));
remaining >= required
}
pub(crate) fn request_budget_exhausted_error(
request_url: &str,
operation: &str,
latency_ms: Option<u64>,
prompt_chars: Option<usize>,
reference_image_count: Option<usize>,
) -> PlatformImageError {
const ERROR_SOURCE: &str = "external request deadline elapsed";
let message = format!("{operation}:外部图片请求执行预算已耗尽");
let audit = build_failure_audit(
request_url,
operation,
"request_budget",
None,
None,
true,
false,
message.as_str(),
Some(ERROR_SOURCE.to_string()),
None,
latency_ms,
prompt_chars,
reference_image_count,
);
tracing::warn!(
provider = VECTOR_ENGINE_PROVIDER,
endpoint = %request_url,
failure_stage = "request_budget",
timeout = true,
elapsed_ms = latency_ms,
prompt_chars,
reference_image_count,
operation,
"VectorEngine 图片请求执行预算已耗尽"
);
PlatformImageError::Request {
provider: VECTOR_ENGINE_PROVIDER,
message,
endpoint: Some(request_url.to_string()),
timeout: true,
connect: false,
request: true,
body: false,
status_code: None,
source: Some(ERROR_SOURCE.to_string()),
audit: Some(audit),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn configured_timeout_is_preserved_without_deadline() {
let now = Instant::now();
assert_eq!(
effective_request_timeout_ms_at(1_000, None, now),
Some(1_000)
);
assert_eq!(effective_request_timeout_ms_at(0, None, now), Some(1));
}
#[test]
fn request_timeout_is_clipped_to_remaining_deadline_budget() {
let now = Instant::now();
assert_eq!(
effective_request_timeout_ms_at(10_000, Some(now + Duration::from_millis(250)), now,),
Some(250)
);
assert_eq!(
effective_request_timeout_ms_at(100, Some(now + Duration::from_millis(250)), now,),
Some(100)
);
}
#[test]
fn exhausted_or_sub_millisecond_deadline_stops_request_before_send() {
let now = Instant::now();
assert_eq!(effective_request_timeout_ms_at(1_000, Some(now), now), None);
assert_eq!(
effective_request_timeout_ms_at(1_000, Some(now + Duration::from_micros(999)), now,),
None
);
}
#[test]
fn retry_requires_budget_for_backoff_and_next_attempt() {
let now = Instant::now();
assert!(retry_delay_fits_request_deadline_at(None, 500, now));
assert!(retry_delay_fits_request_deadline_at(
Some(now + Duration::from_millis(501)),
500,
now,
));
assert!(!retry_delay_fits_request_deadline_at(
Some(now + Duration::from_millis(500)),
500,
now,
));
assert!(!retry_delay_fits_request_deadline_at(Some(now), 0, now));
}
#[test]
fn exhausted_budget_maps_to_timeout_request_error_and_audit() {
let error = request_budget_exhausted_error(
"https://vector.example/v1/images/generations",
"生成图片失败",
Some(900),
Some(12),
Some(1),
);
assert!(matches!(
error,
PlatformImageError::Request { timeout: true, .. }
));
let audit = error.audit().expect("budget error should carry audit");
assert_eq!(audit.failure_stage, "request_budget");
assert!(audit.timeout);
assert_eq!(audit.latency_ms, Some(900));
}
}
@@ -1,10 +1,14 @@
use std::time::{SystemTime, UNIX_EPOCH};
use std::time::{Instant, SystemTime, UNIX_EPOCH};
const VECTOR_ENGINE_SEND_MAX_ATTEMPTS: u32 = 5;
const VECTOR_ENGINE_SEND_RETRY_BASE_DELAY_MS: u64 = 500;
const VECTOR_ENGINE_SEND_RETRY_MAX_JITTER_MS: u64 = 999;
use super::{
budget::{
effective_request_timeout_ms, request_budget_exhausted_error,
retry_delay_fits_request_deadline,
},
constants::{GPT_IMAGE_2_MODEL, VECTOR_ENGINE_PROVIDER},
curl_transport::{
map_curl_error, send_vector_engine_json_request_with_curl,
@@ -63,8 +67,13 @@ pub async fn create_vector_engine_image_generation_with_model(
) -> Result<GeneratedImages, PlatformImageError> {
let model = normalize_vector_engine_image_model(model);
if !reference_images.is_empty() {
let resolved_references =
resolve_reference_images(http_client, reference_images, failure_context).await?;
let resolved_references = resolve_reference_images(
http_client,
reference_images,
failure_context,
settings.request_deadline,
)
.await?;
return create_vector_engine_image_edit_with_references_and_model(
http_client,
settings,
@@ -92,17 +101,28 @@ pub async fn create_vector_engine_image_generation_with_model(
let started_at = std::time::Instant::now();
let mut attempt = 1;
let response = loop {
let Some(attempt_timeout_ms) =
effective_request_timeout_ms(settings.request_timeout_ms, settings.request_deadline)
else {
return Err(request_budget_exhausted_error(
request_url.as_str(),
failure_context,
Some(started_at.elapsed().as_millis() as u64),
Some(prompt.chars().count()),
Some(reference_images.len()),
));
};
match send_vector_engine_json_request_with_curl(
request_url.as_str(),
settings.api_key.as_str(),
&request_body,
settings.request_timeout_ms,
attempt_timeout_ms,
)
.await
{
Ok(response) => {
if should_retry_vector_engine_upstream_status(response.status, attempt) {
retry_vector_engine_upstream_status_after_delay(
if retry_vector_engine_upstream_status_after_delay(
"generation",
request_url.as_str(),
attempt,
@@ -112,16 +132,19 @@ pub async fn create_vector_engine_image_generation_with_model(
Some(prompt.chars().count()),
Some(reference_images.len()),
Some(&request_body),
settings.request_deadline,
)
.await;
attempt += 1;
continue;
.await
{
attempt += 1;
continue;
}
}
break response;
}
Err(error) => {
if should_retry_vector_engine_curl_send_error(&error, attempt) {
retry_vector_engine_send_after_delay(
if retry_vector_engine_send_after_delay(
"generation",
request_url.as_str(),
"request_send",
@@ -135,10 +158,13 @@ pub async fn create_vector_engine_image_generation_with_model(
Some(prompt.chars().count()),
Some(reference_images.len()),
Some(&request_body),
settings.request_deadline,
)
.await;
attempt += 1;
continue;
.await
{
attempt += 1;
continue;
}
}
return Err(map_curl_error(
format!("{failure_context}:创建图片生成任务失败").as_str(),
@@ -179,6 +205,7 @@ pub async fn create_vector_engine_image_generation_with_model(
Some(reference_images.len()),
candidate_count,
"vector-engine",
settings.request_deadline,
)
.await
}
@@ -227,17 +254,28 @@ pub async fn create_vector_engine_nanobanana_generate_content(
let started_at = std::time::Instant::now();
let mut attempt = 1;
let response = loop {
let Some(attempt_timeout_ms) =
effective_request_timeout_ms(settings.request_timeout_ms, settings.request_deadline)
else {
return Err(request_budget_exhausted_error(
request_url.as_str(),
failure_context,
Some(started_at.elapsed().as_millis() as u64),
Some(prompt.chars().count()),
Some(reference_image_count),
));
};
match send_vector_engine_json_request_with_curl(
request_url.as_str(),
settings.api_key.as_str(),
&request_body,
settings.request_timeout_ms,
attempt_timeout_ms,
)
.await
{
Ok(response) => {
if should_retry_vector_engine_upstream_status(response.status, attempt) {
retry_vector_engine_upstream_status_after_delay(
if retry_vector_engine_upstream_status_after_delay(
"nanobanana_generate_content",
request_url.as_str(),
attempt,
@@ -247,16 +285,19 @@ pub async fn create_vector_engine_nanobanana_generate_content(
Some(prompt.chars().count()),
Some(reference_image_count),
Some(&request_params),
settings.request_deadline,
)
.await;
attempt += 1;
continue;
.await
{
attempt += 1;
continue;
}
}
break response;
}
Err(error) => {
if should_retry_vector_engine_curl_send_error(&error, attempt) {
retry_vector_engine_send_after_delay(
if retry_vector_engine_send_after_delay(
"nanobanana_generate_content",
request_url.as_str(),
"request_send",
@@ -270,10 +311,13 @@ pub async fn create_vector_engine_nanobanana_generate_content(
Some(prompt.chars().count()),
Some(reference_image_count),
Some(&request_params),
settings.request_deadline,
)
.await;
attempt += 1;
continue;
.await
{
attempt += 1;
continue;
}
}
return Err(map_curl_error(
format!("{failure_context}:创建 nanobanana2 图片生成任务失败").as_str(),
@@ -317,6 +361,7 @@ pub async fn create_vector_engine_nanobanana_generate_content(
Some(reference_image_count),
1,
"vector-engine-nanobanana",
settings.request_deadline,
)
.await
}
@@ -427,6 +472,17 @@ pub async fn create_vector_engine_image_edit_with_references_and_model(
);
let mut attempt = 1;
let response = loop {
let Some(attempt_timeout_ms) =
effective_request_timeout_ms(settings.request_timeout_ms, settings.request_deadline)
else {
return Err(request_budget_exhausted_error(
request_url.as_str(),
failure_context,
Some(started_at.elapsed().as_millis() as u64),
Some(prompt.chars().count()),
Some(reference_image_count),
));
};
match send_vector_engine_multipart_edit_request_with_curl(
request_url.as_str(),
settings.api_key.as_str(),
@@ -436,13 +492,13 @@ pub async fn create_vector_engine_image_edit_with_references_and_model(
normalized_size.as_str(),
candidate_count,
reference_images,
settings.request_timeout_ms,
attempt_timeout_ms,
)
.await
{
Ok(response) => {
if should_retry_vector_engine_upstream_status(response.status, attempt) {
retry_vector_engine_upstream_status_after_delay(
if retry_vector_engine_upstream_status_after_delay(
"edit",
request_url.as_str(),
attempt,
@@ -452,16 +508,19 @@ pub async fn create_vector_engine_image_edit_with_references_and_model(
Some(prompt.chars().count()),
Some(reference_image_count),
Some(&request_params),
settings.request_deadline,
)
.await;
attempt += 1;
continue;
.await
{
attempt += 1;
continue;
}
}
break response;
}
Err(error) => {
if should_retry_vector_engine_curl_send_error(&error, attempt) {
retry_vector_engine_send_after_delay(
if retry_vector_engine_send_after_delay(
"edit",
request_url.as_str(),
"request_send",
@@ -475,10 +534,13 @@ pub async fn create_vector_engine_image_edit_with_references_and_model(
Some(prompt.chars().count()),
Some(reference_image_count),
Some(&request_params),
settings.request_deadline,
)
.await;
attempt += 1;
continue;
.await
{
attempt += 1;
continue;
}
}
return Err(map_curl_error(
format!("{failure_context}:创建图片编辑任务失败").as_str(),
@@ -520,6 +582,7 @@ pub async fn create_vector_engine_image_edit_with_references_and_model(
Some(reference_image_count),
candidate_count,
"vector-engine-edit",
settings.request_deadline,
)
.await
}
@@ -550,8 +613,31 @@ async fn retry_vector_engine_send_after_delay(
prompt_chars: Option<usize>,
reference_image_count: Option<usize>,
request_params: Option<&serde_json::Value>,
) {
request_deadline: Option<Instant>,
) -> bool {
let delay_ms = vector_engine_send_retry_delay_ms(attempt, vector_engine_send_retry_jitter_ms());
if !retry_delay_fits_request_deadline(request_deadline, delay_ms) {
tracing::warn!(
provider = VECTOR_ENGINE_PROVIDER,
endpoint = %request_url,
request_kind,
failure_stage = "request_budget",
attempt,
max_attempts = VECTOR_ENGINE_SEND_MAX_ATTEMPTS,
retry_delay_ms = delay_ms,
timeout,
connect,
request,
body,
status = 0,
error,
elapsed_ms,
prompt_chars,
reference_image_count,
"VectorEngine 图片请求剩余预算不足,停止重试"
);
return false;
}
tracing::warn!(
provider = VECTOR_ENGINE_PROVIDER,
endpoint = %request_url,
@@ -575,6 +661,7 @@ async fn retry_vector_engine_send_after_delay(
"VectorEngine 图片请求发送失败,准备重试"
);
tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await;
true
}
async fn retry_vector_engine_upstream_status_after_delay(
@@ -587,8 +674,27 @@ async fn retry_vector_engine_upstream_status_after_delay(
prompt_chars: Option<usize>,
reference_image_count: Option<usize>,
request_params: Option<&serde_json::Value>,
) {
request_deadline: Option<Instant>,
) -> bool {
let delay_ms = vector_engine_send_retry_delay_ms(attempt, vector_engine_send_retry_jitter_ms());
if !retry_delay_fits_request_deadline(request_deadline, delay_ms) {
tracing::warn!(
provider = VECTOR_ENGINE_PROVIDER,
endpoint = %request_url,
request_kind,
failure_stage = "request_budget",
attempt,
max_attempts = VECTOR_ENGINE_SEND_MAX_ATTEMPTS,
retry_delay_ms = delay_ms,
status,
retryable = false,
elapsed_ms,
prompt_chars,
reference_image_count,
"VectorEngine 图片请求剩余预算不足,停止上游状态重试"
);
return false;
}
tracing::warn!(
provider = VECTOR_ENGINE_PROVIDER,
endpoint = %request_url,
@@ -609,6 +715,7 @@ async fn retry_vector_engine_upstream_status_after_delay(
"VectorEngine 图片上游状态可重试,准备重试"
);
tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await;
true
}
fn vector_engine_send_retry_delay_ms(attempt: u32, jitter_ms: u64) -> u64 {
@@ -629,6 +736,40 @@ fn vector_engine_send_retry_jitter_ms() -> u64 {
mod tests {
use super::*;
#[tokio::test]
async fn expired_deadline_stops_generation_before_network_send() {
let settings = VectorEngineImageSettings {
base_url: "http://127.0.0.1:9".to_string(),
api_key: "test-key".to_string(),
request_timeout_ms: 1_000,
request_deadline: Some(Instant::now()),
};
let error = create_vector_engine_image_generation(
&reqwest::Client::new(),
&settings,
"测试提示词",
None,
"1024x1024",
1,
&[],
"测试图片生成失败",
)
.await
.expect_err("expired deadline should stop before network send");
assert!(matches!(
error,
PlatformImageError::Request { timeout: true, .. }
));
assert_eq!(
error
.audit()
.expect("budget error should carry audit")
.failure_stage,
"request_budget"
);
}
#[test]
fn vector_engine_send_retry_policy_allows_four_retries_before_final_attempt() {
assert_eq!(VECTOR_ENGINE_SEND_MAX_ATTEMPTS, 5);
@@ -1,7 +1,9 @@
use base64::{Engine as _, engine::general_purpose::STANDARD as BASE64_STANDARD};
use reqwest::header;
use std::time::Instant;
use super::{
budget::request_budget_exhausted_error,
constants::VECTOR_ENGINE_PROVIDER,
error::PlatformImageError,
types::{DownloadedImage, GeneratedImages, ReferenceImage},
@@ -53,18 +55,63 @@ pub async fn download_remote_image(
})
}
async fn download_remote_image_with_deadline(
http_client: &reqwest::Client,
image_url: &str,
request_deadline: Option<Instant>,
operation: &str,
) -> Result<DownloadedImage, PlatformImageError> {
let Some(request_deadline) = request_deadline else {
return download_remote_image(http_client, image_url).await;
};
let started_at = Instant::now();
if request_deadline <= started_at {
return Err(request_budget_exhausted_error(
image_url,
operation,
Some(0),
None,
None,
));
}
tokio::time::timeout_at(
tokio::time::Instant::from_std(request_deadline),
download_remote_image(http_client, image_url),
)
.await
.map_err(|_| {
request_budget_exhausted_error(
image_url,
operation,
Some(started_at.elapsed().as_millis() as u64),
None,
None,
)
})?
}
pub(crate) async fn download_images_from_urls(
http_client: &reqwest::Client,
task_id: String,
image_urls: Vec<String>,
candidate_count: u32,
request_deadline: Option<Instant>,
operation: &str,
) -> Result<GeneratedImages, PlatformImageError> {
let mut images = Vec::with_capacity(candidate_count.clamp(1, 4) as usize);
for image_url in image_urls
.into_iter()
.take(candidate_count.clamp(1, 4) as usize)
{
images.push(download_remote_image(http_client, image_url.as_str()).await?);
images.push(
download_remote_image_with_deadline(
http_client,
image_url.as_str(),
request_deadline,
operation,
)
.await?,
);
}
Ok(GeneratedImages {
task_id,
@@ -77,6 +124,7 @@ pub(crate) async fn resolve_reference_images(
http_client: &reqwest::Client,
reference_images: &[String],
failure_context: &str,
request_deadline: Option<Instant>,
) -> Result<Vec<ReferenceImage>, PlatformImageError> {
let mut resolved = Vec::new();
for (index, source) in reference_images.iter().take(5).enumerate() {
@@ -89,20 +137,14 @@ pub(crate) async fn resolve_reference_images(
continue;
}
if source.starts_with("http://") || source.starts_with("https://") {
let downloaded = download_remote_image(http_client, source)
.await
.map_err(|error| PlatformImageError::Request {
provider: VECTOR_ENGINE_PROVIDER,
message: format!("{failure_context}:下载参考图失败:{error}"),
endpoint: Some(source.to_string()),
timeout: false,
connect: false,
request: false,
body: false,
status_code: None,
source: None,
audit: None,
})?;
let downloaded = download_remote_image_with_deadline(
http_client,
source,
request_deadline,
failure_context,
)
.await
.map_err(|error| contextualize_reference_download_error(error, failure_context))?;
resolved.push(ReferenceImage {
bytes: downloaded.bytes,
mime_type: downloaded.mime_type.clone(),
@@ -246,3 +288,106 @@ fn map_simple_request_error(message: String, endpoint: Option<String>) -> Platfo
audit: None,
}
}
fn contextualize_reference_download_error(
error: PlatformImageError,
failure_context: &str,
) -> PlatformImageError {
match error {
PlatformImageError::Request {
provider,
message,
endpoint,
timeout,
connect,
request,
body,
status_code,
source,
audit,
} => PlatformImageError::Request {
provider,
message: format!("{failure_context}:下载参考图失败:{message}"),
endpoint,
timeout,
connect,
request,
body,
status_code,
source,
audit,
},
error => error,
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::Duration;
use tokio::net::TcpListener;
#[tokio::test]
async fn exhausted_deadline_stops_remote_image_download_before_send() {
let http_client = reqwest::Client::new();
let error = download_remote_image_with_deadline(
&http_client,
"http://127.0.0.1:9/not-called.png",
Some(Instant::now()),
"下载测试图片失败",
)
.await
.expect_err("expired deadline should stop before network send");
assert!(matches!(
error,
PlatformImageError::Request { timeout: true, .. }
));
assert_eq!(
error
.audit()
.expect("budget error should carry audit")
.failure_stage,
"request_budget"
);
}
#[tokio::test]
async fn deadline_cancels_pending_remote_image_download() {
let listener = TcpListener::bind("127.0.0.1:0")
.await
.expect("mock listener should bind");
let address = listener
.local_addr()
.expect("mock listener should expose address");
let server = tokio::spawn(async move {
let (_socket, _) = listener
.accept()
.await
.expect("mock listener should accept request");
tokio::time::sleep(Duration::from_secs(1)).await;
});
let error = download_remote_image_with_deadline(
&reqwest::Client::new(),
format!("http://{address}/pending.png").as_str(),
Some(Instant::now() + Duration::from_millis(100)),
"下载测试图片失败",
)
.await
.expect_err("absolute deadline should cancel pending download");
server.abort();
assert!(matches!(
error,
PlatformImageError::Request { timeout: true, .. }
));
assert_eq!(
error
.audit()
.expect("budget error should carry audit")
.failure_stage,
"request_budget"
);
}
}
@@ -1,4 +1,5 @@
mod audit;
mod budget;
mod client;
mod constants;
mod curl_transport;
@@ -10,6 +10,7 @@ use super::{
types::GeneratedImages,
util::{current_utc_micros, is_timeout_message, truncate_raw},
};
use std::time::Instant;
pub(crate) async fn handle_vector_engine_response(
http_client: &reqwest::Client,
@@ -22,6 +23,7 @@ pub(crate) async fn handle_vector_engine_response(
reference_image_count: Option<usize>,
candidate_count: u32,
task_prefix: &str,
request_deadline: Option<Instant>,
) -> Result<GeneratedImages, PlatformImageError> {
if !(200..=299).contains(&response_status) {
let message = parse_api_error_message(response_text, failure_context);
@@ -101,18 +103,32 @@ pub(crate) async fn handle_vector_engine_response(
task_id,
image_urls,
candidate_count,
request_deadline,
failure_context,
)
.await
{
Ok(generated) => generated,
Err(error) => {
let timeout = matches!(&error, PlatformImageError::Request { timeout: true, .. });
let request_budget_exhausted = error
.audit()
.is_some_and(|audit| audit.failure_stage == "request_budget");
let audit = build_failure_audit(
request_url,
failure_context,
"image_download",
if request_budget_exhausted {
"request_budget"
} else {
"image_download"
},
Some(response_status),
Some("5xx"),
false,
if request_budget_exhausted {
None
} else {
Some("5xx")
},
timeout,
false,
error.message(),
None,
@@ -27,11 +27,13 @@ mod tests {
base_url: "https://vector.example".to_string(),
api_key: "test-key".to_string(),
request_timeout_ms: 1_000,
request_deadline: None,
};
let v1_settings = VectorEngineImageSettings {
base_url: "https://vector.example/v1".to_string(),
api_key: "test-key".to_string(),
request_timeout_ms: 1_000,
request_deadline: None,
};
assert_eq!(
@@ -3,6 +3,7 @@ pub struct VectorEngineImageSettings {
pub base_url: String,
pub api_key: String,
pub request_timeout_ms: u64,
pub request_deadline: Option<std::time::Instant>,
}
#[derive(Clone, Debug)]
@@ -1,7 +1,7 @@
use platform_image::vector_engine::{
GPT_IMAGE_2_MODEL, ReferenceImage, VECTOR_ENGINE_PROVIDER, VectorEngineImageSettings,
build_vector_engine_image_http_client, build_vector_engine_image_request_body,
build_vector_engine_image_request_body_with_model,
GPT_IMAGE_2_MODEL, PlatformImageError, ReferenceImage, VECTOR_ENGINE_PROVIDER,
VectorEngineImageSettings, build_vector_engine_image_http_client,
build_vector_engine_image_request_body, build_vector_engine_image_request_body_with_model,
build_vector_engine_nanobanana_generate_content_request_body, create_vector_engine_image_edit,
create_vector_engine_image_generation, create_vector_engine_nanobanana_generate_content,
vector_engine_images_edit_url, vector_engine_images_generation_url,
@@ -12,7 +12,7 @@ use std::{
Arc,
atomic::{AtomicUsize, Ordering},
},
time::Duration,
time::{Duration, Instant},
};
use tokio::{
io::{AsyncReadExt, AsyncWriteExt},
@@ -25,6 +25,7 @@ fn vector_engine_module_exposes_provider_protocol_helpers() {
base_url: "https://vector.example/v1".to_string(),
api_key: "test-key".to_string(),
request_timeout_ms: 1_000,
request_deadline: None,
};
let body =
@@ -187,6 +188,7 @@ fn nanobanana_generate_content_url_uses_model_path() {
base_url: "https://vector.example/v1".to_string(),
api_key: "test-key".to_string(),
request_timeout_ms: 1_000,
request_deadline: None,
};
assert_eq!(
@@ -235,6 +237,7 @@ async fn vector_engine_image_edit_retries_send_timeout_once_and_succeeds() {
base_url: format!("http://{server_addr}/v1"),
api_key: "test-key".to_string(),
request_timeout_ms: 40,
request_deadline: None,
};
let http_client =
build_vector_engine_image_http_client(&settings).expect("client should build");
@@ -262,6 +265,65 @@ async fn vector_engine_image_edit_retries_send_timeout_once_and_succeeds() {
server.abort();
}
#[tokio::test]
async fn vector_engine_deadline_clips_stalled_attempt_and_prevents_retry() {
let listener = TcpListener::bind("127.0.0.1:0")
.await
.expect("mock server should bind");
let server_addr = listener
.local_addr()
.expect("mock server address should be readable");
let request_count = Arc::new(AtomicUsize::new(0));
let request_count_for_server = Arc::clone(&request_count);
let server = tokio::spawn(async move {
loop {
let Ok((mut stream, _)) = listener.accept().await else {
break;
};
request_count_for_server.fetch_add(1, Ordering::SeqCst);
tokio::spawn(async move {
let mut buffer = [0_u8; 4096];
let _ = stream.read(&mut buffer).await;
tokio::time::sleep(Duration::from_secs(1)).await;
});
}
});
let started_at = Instant::now();
let settings = VectorEngineImageSettings {
base_url: format!("http://{server_addr}/v1"),
api_key: "test-key".to_string(),
request_timeout_ms: 5_000,
request_deadline: Some(started_at + Duration::from_millis(150)),
};
let http_client =
build_vector_engine_image_http_client(&settings).expect("client should build");
let error = create_vector_engine_image_generation(
&http_client,
&settings,
"测试提示词",
None,
"1024x1024",
1,
&[],
"测试 VectorEngine 图片生成失败",
)
.await
.expect_err("stalled request should exhaust the shared deadline");
assert!(matches!(
error,
PlatformImageError::Request { timeout: true, .. }
));
assert!(
started_at.elapsed() < Duration::from_secs(1),
"attempt 应使用剩余 deadline,而不是完整配置 timeout"
);
assert_eq!(request_count.load(Ordering::SeqCst), 1);
server.abort();
}
#[tokio::test]
async fn nanobanana_generate_content_posts_native_body_and_reads_inline_data() {
let listener = TcpListener::bind("127.0.0.1:0")
@@ -307,6 +369,7 @@ async fn nanobanana_generate_content_posts_native_body_and_reads_inline_data() {
base_url: format!("http://{}", server_addr),
api_key: "test-key".to_string(),
request_timeout_ms: 1_000,
request_deadline: None,
};
let client = build_vector_engine_image_http_client(&settings).expect("client should build");
@@ -375,6 +438,7 @@ async fn vector_engine_image_generation_retries_upstream_502_once_and_succeeds()
base_url: format!("http://{server_addr}/v1"),
api_key: "test-key".to_string(),
request_timeout_ms: 1_000,
request_deadline: None,
};
let http_client =
build_vector_engine_image_http_client(&settings).expect("client should build");