Merge branch 'master' into editor-agent-more-tools
# Conflicts: # docs/project-memory/shared-memory/decision-log.md
This commit is contained in:
File diff suppressed because it is too large
Load Diff
@@ -74,7 +74,7 @@ use crate::{
|
||||
apply_editor_screen_background_decision_to_generation_inputs,
|
||||
build_editor_canvas_generated_layer_item, complete_editor_canvas_generation_with_items,
|
||||
persist_editor_generated_media_asset,
|
||||
remove_editor_generated_screen_background_with_bgfilter_with_request_timeout,
|
||||
remove_editor_generated_screen_background_with_bgfilter,
|
||||
resolve_editor_reference_object_key_for_owner,
|
||||
},
|
||||
editor_screen_background_decision::{
|
||||
@@ -118,7 +118,6 @@ const EDITOR_CHARACTER_ANIMATION_MODEL: &str = "seedance2.0-fast";
|
||||
const EDITOR_CHARACTER_ANIMATION_ASSET_KIND: &str = "editor_character_animation";
|
||||
const EDITOR_CHARACTER_ANIMATION_RESOURCE_ASSET_KIND: &str = "character-animation";
|
||||
const EDITOR_CHARACTER_ANIMATION_PROVIDER_SOURCE_SLOT: &str = "provider_source";
|
||||
const EDITOR_CHARACTER_ANIMATION_BGFILTER_TIMEOUT_PER_FRAME_MS: u64 = 2_000;
|
||||
const EDITOR_VIDEO_ASSET_KIND: &str = "editor_video";
|
||||
const EDITOR_VIDEO_ENTITY_KIND: &str = "editor_canvas";
|
||||
const EDITOR_VIDEO_SLOT: &str = "video_preview";
|
||||
@@ -706,6 +705,7 @@ pub(crate) async fn generate_editor_character_animation_for_owner(
|
||||
user_id: Some(owner_user_id.clone()),
|
||||
profile_id: project_id.clone(),
|
||||
request_id: Some(request_context.request_id().to_string()),
|
||||
external_call_deadline: request_context.external_call_deadline(),
|
||||
};
|
||||
|
||||
let result = execute_billable_asset_operation_with_cost(
|
||||
@@ -734,6 +734,7 @@ pub(crate) async fn generate_editor_character_animation_for_owner(
|
||||
user_id: Some(owner_user_id.clone()),
|
||||
profile_id: project_id.clone(),
|
||||
request_id: Some(request_context.request_id().to_string()),
|
||||
external_call_deadline: request_context.external_call_deadline(),
|
||||
},
|
||||
}),
|
||||
)
|
||||
@@ -2295,18 +2296,6 @@ async fn create_editor_ark_image_to_video_task(
|
||||
})
|
||||
}
|
||||
|
||||
fn editor_character_animation_bgfilter_request_timeout_ms(
|
||||
base_timeout_ms: u64,
|
||||
frame_count: usize,
|
||||
) -> u64 {
|
||||
let frame_timeout_increment_ms = u64::try_from(frame_count)
|
||||
.unwrap_or(u64::MAX)
|
||||
.saturating_mul(EDITOR_CHARACTER_ANIMATION_BGFILTER_TIMEOUT_PER_FRAME_MS);
|
||||
base_timeout_ms
|
||||
.saturating_add(frame_timeout_increment_ms)
|
||||
.max(1)
|
||||
}
|
||||
|
||||
async fn extract_and_persist_editor_character_animation_frames(
|
||||
state: &AppState,
|
||||
owner_user_id: &str,
|
||||
@@ -2337,10 +2326,6 @@ async fn extract_and_persist_editor_character_animation_frames(
|
||||
use futures_util::StreamExt as _;
|
||||
|
||||
let frame_count = finalized_frames.len();
|
||||
let bgfilter_request_timeout_ms = editor_character_animation_bgfilter_request_timeout_ms(
|
||||
state.config.editor_bgfilter_request_timeout_ms,
|
||||
frame_count,
|
||||
);
|
||||
let frame_results = futures_util::stream::iter(finalized_frames.into_iter().enumerate().map(
|
||||
|(frame_index, frame)| async move {
|
||||
process_and_persist_editor_character_animation_frame(
|
||||
@@ -2353,15 +2338,14 @@ async fn extract_and_persist_editor_character_animation_frames(
|
||||
request.frame_width,
|
||||
request.frame_height,
|
||||
request.screen_color,
|
||||
bgfilter_request_timeout_ms,
|
||||
audit,
|
||||
)
|
||||
.await
|
||||
.map_err(|error| (frame_index, error))
|
||||
},
|
||||
))
|
||||
// 中文注释:动作帧上限固定为 48。全部帧连续进入在途集合,由 BgFilter 服务端既有进程锁自行排队;
|
||||
// api-server 不再等待前一帧返回,也不因返回乱序产生队头阻塞。
|
||||
// 中文注释:动作帧上限固定为 48。全部帧连续进入在途集合,由唯一内部
|
||||
// bgfilter-worker 的 admission 与 provider semaphore 统一排队和限流;父流程不因返回乱序产生队头阻塞。
|
||||
.buffer_unordered(frame_count.max(1))
|
||||
// 中文注释:不能 try_collect 提前取消。已经发出的请求必须全部排空,避免服务端完成计算后无人接收。
|
||||
.collect::<Vec<_>>()
|
||||
@@ -2417,7 +2401,6 @@ async fn process_and_persist_editor_character_animation_frame(
|
||||
frame_width: u32,
|
||||
frame_height: u32,
|
||||
screen_color: EditorScreenBackgroundColor,
|
||||
bgfilter_request_timeout_ms: u64,
|
||||
audit: &crate::external_api_audit::ExternalApiAuditContext,
|
||||
) -> Result<ProcessedEditorCharacterAnimationFrame, AppError> {
|
||||
// 中文注释:每一帧只要求自己的绿幕源图先落 OSS,不再等待整批源图全部上传完成。
|
||||
@@ -2449,13 +2432,14 @@ async fn process_and_persist_editor_character_animation_frame(
|
||||
)
|
||||
.await?;
|
||||
|
||||
let removed = remove_editor_generated_screen_background_with_bgfilter_with_request_timeout(
|
||||
// 排队与调用预算由共享 helper 按 N/est 公式和父剩余预算派生,动画帧不再携带
|
||||
// 任何按帧数放大的 timeout。
|
||||
let removed = remove_editor_generated_screen_background_with_bgfilter(
|
||||
state,
|
||||
source_put.object_key.as_str(),
|
||||
screen_color,
|
||||
EDITOR_BGFILTER_DEFAULT_SEG_MODEL,
|
||||
EDITOR_BGFILTER_CROSS_CHECK_ENABLED,
|
||||
bgfilter_request_timeout_ms,
|
||||
audit,
|
||||
)
|
||||
.await?;
|
||||
@@ -6431,31 +6415,24 @@ mod tests {
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn editor_character_animation_bgfilter_timeout_scales_with_frame_count() {
|
||||
assert_eq!(
|
||||
editor_character_animation_bgfilter_request_timeout_ms(180_000, 32),
|
||||
244_000
|
||||
);
|
||||
assert_eq!(
|
||||
editor_character_animation_bgfilter_request_timeout_ms(180_000, 40),
|
||||
260_000
|
||||
);
|
||||
assert_eq!(
|
||||
editor_character_animation_bgfilter_request_timeout_ms(180_000, 48),
|
||||
276_000
|
||||
);
|
||||
assert_eq!(
|
||||
editor_character_animation_bgfilter_request_timeout_ms(0, 32),
|
||||
64_000
|
||||
);
|
||||
assert_eq!(
|
||||
editor_character_animation_bgfilter_request_timeout_ms(0, 0),
|
||||
1
|
||||
);
|
||||
assert_eq!(
|
||||
editor_character_animation_bgfilter_request_timeout_ms(u64::MAX - 1, 48),
|
||||
u64::MAX
|
||||
fn editor_character_animation_bgfilter_frames_share_budget_helper_without_per_frame_timeout() {
|
||||
let source = include_str!("character_animation_assets.rs");
|
||||
// 动画帧不携带任何自定 timeout:预算全部由共享 helper 按 N/est 公式派生。
|
||||
assert_function_contains(
|
||||
source,
|
||||
"async fn extract_and_persist_editor_character_animation_frames",
|
||||
"async fn process_and_persist_editor_character_animation_frame",
|
||||
&["process_and_persist_editor_character_animation_frame"],
|
||||
);
|
||||
assert!(!source.contains(concat!(
|
||||
"EDITOR_CHARACTER_ANIMATION_BGFILTER_TIMEOUT_",
|
||||
"PER_FRAME_MS"
|
||||
)));
|
||||
assert!(!source.contains(concat!(
|
||||
"editor_character_animation_bgfilter_",
|
||||
"request_timeout_ms"
|
||||
)));
|
||||
assert!(!source.contains(concat!("bgfilter_request_", "timeout_ms =")));
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -6571,10 +6548,7 @@ mod tests {
|
||||
&[
|
||||
"apply_chroma_key: false",
|
||||
"prepare_for_bgfilter_input: true",
|
||||
"editor_character_animation_bgfilter_request_timeout_ms",
|
||||
"state.config.editor_bgfilter_request_timeout_ms",
|
||||
"process_and_persist_editor_character_animation_frame",
|
||||
"bgfilter_request_timeout_ms",
|
||||
".buffer_unordered(frame_count.max(1))",
|
||||
".collect::<Vec<_>>()",
|
||||
"frame_payloads.sort_by_key",
|
||||
@@ -6604,9 +6578,8 @@ mod tests {
|
||||
&[
|
||||
"put_character_animation_frame_object",
|
||||
"green-screen-frame",
|
||||
"remove_editor_generated_screen_background_with_bgfilter_with_request_timeout",
|
||||
"remove_editor_generated_screen_background_with_bgfilter",
|
||||
"source_put.object_key.as_str()",
|
||||
"bgfilter_request_timeout_ms",
|
||||
"finalize_animation_frame_payload",
|
||||
"put_character_animation_frame_object",
|
||||
"animation_frame",
|
||||
@@ -6620,10 +6593,9 @@ mod tests {
|
||||
"EDITOR_CHARACTER_ANIMATION_ASSET_KIND",
|
||||
"EDITOR_CHARACTER_ANIMATION_PROVIDER_SOURCE_SLOT",
|
||||
"green-screen-frame",
|
||||
"remove_editor_generated_screen_background_with_bgfilter_with_request_timeout",
|
||||
"remove_editor_generated_screen_background_with_bgfilter",
|
||||
"EDITOR_BGFILTER_DEFAULT_SEG_MODEL",
|
||||
"EDITOR_BGFILTER_CROSS_CHECK_ENABLED",
|
||||
"bgfilter_request_timeout_ms",
|
||||
"finalize_animation_frame_payload",
|
||||
],
|
||||
);
|
||||
|
||||
@@ -18,9 +18,11 @@ const DEFAULT_EXTERNAL_GENERATION_WORKER_JOB_TIMEOUT_SECONDS: u64 = 900;
|
||||
const DEFAULT_EXTERNAL_GENERATION_WORKER_LONG_JOB_TIMEOUT_SECONDS: u64 = 1_800;
|
||||
pub(crate) const DEFAULT_VECTOR_ENGINE_IMAGE_REQUEST_TIMEOUT_MS: u64 = 1_000_000;
|
||||
const DEFAULT_EDITOR_BGFILTER_BASE_URL: &str = "http://58.87.105.82/bgfilter";
|
||||
const DEFAULT_EDITOR_BGFILTER_REQUEST_TIMEOUT_MS: u64 = 180_000;
|
||||
const DEFAULT_EDITOR_BGFILTER_SINGLE_IMAGE_ESTIMATE_MS: u64 = 5_000;
|
||||
const BGFILTER_ATTEMPT_SAFETY_FACTOR: u64 = 2;
|
||||
const BGFILTER_WORKER_RESPONSE_WINDOW_MS: u64 = 1_000;
|
||||
const DEFAULT_EDITOR_BGFILTER_CIRCUIT_FAILURE_THRESHOLD: u32 = 3;
|
||||
const DEFAULT_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS: u64 = 300;
|
||||
const DEFAULT_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS: u64 = 120;
|
||||
const DEFAULT_ALIYUN_MATTING_ENDPOINT: &str = "imageseg.cn-shanghai.aliyuncs.com";
|
||||
const DEFAULT_ALIYUN_MATTING_REQUEST_TIMEOUT_MS: u64 = 30_000;
|
||||
|
||||
@@ -32,6 +34,13 @@ pub struct AppConfig {
|
||||
pub listen_backlog: i32,
|
||||
pub worker_threads: Option<usize>,
|
||||
pub process_role: ProcessRole,
|
||||
pub bgfilter_worker_host: String,
|
||||
pub bgfilter_worker_port: u16,
|
||||
pub bgfilter_worker_base_url: String,
|
||||
pub bgfilter_internal_token: Option<String>,
|
||||
pub bgfilter_worker_concurrency: usize,
|
||||
pub bgfilter_worker_max_requests: usize,
|
||||
pub bgfilter_worker_connect_timeout_ms: u64,
|
||||
pub external_generation_mode: ExternalGenerationMode,
|
||||
pub external_generation_worker_id: String,
|
||||
pub external_generation_worker_concurrency: usize,
|
||||
@@ -62,7 +71,7 @@ pub struct AppConfig {
|
||||
pub editor_generation_pricing_override_path: PathBuf,
|
||||
pub editor_bgfilter_base_url: String,
|
||||
pub editor_bgfilter_token: Option<String>,
|
||||
pub editor_bgfilter_request_timeout_ms: u64,
|
||||
pub editor_bgfilter_single_image_estimate_ms: u64,
|
||||
pub aliyun_matting_enabled: bool,
|
||||
pub aliyun_matting_endpoint: String,
|
||||
pub aliyun_matting_access_key_id: Option<String>,
|
||||
@@ -206,6 +215,7 @@ pub struct AppConfig {
|
||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||
pub enum ProcessRole {
|
||||
Api,
|
||||
BgfilterWorker,
|
||||
ExternalGenerationWorker,
|
||||
ExternalGenerationController,
|
||||
All,
|
||||
@@ -234,6 +244,7 @@ impl ProcessRole {
|
||||
pub fn as_str(self) -> &'static str {
|
||||
match self {
|
||||
Self::Api => "api",
|
||||
Self::BgfilterWorker => "bgfilter-worker",
|
||||
Self::ExternalGenerationWorker => "external-generation-worker",
|
||||
Self::ExternalGenerationController => "external-generation-controller",
|
||||
Self::All => "all",
|
||||
@@ -244,6 +255,10 @@ impl ProcessRole {
|
||||
matches!(self, Self::Api | Self::All)
|
||||
}
|
||||
|
||||
pub fn runs_bgfilter_worker(self) -> bool {
|
||||
matches!(self, Self::BgfilterWorker)
|
||||
}
|
||||
|
||||
pub fn runs_external_generation_worker(self) -> bool {
|
||||
matches!(self, Self::ExternalGenerationWorker | Self::All)
|
||||
}
|
||||
@@ -261,6 +276,13 @@ impl Default for AppConfig {
|
||||
listen_backlog: 1024,
|
||||
worker_threads: None,
|
||||
process_role: ProcessRole::Api,
|
||||
bgfilter_worker_host: "127.0.0.1".to_string(),
|
||||
bgfilter_worker_port: 8083,
|
||||
bgfilter_worker_base_url: "http://127.0.0.1:8083".to_string(),
|
||||
bgfilter_internal_token: None,
|
||||
bgfilter_worker_concurrency: 16,
|
||||
bgfilter_worker_max_requests: 2048,
|
||||
bgfilter_worker_connect_timeout_ms: 2_000,
|
||||
external_generation_mode: ExternalGenerationMode::Queue,
|
||||
external_generation_worker_id: default_external_generation_worker_id(),
|
||||
external_generation_worker_concurrency: 2,
|
||||
@@ -299,7 +321,8 @@ impl Default for AppConfig {
|
||||
crate::editor_generation_config::default_editor_generation_pricing_override_path(),
|
||||
editor_bgfilter_base_url: DEFAULT_EDITOR_BGFILTER_BASE_URL.to_string(),
|
||||
editor_bgfilter_token: None,
|
||||
editor_bgfilter_request_timeout_ms: DEFAULT_EDITOR_BGFILTER_REQUEST_TIMEOUT_MS,
|
||||
editor_bgfilter_single_image_estimate_ms:
|
||||
DEFAULT_EDITOR_BGFILTER_SINGLE_IMAGE_ESTIMATE_MS,
|
||||
aliyun_matting_enabled: true,
|
||||
aliyun_matting_endpoint: DEFAULT_ALIYUN_MATTING_ENDPOINT.to_string(),
|
||||
aliyun_matting_access_key_id: None,
|
||||
@@ -454,6 +477,22 @@ impl Default for AppConfig {
|
||||
}
|
||||
|
||||
impl AppConfig {
|
||||
/// 单次 provider attempt 上限:`N × est × 2`。运行时派生,禁止硬编码计算结果;
|
||||
/// worker 角色在 main 中覆盖 `bgfilter_worker_concurrency` 后自动生效。
|
||||
pub fn bgfilter_provider_attempt_timeout_ms(&self) -> u64 {
|
||||
(self.bgfilter_worker_concurrency.max(1) as u64)
|
||||
.saturating_mul(self.editor_bgfilter_single_image_estimate_ms.max(1))
|
||||
.saturating_mul(BGFILTER_ATTEMPT_SAFETY_FACTOR)
|
||||
}
|
||||
|
||||
/// 调用预算:自取得 provider permit 起算,覆盖签名、两次顺序 attempt、
|
||||
/// 结果校验与响应构造;排队时长不消耗它。
|
||||
pub fn bgfilter_call_budget_ms(&self) -> u64 {
|
||||
self.bgfilter_provider_attempt_timeout_ms()
|
||||
.saturating_mul(2)
|
||||
.saturating_add(BGFILTER_WORKER_RESPONSE_WINDOW_MS)
|
||||
}
|
||||
|
||||
pub fn from_env() -> Self {
|
||||
let mut config = Self::default();
|
||||
|
||||
@@ -473,6 +512,36 @@ impl AppConfig {
|
||||
config.bind_port = parsed_port;
|
||||
}
|
||||
|
||||
if let Some(host) = read_first_non_empty_env(&["GENARRATIVE_BGFILTER_WORKER_HOST"]) {
|
||||
config.bgfilter_worker_host = host;
|
||||
}
|
||||
if let Some(port) = read_first_positive_u16_env(&["GENARRATIVE_BGFILTER_WORKER_PORT"]) {
|
||||
config.bgfilter_worker_port = port;
|
||||
}
|
||||
if let Some(base_url) = read_first_non_empty_env(&["GENARRATIVE_BGFILTER_WORKER_BASE_URL"])
|
||||
{
|
||||
config.bgfilter_worker_base_url = base_url;
|
||||
}
|
||||
config.bgfilter_internal_token = read_secret_env_or_file(
|
||||
&["GENARRATIVE_BGFILTER_INTERNAL_TOKEN"],
|
||||
&["GENARRATIVE_BGFILTER_INTERNAL_TOKEN_FILE"],
|
||||
);
|
||||
if let Some(concurrency) =
|
||||
read_first_usize_env(&["GENARRATIVE_BGFILTER_WORKER_CONCURRENCY"])
|
||||
{
|
||||
config.bgfilter_worker_concurrency = concurrency;
|
||||
}
|
||||
if let Some(max_requests) =
|
||||
read_first_usize_env(&["GENARRATIVE_BGFILTER_WORKER_MAX_REQUESTS"])
|
||||
{
|
||||
config.bgfilter_worker_max_requests = max_requests;
|
||||
}
|
||||
if let Some(connect_timeout_ms) =
|
||||
read_first_positive_u64_env(&["GENARRATIVE_BGFILTER_WORKER_CONNECT_TIMEOUT_MS"])
|
||||
{
|
||||
config.bgfilter_worker_connect_timeout_ms = connect_timeout_ms;
|
||||
}
|
||||
|
||||
if let Ok(log_filter) = env::var("GENARRATIVE_API_LOG")
|
||||
&& !log_filter.trim().is_empty()
|
||||
{
|
||||
@@ -491,10 +560,10 @@ impl AppConfig {
|
||||
"GENARRATIVE_EDITOR_BGFILTER_TOKEN",
|
||||
"GENARRATIVE_EDITOR_BACKGROUND_REMOVAL_TOKEN",
|
||||
]);
|
||||
if let Some(timeout_ms) =
|
||||
read_first_positive_u64_env(&["GENARRATIVE_EDITOR_BGFILTER_REQUEST_TIMEOUT_MS"])
|
||||
if let Some(estimate_ms) =
|
||||
read_first_positive_u64_env(&["GENARRATIVE_EDITOR_BGFILTER_SINGLE_IMAGE_ESTIMATE_MS"])
|
||||
{
|
||||
config.editor_bgfilter_request_timeout_ms = timeout_ms;
|
||||
config.editor_bgfilter_single_image_estimate_ms = estimate_ms;
|
||||
}
|
||||
if let Some(enabled) = read_first_bool_env(&["GENARRATIVE_ALIYUN_MATTING_ENABLED"]) {
|
||||
config.aliyun_matting_enabled = enabled;
|
||||
@@ -1370,6 +1439,7 @@ fn default_external_generation_worker_id() -> String {
|
||||
fn parse_process_role(value: &str) -> Option<ProcessRole> {
|
||||
match trim_quoted_env_value(value).to_ascii_lowercase().as_str() {
|
||||
"api" => Some(ProcessRole::Api),
|
||||
"bgfilter-worker" | "bgfilter_worker" => Some(ProcessRole::BgfilterWorker),
|
||||
"external-generation-worker" | "external_generation_worker" | "worker" => {
|
||||
Some(ProcessRole::ExternalGenerationWorker)
|
||||
}
|
||||
@@ -1526,7 +1596,7 @@ mod tests {
|
||||
AppConfig, DEFAULT_EDITOR_BGFILTER_BASE_URL,
|
||||
DEFAULT_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS,
|
||||
DEFAULT_EDITOR_BGFILTER_CIRCUIT_FAILURE_THRESHOLD,
|
||||
DEFAULT_EDITOR_BGFILTER_REQUEST_TIMEOUT_MS,
|
||||
DEFAULT_EDITOR_BGFILTER_SINGLE_IMAGE_ESTIMATE_MS,
|
||||
DEFAULT_EXTERNAL_GENERATION_WORKER_JOB_TIMEOUT_SECONDS,
|
||||
DEFAULT_EXTERNAL_GENERATION_WORKER_LEASE_SECONDS,
|
||||
DEFAULT_EXTERNAL_GENERATION_WORKER_LONG_JOB_TIMEOUT_SECONDS, ExternalGenerationMode,
|
||||
@@ -1543,6 +1613,13 @@ mod tests {
|
||||
fn default_keeps_non_public_model_and_base_url_empty() {
|
||||
let config = AppConfig::default();
|
||||
|
||||
assert_eq!(config.bgfilter_worker_host, "127.0.0.1");
|
||||
assert_eq!(config.bgfilter_worker_port, 8083);
|
||||
assert_eq!(config.bgfilter_worker_base_url, "http://127.0.0.1:8083");
|
||||
assert!(config.bgfilter_internal_token.is_none());
|
||||
assert_eq!(config.bgfilter_worker_concurrency, 16);
|
||||
assert_eq!(config.bgfilter_worker_max_requests, 2048);
|
||||
assert_eq!(config.bgfilter_worker_connect_timeout_ms, 2_000);
|
||||
assert!(config.llm_model.is_empty());
|
||||
assert!(config.llm_base_url.is_empty());
|
||||
// assert!(config.apimart_base_url.is_empty());
|
||||
@@ -1552,9 +1629,11 @@ mod tests {
|
||||
DEFAULT_EDITOR_BGFILTER_BASE_URL
|
||||
);
|
||||
assert_eq!(
|
||||
config.editor_bgfilter_request_timeout_ms,
|
||||
DEFAULT_EDITOR_BGFILTER_REQUEST_TIMEOUT_MS
|
||||
config.editor_bgfilter_single_image_estimate_ms,
|
||||
DEFAULT_EDITOR_BGFILTER_SINGLE_IMAGE_ESTIMATE_MS
|
||||
);
|
||||
assert_eq!(config.bgfilter_provider_attempt_timeout_ms(), 160_000);
|
||||
assert_eq!(config.bgfilter_call_budget_ms(), 321_000);
|
||||
assert_eq!(
|
||||
config.editor_bgfilter_circuit_failure_threshold,
|
||||
DEFAULT_EDITOR_BGFILTER_CIRCUIT_FAILURE_THRESHOLD
|
||||
@@ -1563,6 +1642,7 @@ mod tests {
|
||||
config.editor_bgfilter_circuit_cooldown.as_secs(),
|
||||
DEFAULT_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS
|
||||
);
|
||||
assert_eq!(DEFAULT_EDITOR_BGFILTER_CIRCUIT_COOLDOWN_SECONDS, 120);
|
||||
assert!(config.editor_bgfilter_token.is_none());
|
||||
assert_eq!(
|
||||
config.external_generation_worker_lease.as_secs(),
|
||||
@@ -1607,6 +1687,14 @@ mod tests {
|
||||
#[test]
|
||||
fn process_role_controls_http_and_external_generation_worker_roles() {
|
||||
assert_eq!(parse_process_role("api"), Some(ProcessRole::Api));
|
||||
assert_eq!(
|
||||
parse_process_role("bgfilter-worker"),
|
||||
Some(ProcessRole::BgfilterWorker)
|
||||
);
|
||||
assert_eq!(
|
||||
parse_process_role("'bgfilter_worker'"),
|
||||
Some(ProcessRole::BgfilterWorker)
|
||||
);
|
||||
assert_eq!(
|
||||
parse_process_role("\"external-generation-worker\""),
|
||||
Some(ProcessRole::ExternalGenerationWorker)
|
||||
@@ -1629,17 +1717,26 @@ mod tests {
|
||||
);
|
||||
assert_eq!(parse_process_role("all"), Some(ProcessRole::All));
|
||||
assert_eq!(parse_process_role("unknown"), None);
|
||||
assert_eq!(ProcessRole::BgfilterWorker.as_str(), "bgfilter-worker");
|
||||
|
||||
assert!(ProcessRole::Api.runs_http());
|
||||
assert!(!ProcessRole::Api.runs_bgfilter_worker());
|
||||
assert!(!ProcessRole::Api.runs_external_generation_worker());
|
||||
assert!(!ProcessRole::Api.runs_external_generation_controller());
|
||||
assert!(!ProcessRole::BgfilterWorker.runs_http());
|
||||
assert!(ProcessRole::BgfilterWorker.runs_bgfilter_worker());
|
||||
assert!(!ProcessRole::BgfilterWorker.runs_external_generation_worker());
|
||||
assert!(!ProcessRole::BgfilterWorker.runs_external_generation_controller());
|
||||
assert!(!ProcessRole::ExternalGenerationWorker.runs_http());
|
||||
assert!(!ProcessRole::ExternalGenerationWorker.runs_bgfilter_worker());
|
||||
assert!(ProcessRole::ExternalGenerationWorker.runs_external_generation_worker());
|
||||
assert!(!ProcessRole::ExternalGenerationWorker.runs_external_generation_controller());
|
||||
assert!(!ProcessRole::ExternalGenerationController.runs_http());
|
||||
assert!(!ProcessRole::ExternalGenerationController.runs_bgfilter_worker());
|
||||
assert!(!ProcessRole::ExternalGenerationController.runs_external_generation_worker());
|
||||
assert!(ProcessRole::ExternalGenerationController.runs_external_generation_controller());
|
||||
assert!(ProcessRole::All.runs_http());
|
||||
assert!(!ProcessRole::All.runs_bgfilter_worker());
|
||||
assert!(ProcessRole::All.runs_external_generation_worker());
|
||||
assert!(!ProcessRole::All.runs_external_generation_controller());
|
||||
}
|
||||
@@ -1872,6 +1969,68 @@ mod tests {
|
||||
let _ = fs::remove_file(secret_path);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn from_env_reads_bgfilter_worker_settings_and_secret_from_value_or_file() {
|
||||
let _guard = ENV_LOCK
|
||||
.get_or_init(|| Mutex::new(()))
|
||||
.lock()
|
||||
.expect("env lock should not poison");
|
||||
let secret_path = std::env::temp_dir().join(format!(
|
||||
"genarrative-bgfilter-internal-token-{}.txt",
|
||||
std::process::id()
|
||||
));
|
||||
fs::write(&secret_path, " file-token \n").expect("secret file should write");
|
||||
|
||||
unsafe {
|
||||
std::env::remove_var("GENARRATIVE_BGFILTER_INTERNAL_TOKEN");
|
||||
std::env::set_var("GENARRATIVE_BGFILTER_INTERNAL_TOKEN_FILE", &secret_path);
|
||||
std::env::set_var("GENARRATIVE_BGFILTER_WORKER_HOST", "127.0.0.2");
|
||||
std::env::set_var("GENARRATIVE_BGFILTER_WORKER_PORT", "18083");
|
||||
std::env::set_var(
|
||||
"GENARRATIVE_BGFILTER_WORKER_BASE_URL",
|
||||
"http://127.0.0.2:18083",
|
||||
);
|
||||
std::env::set_var("GENARRATIVE_BGFILTER_WORKER_CONCURRENCY", "6");
|
||||
std::env::set_var("GENARRATIVE_BGFILTER_WORKER_MAX_REQUESTS", "192");
|
||||
std::env::set_var("GENARRATIVE_BGFILTER_WORKER_CONNECT_TIMEOUT_MS", "3500");
|
||||
std::env::set_var("GENARRATIVE_PROCESS_ROLE", "bgfilter_worker");
|
||||
}
|
||||
|
||||
let config = AppConfig::from_env();
|
||||
assert_eq!(config.process_role, ProcessRole::BgfilterWorker);
|
||||
assert_eq!(config.bgfilter_worker_host, "127.0.0.2");
|
||||
assert_eq!(config.bgfilter_worker_port, 18_083);
|
||||
assert_eq!(config.bgfilter_worker_base_url, "http://127.0.0.2:18083");
|
||||
assert_eq!(
|
||||
config.bgfilter_internal_token.as_deref(),
|
||||
Some("file-token")
|
||||
);
|
||||
assert_eq!(config.bgfilter_worker_concurrency, 6);
|
||||
assert_eq!(config.bgfilter_worker_max_requests, 192);
|
||||
assert_eq!(config.bgfilter_worker_connect_timeout_ms, 3_500);
|
||||
|
||||
unsafe {
|
||||
std::env::set_var("GENARRATIVE_BGFILTER_INTERNAL_TOKEN", "direct-token");
|
||||
}
|
||||
assert_eq!(
|
||||
AppConfig::from_env().bgfilter_internal_token.as_deref(),
|
||||
Some("direct-token")
|
||||
);
|
||||
|
||||
unsafe {
|
||||
std::env::remove_var("GENARRATIVE_BGFILTER_INTERNAL_TOKEN");
|
||||
std::env::remove_var("GENARRATIVE_BGFILTER_INTERNAL_TOKEN_FILE");
|
||||
std::env::remove_var("GENARRATIVE_BGFILTER_WORKER_HOST");
|
||||
std::env::remove_var("GENARRATIVE_BGFILTER_WORKER_PORT");
|
||||
std::env::remove_var("GENARRATIVE_BGFILTER_WORKER_BASE_URL");
|
||||
std::env::remove_var("GENARRATIVE_BGFILTER_WORKER_CONCURRENCY");
|
||||
std::env::remove_var("GENARRATIVE_BGFILTER_WORKER_MAX_REQUESTS");
|
||||
std::env::remove_var("GENARRATIVE_BGFILTER_WORKER_CONNECT_TIMEOUT_MS");
|
||||
std::env::remove_var("GENARRATIVE_PROCESS_ROLE");
|
||||
}
|
||||
let _ = fs::remove_file(secret_path);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn from_env_reads_api_runtime_performance_settings() {
|
||||
let _guard = ENV_LOCK
|
||||
@@ -2218,7 +2377,7 @@ mod tests {
|
||||
unsafe {
|
||||
std::env::remove_var("GENARRATIVE_EDITOR_BGFILTER_BASE_URL");
|
||||
std::env::remove_var("GENARRATIVE_EDITOR_BGFILTER_TOKEN");
|
||||
std::env::remove_var("GENARRATIVE_EDITOR_BGFILTER_REQUEST_TIMEOUT_MS");
|
||||
std::env::remove_var("GENARRATIVE_EDITOR_BGFILTER_SINGLE_IMAGE_ESTIMATE_MS");
|
||||
std::env::remove_var("GENARRATIVE_EDITOR_BACKGROUND_REMOVAL_TOKEN");
|
||||
std::env::set_var(
|
||||
"GENARRATIVE_EDITOR_BGFILTER_BASE_URL",
|
||||
@@ -2228,7 +2387,10 @@ mod tests {
|
||||
"GENARRATIVE_EDITOR_BACKGROUND_REMOVAL_TOKEN",
|
||||
"shared-token",
|
||||
);
|
||||
std::env::set_var("GENARRATIVE_EDITOR_BGFILTER_REQUEST_TIMEOUT_MS", "240000");
|
||||
std::env::set_var(
|
||||
"GENARRATIVE_EDITOR_BGFILTER_SINGLE_IMAGE_ESTIMATE_MS",
|
||||
"8000",
|
||||
);
|
||||
}
|
||||
|
||||
let config = AppConfig::from_env();
|
||||
@@ -2237,7 +2399,16 @@ mod tests {
|
||||
config.editor_bgfilter_token.as_deref(),
|
||||
Some("shared-token")
|
||||
);
|
||||
assert_eq!(config.editor_bgfilter_request_timeout_ms, 240_000);
|
||||
assert_eq!(config.editor_bgfilter_single_image_estimate_ms, 8_000);
|
||||
// attempt / callBudget 必须随 N 与 est 联动派生,而不是固定值。
|
||||
assert_eq!(
|
||||
config.bgfilter_provider_attempt_timeout_ms(),
|
||||
config.bgfilter_worker_concurrency as u64 * 8_000 * 2
|
||||
);
|
||||
assert_eq!(
|
||||
config.bgfilter_call_budget_ms(),
|
||||
config.bgfilter_provider_attempt_timeout_ms() * 2 + 1_000
|
||||
);
|
||||
|
||||
unsafe {
|
||||
std::env::set_var("GENARRATIVE_EDITOR_BGFILTER_TOKEN", "bgfilter-token");
|
||||
@@ -2252,7 +2423,7 @@ mod tests {
|
||||
unsafe {
|
||||
std::env::remove_var("GENARRATIVE_EDITOR_BGFILTER_BASE_URL");
|
||||
std::env::remove_var("GENARRATIVE_EDITOR_BGFILTER_TOKEN");
|
||||
std::env::remove_var("GENARRATIVE_EDITOR_BGFILTER_REQUEST_TIMEOUT_MS");
|
||||
std::env::remove_var("GENARRATIVE_EDITOR_BGFILTER_SINGLE_IMAGE_ESTIMATE_MS");
|
||||
std::env::remove_var("GENARRATIVE_EDITOR_BACKGROUND_REMOVAL_TOKEN");
|
||||
}
|
||||
}
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -745,6 +745,7 @@ mod tests {
|
||||
user_id: Some("user-1".to_string()),
|
||||
profile_id: Some("project-1".to_string()),
|
||||
request_id: Some("request-1".to_string()),
|
||||
external_call_deadline: None,
|
||||
},
|
||||
);
|
||||
let tracking = crate::external_api_audit::build_external_api_failure_tracking_draft(&audit);
|
||||
|
||||
@@ -1,3 +1,5 @@
|
||||
use std::time::Instant;
|
||||
|
||||
#[cfg(test)]
|
||||
use axum::http::StatusCode;
|
||||
use module_runtime::RuntimeTrackingScopeKind;
|
||||
@@ -145,6 +147,9 @@ pub(crate) struct ExternalApiAuditContext {
|
||||
pub(crate) user_id: Option<String>,
|
||||
pub(crate) profile_id: Option<String>,
|
||||
pub(crate) request_id: Option<String>,
|
||||
/// 父流程允许外部调用占用到的绝对时刻。该字段只在进程内用于预算截断,
|
||||
/// 不写入 tracking metadata,也不会通过内部协议传递绝对时间。
|
||||
pub(crate) external_call_deadline: Option<Instant>,
|
||||
}
|
||||
|
||||
/// 抠图供应商(BgFilter / 阿里云通用抠图)调用失败的统一失败审计入口。
|
||||
@@ -180,6 +185,45 @@ pub(crate) async fn record_matting_external_api_failure(
|
||||
record_external_api_failure(state, draft).await;
|
||||
}
|
||||
|
||||
/// BgFilter worker 专用入口:保留同一份 OTLP / tracking draft,但只允许写入本进程
|
||||
/// 独立 outbox。outbox 缺失、满载或写盘失败时丢弃,禁止在受限 worker 中逐条同步
|
||||
/// 直写 SpacetimeDB。
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub(crate) async fn record_matting_external_api_failure_outbox_only(
|
||||
state: &AppState,
|
||||
context: &ExternalApiAuditContext,
|
||||
provider: &'static str,
|
||||
endpoint: String,
|
||||
operation: &'static str,
|
||||
failure_stage: &'static str,
|
||||
status_code: Option<u16>,
|
||||
timeout: bool,
|
||||
transport: bool,
|
||||
latency_ms: Option<u64>,
|
||||
error_message: String,
|
||||
raw_excerpt: Option<String>,
|
||||
) {
|
||||
let draft = build_matting_external_api_failure_draft(
|
||||
provider,
|
||||
endpoint,
|
||||
operation,
|
||||
failure_stage,
|
||||
status_code,
|
||||
timeout,
|
||||
transport,
|
||||
latency_ms,
|
||||
error_message,
|
||||
raw_excerpt,
|
||||
context,
|
||||
);
|
||||
record_external_api_failure_with_policy(
|
||||
state,
|
||||
draft,
|
||||
ExternalApiAuditPersistencePolicy::RequireOutboxDropOnFailure,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
/// 构建抠图失败审计 draft。`transport` 必须由调用方从错误结构化字段读取,
|
||||
/// 不能用 `status_code.is_none()` 反推——本地处理失败同样没有上游 HTTP 状态。
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
@@ -347,8 +391,33 @@ pub(crate) fn app_error_status_class(status_code: StatusCode) -> &'static str {
|
||||
status_class(Some(status_code.as_u16()))
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||
enum ExternalApiAuditPersistencePolicy {
|
||||
PreferOutboxThenSync,
|
||||
RequireOutboxDropOnFailure,
|
||||
}
|
||||
|
||||
impl ExternalApiAuditPersistencePolicy {
|
||||
fn allows_sync_fallback(self) -> bool {
|
||||
matches!(self, Self::PreferOutboxThenSync)
|
||||
}
|
||||
}
|
||||
|
||||
/// 中文注释:外部供应商失败同时进入 OTLP 和 tracking_event;失败审计不能反向阻断主业务错误返回。
|
||||
pub(crate) async fn record_external_api_failure(state: &AppState, draft: ExternalApiFailureDraft) {
|
||||
record_external_api_failure_with_policy(
|
||||
state,
|
||||
draft,
|
||||
ExternalApiAuditPersistencePolicy::PreferOutboxThenSync,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
async fn record_external_api_failure_with_policy(
|
||||
state: &AppState,
|
||||
draft: ExternalApiFailureDraft,
|
||||
persistence_policy: ExternalApiAuditPersistencePolicy,
|
||||
) {
|
||||
record_external_api_failure_otlp(&draft);
|
||||
|
||||
let tracking_event = build_external_api_failure_tracking_draft(&draft);
|
||||
@@ -361,6 +430,18 @@ pub(crate) async fn record_external_api_failure(state: &AppState, draft: Externa
|
||||
{
|
||||
Ok(crate::tracking_outbox::TrackingOutboxEnqueueOutcome::Enqueued) => {}
|
||||
Ok(crate::tracking_outbox::TrackingOutboxEnqueueOutcome::Dropped { reason }) => {
|
||||
if !persistence_policy.allows_sync_fallback() {
|
||||
crate::telemetry::record_external_api_audit_dropped(draft.provider, reason);
|
||||
tracing::warn!(
|
||||
provider = draft.provider,
|
||||
endpoint = %draft.endpoint,
|
||||
operation = %draft.operation,
|
||||
failure_stage = draft.failure_stage,
|
||||
reason,
|
||||
"外部 API 失败审计写入专用 outbox 被保护阈值拒绝,已丢弃"
|
||||
);
|
||||
return;
|
||||
}
|
||||
tracing::warn!(
|
||||
provider = draft.provider,
|
||||
endpoint = %draft.endpoint,
|
||||
@@ -377,6 +458,21 @@ pub(crate) async fn record_external_api_failure(state: &AppState, draft: Externa
|
||||
.await;
|
||||
}
|
||||
Err(error) => {
|
||||
if !persistence_policy.allows_sync_fallback() {
|
||||
crate::telemetry::record_external_api_audit_dropped(
|
||||
draft.provider,
|
||||
"outbox_error",
|
||||
);
|
||||
tracing::warn!(
|
||||
provider = draft.provider,
|
||||
endpoint = %draft.endpoint,
|
||||
operation = %draft.operation,
|
||||
failure_stage = draft.failure_stage,
|
||||
error = %error,
|
||||
"外部 API 失败审计写入专用 outbox 失败,已丢弃"
|
||||
);
|
||||
return;
|
||||
}
|
||||
tracing::warn!(
|
||||
provider = draft.provider,
|
||||
endpoint = %draft.endpoint,
|
||||
@@ -396,6 +492,18 @@ pub(crate) async fn record_external_api_failure(state: &AppState, draft: Externa
|
||||
return;
|
||||
}
|
||||
|
||||
if !persistence_policy.allows_sync_fallback() {
|
||||
crate::telemetry::record_external_api_audit_dropped(draft.provider, "outbox_missing");
|
||||
tracing::warn!(
|
||||
provider = draft.provider,
|
||||
endpoint = %draft.endpoint,
|
||||
operation = %draft.operation,
|
||||
failure_stage = draft.failure_stage,
|
||||
"外部 API 失败审计缺少专用 outbox,已丢弃"
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
crate::tracking::record_tracking_event_after_success(
|
||||
state,
|
||||
&audit_request_context(),
|
||||
@@ -567,6 +675,14 @@ mod tests {
|
||||
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn bgfilter_outbox_only_policy_never_allows_sync_fallback() {
|
||||
assert!(ExternalApiAuditPersistencePolicy::PreferOutboxThenSync.allows_sync_fallback());
|
||||
assert!(
|
||||
!ExternalApiAuditPersistencePolicy::RequireOutboxDropOnFailure.allows_sync_fallback()
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn external_api_failure_tracking_draft_uses_module_scope_and_safe_metadata() {
|
||||
let draft = build_external_api_failure_tracking_draft(
|
||||
|
||||
@@ -16,6 +16,7 @@ mod auth_public_user;
|
||||
mod auth_session;
|
||||
mod auth_sessions;
|
||||
mod backpressure;
|
||||
mod bgfilter_worker;
|
||||
mod character_animation_assets;
|
||||
mod character_visual_assets;
|
||||
mod config;
|
||||
@@ -92,6 +93,7 @@ use tracing::{error, info, warn};
|
||||
|
||||
use crate::{
|
||||
app::{build_router, build_spacetime_unavailable_router},
|
||||
bgfilter_worker::{build_bgfilter_worker_router, validate_bgfilter_internal_token},
|
||||
config::{AppConfig, ProcessRole},
|
||||
external_generation_worker::run_external_generation_worker,
|
||||
external_generation_worker_controller::run_external_generation_worker_controller,
|
||||
@@ -141,6 +143,7 @@ fn main() -> Result<(), io::Error> {
|
||||
}
|
||||
|
||||
async fn run_server(config: AppConfig) -> Result<(), io::Error> {
|
||||
validate_bgfilter_internal_token_for_startup(&config).map_err(io::Error::other)?;
|
||||
init_tracing(
|
||||
&config.log_filter,
|
||||
OtelConfig {
|
||||
@@ -150,6 +153,10 @@ async fn run_server(config: AppConfig) -> Result<(), io::Error> {
|
||||
process_metrics::register_process_metrics();
|
||||
telemetry::register_http_runtime_metrics();
|
||||
|
||||
if config.process_role.runs_bgfilter_worker() {
|
||||
return run_bgfilter_worker_role(config).await;
|
||||
}
|
||||
|
||||
if !config.process_role.runs_http() {
|
||||
return run_worker_only(config).await;
|
||||
}
|
||||
@@ -157,6 +164,159 @@ async fn run_server(config: AppConfig) -> Result<(), io::Error> {
|
||||
run_http_role(config).await
|
||||
}
|
||||
|
||||
fn validate_bgfilter_internal_token_for_startup(config: &AppConfig) -> Result<(), String> {
|
||||
if should_validate_bgfilter_internal_token_for_startup(config.process_role) {
|
||||
validate_bgfilter_internal_token(config.bgfilter_internal_token.as_deref())?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn should_validate_bgfilter_internal_token_for_startup(process_role: ProcessRole) -> bool {
|
||||
matches!(
|
||||
process_role,
|
||||
ProcessRole::Api
|
||||
| ProcessRole::BgfilterWorker
|
||||
| ProcessRole::ExternalGenerationWorker
|
||||
| ProcessRole::All
|
||||
)
|
||||
}
|
||||
|
||||
async fn run_bgfilter_worker_role(mut config: AppConfig) -> Result<(), io::Error> {
|
||||
let (concurrency, single_image_estimate_ms, max_requests) =
|
||||
required_bgfilter_worker_capacity_from_env()?;
|
||||
config.bgfilter_worker_concurrency = concurrency;
|
||||
config.editor_bgfilter_single_image_estimate_ms = single_image_estimate_ms;
|
||||
config.bgfilter_worker_max_requests = max_requests;
|
||||
let bind_address = format!(
|
||||
"{}:{}",
|
||||
config.bgfilter_worker_host, config.bgfilter_worker_port
|
||||
)
|
||||
.parse::<SocketAddr>()
|
||||
.map_err(|error| io::Error::other(format!("bgfilter-worker 监听地址无效:{error}")))?;
|
||||
if !bind_address.ip().is_loopback() {
|
||||
return Err(io::Error::other(format!(
|
||||
"bgfilter-worker 首版只允许监听 loopback,当前地址为 {bind_address}"
|
||||
)));
|
||||
}
|
||||
let listen_backlog = config.listen_backlog;
|
||||
let outbox_flush_timeout = config.shutdown_outbox_flush_timeout;
|
||||
let listener = build_tcp_listener(bind_address, listen_backlog)?;
|
||||
|
||||
configure_bgfilter_worker_outboxes(&mut config);
|
||||
let state = AppState::new_with_empty_auth_store(config)
|
||||
.map_err(|error| io::Error::other(format!("初始化 bgfilter-worker 状态失败:{error}")))?;
|
||||
let (router, task_tracker) = build_bgfilter_worker_router(state.clone())
|
||||
.map_err(|error| io::Error::other(format!("初始化 bgfilter-worker 路由失败:{error}")))?;
|
||||
let tracking_outbox = state.tracking_outbox();
|
||||
if let Some(outbox) = tracking_outbox.clone() {
|
||||
outbox.spawn_worker();
|
||||
}
|
||||
let shutdown_context = ShutdownContext {
|
||||
app_state: Some(state),
|
||||
tracking_outbox,
|
||||
wallet_refund_outbox: None,
|
||||
outbox_flush_timeout,
|
||||
};
|
||||
info!(
|
||||
%bind_address,
|
||||
listen_backlog,
|
||||
process_role = ProcessRole::BgfilterWorker.as_str(),
|
||||
"bgfilter-worker 已开始监听内部 HTTP"
|
||||
);
|
||||
let shutdown_tracker = task_tracker.clone();
|
||||
let shutdown_context_for_signal = shutdown_context.clone();
|
||||
let result = axum::serve(listener, router)
|
||||
.with_graceful_shutdown(async move {
|
||||
shutdown_signal(shutdown_context_for_signal).await;
|
||||
shutdown_tracker.close();
|
||||
})
|
||||
.await;
|
||||
task_tracker.close();
|
||||
task_tracker.wait_for_drain().await;
|
||||
finalize_shutdown(shutdown_context).await;
|
||||
result
|
||||
}
|
||||
|
||||
fn configure_bgfilter_worker_outboxes(config: &mut AppConfig) {
|
||||
// 多进程不能操作同一个 active 文件;worker 从共享基础目录派生自己的持久子目录。
|
||||
config.tracking_outbox_enabled = true;
|
||||
config.tracking_outbox_dir = config.tracking_outbox_dir.join("bgfilter-worker");
|
||||
config.wallet_refund_outbox_enabled = false;
|
||||
}
|
||||
|
||||
const DEFAULT_BGFILTER_WORKER_MAX_REQUESTS_FUSE: usize = 2_048;
|
||||
|
||||
fn required_bgfilter_worker_capacity_from_env() -> Result<(usize, u64, usize), io::Error> {
|
||||
let concurrency = env::var("GENARRATIVE_BGFILTER_WORKER_CONCURRENCY").ok();
|
||||
let single_image_estimate_ms =
|
||||
env::var("GENARRATIVE_EDITOR_BGFILTER_SINGLE_IMAGE_ESTIMATE_MS").ok();
|
||||
let max_requests = env::var("GENARRATIVE_BGFILTER_WORKER_MAX_REQUESTS").ok();
|
||||
parse_required_bgfilter_worker_capacity(
|
||||
concurrency.as_deref(),
|
||||
single_image_estimate_ms.as_deref(),
|
||||
max_requests.as_deref(),
|
||||
)
|
||||
}
|
||||
|
||||
fn parse_required_bgfilter_worker_capacity(
|
||||
concurrency: Option<&str>,
|
||||
single_image_estimate_ms: Option<&str>,
|
||||
max_requests: Option<&str>,
|
||||
) -> Result<(usize, u64, usize), io::Error> {
|
||||
fn parse_required_positive(name: &str, raw: Option<&str>) -> Result<usize, io::Error> {
|
||||
let value = raw
|
||||
.map(strip_env_value)
|
||||
.map(|value| value.trim().to_string())
|
||||
.filter(|value| !value.is_empty())
|
||||
.ok_or_else(|| io::Error::other(format!("bgfilter-worker 启动必须显式配置 {name}")))?;
|
||||
let parsed = value.parse::<usize>().map_err(|error| {
|
||||
io::Error::other(format!(
|
||||
"bgfilter-worker 配置 {name} 不是有效正整数:{error}"
|
||||
))
|
||||
})?;
|
||||
if parsed == 0 {
|
||||
return Err(io::Error::other(format!(
|
||||
"bgfilter-worker 配置 {name} 必须大于 0"
|
||||
)));
|
||||
}
|
||||
Ok(parsed)
|
||||
}
|
||||
|
||||
let concurrency =
|
||||
parse_required_positive("GENARRATIVE_BGFILTER_WORKER_CONCURRENCY", concurrency)?;
|
||||
let single_image_estimate_ms = parse_required_positive(
|
||||
"GENARRATIVE_EDITOR_BGFILTER_SINGLE_IMAGE_ESTIMATE_MS",
|
||||
single_image_estimate_ms,
|
||||
)? as u64;
|
||||
// Q 已降级为 admission 保险丝:可缺省(默认 2048),显式配置时仍必须为正且不小于 N。
|
||||
let max_requests = match max_requests
|
||||
.map(strip_env_value)
|
||||
.map(|value| value.trim().to_string())
|
||||
.filter(|value| !value.is_empty())
|
||||
{
|
||||
None => DEFAULT_BGFILTER_WORKER_MAX_REQUESTS_FUSE,
|
||||
Some(raw) => {
|
||||
let parsed = raw.parse::<usize>().map_err(|error| {
|
||||
io::Error::other(format!(
|
||||
"bgfilter-worker 配置 GENARRATIVE_BGFILTER_WORKER_MAX_REQUESTS 不是有效正整数:{error}"
|
||||
))
|
||||
})?;
|
||||
if parsed == 0 {
|
||||
return Err(io::Error::other(
|
||||
"bgfilter-worker 配置 GENARRATIVE_BGFILTER_WORKER_MAX_REQUESTS 必须大于 0",
|
||||
));
|
||||
}
|
||||
parsed
|
||||
}
|
||||
};
|
||||
if max_requests < concurrency {
|
||||
return Err(io::Error::other(
|
||||
"GENARRATIVE_BGFILTER_WORKER_MAX_REQUESTS 不能小于 GENARRATIVE_BGFILTER_WORKER_CONCURRENCY",
|
||||
));
|
||||
}
|
||||
Ok((concurrency, single_image_estimate_ms, max_requests))
|
||||
}
|
||||
|
||||
async fn run_worker_only(config: AppConfig) -> Result<(), io::Error> {
|
||||
let process_role = config.process_role;
|
||||
let state = build_non_http_app_state_for_startup(config).map_err(|error| {
|
||||
@@ -545,12 +705,14 @@ fn is_valid_env_key(key: &str) -> bool {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::{
|
||||
AUTH_STORE_STARTUP_RETRY_INTERVAL, is_valid_env_key, protected_env_keys_from,
|
||||
AUTH_STORE_STARTUP_RETRY_INTERVAL, configure_bgfilter_worker_outboxes, is_valid_env_key,
|
||||
parse_required_bgfilter_worker_capacity, protected_env_keys_from,
|
||||
should_initialize_editor_generation_pricing_for_startup,
|
||||
should_restore_auth_store_for_startup, should_start_profile_recharge_expiration_listener,
|
||||
strip_env_value,
|
||||
should_validate_bgfilter_internal_token_for_startup, strip_env_value,
|
||||
validate_bgfilter_internal_token_for_startup,
|
||||
};
|
||||
use crate::config::ProcessRole;
|
||||
use crate::config::{AppConfig, ProcessRole};
|
||||
|
||||
#[test]
|
||||
fn strip_env_value_removes_wrapping_quotes() {
|
||||
@@ -559,6 +721,51 @@ mod tests {
|
||||
assert_eq!(strip_env_value("plain\r"), "plain");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn bgfilter_worker_capacity_must_be_explicit_positive_and_bounded_by_q() {
|
||||
assert_eq!(
|
||||
parse_required_bgfilter_worker_capacity(Some("'16'"), Some("5000"), Some(" 128 "))
|
||||
.expect("valid explicit N/est/Q"),
|
||||
(16, 5_000, 128)
|
||||
);
|
||||
// Q 是可选保险丝:缺省时取默认值 2048,N 与 est 仍必须显式且为正。
|
||||
assert_eq!(
|
||||
parse_required_bgfilter_worker_capacity(Some("16"), Some("5000"), None)
|
||||
.expect("missing Q falls back to fuse default"),
|
||||
(16, 5_000, 2_048)
|
||||
);
|
||||
for (concurrency, estimate, max_requests) in [
|
||||
(None, Some("5000"), Some("128")),
|
||||
(Some("16"), None, Some("128")),
|
||||
(Some(""), Some("5000"), Some("128")),
|
||||
(Some("0"), Some("5000"), Some("128")),
|
||||
(Some("16"), Some("0"), Some("128")),
|
||||
(Some("four"), Some("5000"), Some("128")),
|
||||
(Some("8"), Some("5000"), Some("4")),
|
||||
(Some("8"), Some("5000"), Some("0")),
|
||||
] {
|
||||
assert!(
|
||||
parse_required_bgfilter_worker_capacity(concurrency, estimate, max_requests)
|
||||
.is_err(),
|
||||
"invalid N/est/Q should fail closed: N={concurrency:?}, est={estimate:?}, Q={max_requests:?}"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn bgfilter_worker_uses_its_own_tracking_outbox_directory() {
|
||||
let mut config = AppConfig::default();
|
||||
let base_dir = config.tracking_outbox_dir.clone();
|
||||
config.tracking_outbox_enabled = false;
|
||||
config.wallet_refund_outbox_enabled = true;
|
||||
|
||||
configure_bgfilter_worker_outboxes(&mut config);
|
||||
|
||||
assert!(config.tracking_outbox_enabled);
|
||||
assert_eq!(config.tracking_outbox_dir, base_dir.join("bgfilter-worker"));
|
||||
assert!(!config.wallet_refund_outbox_enabled);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn load_env_key_can_strip_utf8_bom_prefix() {
|
||||
let key = "\u{feff}SMS_AUTH_ENABLED"
|
||||
@@ -597,10 +804,39 @@ mod tests {
|
||||
assert_eq!(AUTH_STORE_STARTUP_RETRY_INTERVAL.as_secs(), 5);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn bgfilter_internal_token_startup_validation_is_limited_to_consumers() {
|
||||
for role in [
|
||||
ProcessRole::Api,
|
||||
ProcessRole::BgfilterWorker,
|
||||
ProcessRole::ExternalGenerationWorker,
|
||||
ProcessRole::All,
|
||||
] {
|
||||
assert!(should_validate_bgfilter_internal_token_for_startup(role));
|
||||
let mut config = AppConfig::default();
|
||||
config.process_role = role;
|
||||
config.bgfilter_internal_token = Some("invalid token".to_string());
|
||||
assert!(validate_bgfilter_internal_token_for_startup(&config).is_err());
|
||||
}
|
||||
assert!(!should_validate_bgfilter_internal_token_for_startup(
|
||||
ProcessRole::ExternalGenerationController
|
||||
));
|
||||
let mut controller_config = AppConfig::default();
|
||||
controller_config.process_role = ProcessRole::ExternalGenerationController;
|
||||
controller_config.bgfilter_internal_token = Some("unused invalid token".to_string());
|
||||
assert!(validate_bgfilter_internal_token_for_startup(&controller_config).is_ok());
|
||||
|
||||
let missing_config = AppConfig::default();
|
||||
assert!(validate_bgfilter_internal_token_for_startup(&missing_config).is_ok());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn auth_store_startup_restore_is_limited_to_http_roles() {
|
||||
assert!(should_restore_auth_store_for_startup(ProcessRole::Api));
|
||||
assert!(should_restore_auth_store_for_startup(ProcessRole::All));
|
||||
assert!(!should_restore_auth_store_for_startup(
|
||||
ProcessRole::BgfilterWorker
|
||||
));
|
||||
assert!(!should_restore_auth_store_for_startup(
|
||||
ProcessRole::ExternalGenerationWorker
|
||||
));
|
||||
@@ -617,6 +853,9 @@ mod tests {
|
||||
assert!(should_initialize_editor_generation_pricing_for_startup(
|
||||
ProcessRole::All
|
||||
));
|
||||
assert!(!should_initialize_editor_generation_pricing_for_startup(
|
||||
ProcessRole::BgfilterWorker
|
||||
));
|
||||
assert!(!should_initialize_editor_generation_pricing_for_startup(
|
||||
ProcessRole::ExternalGenerationWorker
|
||||
));
|
||||
@@ -633,6 +872,9 @@ mod tests {
|
||||
assert!(should_start_profile_recharge_expiration_listener(
|
||||
ProcessRole::All
|
||||
));
|
||||
assert!(!should_start_profile_recharge_expiration_listener(
|
||||
ProcessRole::BgfilterWorker
|
||||
));
|
||||
assert!(!should_start_profile_recharge_expiration_listener(
|
||||
ProcessRole::ExternalGenerationWorker
|
||||
));
|
||||
|
||||
@@ -51,6 +51,9 @@ 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;
|
||||
// P=8:父侧成功图片读取/解码槽。配 N=16 是内存与出口吞吐的折中,
|
||||
// 极端完整 body 内存按 (N + P) × 32 MiB 评估(见调度方案 §9.2)。
|
||||
pub(crate) const BGFILTER_IMAGE_VALIDATION_MAX_CONCURRENCY: usize = 8;
|
||||
|
||||
pub type HttpRequestPermitPool = Semaphore;
|
||||
|
||||
@@ -228,6 +231,9 @@ pub struct AppStateInner {
|
||||
#[allow(dead_code)]
|
||||
pub config: AppConfig,
|
||||
ready: AtomicBool,
|
||||
/// 本进程是否已至少一次连通 BgFilter worker(收到任意 HTTP 响应即算)。
|
||||
/// 连接失败重试用它区分冷启动窗口(开机/首连,宽限退避)与运行中途故障(快速收口)。
|
||||
bgfilter_worker_reached: AtomicBool,
|
||||
http_request_permit_pools: HttpRequestPermitPools,
|
||||
auth_jwt_config: JwtConfig,
|
||||
admin_runtime: Option<AdminRuntime>,
|
||||
@@ -261,7 +267,9 @@ pub struct AppStateInner {
|
||||
llm_client: Option<LlmClient>,
|
||||
editor_agent_llm_client: Option<LlmClient>,
|
||||
matting_client: Option<MattingClient>,
|
||||
editor_bgfilter_http_client: reqwest::Client,
|
||||
bgfilter_provider_http_client: reqwest::Client,
|
||||
bgfilter_worker_http_client: reqwest::Client,
|
||||
bgfilter_image_validation_limiter: Arc<Semaphore>,
|
||||
character_animation_oss_http_client: reqwest::Client,
|
||||
character_animation_oss_io_limiter: Arc<Semaphore>,
|
||||
#[cfg(any())]
|
||||
@@ -506,7 +514,18 @@ impl AppState {
|
||||
let llm_client = build_llm_client(&config)?;
|
||||
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 bgfilter_provider_http_client = build_bgfilter_provider_http_client(&config)?;
|
||||
let bgfilter_worker_http_client = build_bgfilter_worker_http_client(&config)?;
|
||||
let bgfilter_image_validation_concurrency = if config.process_role.runs_bgfilter_worker() {
|
||||
// 子 worker 已由 provider N 限流;图片校验槽与 N 对齐,避免引入第二个隐藏吞吐上限。
|
||||
config.bgfilter_worker_concurrency.max(1)
|
||||
} else {
|
||||
// 父 API / external-generation worker 固定限制解码并发,避免动画响应同时进入 blocking pool。
|
||||
BGFILTER_IMAGE_VALIDATION_MAX_CONCURRENCY
|
||||
};
|
||||
let bgfilter_image_validation_limiter = Arc::new(Semaphore::new(
|
||||
bgfilter_image_validation_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));
|
||||
@@ -516,6 +535,7 @@ impl AppState {
|
||||
Ok(Self(Arc::new(AppStateInner {
|
||||
config,
|
||||
ready: AtomicBool::new(true),
|
||||
bgfilter_worker_reached: AtomicBool::new(false),
|
||||
http_request_permit_pools,
|
||||
auth_jwt_config,
|
||||
admin_runtime,
|
||||
@@ -551,7 +571,9 @@ impl AppState {
|
||||
llm_client,
|
||||
editor_agent_llm_client,
|
||||
matting_client,
|
||||
editor_bgfilter_http_client,
|
||||
bgfilter_provider_http_client,
|
||||
bgfilter_worker_http_client,
|
||||
bgfilter_image_validation_limiter,
|
||||
character_animation_oss_http_client,
|
||||
character_animation_oss_io_limiter,
|
||||
#[cfg(any())]
|
||||
@@ -1266,8 +1288,24 @@ impl AppState {
|
||||
self.matting_client.as_ref()
|
||||
}
|
||||
|
||||
pub fn editor_bgfilter_http_client(&self) -> &reqwest::Client {
|
||||
&self.editor_bgfilter_http_client
|
||||
pub fn bgfilter_provider_http_client(&self) -> &reqwest::Client {
|
||||
&self.bgfilter_provider_http_client
|
||||
}
|
||||
|
||||
pub fn bgfilter_worker_http_client(&self) -> &reqwest::Client {
|
||||
&self.bgfilter_worker_http_client
|
||||
}
|
||||
|
||||
pub fn bgfilter_worker_reached(&self) -> bool {
|
||||
self.bgfilter_worker_reached.load(Ordering::Relaxed)
|
||||
}
|
||||
|
||||
pub fn mark_bgfilter_worker_reached(&self) {
|
||||
self.bgfilter_worker_reached.store(true, Ordering::Relaxed);
|
||||
}
|
||||
|
||||
pub fn bgfilter_image_validation_limiter(&self) -> Arc<Semaphore> {
|
||||
self.bgfilter_image_validation_limiter.clone()
|
||||
}
|
||||
|
||||
pub fn character_animation_oss_http_client(&self) -> &reqwest::Client {
|
||||
@@ -1992,12 +2030,13 @@ fn build_matting_client(config: &AppConfig) -> Result<Option<MattingClient>, App
|
||||
.map_err(|error| AppStateInitError::DependencyUnavailable(error.to_string()))
|
||||
}
|
||||
|
||||
fn build_editor_bgfilter_http_client(
|
||||
fn build_bgfilter_provider_http_client(
|
||||
config: &AppConfig,
|
||||
) -> Result<reqwest::Client, AppStateInitError> {
|
||||
reqwest::Client::builder()
|
||||
.redirect(reqwest::redirect::Policy::none())
|
||||
.timeout(std::time::Duration::from_millis(
|
||||
config.editor_bgfilter_request_timeout_ms.max(1),
|
||||
config.bgfilter_provider_attempt_timeout_ms().max(1),
|
||||
))
|
||||
.connect_timeout(std::time::Duration::from_secs(30))
|
||||
.pool_idle_timeout(std::time::Duration::from_secs(300))
|
||||
@@ -2006,7 +2045,26 @@ fn build_editor_bgfilter_http_client(
|
||||
.build()
|
||||
.map_err(|error| {
|
||||
AppStateInitError::DependencyUnavailable(format!(
|
||||
"构建共享 BgFilter HTTP 客户端失败:{error}"
|
||||
"构建 BgFilter provider HTTP 客户端失败:{error}"
|
||||
))
|
||||
})
|
||||
}
|
||||
|
||||
fn build_bgfilter_worker_http_client(
|
||||
config: &AppConfig,
|
||||
) -> Result<reqwest::Client, AppStateInitError> {
|
||||
reqwest::Client::builder()
|
||||
.redirect(reqwest::redirect::Policy::none())
|
||||
.connect_timeout(std::time::Duration::from_millis(
|
||||
config.bgfilter_worker_connect_timeout_ms.max(1),
|
||||
))
|
||||
.pool_idle_timeout(std::time::Duration::from_secs(300))
|
||||
.pool_max_idle_per_host(128)
|
||||
.tcp_keepalive(std::time::Duration::from_secs(60))
|
||||
.build()
|
||||
.map_err(|error| {
|
||||
AppStateInitError::DependencyUnavailable(format!(
|
||||
"构建 BgFilter 内部 worker HTTP 客户端失败:{error}"
|
||||
))
|
||||
})
|
||||
}
|
||||
@@ -2173,6 +2231,30 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn bgfilter_image_validation_limiter_is_bounded_per_process_role() {
|
||||
let parent = AppState::new(AppConfig::default()).expect("parent state should build");
|
||||
assert_eq!(
|
||||
parent
|
||||
.bgfilter_image_validation_limiter()
|
||||
.available_permits(),
|
||||
BGFILTER_IMAGE_VALIDATION_MAX_CONCURRENCY
|
||||
);
|
||||
|
||||
let worker = AppState::new(AppConfig {
|
||||
process_role: crate::config::ProcessRole::BgfilterWorker,
|
||||
bgfilter_worker_concurrency: 6,
|
||||
..AppConfig::default()
|
||||
})
|
||||
.expect("worker state should build");
|
||||
assert_eq!(
|
||||
worker
|
||||
.bgfilter_image_validation_limiter()
|
||||
.available_permits(),
|
||||
6
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn editor_generation_pricing_typed_record_round_trips() {
|
||||
let expected = crate::editor_generation_config::parse_editor_generation_pricing_json(
|
||||
|
||||
@@ -180,6 +180,16 @@ pub(crate) fn record_external_api_failure(
|
||||
);
|
||||
}
|
||||
|
||||
pub(crate) fn record_external_api_audit_dropped(provider: &'static str, reason: &'static str) {
|
||||
external_api_metrics().audit_dropped.add(
|
||||
1,
|
||||
&[
|
||||
KeyValue::new("provider", provider),
|
||||
KeyValue::new("reason", reason),
|
||||
],
|
||||
);
|
||||
}
|
||||
|
||||
fn track_response_body_in_flight(response: Response<Body>) -> Response<Body> {
|
||||
response.map(|body| {
|
||||
HTTP_RESPONSE_BODY_IN_FLIGHT.fetch_add(1, Ordering::Relaxed);
|
||||
@@ -219,6 +229,7 @@ struct TrackingOutboxMetrics {
|
||||
|
||||
struct ExternalApiMetrics {
|
||||
failures: Counter<u64>,
|
||||
audit_dropped: Counter<u64>,
|
||||
}
|
||||
|
||||
struct HttpRequestPermitsAvailableGauges {
|
||||
@@ -363,6 +374,12 @@ fn external_api_metrics() -> &'static ExternalApiMetrics {
|
||||
"External API call failures grouped by provider and failure stage",
|
||||
)
|
||||
.build(),
|
||||
audit_dropped: meter
|
||||
.u64_counter("genarrative.external_api.audit.dropped")
|
||||
.with_description(
|
||||
"External API failure audit records dropped when synchronous fallback is disabled",
|
||||
)
|
||||
.build(),
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user