动画生成过程中的OSS上传行为优化 (#93)

Reviewed-on: https://git.genarrative.world/git/GenarrativeAI/Genarrative/pulls/93
Reviewed-by: 段舒康 <kdletters@qq.com>
Co-authored-by: Linghong <ink29535@proton.me>
Co-committed-by: Linghong <ink29535@proton.me>
This commit was merged in pull request #93.
This commit is contained in:
2026-07-20 19:50:37 +08:00
committed by 段舒康
parent 8965a53fb8
commit a4b3898ff4
8 changed files with 1317 additions and 27 deletions
@@ -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};
@@ -2417,8 +2417,10 @@ async fn process_and_persist_editor_character_animation_frame(
audit: &crate::external_api_audit::ExternalApiAuditContext,
) -> Result<ProcessedEditorCharacterAnimationFrame, AppError> {
// 中文注释:每一帧只要求自己的绿幕源图先落 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(),
@@ -2466,8 +2468,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(),
@@ -2494,6 +2498,7 @@ async fn process_and_persist_editor_character_animation_frame(
task_id,
put_result.object_key.clone(),
content_type,
frame_index + 1,
)
.await?;
@@ -2751,6 +2756,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<String>,
file_name: String,
content_type: String,
body: Vec<u8>,
metadata: BTreeMap<String, String>,
) -> Result<platform_oss::OssPutObjectResponse, AppError> {
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,
@@ -2883,13 +2921,27 @@ async fn confirm_editor_character_animation_frame_asset_object(
task_id: &str,
object_key: String,
content_type: String,
frame_index: usize,
) -> Result<module_assets::ConfirmAssetObjectResult, AppError> {
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,
)
@@ -2910,6 +2962,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<module_assets::ConfirmAssetObjectResult, AppError> {
let now_micros = current_utc_micros();
let record = state
.spacetime_client()
@@ -6214,13 +6287,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",
],
);
+47
View File
@@ -54,6 +54,7 @@ use crate::work_author::{
};
const ADMIN_ROLE: &str = "admin";
pub(crate) const CHARACTER_ANIMATION_OSS_MAX_CONCURRENCY: usize = 8;
pub type HttpRequestPermitPool = Semaphore;
@@ -277,6 +278,8 @@ pub struct AppStateInner {
creative_agent_gpt5_client: Option<LlmClient>,
matting_client: Option<MattingClient>,
editor_bgfilter_http_client: reqwest::Client,
character_animation_oss_http_client: reqwest::Client,
character_animation_oss_io_limiter: Arc<Semaphore>,
creative_agent_executor: Arc<MockLangChainRustAgentExecutor>,
// Phase 1 任务 E 的 creative session facade 暂存在 api-server。
// creative_agent_* 表由任务 D 收口后,这里只保留读写 facade。
@@ -517,6 +520,9 @@ impl AppState {
let creative_agent_gpt5_client = build_creative_agent_gpt5_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);
@@ -558,6 +564,8 @@ impl AppState {
creative_agent_gpt5_client,
matting_client,
editor_bgfilter_http_client,
character_animation_oss_http_client,
character_animation_oss_io_limiter,
creative_agent_executor: Arc::new(MockLangChainRustAgentExecutor),
creative_agent_sessions: Arc::new(Mutex::new(HashMap::new())),
profile_recharge_order_updates,
@@ -1262,6 +1270,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<Semaphore> {
self.character_animation_oss_io_limiter.clone()
}
pub fn creative_agent_executor(&self) -> Arc<MockLangChainRustAgentExecutor> {
self.creative_agent_executor.clone()
}
@@ -1980,6 +1996,21 @@ fn build_editor_bgfilter_http_client(
})
}
fn build_character_animation_oss_http_client() -> Result<reqwest::Client, AppStateInitError> {
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(),
@@ -2112,6 +2143,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(