diff --git a/docs/project-memory/shared-memory/decision-log.md b/docs/project-memory/shared-memory/decision-log.md index 67698e4a2..57a8898bf 100644 --- a/docs/project-memory/shared-memory/decision-log.md +++ b/docs/project-memory/shared-memory/decision-log.md @@ -16,6 +16,16 @@ --- +## 2026-07-20 角色动画帧 OSS 请求使用专用连接池、并发保护与结构化重试 + +- 背景:角色动作逐帧流水线会同时发起源帧 PUT、透明帧 PUT 和最终帧 HEAD;原路径每次请求新建 `reqwest::Client`,且 OSS 请求错误丢失 HTTP 状态和 timeout/connect/transport 分类,多个动画任务叠加时无法在进程级限制 OSS 在途请求,也无法安全区分 PUT 与 HEAD 的失败。 +- 决策:`AppState` 仅为角色动画帧初始化一次 OSS HTTP Client 和 8 路 `Semaphore`。全帧 Future 仍保持 `buffer_unordered(frame_count.max(1))`,BgFilter、阿里云抠图和本地处理不占 OSS permit;每次 PUT/HEAD 网络 attempt 单独获取 permit,退避期间释放。`platform-oss` 保留 `OssErrorKind::Request`,但在 `OssError::Request` 中保留 operation、status、timeout、connect、transport、OSS code、OSS request-id 和原脱敏 message,并为动画帧提供 3 次 attempt、250ms/500ms 退避的 PUT/HEAD 独立重试。仅无响应传输错误、timeout、OSS PutObject 的 `400 + RequestTimeout`、PUT 400 错误体读取失败(未解析出 `Code`,按 timeout/transport 归类)、408、429 和 5xx 可重试;动作帧 PUT 对 400 错误体最多读取 16 KiB,只提取 `Code` 和响应头优先的 `x-oss-request-id`,不记录完整 XML;错误体读取超时/断流时保留已读字节,已解析出的 `Code` 优先生效。除 `RequestTimeout` 与该错误体读取失败情形外的确定性 4xx、配置、签名、URL、空请求体、抠图和素材登记错误不重试。重试体在 platform-oss 内一次转为可复用 `Bytes`,每次重新签名和构造 Request,不复制整帧字节。 +- 失败语义:最终帧 PUT 成功后才执行 HEAD;HEAD 失败只重试 HEAD,不重复 PUT。任一帧最终失败仍排空已启动的 Future、整段动作退款并禁止发布缺帧动画,帧结果继续按原始序号排序。 +- 影响范围:`state.rs`、`platform-oss/lib.rs`、`character_animation_assets.rs`、对应 Cargo 依赖和架构 / 运维文档;不改变其他 OSS 调用方、BgFilter/阿里云降级、worker、计费退款、SpacetimeDB schema/DTO 或前端接口。 +- 验证方式:`cargo test -p platform-oss --manifest-path server-rs/Cargo.toml`、`cargo test -p api-server character_animation --manifest-path server-rs/Cargo.toml`、`cargo check -p api-server --manifest-path server-rs/Cargo.toml`、`npm run check:encoding`、`git diff --check`。 + +--- + ## 2026-07-18 图片生成 K 档由 provider 直接生成 - 背景:旧 gpt-image-2 尺寸表会把 2K 竖版回落到 `1024x1536`,图标入口又使用固定 `360x360 / 512x512` 占位;角色去背景结果变小时还会直接放大整张透明成品,导致 UI 显示的 2K 与模型实际生成清晰度不一致。 diff --git a/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md b/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md index 63037e502..d9d38b5e9 100644 --- a/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md +++ b/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md @@ -243,7 +243,7 @@ npm run check:server-rs-ddd - 抠图输入以私有 OSS 作为内存生命周期边界:生成原图和角色动作抽取帧上传时消费图片字节所有权,上传完成后不保留原图缓冲;手动去背景直接解析并校验已有 OSS object key,不下载原图。BgFilter 必须为 object key 签发 600 秒 GET URL 并通过 multipart `image_url` 提交,不用 `file` 重传;flat 链路进入阿里云 fallback 时由 `platform-matting` URL 接口单独下载并上传 `AuthorizeFileUpload` 临时对象,在推理前释放下载缓冲,继续 fallback 到本地键色时再单独下载一次原图,本地产出后释放本次原图下载缓冲。签名 URL 不得写入日志、审计或持久化。 - 阿里云通用抠图的非上海地域输入不得使用 `viapiutils/GetOssStsToken`、固定 `viapi-customer-temp` 或 OSS V1 PUT。`platform-matting` 必须按官方新版 SDK Advance 协议调用 `AuthorizeFileUpload`,使用动态返回的单对象 Policy 执行 multipart POST,再把临时上海 OSS URL 交给 `SegmentCommonImage`;输入归一化、结果下载与原尺寸 Alpha 回贴继续留在同一适配器内。该协议仍上传图片字节,不等同于阿里云服务端直接抓取任意公网 URL,也不改变上层 BgFilter → 阿里云 → 本地降级顺序。 - 编辑器抠图服务:手动 `POST /api/editor/images/background-removals` 与角色形象生成、图标 spritesheet 生成、UI 设计图素材提取、角色动作抽帧后的透明化统一走 BgFilter,配置为 `GENARRATIVE_EDITOR_BGFILTER_BASE_URL`、`GENARRATIVE_EDITOR_BGFILTER_TOKEN` 和 `GENARRATIVE_EDITOR_BGFILTER_REQUEST_TIMEOUT_MS`,默认 base URL 为 `http://58.87.105.82/bgfilter`,默认请求超时为 `180000ms`(BgFilter 当前为 CPU 推理,单次抠图较慢,必须留足超时);旧 `GENARRATIVE_EDITOR_BACKGROUND_REMOVAL_TOKEN` 只作为 BgFilter token 的兼容回退别名,原手动去背景专用 base URL / timeout 配置已经删除。手动去背景固定传 `image_url`、`background_mode=complex`、`seg_model=birefnet`、`cross_check=off`,不传 `file` 或 `screen_color`;该 API 接收 `objectKey`、`resourceId` 或 `assetId` 候选引用;BFF 入队前统一拒绝 `data:` / `blob:`;worker 不重复入口校验,只调用 `resolve_editor_reference_object_key_for_owner`,底层 resolver 在解析引用前拒绝内联媒体,并在签名前完成登记状态和 owner 校验;直接签发 OSS URL,不下载原图。标准纯色背景四条链路固定传 `background_mode=flat`,并显式传 `image_url`、`screen_color=`、`seg_model=` 和 `cross_check=`,其中角色形象生成和角色动作逐帧去背传 `cross_check=on`,图标 spritesheet 生成和 UI 设计图素材提取传 `cross_check=off`。前端用户路径不展示抠图模型、模式或 cross-check,固定提交默认 `birefnet`,后端仍识别内部保留的 `anime-seg`;这些参数只属于后端内部供应商策略,不进入前端或外部 OpenAPI。标准纯色背景 BgFilter 调用失败,或连续失败达到 `GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_FAILURE_THRESHOLD`(默认 `3`)并在 `GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS`(默认 `300`)内打开熔断时,继续复用“阿里云通用抠图 → 本地 `editor_green_screen` 键色扣除”兜底链,熔断期不得直接退化到本地兜底。角色动作视频生成的背景色已与生图链路统一:`screenColor=auto` 时由视觉 LLM(`gpt-5-mini`,Responses 协议、low 推理档)读源角色图自动决策,并经硬过滤器剔除与前景 / 皮肤撞色的候选,手动 hex 则尊重用户选择;透明源角色图在提交 Ark 图生视频前先合成到选定背景色实色,使视频背景等于抠图键色;抽帧后每帧先上传私有 OSS 并释放原帧缓冲,再以该 object key 的签名 URL 固定使用 `seg_model=birefnet`、`cross_check=on` 进入上述三段式链路。阿里云通用抠图配置为 `GENARRATIVE_ALIYUN_MATTING_ENABLED`、`GENARRATIVE_ALIYUN_MATTING_ENDPOINT`、`GENARRATIVE_ALIYUN_MATTING_ACCESS_KEY_ID`、`GENARRATIVE_ALIYUN_MATTING_ACCESS_KEY_SECRET` 和 `GENARRATIVE_ALIYUN_MATTING_REQUEST_TIMEOUT_MS`;未配置专用 AK/SK 时可复用 `ALIBABA_CLOUD_ACCESS_KEY_ID` / `ALIBABA_CLOUD_ACCESS_KEY_SECRET`,默认 endpoint 为 `imageseg.cn-shanghai.aliyuncs.com`。标准纯色背景链路中,BgFilter 调用失败和阿里云抠图链路已开始后的失败(包括源 OSS GET 成功后的解码、尺寸校验和归一化失败)都写入 `external_api_call_failure` 审计;真正开始外部调用前的本地预检不写该审计,并在 `failureStage` 中保留 `source_decode`、`source_validate` 等阶段。 -- BgFilter 连接复用、重试与动作帧流水线:api-server 必须在 `AppState` 复用同一个 BgFilter HTTP Client 及 keep-alive 连接池。`GENARRATIVE_EDITOR_BGFILTER_REQUEST_TIMEOUT_MS` 是所有路径的基准请求超时;角色动作逐帧 BgFilter 的每一次 HTTP attempt 使用“基准超时 + `2000ms × 本次实际帧数`”,默认 `32 / 40 / 48` 帧分别为 `244000 / 260000 / 276000ms`,角色形象单图、图标、UI 和手动去背景仍使用基准值。api-server 只在共享 Client 的单次 RequestBuilder 上覆盖该值;它覆盖从请求发起到响应体读取完成,是单次 attempt 的总 deadline,不是整批帧或 worker job 超时,重试会重新签发 600 秒 OSS URL 并获得同样的 request deadline,整项任务仍受 worker long-job 预算约束。flat 与 complex 请求首次失败后都立即重试 `1` 次;标准纯色背景 flat 请求第二次仍失败才进入“阿里云通用抠图 → 本地键色”降级链,手动 complex 请求第二次仍失败则返回最终错误,不接入依赖纯色键值的降级链,也不改变 flat 路径的熔断状态。角色动作全部 `32 / 40 / 48` 帧按“单帧绿幕源图 owned 上传 OSS 并释放原帧 → 以签名 URL 调 BgFilter/按 object key 降级 → 透明帧落 OSS”独立流水化,使用覆盖本次全部帧的无序在途集合连续发射,不在 api-server 增加供应商进程锁或固定小并发窗口;返回结果携带原始帧序并在收口时排序。任一帧最终失败时必须先排空全部已启动 Future,再让整个动作任务失败退款,不能发布缺帧动画。 +- BgFilter 连接复用、重试与动作帧流水线:api-server 必须在 `AppState` 复用同一个 BgFilter HTTP Client 及 keep-alive 连接池。`GENARRATIVE_EDITOR_BGFILTER_REQUEST_TIMEOUT_MS` 是所有路径的基准请求超时;角色动作逐帧 BgFilter 的每一次 HTTP attempt 使用“基准超时 + `2000ms × 本次实际帧数`”,默认 `32 / 40 / 48` 帧分别为 `244000 / 260000 / 276000ms`,角色形象单图、图标、UI 和手动去背景仍使用基准值。api-server 只在共享 Client 的单次 RequestBuilder 上覆盖该值;它覆盖从请求发起到响应体读取完成,是单次 attempt 的总 deadline,不是整批帧或 worker job 超时,重试会重新签发 600 秒 OSS URL 并获得同样的 request deadline,整项任务仍受 worker long-job 预算约束。flat 与 complex 请求首次失败后都立即重试 `1` 次;标准纯色背景 flat 请求第二次仍失败才进入“阿里云通用抠图 → 本地键色”降级链,手动 complex 请求第二次仍失败则返回最终错误,不接入依赖纯色键值的降级链,也不改变 flat 路径的熔断状态。角色动作全部 `32 / 40 / 48` 帧按“单帧绿幕源图 owned 上传 OSS 并释放原帧 → 以签名 URL 调 BgFilter/按 object key 降级 → 透明帧落 OSS”独立流水化,使用覆盖本次全部帧的无序在途集合连续发射;不限制 BgFilter、阿里云或本地处理,但角色动画源帧 PUT、透明帧 PUT 和最终帧 HEAD 统一复用 `AppState` 内初始化一次的 OSS HTTP Client(连接池参数为 connect 30 秒、request 60 秒、idle 300 秒、每 host 8 个 idle 连接、TCP keepalive 60 秒),并受进程级 8 路 OSS semaphore 限制。每个 OSS 网络 attempt 单独获取 permit,退避期间释放;PUT/HEAD 动画帧请求最多 3 次(250ms、500ms 退避),只重试无 HTTP 响应的传输错误、timeout、OSS PutObject 的 `400 + RequestTimeout`、PUT `400` 错误体读取失败(未解析出 `Code`,按 timeout/transport 归类)、408、429 和 5xx。动作帧 PUT 只在 400 响应中有界读取最多 16 KiB OSS 错误 XML,并保留 `Code` 与响应头优先的 `x-oss-request-id`;错误体读取超时/断流时保留已读字节,已解析出的 `Code` 优先生效,未解析出 `Code` 则按 timeout/transport 归类重试;除 `RequestTimeout` 与该错误体读取失败情形外的其他 400、401/403/404、配置、URL/签名和空请求体错误不重试。返回结果携带原始帧序并在收口时排序。任一帧最终失败时必须先排空全部已启动 Future,再让整个动作任务失败退款,不能发布缺帧动画。 - Match3D 物品 sheet:关卡整图完成后走 VectorEngine `/v1/images/edits` multipart `image`,模型为 `gpt-image-2`,`2K 1:1` 输出 `10*10` spritesheet;物品 sheet prompt 固定要求单一纯绿色 `#00FF00 / RGB(0,255,0)` 绿幕背景,后端上传 OSS 前必须把绿幕扣成透明 PNG,并把透明整图写入 `itemSpritesheetImageSrc/itemSpritesheetImageObjectKey`。后端优先按透明 alpha 连通域从该 sheet 识别真实素材矩形并持久化 20 个物品、每个 5 个形态;识别数量不足时才回退 `10*10` 固定网格。通用系列素材图集的行列索引按每行 2 个物品计算,必须落在 `1..=10`,难度只决定运行态加载 3 / 9 / 15 / 20 种。 - Match3D UI spritesheet 和背景派生图:关卡整图作为参考图并发生成 `1K 1:1` UI spritesheet 与 `1K 9:16` 背景图,模型均为 `gpt-image-2`。UI spritesheet prompt 固定要求单一纯绿色 `#00FF00 / RGB(0,255,0)` 绿幕背景,后端上传 OSS 前必须把绿幕扣成透明 PNG;背景图必须合成为全画幅不透明 PNG。 - Match3D 1:1 容器 UI:VectorEngine `/v1/images/edits` multipart 参考图。该容器参考图是后端生图协议输入,必须通过 `include_bytes!` 随 `api-server` 编译进二进制,避免 API 单独发布或运行目录缺少 `public/` 时生成失败。 diff --git a/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md b/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md index 9abb5d31c..fd027403b 100644 --- a/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md +++ b/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md @@ -407,6 +407,10 @@ curl -fsS --max-time 5 http://127.0.0.1/ >/dev/null curl -fsS --max-time 5 http://127.0.0.1/api/editor/showcase/resources >/dev/null ``` +### 角色动画帧 OSS 排障 + +角色动画源帧 PUT、透明帧 PUT 和最终帧 HEAD 使用 `AppState` 内同一个 OSS HTTP Client/连接池,并受进程级 8 路 OSS permit 保护;BgFilter、阿里云抠图和本地处理不占用该 permit。每个 OSS attempt 最多 3 次(首次 + 2 次重试),退避为 250ms、500ms;只重试 timeout、无 HTTP 响应传输错误、OSS PutObject 的 `400 + RequestTimeout`、PUT 400 错误体读取失败(未解析出 `Code`,按 timeout/transport 归类)、408、429 和 500–599。动作帧 PUT 收到 400 时只读取最多 16 KiB OSS 错误 XML,提取 `Code` 和 `RequestId`;`oss_request_id` 优先使用响应头 `x-oss-request-id`,XML 字段只作回退。错误体读取超时/断流不再按确定性 400 处理:已解析出的 `Code` 优先生效;未解析出 `Code` 时按读取失败原因置 `timeout`/`transport` 并重试,message 追加「错误响应体读取失败」。日志字段包括 `frame_index`、`object_key`、`operation=source_put|final_put|final_head`、`attempt`、`max_attempts`、`retryable`、`will_retry`、`retry_delay_ms`、`permit_wait_ms`、`timeout`、`connect`、`transport`、`oss_code`、`oss_request_id`、`status` 和 `elapsed_ms`。`请求 OSS 失败` 时,`timeout/connect/transport=true` 表示传输类失败;`status=400, oss_code=RequestTimeout, timeout=true`、`status=429` 或 `500–599` 表示暂时性失败,PUT 的 `status=400`、`oss_code` 为空且 `timeout=true` 或 `transport=true`(message 含「错误响应体读取失败」)同样是暂时性失败。除 `RequestTimeout` 和该错误体读取失败两类例外外,其他 400、401/403/404、配置、URL 和签名错误是确定性失败,不会重试。最终帧 HEAD 失败只会重试 HEAD,不会重复 PUT;如果任一帧最终失败,确认整段动作已排空已启动 Future,并检查任务按现有契约退款且没有发布缺帧动画。 + ## 生产运维 生产部署当前口径: @@ -600,7 +604,7 @@ OpenTelemetry 现阶段默认开启 OTLP traces / metrics / logs,但本地日 - api-server 会随 metrics 发送进程级指标:`process.memory.usage`、`process.memory.virtual`、`process.cpu.time`、`genarrative.process.cpu.usage_percent`、`process.thread.count`、`genarrative.process.memory.private`;Windows 额外发送 `process.windows.handle.count`,Linux 额外发送 `process.unix.file_descriptor.count`。这些指标只描述当前进程,不携带请求、用户或作品 label。 - HTTP 运行态补充发送 `genarrative.http.server.response_bodies.in_flight` 与 `genarrative.http.server.request_permits.available`,后者带低基数 `pool=default|gallery|detail|admin` label,用于区分业务 handler / 背压 permit 是否仍被占用;拼图广场热点缓存补充发送 `genarrative.puzzle_gallery.cache.*` 指标,记录 fresh hit、stale hit、未命中、后台刷新开始 / 失败、重建耗时和预序列化 data JSON 字节数。 - 外部 API 失败统一发送 OTLP 并落库。当前 VectorEngine `gpt-image-2` 图片生成 / 编辑失败由 `platform-image` provider 输出结构化日志字段,字段包括 provider、endpoint、failure_stage、status、source、source_chain、source_chain_depth、timeout、retryable、latency_ms、prompt_chars、reference_image_count、image_model、request_params 和 raw_excerpt;图片编辑请求参数日志还会带 reference_image_bytes_total,并在 request_params.referenceImages 中记录每个 multipart `image` part 的 fileName、mimeType 和 bytes,不记录 API key 或原始图片 bytes;`api-server` 再记录指标 `genarrative.external_api.failures{provider,failure_stage,status_class,retryable}`,并写入 `tracking_event`,`event_key = external_api_call_failure`、`module_key = external-api`、`scope_kind = module`、`scope_id = provider`。调用方能拿到身份上下文时,失败事件还会在行级 `user_id` / `owner_user_id` / `profile_id` 和 `metadata_json.userId` / `metadata_json.profileId` / `metadata_json.requestId` / `metadata_json.errorSource` 中记录触发者、草稿 / 作品作用域、请求标识和传输错误链。排障时先按 provider / failureStage 聚合,再下钻 userId / profileId,最后结合 request 日志、errorSource 和上游响应 excerpt 判断是限流、超时、解析失败还是未返回图片。 -- OSS 平台适配器也输出结构化日志,覆盖 `sign_post_object`、`sign_get_object_url`、`head_object` 和 `put_object`。排查资产签名、上传或确认失败时,先按 `provider=aliyun-oss` 与 `operation` 过滤,再看 `object_key` / `key_prefix`、`status`、`status_class`、`error_kind`、`content_length`、`content_type` 和 `elapsed_ms`;日志不得包含 AccessKey、policy、signature、Authorization header 或完整 signed URL。排查 generated 图片重复下载时,先确认前端输入是否为 `/generated-*` legacy path 或可归一化的 `https://*.oss-*.aliyuncs.com/generated-*`;正确链路应先调 `/api/assets/read-url`,再由浏览器请求 signed URL,且同一路径、同一 `refreshKey` 版本和未临近过期的 signed URL 应复用。新上传 generated 私有对象应带 `Cache-Control: public, max-age=31536000, immutable`;旧对象若只有 `ETag` / `Last-Modified`,浏览器会走 304 协商缓存而不是长期强缓存,可通过刷新 OSS 元数据或 CDN 配置补齐。 +- OSS 平台适配器也输出结构化日志,覆盖 `sign_post_object`、`sign_get_object_url`、`head_object` 和 `put_object`。排查资产签名、上传或确认失败时,先按 `provider=aliyun-oss` 与 `operation` 过滤,再看 `object_key` / `key_prefix`、`status`、`status_class`、`error_kind`、`content_length`、`content_type` 和 `elapsed_ms`;角色动画逐帧额外按 `frame_index`、`operation=source_put|final_put|final_head`、`attempt/max_attempts`、`will_retry`、`oss_code` 和 `oss_request_id` 对齐同一对象的请求尝试。`请求 OSS 失败` 时,`timeout/connect/transport=true` 表示传输类失败,OSS PutObject 的 `status=400, oss_code=RequestTimeout, timeout=true`、`status=429` 或 `500–599` 表示暂时性失败,PUT 的 `status=400`、`oss_code` 为空且 `timeout=true` 或 `transport=true`(message 含「错误响应体读取失败」,即 400 错误体读取超时/断流)也会重试;除这两类例外外,其他 400、401/403/404、配置、URL 和签名错误是确定性失败,不会重试。最终帧 HEAD 失败只会重试 HEAD,不会重复 PUT。日志不得包含 AccessKey、policy、signature、Authorization header、完整 signed URL 或 OSS 错误响应体;`oss_request_id` 只用于关联 OSS 服务端排障。排查 generated 图片重复下载时,先确认前端输入是否为 `/generated-*` legacy path 或可归一化的 `https://*.oss-*.aliyuncs.com/generated-*`;正确链路应先调 `/api/assets/read-url`,再由浏览器请求 signed URL,且同一路径、同一 `refreshKey` 版本和未临近过期的 signed URL 应复用。新上传 generated 私有对象应带 `Cache-Control: public, max-age=31536000, immutable`;旧对象若只有 `ETag` / `Last-Modified`,浏览器会走 304 协商缓存而不是长期强缓存,可通过刷新 OSS 元数据或 CDN 配置补齐。 - SpacetimeDB 观测分为两类:procedure / reducer 调用继续用 `genarrative.spacetime.procedure.*`,订阅本地 cache 读使用 `genarrative.spacetime.read.*`。`read=list_puzzle_gallery` 表示拼图广场当前从 `puzzle_gallery_card_view` 本地 cache 读取,不再每个 HTTP 请求调用 `list_puzzle_gallery` procedure。 - 本地 Windows 直连压测的内存高水位要结合 K6 VU / 连接数解释。250 RPS 下过高 `PREALLOCATED_VUS` 可能让 300 个本地 Established 连接把 `api-server` private memory 瞬时推到 GB 级,且 `/healthz` 小响应也能复现;若压测结束后回落、`response_bodies.in_flight` 和背压 permit 未显示业务积压,应优先按连接 / 发送链路高水位处理,而不是判断为 SpacetimeDB 或 JSON 缓存泄漏。 - Rider 的 Logs 面板只展示 log event 自身字段,不会自动展开父 span 的全部 attributes;请求完成日志会直接带 `request_id`、`http.request.method`、`http.route`、`url.scheme`、`url.path`、`http.response.status_code`、`status_class`、`latency_ms` 和 `slow_request`,完整链路继续到 Traces 面板按 trace/span 查看。 diff --git a/server-rs/Cargo.lock b/server-rs/Cargo.lock index 28712646b..0dc45c3b1 100644 --- a/server-rs/Cargo.lock +++ b/server-rs/Cargo.lock @@ -4125,6 +4125,7 @@ name = "platform-oss" version = "0.1.0" dependencies = [ "base64 0.22.1", + "bytes", "hmac", "reqwest", "serde", diff --git a/server-rs/crates/api-server/src/character_animation_assets.rs b/server-rs/crates/api-server/src/character_animation_assets.rs index bb86bd9af..83b3e074c 100644 --- a/server-rs/crates/api-server/src/character_animation_assets.rs +++ b/server-rs/crates/api-server/src/character_animation_assets.rs @@ -27,7 +27,7 @@ use module_assets::{ }; use platform_oss::{ LegacyAssetPrefix, OssHeadObjectRequest, OssObjectAccess, OssPutObjectRequest, - OssSignedGetObjectUrlRequest, + OssRequestAttemptContext, OssSignedGetObjectUrlRequest, }; use serde::Deserialize; use serde_json::{Value, json}; @@ -2420,8 +2420,10 @@ async fn process_and_persist_editor_character_animation_frame( audit: &crate::external_api_audit::ExternalApiAuditContext, ) -> Result { // 中文注释:每一帧只要求自己的绿幕源图先落 OSS,不再等待整批源图全部上传完成。 - let source_put = put_character_animation_object( + let source_put = put_character_animation_frame_object( state, + frame_index + 1, + "source_put", LegacyAssetPrefix::Animations, vec![ "editor".to_string(), @@ -2469,8 +2471,10 @@ async fn process_and_persist_editor_character_animation_frame( false, )?; let content_type = finalized.mime_type.clone(); - let put_result = put_character_animation_object( + let put_result = put_character_animation_frame_object( state, + frame_index + 1, + "final_put", LegacyAssetPrefix::Animations, vec![ "editor".to_string(), @@ -2497,6 +2501,7 @@ async fn process_and_persist_editor_character_animation_frame( task_id, put_result.object_key.clone(), content_type, + frame_index + 1, ) .await?; @@ -2754,6 +2759,39 @@ async fn publish_single_animation_action( }) } +async fn put_character_animation_frame_object( + state: &AppState, + frame_index: usize, + operation: &'static str, + prefix: LegacyAssetPrefix, + path_segments: Vec, + file_name: String, + content_type: String, + body: Vec, + metadata: BTreeMap, +) -> Result { + require_oss_client(state)? + .put_object_with_retry( + state.character_animation_oss_http_client(), + OssPutObjectRequest { + prefix, + path_segments, + file_name, + content_type: Some(content_type), + access: OssObjectAccess::Private, + metadata, + body, + }, + state.character_animation_oss_io_limiter(), + OssRequestAttemptContext { + frame_index, + operation, + }, + ) + .await + .map_err(map_character_animation_oss_error) +} + async fn put_character_animation_object( state: &AppState, prefix: LegacyAssetPrefix, @@ -2886,13 +2924,27 @@ async fn confirm_editor_character_animation_frame_asset_object( task_id: &str, object_key: String, content_type: String, + frame_index: usize, ) -> Result { - confirm_editor_character_animation_asset_object( + let oss_client = require_oss_client(state)?; + let head = oss_client + .head_object_with_retry( + state.character_animation_oss_http_client(), + OssHeadObjectRequest { object_key }, + state.character_animation_oss_io_limiter(), + OssRequestAttemptContext { + frame_index, + operation: "final_head", + }, + ) + .await + .map_err(map_character_animation_oss_error)?; + confirm_editor_character_animation_asset_object_from_head( state, owner_user_id, source_layer_id, task_id, - object_key, + head, content_type, EDITOR_CHARACTER_ANIMATION_ASSET_KIND, ) @@ -2913,6 +2965,27 @@ async fn confirm_editor_character_animation_asset_object( .head_object(&reqwest::Client::new(), OssHeadObjectRequest { object_key }) .await .map_err(map_character_animation_oss_error)?; + confirm_editor_character_animation_asset_object_from_head( + state, + owner_user_id, + source_layer_id, + task_id, + head, + content_type, + asset_kind, + ) + .await +} + +async fn confirm_editor_character_animation_asset_object_from_head( + state: &AppState, + owner_user_id: &str, + source_layer_id: &str, + task_id: &str, + head: platform_oss::OssHeadObjectResponse, + content_type: String, + asset_kind: &str, +) -> Result { let now_micros = current_utc_micros(); let record = state .spacetime_client() @@ -6218,13 +6291,13 @@ mod tests { "async fn process_and_persist_editor_character_animation_frame", "async fn publish_animation_set", &[ - "put_character_animation_object", + "put_character_animation_frame_object", "green-screen-frame", "remove_editor_generated_screen_background_with_bgfilter_with_request_timeout", "source_put.object_key.as_str()", "bgfilter_request_timeout_ms", "finalize_animation_frame_payload", - "put_character_animation_object", + "put_character_animation_frame_object", "animation_frame", ], ); diff --git a/server-rs/crates/api-server/src/state.rs b/server-rs/crates/api-server/src/state.rs index f5c3a3495..6d0b18d9a 100644 --- a/server-rs/crates/api-server/src/state.rs +++ b/server-rs/crates/api-server/src/state.rs @@ -48,6 +48,7 @@ use crate::work_author::{ }; const ADMIN_ROLE: &str = "admin"; +pub(crate) const CHARACTER_ANIMATION_OSS_MAX_CONCURRENCY: usize = 8; pub type HttpRequestPermitPool = Semaphore; @@ -259,6 +260,8 @@ pub struct AppStateInner { editor_agent_llm_client: Option, matting_client: Option, editor_bgfilter_http_client: reqwest::Client, + character_animation_oss_http_client: reqwest::Client, + character_animation_oss_io_limiter: Arc, #[cfg(any())] creative_agent_executor: Arc, // Phase 1 任务 E 的 creative session facade 暂存在 api-server。 @@ -502,6 +505,9 @@ impl AppState { let editor_agent_llm_client = build_editor_agent_llm_client(&config)?; let matting_client = build_matting_client(&config)?; let editor_bgfilter_http_client = build_editor_bgfilter_http_client(&config)?; + let character_animation_oss_http_client = build_character_animation_oss_http_client()?; + let character_animation_oss_io_limiter = + Arc::new(Semaphore::new(CHARACTER_ANIMATION_OSS_MAX_CONCURRENCY)); let http_request_permit_pools = HttpRequestPermitPools::from_config(&config); let (profile_recharge_order_updates, _) = broadcast::channel(128); @@ -544,6 +550,8 @@ impl AppState { editor_agent_llm_client, matting_client, editor_bgfilter_http_client, + character_animation_oss_http_client, + character_animation_oss_io_limiter, #[cfg(any())] creative_agent_executor: Arc::new(MockLangChainRustAgentExecutor), #[cfg(any())] @@ -1260,6 +1268,14 @@ impl AppState { &self.editor_bgfilter_http_client } + pub fn character_animation_oss_http_client(&self) -> &reqwest::Client { + &self.character_animation_oss_http_client + } + + pub fn character_animation_oss_io_limiter(&self) -> Arc { + self.character_animation_oss_io_limiter.clone() + } + #[cfg(any())] pub fn creative_agent_executor(&self) -> Arc { self.creative_agent_executor.clone() @@ -1991,6 +2007,21 @@ fn build_editor_bgfilter_http_client( }) } +fn build_character_animation_oss_http_client() -> Result { + reqwest::Client::builder() + .connect_timeout(std::time::Duration::from_secs(30)) + .timeout(std::time::Duration::from_secs(60)) + .pool_idle_timeout(std::time::Duration::from_secs(300)) + .pool_max_idle_per_host(8) + .tcp_keepalive(std::time::Duration::from_secs(60)) + .build() + .map_err(|error| { + AppStateInitError::DependencyUnavailable(format!( + "初始化角色动画 OSS HTTP Client 失败:{error}" + )) + }) +} + fn build_wechat_client(config: &AppConfig) -> WechatClient { WechatClient::new(WechatConfig { app_id: config.wechat_mini_program_app_id.clone(), @@ -2120,6 +2151,22 @@ mod tests { use super::*; + #[test] + fn app_state_reuses_character_animation_oss_client_and_eight_permits() { + let state = AppState::new(AppConfig::default()).expect("state should build"); + + assert!(std::ptr::eq( + state.character_animation_oss_http_client(), + state.character_animation_oss_http_client(), + )); + assert_eq!( + state + .character_animation_oss_io_limiter() + .available_permits(), + CHARACTER_ANIMATION_OSS_MAX_CONCURRENCY + ); + } + #[test] fn editor_generation_pricing_typed_record_round_trips() { let expected = crate::editor_generation_config::parse_editor_generation_pricing_json( diff --git a/server-rs/crates/platform-oss/Cargo.toml b/server-rs/crates/platform-oss/Cargo.toml index 2b8a9b9ea..e514611cd 100644 --- a/server-rs/crates/platform-oss/Cargo.toml +++ b/server-rs/crates/platform-oss/Cargo.toml @@ -6,13 +6,15 @@ license.workspace = true [dependencies] base64 = { workspace = true } +bytes = { workspace = true } hmac = { workspace = true } reqwest = { workspace = true, features = ["rustls-tls"] } serde = { workspace = true } serde_json = { workspace = true } sha2 = { workspace = true } time = { workspace = true, features = ["formatting"] } +tokio = { workspace = true, features = ["sync", "time"] } tracing = { workspace = true } [dev-dependencies] -tokio = { workspace = true, features = ["macros", "rt"] } +tokio = { workspace = true, features = ["macros", "rt", "net", "io-util"] } diff --git a/server-rs/crates/platform-oss/src/lib.rs b/server-rs/crates/platform-oss/src/lib.rs index 37dd079a1..c563ee158 100644 --- a/server-rs/crates/platform-oss/src/lib.rs +++ b/server-rs/crates/platform-oss/src/lib.rs @@ -1,12 +1,14 @@ -use std::{collections::BTreeMap, error::Error, fmt, time::Instant}; +use std::{collections::BTreeMap, error::Error, fmt, future::Future, sync::Arc, time::Instant}; use base64::{Engine as _, engine::general_purpose::STANDARD as BASE64_STANDARD}; +use bytes::Bytes; use hmac::{Hmac, Mac}; use reqwest::Method; use serde::{Deserialize, Serialize}; use serde_json::{Value, json}; use sha2::{Digest, Sha256}; use time::{Duration, OffsetDateTime, format_description::well_known::Rfc3339}; +use tokio::{sync::Semaphore, time::sleep}; use tracing::{info, warn}; type HmacSha256 = Hmac; @@ -212,12 +214,30 @@ pub struct OssClient { config: OssConfig, } +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum OssRequestOperation { + Put, + Head, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct OssRequestError { + pub status: Option, + pub timeout: bool, + pub connect: bool, + pub transport: bool, + pub oss_code: Option, + pub oss_request_id: Option, + pub operation: OssRequestOperation, + pub message: String, +} + #[derive(Debug, PartialEq, Eq)] pub enum OssError { InvalidConfig(String), InvalidRequest(String), ObjectNotFound(String), - Request(String), + Request(OssRequestError), SerializePolicy(String), Sign(String), } @@ -233,6 +253,33 @@ pub enum OssErrorKind { Sign, } +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct OssRequestAttemptContext { + pub frame_index: usize, + pub operation: &'static str, +} + +const CHARACTER_ANIMATION_OSS_MAX_ATTEMPTS: u32 = 3; +const CHARACTER_ANIMATION_OSS_RETRY_DELAYS_MS: [u64; 2] = [250, 500]; +const CHARACTER_ANIMATION_OSS_ERROR_BODY_MAX_BYTES: usize = 16 * 1024; +const OSS_ERROR_CODE_MAX_BYTES: usize = 128; +const OSS_REQUEST_ID_MAX_BYTES: usize = 256; + +struct PreparedHeadObject { + object_key: String, + target_url: reqwest::Url, +} + +struct PreparedPutObject { + object_key: String, + target_url: reqwest::Url, + content_type: Option, + headers: BTreeMap, + content_length: u64, + access: OssObjectAccess, + body: Bytes, +} + impl LegacyAssetPrefix { pub fn parse(raw: &str) -> Option { let normalized = raw @@ -697,7 +744,12 @@ impl OssClient { let object_key = normalize_object_key(&request.object_key)?; let target_url = build_object_url(&self.config.bucket, &self.config.endpoint, &object_key).map_err( - |error| OssError::Request(format!("构造 OSS 对象 URL 失败:{error}")), + |error| { + request_error( + OssRequestOperation::Head, + &format!("构造 OSS 对象 URL 失败:{error}"), + ) + }, )?; let response = send_signed_request( client, @@ -705,6 +757,7 @@ impl OssClient { Method::HEAD, Some(&object_key), target_url, + OssRequestOperation::Head, ) .await?; response_status = Some(response.status().as_u16()); @@ -717,10 +770,11 @@ impl OssClient { } if !response.status().is_success() { - return Err(OssError::Request(format!( - "OSS HEAD Object 失败,状态码:{}", - response.status() - ))); + return Err(request_status_error( + OssRequestOperation::Head, + response.status().as_u16(), + format!("OSS HEAD Object 失败,状态码:{}", response.status()), + )); } let headers = response.headers(); @@ -824,7 +878,12 @@ impl OssClient { let headers = build_put_object_headers(request.metadata)?; let target_url = build_object_url(&self.config.bucket, &self.config.endpoint, &object_key).map_err( - |error| OssError::Request(format!("构造 OSS 对象 URL 失败:{error}")), + |error| { + request_error( + OssRequestOperation::Put, + &format!("构造 OSS 对象 URL 失败:{error}"), + ) + }, )?; let content_length = u64::try_from(request.body.len()) .map_err(|_| OssError::InvalidRequest("上传对象大小超出可支持范围".to_string()))?; @@ -843,14 +902,15 @@ impl OssClient { let response = builder .send() .await - .map_err(|error| OssError::Request(format!("请求 OSS 失败:{error}")))?; + .map_err(|error| request_error_from_reqwest(OssRequestOperation::Put, error))?; response_status = Some(response.status().as_u16()); if !response.status().is_success() { - return Err(OssError::Request(format!( - "OSS PutObject 失败,状态码:{}", - response.status() - ))); + return Err(request_status_error( + OssRequestOperation::Put, + response.status().as_u16(), + format!("OSS PutObject 失败,状态码:{}", response.status()), + )); } let headers = response.headers(); @@ -916,6 +976,533 @@ impl OssClient { result } + + /// 角色动画帧专用的可重试 PUT。调用方传入进程级并发限制器,单次网络 attempt + /// 独占一个 permit,退避等待期间不会占用 permit。 + pub async fn put_object_with_retry( + &self, + client: &reqwest::Client, + request: OssPutObjectRequest, + io_limiter: Arc, + attempt_context: OssRequestAttemptContext, + ) -> Result { + let prepared = self.prepare_put_object(request)?; + let object_key = prepared.object_key.clone(); + + run_animation_request_with_retry(io_limiter, attempt_context, &object_key, || { + self.put_object_once(client, &prepared) + }) + .await + } + + /// 角色动画帧专用的可重试 HEAD。HEAD 与 PUT 独立重试,HEAD 失败不会重新上传 PUT。 + pub async fn head_object_with_retry( + &self, + client: &reqwest::Client, + request: OssHeadObjectRequest, + io_limiter: Arc, + attempt_context: OssRequestAttemptContext, + ) -> Result { + let prepared = self.prepare_head_object(request)?; + let object_key = prepared.object_key.clone(); + + run_animation_request_with_retry(io_limiter, attempt_context, &object_key, || { + self.head_object_once(client, &prepared) + }) + .await + } + + fn prepare_head_object( + &self, + request: OssHeadObjectRequest, + ) -> Result { + let object_key = normalize_object_key(&request.object_key)?; + let target_url = build_object_url(&self.config.bucket, &self.config.endpoint, &object_key) + .map_err(|error| { + request_error( + OssRequestOperation::Head, + &format!("构造 OSS 对象 URL 失败:{error}"), + ) + })?; + Ok(PreparedHeadObject { + object_key, + target_url, + }) + } + + async fn head_object_once( + &self, + client: &reqwest::Client, + prepared: &PreparedHeadObject, + ) -> Result<(OssHeadObjectResponse, u16), OssError> { + let response = send_signed_request( + client, + &self.config, + Method::HEAD, + Some(&prepared.object_key), + prepared.target_url.clone(), + OssRequestOperation::Head, + ) + .await?; + + if response.status() == reqwest::StatusCode::NOT_FOUND { + return Err(OssError::ObjectNotFound(format!( + "OSS 对象不存在:{}", + prepared.object_key + ))); + } + if !response.status().is_success() { + return Err(request_status_error( + OssRequestOperation::Head, + response.status().as_u16(), + format!("OSS HEAD Object 失败,状态码:{}", response.status()), + )); + } + + let headers = response.headers(); + let content_length = headers + .get(reqwest::header::CONTENT_LENGTH) + .and_then(|value| value.to_str().ok()) + .and_then(|value| value.parse::().ok()) + .unwrap_or(0); + let content_type = headers + .get(reqwest::header::CONTENT_TYPE) + .and_then(|value| value.to_str().ok()) + .map(str::to_string); + let etag = headers + .get(reqwest::header::ETAG) + .and_then(|value| value.to_str().ok()) + .map(|value| value.trim_matches('"').to_string()); + let last_modified = headers + .get(reqwest::header::LAST_MODIFIED) + .and_then(|value| value.to_str().ok()) + .map(str::to_string); + + Ok(( + OssHeadObjectResponse { + bucket: self.config.bucket.clone(), + object_key: prepared.object_key.clone(), + content_length, + content_type, + etag, + last_modified, + }, + response.status().as_u16(), + )) + } + + fn prepare_put_object( + &self, + request: OssPutObjectRequest, + ) -> Result { + if request.body.is_empty() { + return Err(OssError::InvalidRequest( + "服务端上传对象内容不能为空".to_string(), + )); + } + + let sanitized_segments = request + .path_segments + .iter() + .map(|segment| sanitize_path_segment(segment)) + .filter(|segment| !segment.is_empty()) + .collect::>(); + let file_name = sanitize_file_name(&request.file_name)?; + let object_key = build_object_key(request.prefix, &sanitized_segments, &file_name); + let content_type = normalize_optional_value(request.content_type); + let headers = build_put_object_headers(request.metadata)?; + let target_url = build_object_url(&self.config.bucket, &self.config.endpoint, &object_key) + .map_err(|error| { + request_error( + OssRequestOperation::Put, + &format!("构造 OSS 对象 URL 失败:{error}"), + ) + })?; + let content_length = u64::try_from(request.body.len()) + .map_err(|_| OssError::InvalidRequest("上传对象大小超出可支持范围".to_string()))?; + + Ok(PreparedPutObject { + object_key, + target_url, + content_type, + headers, + content_length, + access: request.access, + body: Bytes::from(request.body), + }) + } + + async fn put_object_once( + &self, + client: &reqwest::Client, + prepared: &PreparedPutObject, + ) -> Result<(OssPutObjectResponse, u16), OssError> { + let response = signed_request_builder( + client, + &self.config, + Method::PUT, + Some(&prepared.object_key), + prepared.target_url.clone(), + prepared.content_type.as_deref(), + &prepared.headers, + )? + .header(reqwest::header::CONTENT_LENGTH, prepared.content_length) + .body(prepared.body.clone()) + .send() + .await + .map_err(|error| request_error_from_reqwest(OssRequestOperation::Put, error))?; + + if !response.status().is_success() { + return Err(request_status_error_from_oss_put_response(response).await); + } + + let headers = response.headers(); + let etag = headers + .get(reqwest::header::ETAG) + .and_then(|value| value.to_str().ok()) + .map(|value| value.trim_matches('"').to_string()); + let last_modified = headers + .get(reqwest::header::LAST_MODIFIED) + .and_then(|value| value.to_str().ok()) + .map(str::to_string); + + Ok(( + OssPutObjectResponse { + provider: OSS_PROVIDER, + bucket: self.config.bucket.clone(), + endpoint: self.config.endpoint.clone(), + host: self.config.upload_host(), + legacy_public_path: format!("/{}", prepared.object_key), + object_key: prepared.object_key.clone(), + content_type: prepared.content_type.clone(), + content_length: prepared.content_length, + access: prepared.access, + etag, + last_modified, + }, + response.status().as_u16(), + )) + } +} + +fn request_error(operation: OssRequestOperation, message: &str) -> OssError { + OssError::Request(OssRequestError { + status: None, + timeout: false, + connect: false, + transport: false, + oss_code: None, + oss_request_id: None, + operation, + message: message.to_string(), + }) +} + +async fn request_status_error_from_oss_put_response(mut response: reqwest::Response) -> OssError { + let status = response.status(); + let header_request_id = response + .headers() + .get("x-oss-request-id") + .and_then(|value| value.to_str().ok()) + .and_then(|value| normalize_oss_error_field(value, OSS_REQUEST_ID_MAX_BYTES)); + let body_read = if status == reqwest::StatusCode::BAD_REQUEST { + read_bounded_oss_error_body(&mut response).await + } else { + OssErrorBodyRead { + body: Vec::new(), + read_failure: None, + } + }; + + request_status_error_from_oss_parts( + OssRequestOperation::Put, + status.as_u16(), + header_request_id, + &body_read.body, + body_read.read_failure, + ) +} + +/// 400 错误响应体读取失败的原因。仅在部分响应体尚未解析出确定性 +/// OSS 错误码时参与重试判定,否则只体现在 message 里。 +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +struct OssErrorBodyReadFailure { + timeout: bool, +} + +struct OssErrorBodyRead { + body: Vec, + read_failure: Option, +} + +async fn read_bounded_oss_error_body(response: &mut reqwest::Response) -> OssErrorBodyRead { + let mut body = Vec::new(); + let mut read_failure = None; + while body.len() < CHARACTER_ANIMATION_OSS_ERROR_BODY_MAX_BYTES { + let chunk = match response.chunk().await { + Ok(Some(chunk)) => chunk, + Ok(None) => break, + Err(error) => { + // 断流/超时不丢弃已读字节:部分响应体可能已含确定性错误码。 + read_failure = Some(OssErrorBodyReadFailure { + timeout: error.is_timeout(), + }); + break; + } + }; + let remaining = CHARACTER_ANIMATION_OSS_ERROR_BODY_MAX_BYTES - body.len(); + body.extend_from_slice(&chunk[..chunk.len().min(remaining)]); + if chunk.len() > remaining { + break; + } + } + OssErrorBodyRead { body, read_failure } +} + +fn request_status_error_from_oss_parts( + operation: OssRequestOperation, + status: u16, + header_request_id: Option, + body: &[u8], + body_read_failure: Option, +) -> OssError { + let header_request_id = header_request_id + .as_deref() + .and_then(|value| normalize_oss_error_field(value, OSS_REQUEST_ID_MAX_BYTES)); + let bounded_body = &body[..body.len().min(CHARACTER_ANIMATION_OSS_ERROR_BODY_MAX_BYTES)]; + let oss_code = extract_oss_error_xml_field(bounded_body, "Code", OSS_ERROR_CODE_MAX_BYTES); + let xml_request_id = + extract_oss_error_xml_field(bounded_body, "RequestId", OSS_REQUEST_ID_MAX_BYTES); + let oss_request_id = header_request_id.or(xml_request_id); + // 确定性错误码优先:已解析出 Code 时,响应体读取失败只保留在 message 里, + // 不改变重试语义;错误码缺失时才按读取失败归类为可重试的超时/传输错误。 + let unclassified_read_failure = if oss_code.is_some() { + None + } else { + body_read_failure + }; + let timeout = (status == reqwest::StatusCode::BAD_REQUEST.as_u16() + && oss_code.as_deref() == Some("RequestTimeout")) + || unclassified_read_failure.is_some_and(|failure| failure.timeout); + let transport = unclassified_read_failure.is_some_and(|failure| !failure.timeout); + let mut message = format!("OSS PutObject 失败,状态码:{status}"); + if let Some(oss_code) = oss_code.as_deref() { + message.push_str(&format!(",OSS 错误码:{oss_code}")); + } + if let Some(oss_request_id) = oss_request_id.as_deref() { + message.push_str(&format!(",OSS Request ID:{oss_request_id}")); + } + if body_read_failure.is_some() { + message.push_str(",错误响应体读取失败"); + } + + OssError::Request(OssRequestError { + status: Some(status), + timeout, + connect: false, + transport, + oss_code, + oss_request_id, + operation, + message, + }) +} + +fn extract_oss_error_xml_field(body: &[u8], field: &str, max_bytes: usize) -> Option { + let body = std::str::from_utf8(body).ok()?; + let start_tag = format!("<{field}>"); + let end_tag = format!(""); + let value_start = body.find(&start_tag)? + start_tag.len(); + let value_end = value_start + body[value_start..].find(&end_tag)?; + normalize_oss_error_field(&body[value_start..value_end], max_bytes) +} + +fn normalize_oss_error_field(value: &str, max_bytes: usize) -> Option { + let value = value.trim(); + if value.is_empty() || value.len() > max_bytes || value.chars().any(char::is_control) { + return None; + } + Some(value.to_string()) +} + +async fn run_animation_request_with_retry( + io_limiter: Arc, + attempt_context: OssRequestAttemptContext, + object_key: &str, + mut attempt_request: F, +) -> Result +where + F: FnMut() -> Fut, + Fut: Future>, +{ + for attempt in 1..=CHARACTER_ANIMATION_OSS_MAX_ATTEMPTS { + let attempt_started_at = Instant::now(); + let permit_wait_started_at = Instant::now(); + let permit = io_limiter + .acquire() + .await + .map_err(|_| OssError::InvalidConfig("角色动画 OSS 并发限制器已关闭".to_string()))?; + let permit_wait_ms = elapsed_ms(permit_wait_started_at); + let result = attempt_request().await; + drop(permit); + + let retryable = result.as_ref().err().is_some_and(oss_error_is_retryable); + let will_retry = retryable && attempt < CHARACTER_ANIMATION_OSS_MAX_ATTEMPTS; + let retry_delay_ms = if will_retry { + CHARACTER_ANIMATION_OSS_RETRY_DELAYS_MS[(attempt - 1) as usize] + } else { + 0 + }; + let success_status = result.as_ref().ok().map(|(_, status)| *status); + log_animation_request_attempt( + attempt_context, + object_key, + attempt, + retryable, + will_retry, + retry_delay_ms, + permit_wait_ms, + elapsed_ms(attempt_started_at), + success_status, + result.as_ref().err(), + ); + + match result { + Ok((response, _status)) => return Ok(response), + Err(_error) if will_retry => { + sleep(std::time::Duration::from_millis(retry_delay_ms)).await + } + Err(error) => return Err(error), + } + } + + unreachable!("角色动画 OSS 重试循环必须返回结果") +} + +fn request_error_from_reqwest(operation: OssRequestOperation, error: reqwest::Error) -> OssError { + let status = error.status().map(|status| status.as_u16()); + let timeout = error.is_timeout(); + let connect = error.is_connect(); + let transport = !timeout && !connect && (error.is_request() || error.is_body()); + + OssError::Request(OssRequestError { + status, + timeout, + connect, + transport, + oss_code: None, + oss_request_id: None, + operation, + message: format!("请求 OSS 失败:{error}"), + }) +} + +fn request_status_error(operation: OssRequestOperation, status: u16, message: String) -> OssError { + OssError::Request(OssRequestError { + status: Some(status), + timeout: false, + connect: false, + transport: false, + oss_code: None, + oss_request_id: None, + operation, + message, + }) +} + +fn oss_error_is_retryable(error: &OssError) -> bool { + let OssError::Request(request_error) = error else { + return false; + }; + + match (request_error.status, request_error.oss_code.as_deref()) { + (Some(400), Some("RequestTimeout")) => true, + // 400 是唯一会读取错误响应体的状态码:错误码缺失且响应体读取 + // 超时/断流时,无法证明是确定性 400,按传输错误重试。 + (Some(400), None) if request_error.timeout || request_error.transport => true, + (Some(408 | 429 | 500..=599), _) => true, + (Some(_), _) => false, + (None, _) => request_error.timeout || request_error.connect || request_error.transport, + } +} + +fn request_error_details( + error: Option<&OssError>, +) -> (Option, bool, bool, bool, Option<&str>, Option<&str>) { + match error { + Some(OssError::Request(request_error)) => ( + request_error.status, + request_error.timeout, + request_error.connect, + request_error.transport, + request_error.oss_code.as_deref(), + request_error.oss_request_id.as_deref(), + ), + Some(OssError::ObjectNotFound(_)) => (Some(404), false, false, false, None, None), + _ => (None, false, false, false, None, None), + } +} + +fn log_animation_request_attempt( + context: OssRequestAttemptContext, + object_key: &str, + attempt: u32, + retryable: bool, + will_retry: bool, + retry_delay_ms: u64, + permit_wait_ms: u64, + elapsed_ms: u64, + success_status: Option, + error: Option<&OssError>, +) { + let (error_status, timeout, connect, transport, oss_code, oss_request_id) = + request_error_details(error); + let status = success_status.or(error_status); + if error.is_none() { + info!( + provider = OSS_PROVIDER, + frame_index = context.frame_index, + object_key, + operation = context.operation, + attempt, + max_attempts = CHARACTER_ANIMATION_OSS_MAX_ATTEMPTS, + retryable, + will_retry, + retry_delay_ms, + permit_wait_ms, + timeout, + connect, + transport, + oss_code = oss_code.unwrap_or_default(), + oss_request_id = oss_request_id.unwrap_or_default(), + status = status.unwrap_or_default(), + elapsed_ms, + "角色动画 OSS 请求 attempt 完成" + ); + return; + } + warn!( + provider = OSS_PROVIDER, + frame_index = context.frame_index, + object_key, + operation = context.operation, + attempt, + max_attempts = CHARACTER_ANIMATION_OSS_MAX_ATTEMPTS, + retryable, + will_retry, + retry_delay_ms, + permit_wait_ms, + timeout, + connect, + transport, + oss_code = oss_code.unwrap_or_default(), + oss_request_id = oss_request_id.unwrap_or_default(), + status = status.unwrap_or_default(), + elapsed_ms, + error_kind = error.map(oss_error_kind_label), + message = error.map(ToString::to_string).unwrap_or_default(), + "角色动画 OSS 请求 attempt 完成" + ); } impl fmt::Display for OssError { @@ -923,10 +1510,10 @@ impl fmt::Display for OssError { match self { Self::InvalidConfig(message) | Self::InvalidRequest(message) - | Self::ObjectNotFound(message) - | Self::Request(message) | Self::SerializePolicy(message) | Self::Sign(message) => f.write_str(message), + Self::ObjectNotFound(message) => f.write_str(message), + Self::Request(error) => f.write_str(&error.message), } } } @@ -1325,6 +1912,7 @@ async fn send_signed_request( method: Method, object_key: Option<&str>, target_url: reqwest::Url, + operation: OssRequestOperation, ) -> Result { signed_request_builder( client, @@ -1337,7 +1925,7 @@ async fn send_signed_request( )? .send() .await - .map_err(|error| OssError::Request(format!("请求 OSS 失败:{error}"))) + .map_err(|error| request_error_from_reqwest(operation, error)) } fn signed_request_builder( @@ -1609,6 +2197,32 @@ fn encode_url_query_value(value: &str) -> String { #[cfg(test)] mod tests { use super::*; + use std::sync::atomic::{AtomicUsize, Ordering}; + + fn retryable_request_error( + status: Option, + timeout: bool, + connect: bool, + transport: bool, + ) -> OssError { + OssError::Request(OssRequestError { + status, + timeout, + connect, + transport, + oss_code: None, + oss_request_id: None, + operation: OssRequestOperation::Put, + message: "mock request failure".to_string(), + }) + } + + fn animation_attempt_context(operation: &'static str) -> OssRequestAttemptContext { + OssRequestAttemptContext { + frame_index: 1, + operation, + } + } #[test] fn oss_error_kind_is_stable_for_adapter_mapping() { @@ -1621,11 +2235,550 @@ mod tests { OssErrorKind::ObjectNotFound ); assert_eq!( - OssError::Request("network".to_string()).kind(), + OssError::Request(OssRequestError { + status: None, + timeout: false, + connect: false, + transport: true, + oss_code: None, + oss_request_id: None, + operation: OssRequestOperation::Put, + message: "network".to_string(), + }) + .kind(), OssErrorKind::Request ); } + #[test] + fn object_not_found_attempt_log_details_keep_404_status() { + let error = OssError::ObjectNotFound("missing".to_string()); + assert_eq!( + request_error_details(Some(&error)), + (Some(404), false, false, false, None, None) + ); + assert_eq!(oss_error_kind_label(&error), "object_not_found"); + assert!(!oss_error_is_retryable(&error)); + } + + #[test] + fn oss_request_timeout_400_is_retryable_and_prefers_header_request_id() { + let body = br#" +RequestTimeoutxml-request-id"#; + let error = request_status_error_from_oss_parts( + OssRequestOperation::Put, + 400, + Some("header-request-id".to_string()), + body, + None, + ); + let OssError::Request(request_error) = &error else { + panic!("OSS status failure should remain a request error"); + }; + + assert_eq!(request_error.status, Some(400)); + assert!(request_error.timeout); + assert_eq!(request_error.oss_code.as_deref(), Some("RequestTimeout")); + assert_eq!( + request_error.oss_request_id.as_deref(), + Some("header-request-id") + ); + assert!(oss_error_is_retryable(&error)); + } + + #[test] + fn oss_request_timeout_400_uses_xml_request_id_when_header_is_missing() { + let body = br#" +RequestTimeoutxml-request-id +"#; + let error = + request_status_error_from_oss_parts(OssRequestOperation::Put, 400, None, body, None); + let OssError::Request(request_error) = &error else { + panic!("OSS status failure should remain a request error"); + }; + + assert_eq!( + request_error.oss_request_id.as_deref(), + Some("xml-request-id") + ); + assert!(oss_error_is_retryable(&error)); + } + + #[test] + fn other_oss_400_errors_and_malformed_xml_are_not_retryable() { + for body in [ + b"InvalidArgument".as_slice(), + b"RequestTimeout".as_slice(), + b"not xml".as_slice(), + ] { + let error = request_status_error_from_oss_parts( + OssRequestOperation::Put, + 400, + None, + body, + None, + ); + assert!(!oss_error_is_retryable(&error)); + } + + let error = request_status_error_from_oss_parts( + OssRequestOperation::Put, + 403, + None, + b"RequestTimeout", + None, + ); + assert!(!oss_error_is_retryable(&error)); + } + + #[test] + fn oss_error_xml_fields_beyond_body_limit_are_ignored() { + let mut body = vec![b' '; CHARACTER_ANIMATION_OSS_ERROR_BODY_MAX_BYTES]; + body.extend_from_slice( + b"RequestTimeoutlate", + ); + let error = + request_status_error_from_oss_parts(OssRequestOperation::Put, 400, None, &body, None); + let OssError::Request(request_error) = &error else { + panic!("OSS status failure should remain a request error"); + }; + + assert_eq!(request_error.oss_code, None); + assert_eq!(request_error.oss_request_id, None); + assert!(!request_error.timeout); + assert!(!oss_error_is_retryable(&error)); + } + + #[test] + fn oss_400_without_code_and_broken_body_read_is_retryable() { + for (read_timeout, expect_transport) in [(true, false), (false, true)] { + let error = request_status_error_from_oss_parts( + OssRequestOperation::Put, + 400, + None, + b"Request", + Some(OssErrorBodyReadFailure { + timeout: read_timeout, + }), + ); + let OssError::Request(request_error) = &error else { + panic!("OSS status failure should remain a request error"); + }; + + assert_eq!(request_error.status, Some(400)); + assert_eq!(request_error.oss_code, None); + assert_eq!(request_error.timeout, read_timeout); + assert_eq!(request_error.transport, expect_transport); + assert!(request_error.message.contains("错误响应体读取失败")); + assert!(oss_error_is_retryable(&error)); + } + } + + #[test] + fn oss_400_with_parsed_code_keeps_deterministic_semantics_on_read_failure() { + let error = request_status_error_from_oss_parts( + OssRequestOperation::Put, + 400, + None, + b"InvalidArgumentpartial", + Some(OssErrorBodyReadFailure { timeout: true }), + ); + let OssError::Request(request_error) = &error else { + panic!("OSS status failure should remain a request error"); + }; + + assert_eq!(request_error.oss_code.as_deref(), Some("InvalidArgument")); + assert!(!request_error.timeout); + assert!(!request_error.transport); + assert!(request_error.message.contains("错误响应体读取失败")); + assert!(!oss_error_is_retryable(&error)); + + let error = request_status_error_from_oss_parts( + OssRequestOperation::Put, + 400, + None, + b"RequestTimeout", + Some(OssErrorBodyReadFailure { timeout: false }), + ); + assert!(oss_error_is_retryable(&error)); + } + + const MOCK_PUT_BODY: &[u8] = b"animation-frame-bytes"; + + /// 极简 HTTP/1.1 mock:读完整个 PUT 请求后返回 400 与部分 XML 响应体 + /// (Content-Length 大于实际发送字节),`stall_before_close` 决定挂住 + /// 连接触发客户端读超时,还是直接断开触发传输错误。 + async fn spawn_broken_error_body_server(stall_before_close: bool) -> std::net::SocketAddr { + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + + let listener = tokio::net::TcpListener::bind("127.0.0.1:0") + .await + .expect("mock server should bind"); + let addr = listener + .local_addr() + .expect("mock server should expose its addr"); + tokio::spawn(async move { + let Ok((mut socket, _)) = listener.accept().await else { + return; + }; + let mut received = Vec::new(); + let mut buffer = [0u8; 4096]; + while !received.ends_with(MOCK_PUT_BODY) { + match socket.read(&mut buffer).await { + Ok(0) | Err(_) => return, + Ok(read) => received.extend_from_slice(&buffer[..read]), + } + } + let response = "HTTP/1.1 400 Bad Request\r\n\ + x-oss-request-id: mock-request-id\r\n\ + Content-Type: application/xml\r\n\ + Content-Length: 4096\r\n\ + \r\n\ + Request"; + if socket.write_all(response.as_bytes()).await.is_err() { + return; + } + let _ = socket.flush().await; + if stall_before_close { + tokio::time::sleep(std::time::Duration::from_secs(5)).await; + } + }); + addr + } + + #[tokio::test] + async fn oss_400_with_broken_error_body_stream_is_retryable_transport() { + let addr = spawn_broken_error_body_server(false).await; + let response = reqwest::Client::new() + .put(format!("http://{addr}/generated-animations/frame01.png")) + .body(MOCK_PUT_BODY.to_vec()) + .send() + .await + .expect("response headers should arrive before the body breaks"); + assert_eq!(response.status(), reqwest::StatusCode::BAD_REQUEST); + + let error = request_status_error_from_oss_put_response(response).await; + let OssError::Request(request_error) = &error else { + panic!("OSS status failure should remain a request error"); + }; + + assert_eq!(request_error.status, Some(400)); + assert_eq!(request_error.oss_code, None); + assert!(!request_error.timeout); + assert!(request_error.transport); + assert_eq!( + request_error.oss_request_id.as_deref(), + Some("mock-request-id") + ); + assert!(request_error.message.contains("错误响应体读取失败")); + assert!(oss_error_is_retryable(&error)); + } + + #[tokio::test] + async fn oss_400_with_stalled_error_body_stream_is_retryable_timeout() { + let addr = spawn_broken_error_body_server(true).await; + let client = reqwest::Client::builder() + .timeout(std::time::Duration::from_millis(300)) + .build() + .expect("test client should build"); + let response = client + .put(format!("http://{addr}/generated-animations/frame01.png")) + .body(MOCK_PUT_BODY.to_vec()) + .send() + .await + .expect("response headers should arrive before the body stalls"); + assert_eq!(response.status(), reqwest::StatusCode::BAD_REQUEST); + + let error = request_status_error_from_oss_put_response(response).await; + let OssError::Request(request_error) = &error else { + panic!("OSS status failure should remain a request error"); + }; + + assert_eq!(request_error.status, Some(400)); + assert_eq!(request_error.oss_code, None); + assert!(request_error.timeout); + assert!(!request_error.transport); + assert!(oss_error_is_retryable(&error)); + } + + #[tokio::test] + async fn reqwest_builder_error_is_not_retryable_transport() { + let error = reqwest::Client::new() + .put("https://example.com") + .header("x-oss-meta-invalid", "first line\nsecond line") + .send() + .await + .expect_err("invalid header must fail while building the request"); + assert!(error.is_builder()); + + let error = request_error_from_reqwest(OssRequestOperation::Put, error); + let OssError::Request(request_error) = &error else { + panic!("builder failure should remain an OSS request error"); + }; + + assert_eq!(request_error.status, None); + assert!(!request_error.timeout); + assert!(!request_error.connect); + assert!(!request_error.transport); + assert!(!oss_error_is_retryable(&error)); + } + + #[tokio::test] + async fn animation_retry_retries_transport_then_succeeds() { + let attempts = Arc::new(AtomicUsize::new(0)); + let attempts_for_request = attempts.clone(); + let result = run_animation_request_with_retry( + Arc::new(Semaphore::new(8)), + animation_attempt_context("source_put"), + "generated-animations/editor/layer/task/green-screen-frame01.png", + move || { + let attempts = attempts_for_request.clone(); + async move { + if attempts.fetch_add(1, Ordering::SeqCst) == 0 { + Err(retryable_request_error(None, false, false, true)) + } else { + Ok(("uploaded", 201)) + } + } + }, + ) + .await; + + assert_eq!(result.expect("second attempt should succeed"), "uploaded"); + assert_eq!(attempts.load(Ordering::SeqCst), 2); + } + + #[tokio::test] + async fn animation_retry_retries_timeout_and_retryable_statuses() { + for (status, timeout, connect, transport) in [ + (None, true, false, false), + (Some(408), false, false, false), + (Some(429), false, false, false), + (Some(500), false, false, false), + (Some(502), false, false, false), + (Some(503), false, false, false), + (Some(504), false, false, false), + ] { + let attempts = Arc::new(AtomicUsize::new(0)); + let attempts_for_request = attempts.clone(); + let result = run_animation_request_with_retry( + Arc::new(Semaphore::new(8)), + animation_attempt_context("final_put"), + "generated-animations/editor/layer/task/frame01.png", + move || { + let attempts = attempts_for_request.clone(); + async move { + if attempts.fetch_add(1, Ordering::SeqCst) == 0 { + Err(retryable_request_error(status, timeout, connect, transport)) + } else { + Ok(((), 204)) + } + } + }, + ) + .await; + + assert!(result.is_ok(), "status={status:?} should retry"); + assert_eq!(attempts.load(Ordering::SeqCst), 2, "status={status:?}"); + } + } + + #[tokio::test] + async fn animation_retry_does_not_retry_deterministic_statuses() { + for status in [400, 403, 404] { + let attempts = Arc::new(AtomicUsize::new(0)); + let attempts_for_request = attempts.clone(); + let result = run_animation_request_with_retry( + Arc::new(Semaphore::new(8)), + animation_attempt_context("final_head"), + "generated-animations/editor/layer/task/frame01.png", + move || { + let attempts = attempts_for_request.clone(); + async move { + attempts.fetch_add(1, Ordering::SeqCst); + Err::<((), u16), _>(retryable_request_error( + Some(status), + false, + false, + false, + )) + } + }, + ) + .await; + + assert!(result.is_err(), "status={status} should fail"); + assert_eq!(attempts.load(Ordering::SeqCst), 1, "status={status}"); + } + } + + #[tokio::test] + async fn animation_retry_returns_final_error_after_three_attempts() { + let attempts = Arc::new(AtomicUsize::new(0)); + let attempts_for_request = attempts.clone(); + let result = run_animation_request_with_retry( + Arc::new(Semaphore::new(8)), + animation_attempt_context("final_put"), + "generated-animations/editor/layer/task/frame01.png", + move || { + let attempts = attempts_for_request.clone(); + async move { + attempts.fetch_add(1, Ordering::SeqCst); + Err::<((), u16), _>(retryable_request_error(Some(503), false, false, false)) + } + }, + ) + .await; + + assert!(result.is_err()); + assert_eq!(attempts.load(Ordering::SeqCst), 3); + } + + #[tokio::test] + async fn animation_retry_keeps_head_retry_independent_from_successful_put() { + let put_attempts = Arc::new(AtomicUsize::new(0)); + let put_attempts_for_request = put_attempts.clone(); + run_animation_request_with_retry( + Arc::new(Semaphore::new(8)), + animation_attempt_context("final_put"), + "generated-animations/editor/layer/task/frame01.png", + move || { + put_attempts_for_request.fetch_add(1, Ordering::SeqCst); + async { Ok::<_, OssError>(((), 204)) } + }, + ) + .await + .expect("PUT should succeed once"); + + let head_attempts = Arc::new(AtomicUsize::new(0)); + let head_attempts_for_request = head_attempts.clone(); + run_animation_request_with_retry( + Arc::new(Semaphore::new(8)), + animation_attempt_context("final_head"), + "generated-animations/editor/layer/task/frame01.png", + move || { + let head_attempts = head_attempts_for_request.clone(); + async move { + if head_attempts.fetch_add(1, Ordering::SeqCst) == 0 { + Err(retryable_request_error(Some(503), false, false, false)) + } else { + Ok(((), 204)) + } + } + }, + ) + .await + .expect("HEAD should succeed on its retry"); + + assert_eq!(put_attempts.load(Ordering::SeqCst), 1); + assert_eq!(head_attempts.load(Ordering::SeqCst), 2); + } + + #[tokio::test] + async fn animation_retry_keeps_permit_bound_at_eight_in_flight_requests() { + let limiter = Arc::new(Semaphore::new(8)); + let current = Arc::new(AtomicUsize::new(0)); + let maximum = Arc::new(AtomicUsize::new(0)); + let mut tasks = Vec::new(); + + for frame_index in 0..32 { + let limiter = limiter.clone(); + let current_for_request = current.clone(); + let maximum_for_request = maximum.clone(); + tasks.push(tokio::spawn(async move { + run_animation_request_with_retry( + limiter, + OssRequestAttemptContext { + frame_index: frame_index + 1, + operation: "source_put", + }, + "generated-animations/editor/layer/task/green-screen-frame01.png", + move || { + let current = current_for_request.clone(); + let maximum = maximum_for_request.clone(); + async move { + let in_flight = current.fetch_add(1, Ordering::SeqCst) + 1; + maximum.fetch_max(in_flight, Ordering::SeqCst); + tokio::time::sleep(std::time::Duration::from_millis(5)).await; + current.fetch_sub(1, Ordering::SeqCst); + Ok::<_, OssError>(((), 204)) + } + }, + ) + .await + })); + } + + for task in tasks { + task.await + .expect("mock request task should join") + .expect("request should pass"); + } + assert!(maximum.load(Ordering::SeqCst) <= 8); + } + + #[tokio::test] + async fn animation_retry_reuses_object_key_and_body_across_attempts() { + let client = OssClient::new( + OssConfig::new( + "bucket".to_string(), + "oss-cn-shanghai.aliyuncs.com".to_string(), + "access-key".to_string(), + "access-secret".to_string(), + 60, + 60, + 1024, + 204, + ) + .expect("test OSS config should be valid"), + ); + let prepared = client + .prepare_put_object(OssPutObjectRequest { + prefix: LegacyAssetPrefix::Animations, + path_segments: vec![ + "editor".to_string(), + "layer".to_string(), + "task".to_string(), + ], + file_name: "frame01.png".to_string(), + content_type: Some("image/png".to_string()), + access: OssObjectAccess::Private, + metadata: BTreeMap::new(), + body: vec![1, 2, 3, 4], + }) + .expect("test request should be prepared"); + let expected_key = prepared.object_key.clone(); + let expected_body = prepared.body.clone(); + let expected_key_for_request = expected_key.clone(); + let attempts = Arc::new(AtomicUsize::new(0)); + let attempts_for_request = attempts.clone(); + let result = run_animation_request_with_retry( + Arc::new(Semaphore::new(8)), + animation_attempt_context("final_put"), + &expected_key_for_request, + move || { + let attempts = attempts_for_request.clone(); + let key = prepared.object_key.clone(); + let body = prepared.body.clone(); + let expected_key = expected_key.clone(); + let expected_body = expected_body.clone(); + async move { + assert_eq!(key, expected_key); + assert_eq!(body, expected_body); + if attempts.fetch_add(1, Ordering::SeqCst) == 0 { + Err(retryable_request_error(Some(503), false, false, false)) + } else { + Ok(((), 204)) + } + } + }, + ) + .await; + + assert!(result.is_ok()); + assert_eq!(attempts.load(Ordering::SeqCst), 2); + } + #[test] fn structured_log_labels_are_stable() { assert_eq!(