From 69f850ff827fe64ec1c7f7a325e34ecf9c338252 Mon Sep 17 00:00:00 2001 From: Linghong Date: Fri, 7 Aug 2026 08:27:31 +0000 Subject: [PATCH] =?UTF-8?q?=E9=9F=B3=E6=95=88=E5=A4=96=E9=83=A8=E8=B0=83?= =?UTF-8?q?=E7=94=A8=E7=BA=B3=E5=85=A5=20worker=20provider=20=E9=A2=84?= =?UTF-8?q?=E7=AE=97?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ElevenLabs 设置新增 request_deadline,单次请求超时按剩余预算收窄,预算耗尽时不再发出请求。 音效翻译按同一预算逐轮判定,超期直接失败而不是继续占用终态写回时间。 音效生成把 request_context 上的 provider 截止时刻分别传给翻译与 ElevenLabs 调用。 新增预算耗尽时 provider 与 LLM 请求数均为 0 的回归测试及各自的对照用例。 --- .../generation.rs | 29 +++-- .../settings.rs | 2 + .../sound_effect_translation.rs | 100 ++++++++++++++-- .../crates/platform-audio/src/elevenlabs.rs | 110 +++++++++++++++++- .../crates/platform-audio/tests/elevenlabs.rs | 76 ++++++++++++ 5 files changed, 294 insertions(+), 23 deletions(-) diff --git a/server-rs/crates/api-server/src/vector_engine_audio_generation/generation.rs b/server-rs/crates/api-server/src/vector_engine_audio_generation/generation.rs index 51a33d48b..5ca2e28e0 100644 --- a/server-rs/crates/api-server/src/vector_engine_audio_generation/generation.rs +++ b/server-rs/crates/api-server/src/vector_engine_audio_generation/generation.rs @@ -294,8 +294,12 @@ pub(crate) async fn generate_editor_sound_effect_for_owner( .map_err(|error| error.into_response_with_context(Some(&request_context)))?; let normalized = normalize_editor_sound_effect_request_with_pricing(payload.clone(), &pricing) .map_err(|error| error.into_response_with_context(Some(&request_context)))?; - let settings = require_elevenlabs_audio_settings(&state) + // Worker 已经把「留出终态写回时间」后的 provider 截止时刻挂在 request_context 上, + // 翻译和 ElevenLabs 两段外部调用都必须落在这个预算内。 + let request_deadline = request_context.external_call_deadline(); + let mut settings = require_elevenlabs_audio_settings(&state) .map_err(|error| error.into_response_with_context(Some(&request_context)))?; + settings.request_deadline = request_deadline; let http_client = platform_audio::build_elevenlabs_audio_http_client(&settings) .map_err(map_platform_audio_error) .map_err(|error| error.into_response_with_context(Some(&request_context)))?; @@ -328,16 +332,19 @@ pub(crate) async fn generate_editor_sound_effect_for_owner( billing_asset_id.as_str(), u64::from(normalized.price_mud_points), async { - let actual_prompt = - translate_sound_effect_prompt_for_worker(llm_client, normalized.prompt.as_str()) - .await - .map_err(|error| { - AppError::from_status(StatusCode::BAD_GATEWAY).with_details(json!({ - "provider": "editor-sound-effect-translation", - "reason": error.reason_code(), - "message": error.to_string(), - })) - })?; + let actual_prompt = translate_sound_effect_prompt_for_worker( + llm_client, + normalized.prompt.as_str(), + request_deadline, + ) + .await + .map_err(|error| { + AppError::from_status(StatusCode::BAD_GATEWAY).with_details(json!({ + "provider": "editor-sound-effect-translation", + "reason": error.reason_code(), + "message": error.to_string(), + })) + })?; let generated = platform_audio::generate_elevenlabs_sound_effect( &http_client, &settings, diff --git a/server-rs/crates/api-server/src/vector_engine_audio_generation/settings.rs b/server-rs/crates/api-server/src/vector_engine_audio_generation/settings.rs index 00e828dc5..7923a9a9c 100644 --- a/server-rs/crates/api-server/src/vector_engine_audio_generation/settings.rs +++ b/server-rs/crates/api-server/src/vector_engine_audio_generation/settings.rs @@ -77,6 +77,8 @@ pub(super) fn require_elevenlabs_audio_settings( base_url: base_url.to_string(), api_key: api_key.to_string(), request_timeout_ms: state.config.elevenlabs_request_timeout_ms.max(1), + // 由调用方按本次请求的 provider 预算覆盖;配置本身不携带截止时刻。 + request_deadline: None, }) } diff --git a/server-rs/crates/api-server/src/vector_engine_audio_generation/sound_effect_translation.rs b/server-rs/crates/api-server/src/vector_engine_audio_generation/sound_effect_translation.rs index 8a849c674..89edba88b 100644 --- a/server-rs/crates/api-server/src/vector_engine_audio_generation/sound_effect_translation.rs +++ b/server-rs/crates/api-server/src/vector_engine_audio_generation/sound_effect_translation.rs @@ -1,4 +1,4 @@ -use std::{error::Error, fmt}; +use std::{error::Error, fmt, time::Instant}; use platform_llm::{ LlmClient, LlmMessage, LlmResponseReasoningEffort, LlmRunRequest, LlmRunResponse, @@ -25,6 +25,7 @@ struct SoundEffectPromptTranslationEnvelope { #[derive(Debug)] pub(super) enum SoundEffectPromptTranslationError { InvalidResponse, + BudgetExhausted, Upstream(platform_llm::LlmError), } @@ -32,6 +33,7 @@ impl SoundEffectPromptTranslationError { pub(super) fn reason_code(&self) -> &'static str { match self { Self::InvalidResponse => "translation_invalid", + Self::BudgetExhausted => "translation_budget_exhausted", Self::Upstream(_) => "translation_upstream_failed", } } @@ -41,6 +43,7 @@ impl fmt::Display for SoundEffectPromptTranslationError { fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { formatter.write_str(match self { Self::InvalidResponse => "音效 Prompt 英文化结果未通过校验", + Self::BudgetExhausted => "音效 Prompt 英文化超出本次任务的调用预算", Self::Upstream(_) => "音效 Prompt 英文化服务暂时不可用", }) } @@ -49,26 +52,29 @@ impl fmt::Display for SoundEffectPromptTranslationError { impl Error for SoundEffectPromptTranslationError { fn source(&self) -> Option<&(dyn Error + 'static)> { match self { - Self::InvalidResponse => None, + Self::InvalidResponse | Self::BudgetExhausted => None, Self::Upstream(error) => Some(error), } } } /// 仅供正式音效生成流水线在 Worker 执行阶段调用;不注册同步翻译 HTTP 路由。 +/// +/// `request_deadline` 是 worker 为本次任务留出的 provider 调用截止时刻。翻译最多两个业务 +/// 语义轮,每轮都可能吃满 LLM 固定超时;不按剩余预算截断,慢翻译会把预算耗尽后仍然继续, +/// 后面的 ElevenLabs 调用就落在预算之外了。 pub(super) async fn translate_sound_effect_prompt_for_worker( llm_client: &LlmClient, user_prompt: &str, + request_deadline: Option, ) -> Result { let user_prompt = platform_audio::validate_sound_effect_prompt(user_prompt) .map_err(|_| SoundEffectPromptTranslationError::InvalidResponse)? .prompt .to_string(); - let first_response = llm_client - .run(build_sound_effect_translation_llm_request(&user_prompt)) - .await - .map_err(SoundEffectPromptTranslationError::Upstream)?; + let first_response = + run_sound_effect_translation_round(llm_client, &user_prompt, request_deadline).await?; if has_finish_reason(&first_response, "content_filter") { return Err(SoundEffectPromptTranslationError::InvalidResponse); } @@ -77,14 +83,38 @@ pub(super) async fn translate_sound_effect_prompt_for_worker( } // 第二业务语义轮始终重新使用冻结的原始 userPrompt;首轮候选不能成为事实源。 - let second_response = llm_client - .run(build_sound_effect_translation_llm_request(&user_prompt)) - .await - .map_err(SoundEffectPromptTranslationError::Upstream)?; + let second_response = + run_sound_effect_translation_round(llm_client, &user_prompt, request_deadline).await?; inspect_sound_effect_translation_candidate(&second_response) .ok_or(SoundEffectPromptTranslationError::InvalidResponse) } +/// 每轮开始前先判定预算:`timeout_at` 会先 poll 一次 future,已过期时请求其实已经发出去了, +/// 所以必须显式前置检查,不能只靠超时包裹。 +async fn run_sound_effect_translation_round( + llm_client: &LlmClient, + user_prompt: &str, + request_deadline: Option, +) -> Result { + let request = build_sound_effect_translation_llm_request(user_prompt); + let Some(request_deadline) = request_deadline else { + return llm_client + .run(request) + .await + .map_err(SoundEffectPromptTranslationError::Upstream); + }; + if request_deadline <= Instant::now() { + return Err(SoundEffectPromptTranslationError::BudgetExhausted); + } + tokio::time::timeout_at( + tokio::time::Instant::from_std(request_deadline), + llm_client.run(request), + ) + .await + .map_err(|_| SoundEffectPromptTranslationError::BudgetExhausted)? + .map_err(SoundEffectPromptTranslationError::Upstream) +} + fn build_sound_effect_translation_llm_request(user_prompt: &str) -> LlmRunRequest { LlmRunRequest::new(vec![ LlmMessage::system(sound_effect_translation_system_prompt()), @@ -325,7 +355,7 @@ mod tests { for (input, expected) in cases { assert_eq!( - translate_sound_effect_prompt_for_worker(client, input) + translate_sound_effect_prompt_for_worker(client, input, None) .await .expect("translation should pass"), expected @@ -373,6 +403,7 @@ mod tests { .editor_agent_llm_client() .expect("client should exist"), "\u{feff}金币拾取\u{2003}", + None, ) .await .expect("second business round should pass"); @@ -406,6 +437,7 @@ mod tests { .editor_agent_llm_client() .expect("client should exist"), "金币拾取", + None, ) .await .expect_err("second length should fail"); @@ -429,6 +461,7 @@ mod tests { .editor_agent_llm_client() .expect("client should exist"), "金币拾取", + None, ) .await .expect_err("content filter should fail"); @@ -436,6 +469,49 @@ mod tests { assert_eq!(mock.finish().len(), 1); } + #[tokio::test] + async fn an_exhausted_provider_budget_stops_before_the_first_llm_request() { + let mock = spawn_mock_llm_server(Vec::new()); + let state = editor_llm_test_state(mock.base_url.clone(), 0); + let error = translate_sound_effect_prompt_for_worker( + state + .editor_agent_llm_client() + .expect("client should exist"), + "金币拾取", + Some(Instant::now() - std::time::Duration::from_millis(1)), + ) + .await + .expect_err("an exhausted provider budget should fail before calling the LLM"); + + // 预算耗尽时翻译一次都不能发出去,否则后面的 ElevenLabs 调用必然落在预算之外。 + assert_eq!(error.reason_code(), "translation_budget_exhausted"); + assert_eq!(mock.finish().len(), 0); + } + + #[tokio::test] + async fn a_remaining_provider_budget_still_runs_the_business_round() { + let mock = spawn_mock_llm_server(vec![MockLlmResponse { + status_line: "200 OK", + body: chat_text_response( + valid_translation_candidate("Bright metallic coin pickup."), + Some("stop"), + ), + }]); + let state = editor_llm_test_state(mock.base_url.clone(), 0); + let actual = translate_sound_effect_prompt_for_worker( + state + .editor_agent_llm_client() + .expect("client should exist"), + "金币拾取", + Some(Instant::now() + std::time::Duration::from_secs(30)), + ) + .await + .expect("a remaining provider budget should still translate"); + + assert_eq!(actual, "Bright metallic coin pickup."); + assert_eq!(mock.finish().len(), 1); + } + #[tokio::test] async fn transport_failure_does_not_start_a_content_retry_round() { let mock = spawn_mock_llm_server(vec![MockLlmResponse { @@ -448,6 +524,7 @@ mod tests { .editor_agent_llm_client() .expect("client should exist"), "金币拾取", + None, ) .await .expect_err("transport failure should fail"); @@ -477,6 +554,7 @@ mod tests { .editor_agent_llm_client() .expect("client should exist"), "金币拾取", + None, ) .await .expect("transport retry should recover inside the first business round"); diff --git a/server-rs/crates/platform-audio/src/elevenlabs.rs b/server-rs/crates/platform-audio/src/elevenlabs.rs index 0904820d5..e3fc66785 100644 --- a/server-rs/crates/platform-audio/src/elevenlabs.rs +++ b/server-rs/crates/platform-audio/src/elevenlabs.rs @@ -1,4 +1,4 @@ -use std::error::Error; +use std::{error::Error, time::Instant}; use bytes::BytesMut; use reqwest::header; @@ -20,6 +20,9 @@ pub struct ElevenLabsAudioSettings { pub base_url: String, pub api_key: String, pub request_timeout_ms: u64, + /// Worker 为本次任务留出的 provider 调用截止时刻。超出后不得再发起 provider 请求, + /// 否则终态写回预留时间会被 provider 调用吃掉。为 `None` 时只受固定超时约束。 + pub request_deadline: Option, } impl std::fmt::Debug for ElevenLabsAudioSettings { @@ -29,10 +32,51 @@ impl std::fmt::Debug for ElevenLabsAudioSettings { .field("base_url", &self.base_url) .field("api_key", &"[redacted]") .field("request_timeout_ms", &self.request_timeout_ms) + .field("request_deadline", &self.request_deadline) .finish() } } +/// 单次 provider 请求的实际超时:固定超时与剩余预算取更严格的一个。 +/// 返回 `None` 表示预算已经耗尽,调用方必须在发出请求之前失败。 +fn effective_request_timeout_ms( + configured_timeout_ms: u64, + request_deadline: Option, +) -> Option { + 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, + now: Instant, +) -> Option { + let configured_timeout_ms = configured_timeout_ms.max(1); + let Some(request_deadline) = request_deadline else { + return Some(configured_timeout_ms); + }; + let remaining_ms = u64::try_from(request_deadline.checked_duration_since(now)?.as_millis()) + .unwrap_or(u64::MAX); + if remaining_ms == 0 { + return None; + } + Some(configured_timeout_ms.min(remaining_ms)) +} + +fn elevenlabs_budget_exhausted_error(endpoint: &str) -> AudioError { + AudioError::request_for( + ELEVENLABS_PROVIDER, + "ElevenLabs 音效生成超出本次任务的 provider 调用预算".to_string(), + Some(endpoint.to_string()), + true, + false, + true, + false, + None, + None, + ) +} + #[derive(Clone, Debug, PartialEq)] pub struct ElevenLabsSoundEffectRequest { pub text: String, @@ -99,8 +143,16 @@ pub async fn generate_elevenlabs_sound_effect( ) -> Result { let endpoint = elevenlabs_sound_generation_endpoint(&settings.base_url); let body = build_elevenlabs_sound_effect_body(&request)?; + // 预算耗尽时必须在发出请求之前失败:翻译轮次拖长后再发 provider 请求,会把 worker + // 留给终态写回的时间吃掉,且这次调用一定来不及被本 attempt 使用。 + let Some(attempt_timeout_ms) = + effective_request_timeout_ms(settings.request_timeout_ms, settings.request_deadline) + else { + return Err(elevenlabs_budget_exhausted_error(endpoint.as_str())); + }; let response = http_client .post(endpoint.as_str()) + .timeout(std::time::Duration::from_millis(attempt_timeout_ms)) .query(&[("output_format", ELEVENLABS_SOUND_EFFECT_OUTPUT_FORMAT)]) .header("xi-api-key", settings.api_key.as_str()) .header( @@ -269,6 +321,7 @@ mod tests { base_url: "https://api.elevenlabs.test".to_string(), api_key: "elevenlabs-secret-test-key".to_string(), request_timeout_ms: 1_000, + request_deadline: None, }; let debug = format!("{settings:?}"); @@ -276,6 +329,61 @@ mod tests { assert!(!debug.contains(settings.api_key.as_str())); } + #[test] + fn configured_timeout_is_preserved_without_a_request_deadline() { + let now = Instant::now(); + + assert_eq!( + effective_request_timeout_ms_at(180_000, None, now), + Some(180_000) + ); + // 0 会被 reqwest 当成非法超时,最低收敛到 1ms。 + assert_eq!(effective_request_timeout_ms_at(0, None, now), Some(1)); + } + + #[test] + fn request_timeout_is_clipped_to_the_remaining_provider_budget() { + let now = Instant::now(); + let deadline = now + std::time::Duration::from_millis(2_000); + + assert_eq!( + effective_request_timeout_ms_at(180_000, Some(deadline), now), + Some(2_000) + ); + // 预算比固定超时宽松时不放大固定超时。 + assert_eq!( + effective_request_timeout_ms_at( + 500, + Some(now + std::time::Duration::from_secs(600)), + now + ), + Some(500) + ); + } + + #[test] + fn exhausted_budget_stops_the_provider_request_before_it_is_sent() { + let now = Instant::now(); + + assert_eq!( + effective_request_timeout_ms_at( + 180_000, + Some(now - std::time::Duration::from_millis(1)), + now + ), + None + ); + assert_eq!( + effective_request_timeout_ms_at(180_000, Some(now), now), + None + ); + + let error = + elevenlabs_budget_exhausted_error("https://api.elevenlabs.test/v1/sound-generation"); + assert_eq!(error.provider(), ELEVENLABS_PROVIDER); + assert!(matches!(error, AudioError::Request { timeout: true, .. })); + } + #[test] fn request_validation_errors_are_attributed_to_elevenlabs() { let error = build_elevenlabs_sound_effect_body(&ElevenLabsSoundEffectRequest { diff --git a/server-rs/crates/platform-audio/tests/elevenlabs.rs b/server-rs/crates/platform-audio/tests/elevenlabs.rs index eaea9a8c2..c9df08581 100644 --- a/server-rs/crates/platform-audio/tests/elevenlabs.rs +++ b/server-rs/crates/platform-audio/tests/elevenlabs.rs @@ -31,6 +31,7 @@ fn settings(base_url: String, request_timeout_ms: u64) -> ElevenLabsAudioSetting base_url, api_key: TEST_API_KEY.to_string(), request_timeout_ms, + request_deadline: None, } } @@ -142,6 +143,35 @@ fn spawn_counting_response_server( (format!("http://{address}"), request_count, server) } +/// 用于断言「provider 一次都没被调用」:只做非阻塞轮询,不会在零请求时把 join 卡死。 +fn spawn_never_answering_server() -> (String, Arc, thread::JoinHandle<()>) { + let listener = TcpListener::bind("127.0.0.1:0").expect("mock server should bind"); + let address = listener + .local_addr() + .expect("mock address should be readable"); + let request_count = Arc::new(AtomicUsize::new(0)); + let request_count_for_server = Arc::clone(&request_count); + let server = thread::spawn(move || { + listener + .set_nonblocking(true) + .expect("mock listener should become nonblocking"); + let deadline = Instant::now() + Duration::from_millis(200); + while Instant::now() < deadline { + match listener.accept() { + Ok((mut stream, _)) => { + request_count_for_server.fetch_add(1, Ordering::SeqCst); + let _ = read_http_request(&mut stream); + } + Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => { + thread::sleep(Duration::from_millis(5)); + } + Err(error) => panic!("mock listener failed: {error}"), + } + } + }); + (format!("http://{address}"), request_count, server) +} + fn success_response(content_type: Option<&str>, body: &[u8]) -> Vec { let content_type_header = content_type .map(|value| format!("Content-Type: {value}\r\n")) @@ -299,6 +329,52 @@ fn timeout_does_not_retry_the_provider_post() { assert_eq!(request_count.load(Ordering::SeqCst), 1); } +#[test] +fn an_exhausted_provider_budget_fails_before_any_provider_post() { + let (base_url, request_count, server) = spawn_never_answering_server(); + let mut settings = settings(base_url, 180_000); + settings.request_deadline = Some(Instant::now() - Duration::from_millis(1)); + let client = + build_elevenlabs_audio_http_client(&settings).expect("ElevenLabs HTTP client should build"); + + let error = runtime() + .block_on(generate_elevenlabs_sound_effect( + &client, + &settings, + request(), + )) + .expect_err("an exhausted provider budget should fail before sending"); + server.join().expect("mock server should finish"); + + // 预算耗尽时 provider 一次都不能被调用,否则慢翻译会白白消耗一次计费请求。 + assert_eq!(request_count.load(Ordering::SeqCst), 0); + assert_eq!(error.provider(), ELEVENLABS_PROVIDER); + assert!(matches!(error, AudioError::Request { timeout: true, .. })); +} + +#[test] +fn a_remaining_provider_budget_still_sends_exactly_one_provider_post() { + let response = success_response(Some("audio/mpeg"), TEST_MP3); + let (base_url, request_count, server) = + spawn_counting_response_server(response, Duration::ZERO); + let mut settings = settings(base_url, 180_000); + settings.request_deadline = Some(Instant::now() + Duration::from_secs(30)); + let client = + build_elevenlabs_audio_http_client(&settings).expect("ElevenLabs HTTP client should build"); + + let generated = runtime() + .block_on(generate_elevenlabs_sound_effect( + &client, + &settings, + request(), + )) + .expect("a remaining provider budget should still generate"); + server.join().expect("mock server should finish"); + + assert_eq!(request_count.load(Ordering::SeqCst), 1); + assert!(generated.duration_seconds > 0.0); +} + #[test] fn content_length_precheck_accepts_the_limit_and_rejects_limit_plus_one() { for (content_length, should_be_size_error) in [