From 0944667c370747f49c9fe1b2aae61f271449d10f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=8E=8B=E5=BE=B7=E5=AE=87?= Date: Sat, 12 Sep 2026 16:22:06 +0800 Subject: [PATCH] =?UTF-8?q?=E9=99=90=E5=88=B6=E5=8E=9F=E5=A7=8B=E5=9B=BE?= =?UTF-8?q?=E7=89=87=E8=A7=A3=E7=A0=81=E5=B9=B6=E5=8F=91?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 为 raw image PNG 预处理增加 AppState 专用 semaphore。 按请求截止时间限制槽位等待和 blocking 解码等待,并让 permit 持有到任务结束。 --- server-rs/crates/api-server/src/raw_image.rs | 59 ++++++++++++++++++-- server-rs/crates/api-server/src/state.rs | 11 ++++ 2 files changed, 64 insertions(+), 6 deletions(-) diff --git a/server-rs/crates/api-server/src/raw_image.rs b/server-rs/crates/api-server/src/raw_image.rs index 6f3b55c52..c6bdf1d65 100644 --- a/server-rs/crates/api-server/src/raw_image.rs +++ b/server-rs/crates/api-server/src/raw_image.rs @@ -13,6 +13,7 @@ use platform_image::{ use serde::Serialize; use serde_json::json; use std::io::Cursor; +use std::time::{Duration, Instant}; use crate::{ asset_billing::{ @@ -25,7 +26,7 @@ use crate::{ require_openai_image_settings, }, request_context::RequestContext, - state::AppState, + state::{AppState, RAW_IMAGE_DECODE_MAX_CONCURRENCY}, tracking::record_external_generation_run_after_success, }; use time::OffsetDateTime; @@ -60,6 +61,7 @@ pub(crate) struct RawImageEditResponse { } const RAW_IMAGE_MAX_TEXT_FIELD_BYTES: usize = 16 * 1024; +const RAW_IMAGE_PREPARE_TIMEOUT: Duration = Duration::from_secs(30); pub(crate) async fn edit_raw_image( State(state): State, @@ -68,11 +70,47 @@ pub(crate) async fn edit_raw_image( multipart: Multipart, ) -> Result, AppError> { let payload = parse_multipart_request(multipart).await?; - let prepared = tokio::task::spawn_blocking(move || prepare_request(payload)) - .await - .map_err(|error| { - AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR).with_message(error.to_string()) - })??; + let local_deadline = Instant::now() + .checked_add(RAW_IMAGE_PREPARE_TIMEOUT) + .unwrap_or_else(Instant::now); + let processing_deadline = request_context + .external_call_deadline() + .map(|deadline| deadline.min(local_deadline)) + .unwrap_or(local_deadline); + let permit = match tokio::time::timeout_at( + tokio::time::Instant::from_std(processing_deadline), + state.raw_image_decode_limiter().acquire_owned(), + ) + .await + { + Ok(Ok(permit)) => permit, + Ok(Err(error)) => { + return Err( + AppError::from_status(StatusCode::SERVICE_UNAVAILABLE).with_details(json!({ + "provider": "raw-image-edit", + "code": "RAW_IMAGE_DECODE_LIMITER_UNAVAILABLE", + "message": format!("raw 图片解码并发控制器不可用:{error}"), + })), + ); + } + Err(_) => return Err(raw_image_prepare_timeout_error("等待 raw 图片解码槽位超时")), + }; + let worker = tokio::task::spawn_blocking(move || { + // 超时只能停止 async 等待,permit 必须由 blocking closure 持有到解码真正结束。 + let _permit = permit; + prepare_request(payload) + }); + let prepared = + match tokio::time::timeout_at(tokio::time::Instant::from_std(processing_deadline), worker) + .await + { + Ok(Ok(result)) => result?, + Ok(Err(error)) => { + return Err(AppError::from_status(StatusCode::INTERNAL_SERVER_ERROR) + .with_message(error.to_string())); + } + Err(_) => return Err(raw_image_prepare_timeout_error("raw 图片解码处理超时")), + }; let settings = require_openai_image_settings(&state)?.with_external_api_audit_context( &request_context, Some(authenticated.claims().user_id().to_string()), @@ -472,6 +510,15 @@ fn bad_request(message: impl Into) -> AppError { })) } +fn raw_image_prepare_timeout_error(message: &str) -> AppError { + AppError::from_status(StatusCode::GATEWAY_TIMEOUT).with_details(json!({ + "provider": "raw-image-edit", + "code": "RAW_IMAGE_PREPARE_TIMEOUT", + "message": message, + "maxConcurrency": RAW_IMAGE_DECODE_MAX_CONCURRENCY, + })) +} + #[cfg(test)] mod tests { use super::*; diff --git a/server-rs/crates/api-server/src/state.rs b/server-rs/crates/api-server/src/state.rs index dc53af8cc..dc7b08503 100644 --- a/server-rs/crates/api-server/src/state.rs +++ b/server-rs/crates/api-server/src/state.rs @@ -56,6 +56,8 @@ const ADMIN_ROLE: &str = "admin"; const EDITOR_AGENT_LLM_MAX_RETRIES: u32 = 1; const EDITOR_AGENT_LLM_MAX_RETRY_BACKOFF_MS: u64 = 60_000; pub(crate) const CHARACTER_ANIMATION_OSS_MAX_CONCURRENCY: usize = 8; +// Raw image PNG 解码会进入 Tokio blocking pool;单独限流,避免图片请求挤占其它 blocking 工作。 +pub(crate) const RAW_IMAGE_DECODE_MAX_CONCURRENCY: usize = 4; // P=8:父侧成功图片读取/解码槽。配 N=16 是内存与出口吞吐的折中, // 极端完整 body 内存按 (N + P) × 32 MiB 评估(见调度方案 §9.2)。 pub(crate) const BGFILTER_IMAGE_VALIDATION_MAX_CONCURRENCY: usize = 8; @@ -304,6 +306,7 @@ pub struct AppStateInner { matting_client: Option, bgfilter_provider_http_client: reqwest::Client, bgfilter_worker_http_client: reqwest::Client, + raw_image_decode_limiter: Arc, bgfilter_image_validation_limiter: Arc, character_animation_oss_http_client: reqwest::Client, character_animation_oss_io_limiter: Arc, @@ -612,6 +615,9 @@ impl AppState { let bgfilter_image_validation_limiter = Arc::new(Semaphore::new( bgfilter_image_validation_concurrency.min(Semaphore::MAX_PERMITS), )); + let raw_image_decode_limiter = Arc::new(Semaphore::new( + RAW_IMAGE_DECODE_MAX_CONCURRENCY.min(Semaphore::MAX_PERMITS), + )); 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)); @@ -674,6 +680,7 @@ impl AppState { matting_client, bgfilter_provider_http_client, bgfilter_worker_http_client, + raw_image_decode_limiter, bgfilter_image_validation_limiter, character_animation_oss_http_client, character_animation_oss_io_limiter, @@ -1584,6 +1591,10 @@ impl AppState { &self.bgfilter_worker_http_client } + pub fn raw_image_decode_limiter(&self) -> Arc { + self.raw_image_decode_limiter.clone() + } + pub fn bgfilter_worker_reached(&self) -> bool { self.bgfilter_worker_reached.load(Ordering::Relaxed) }