BgFilter 父侧连接失败有界重试,跨过 worker 重启窗口
主机重启 / worker 崩溃拉起期间,父侧对 127.0.0.1 worker 的 TCP 连接 失败此前直接映射 internal_error:complex(max_attempts=1)终态失败不可 自愈,flat 被迫降级。现仅对连接从未建立的失败(无副作用、天然幂等)做 有界退避重试:每轮按公式重算 maxQueueWaitMs,不突破「预算不足不发送」 不变量;complex 用完整序列(约 22.5s,覆盖 RestartSec=5s + 启动窗), flat 只取前 2 项(≤1.5s,不侵蚀 39s fallback 预留)。收到任何 HTTP 响应(含 5xx)立即停止重试;新增 bgfilter_internal_connect_retry_total 指标。设计文档与数据契约同步刻出「已建立连接后不重试」禁令的精确例外。 Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
@@ -227,6 +227,8 @@ parent client timeout = maxQueueWaitMs + callBudgetMs + 2s 传输窗
|
||||
|
||||
`maxQueueWaitMs <= 0` 时父侧不得发送请求:flat 直接进入既有“阿里云 → 本地”fallback,complex 直接失败;不允许把注定超时的请求塞进队列。flat 扣除的 `39s` 父侧预留由 `37s` fallback 窗口和 `2s` 内部响应传输窗组成;complex 没有 fallback,只保留 `2s` 传输窗。
|
||||
|
||||
父侧对「TCP 连接从未建立」的失败(worker 重启、主机开机排序窗口内的连接拒绝 / 不可达 / connect 阶段超时)做有界自动重试:这类请求从未进入 worker admission,无副作用、天然幂等。每轮重试前按上式重算 `maxQueueWaitMs`,重试消耗的是父预算的自然余量,不突破「预算不足不发送」的不变量;退避序列本身有界——complex 用完整序列(总额约 `22.5s`,按上界覆盖 systemd `RestartSec=5s` + 进程启动窗口,更长的停机应快速失败而非挂住父 job),flat 只取前 2 项(额外延迟 `≤1.5s`,不侵蚀 `39s` fallback 预留)——父无绝对 deadline 时也不会无限等待。收到任何 HTTP 响应(含 5xx)或其它错误类别一律不重试,边界见 §7.1。
|
||||
|
||||
queue job 的总预算从父 job 开始执行时起算,不从开始申请 BgFilter 时重新计时。现有父 worker 还会把 provider deadline 设在 job deadline 前 `60s`,为最终写回和终态保留时间。同步 RPC 实现必须显式读取父侧剩余 provider budget 派生 `maxQueueWaitMs`,不能忽略 `RequestContext` deadline。
|
||||
|
||||
本版没有 BgFilter 子任务等待 claim 的阶段。几个起算点必须区分:父 job 在数据库中尚未被 claim 的等待不消耗 job 执行预算;父 job 开始实际执行后,生图及 BgFilter 之前的耗时都会消耗父总预算;父内部 HTTP client timeout 从开始发送请求起覆盖 loopback 传输、worker admission、排队、provider 和回包;`callBudgetMs` 计时只在子 worker 取得 provider permit 后启动,attempt timer 只在真正开始一次 provider HTTP 时启动。父侧不是放弃超时,而是不再直接执行 provider 单次 attempt 的计时器。
|
||||
@@ -335,7 +337,7 @@ flat / complex 统一使用的 `GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_FAILURE_THRE
|
||||
|
||||
### 7.1 内部 RPC 断连
|
||||
|
||||
- 父侧不得自动重试整次内部 HTTP。断连时结果未知,重试可能让一次逻辑调用从最多两次 provider attempt 扩大为四次,并可能突破瞬时并发预期。
|
||||
- 连接已建立后的断连不得自动重试整次内部 HTTP:此时结果未知,重试可能让一次逻辑调用从最多两次 provider attempt 扩大为四次,并可能突破瞬时并发预期。唯一例外是 TCP 连接从未建立的失败(连接拒绝 / 不可达 / connect 阶段超时):请求从未进入 worker admission,结果确定为「未发生」,父侧按 §5.1 的预算约束有界退避重试,用于跨过 worker 重启与主机开机排序窗口;收到任何 HTTP 响应后即回到本条禁令。
|
||||
- 尚在等待 permit 的请求到达自身 deadline 后必须取消,不再发送 provider 请求;运行时若能可靠观察客户端断连,也可提前取消,但正确性不能只依赖断连事件。
|
||||
- 已经开始的 provider attempt 必须继续读取到完成或本次 attempt timeout,并持有 permit;可观察到的 handler / client drop 只丢弃最终结果,不能让已启动请求变成无人管理的本地 future。
|
||||
- deadline 已被子 worker 观察到后,不再开始第二次 attempt。首版没有显式 cancellation signal 通道;单纯 TCP 断连只能 best-effort 阻止二试(handler future 被 drop 后自然不再开始新 attempt),Axum / Hyper 不保证立刻通知 handler,因此不能承诺所有断连都阻止二试,排队阶段的 queue timeout 与 permit 后的 `callBudgetMs` deadline 是最终可靠的停止条件。
|
||||
@@ -345,7 +347,7 @@ flat / complex 统一使用的 `GENARRATIVE_EDITOR_BGFILTER_CIRCUIT_FAILURE_THRE
|
||||
### 7.2 进程崩溃
|
||||
|
||||
- 父进程崩溃:内部连接最终断开,父 job 按现有 heartbeat、lease、`max_attempts = 1`、失败和退款语义收口。
|
||||
- 子 worker 崩溃或重启:当前内部 RPC 失败;父侧不查询、不恢复、不重发同一次 RPC。
|
||||
- 子 worker 崩溃或重启:已在途的内部 RPC 失败;父侧不查询、不恢复、不重发同一次 RPC。重启窗口内连接从未建立的新调用按 §7.1 的例外有界重试。
|
||||
- 子 worker 成功但响应在网络中丢失:结果视为未知;flat 进入原 fallback,complex 失败。
|
||||
- 子 worker 不得反向 complete / fail 父 job,也不得写画布、业务资源或账单。
|
||||
|
||||
@@ -406,6 +408,7 @@ BgFilter 成功二进制不是一份新的业务资产:
|
||||
- `bgfilter_internal_waiting_requests`(即 admission 后等待 `N` permit 的队长)
|
||||
- `bgfilter_internal_queue_timeout_total{bound}`(排队超时按触发边界 `estimate | parent` 分维度)
|
||||
- `bgfilter_internal_call_budget_drift_total{mode}`(请求 `callBudgetMs` 与本进程公式值不一致;发布窗口内短暂非零正常,持续增长说明父子 N / est 真漂移)
|
||||
- `bgfilter_internal_connect_retry_total{mode}`(父侧连接失败重试次数;worker 重启窗口内短暂非零正常,持续增长说明 worker 长期不可达)
|
||||
- `bgfilter_internal_in_flight`
|
||||
- `bgfilter_internal_request_seconds{mode,outcome}`
|
||||
- `bgfilter_provider_http_seconds{mode,attempt,outcome}`
|
||||
@@ -466,7 +469,7 @@ flat / complex 的每次 provider 失败审计都必须留在子 worker,保留
|
||||
- `maxQueueWaitMs <= 0` 时父侧不发送请求:flat 直接 fallback,complex 直接失败。
|
||||
- `callBudgetMs` 与 worker 本进程公式值不一致时不拒绝:worker 以自身公式值执行,记 warn 并递增漂移指标;发布调优 N / est 的新旧进程共存窗口内,在途 flat 任务仍能正常执行或走既有 fallback,不得因瞬态漂移触发 `invalid_request`(该码禁止 fallback)。attempt、callBudget、client timeout 全部由 `N / est` 运行时派生,代码不存在硬编码结果值。
|
||||
- parent client timeout 精确取 `maxQueueWaitMs + callBudgetMs + 2s`,helper 保持 infallible;父绝对预算通过 `maxQueueWaitMs` 的派生公式预先约束,结果校验等待也必须 deadline-aware,不能只在校验完成后事后判超时。
|
||||
- 第一次失败后预算不足时不开始第二次;父侧从不重试整次内部 RPC。
|
||||
- 第一次失败后预算不足时不开始第二次;父侧从不重试已建立连接的内部 RPC。连接从未建立的失败按 §5.1 有界退避重试:仅 connect 类失败重入、收到任何 HTTP 响应立即停止、每轮重算 `maxQueueWaitMs`、complex / flat 各自的退避配额封顶(Rust 集成测试覆盖「重试跨过监听空窗后停在首个 HTTP 响应」「配额耗尽返回 connect 失败」「deadline 放不下下一轮时不空睡」三条路径)。
|
||||
- 父业务预算仍有效时,flat 两次失败、熔断、overload、内部 RPC deadline 或断连仍走“阿里云 → 本地”;complex 任意失败或自身熔断都直接失败,不接 flat fallback。
|
||||
- flat / complex 分别按自身真实失败 attempt 计数且状态互不影响;由剩余业务预算截短的 timeout 不计入。两种模式都在 permit 前二次检查;已获准调用可完成第二次,后续同模式排队请求快速 `circuit_open`。
|
||||
- `cancelled`(仅验证父侧映射,保留码首版不产生)、父 cancellation / 绝对 deadline、`invalid_request` 和 `unauthorized` 不启动 flat fallback;其它 flat 错误只在父业务预算仍有效时进入 fallback。
|
||||
|
||||
File diff suppressed because one or more lines are too long
@@ -51,6 +51,13 @@ const BGFILTER_PROVIDER_ATTEMPT_RESERVE: Duration = Duration::from_secs(1);
|
||||
const BGFILTER_INTERNAL_CLIENT_RESPONSE_RESERVE: Duration = Duration::from_secs(2);
|
||||
const BGFILTER_SOURCE_URL_EXPIRE_SECONDS: u64 = 600;
|
||||
const BGFILTER_PROVIDER_TOKEN_HEADER: &str = "X-Genarrative-Image-Token";
|
||||
/// 连接失败(TCP 从未建立)的重试退避序列。总长约 22.5s,按上界覆盖
|
||||
/// worker 的 systemd 自动拉起窗口(RestartSec=5s + 进程启动)与主机重启的
|
||||
/// 并发启动窗口;更长的 worker 故障属于真实停机,应当快速失败而不是挂住父 job。
|
||||
const BGFILTER_CONNECT_RETRY_BACKOFF_MS: [u64; 7] = [500, 1_000, 2_000, 4_000, 5_000, 5_000, 5_000];
|
||||
/// flat 有完整 fallback 链,重试只为跨过最短窗口,不得侵蚀 fallback 预留:
|
||||
/// 只取退避序列前 2 项(额外延迟 ≤1.5s)。
|
||||
const BGFILTER_FLAT_CONNECT_RETRY_LIMIT: usize = 2;
|
||||
|
||||
static BGFILTER_FLAT_CIRCUIT: OnceLock<Mutex<BgfilterCircuitState>> = OnceLock::new();
|
||||
static BGFILTER_COMPLEX_CIRCUIT: OnceLock<Mutex<BgfilterCircuitState>> = OnceLock::new();
|
||||
@@ -60,6 +67,7 @@ struct BgfilterMetrics {
|
||||
waiting_requests: UpDownCounter<i64>,
|
||||
queue_timeout_total: Counter<u64>,
|
||||
call_budget_drift_total: Counter<u64>,
|
||||
connect_retry_total: Counter<u64>,
|
||||
in_flight: UpDownCounter<i64>,
|
||||
internal_request_seconds: Histogram<f64>,
|
||||
provider_http_seconds: Histogram<f64>,
|
||||
@@ -110,6 +118,13 @@ fn bgfilter_metrics() -> &'static BgfilterMetrics {
|
||||
"Requests whose callBudgetMs fingerprint mismatched the worker formula; transient during deploys, sustained growth means real N/est drift",
|
||||
)
|
||||
.build(),
|
||||
connect_retry_total: meter
|
||||
.u64_counter("bgfilter_internal_connect_retry_total")
|
||||
.with_unit("{attempt}")
|
||||
.with_description(
|
||||
"Parent-side retries after the worker TCP connect failed; transient spikes cover worker restarts, sustained growth means the worker is down",
|
||||
)
|
||||
.build(),
|
||||
in_flight: meter
|
||||
.i64_up_down_counter("bgfilter_internal_in_flight")
|
||||
.with_unit("{request}")
|
||||
@@ -171,6 +186,9 @@ pub(crate) struct BgfilterClientError {
|
||||
status: StatusCode,
|
||||
timeout: bool,
|
||||
transport: bool,
|
||||
/// TCP 连接从未建立(拒绝/不可达/connect 阶段失败):请求未进入 worker
|
||||
/// 准入队列,无副作用,是唯一可以安全自动重试的失败类别。
|
||||
connect: bool,
|
||||
}
|
||||
|
||||
impl BgfilterClientError {
|
||||
@@ -184,6 +202,10 @@ impl BgfilterClientError {
|
||||
!matches!(self.code, "invalid_request" | "unauthorized" | "cancelled")
|
||||
}
|
||||
|
||||
fn is_connect_failure(&self) -> bool {
|
||||
self.connect
|
||||
}
|
||||
|
||||
pub(crate) fn into_app_error(self) -> AppError {
|
||||
let status = match self.code {
|
||||
"invalid_request" => StatusCode::BAD_REQUEST,
|
||||
@@ -210,6 +232,7 @@ impl BgfilterClientError {
|
||||
status: StatusCode::BAD_GATEWAY,
|
||||
timeout: code == "deadline_exceeded",
|
||||
transport: false,
|
||||
connect: false,
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -2065,6 +2088,95 @@ fn worker_image_response(image: BgfilterImage, request_id: &str) -> Response {
|
||||
response
|
||||
}
|
||||
|
||||
/// 决定第 `attempt` 次连接失败后是否还允许重试及退避时长;`None` 表示配额用尽。
|
||||
fn connect_retry_backoff(mode: BgfilterBackgroundMode, attempt: usize) -> Option<Duration> {
|
||||
let limit = match mode {
|
||||
BgfilterBackgroundMode::Flat => BGFILTER_FLAT_CONNECT_RETRY_LIMIT,
|
||||
BgfilterBackgroundMode::Complex => BGFILTER_CONNECT_RETRY_BACKOFF_MS.len(),
|
||||
};
|
||||
if attempt >= limit {
|
||||
return None;
|
||||
}
|
||||
BGFILTER_CONNECT_RETRY_BACKOFF_MS
|
||||
.get(attempt)
|
||||
.copied()
|
||||
.map(Duration::from_millis)
|
||||
}
|
||||
|
||||
/// 围绕 `request_bgfilter_worker` 的连接失败有界重试:覆盖 worker 重启与主机
|
||||
/// 开机排序窗口。只重试 TCP 连接从未建立的失败(请求未进入 worker 准入队列,
|
||||
/// 无副作用、天然幂等);收到任何 HTTP 响应(含 5xx)或其它错误类别都按原语义
|
||||
/// 立刻返回。每轮从父剩余预算按现有公式重算 `maxQueueWaitMs`,重试消耗的是
|
||||
/// 预算自然余量,不突破「预算不足不发送」的协议不变量;退避序列本身有界,
|
||||
/// 父无绝对 deadline 时也不会无限等待。
|
||||
pub(crate) async fn request_bgfilter_worker_with_connect_retry(
|
||||
state: &AppState,
|
||||
source_object_key: &str,
|
||||
background_mode: BgfilterBackgroundMode,
|
||||
screen_color: Option<&str>,
|
||||
seg_model: &str,
|
||||
cross_check: bool,
|
||||
deadline_reserve: Duration,
|
||||
audit: &ExternalApiAuditContext,
|
||||
) -> Result<BgfilterImage, BgfilterClientError> {
|
||||
let call_budget_ms = state.config.bgfilter_call_budget_ms();
|
||||
let mut connect_failures = 0usize;
|
||||
loop {
|
||||
let max_queue_wait_ms = max_queue_wait_ms(
|
||||
call_budget_ms,
|
||||
audit.external_call_deadline,
|
||||
deadline_reserve,
|
||||
);
|
||||
let error = match request_bgfilter_worker(
|
||||
state,
|
||||
source_object_key,
|
||||
background_mode,
|
||||
screen_color,
|
||||
seg_model,
|
||||
cross_check,
|
||||
max_queue_wait_ms,
|
||||
audit,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(image) => return Ok(image),
|
||||
Err(error) => error,
|
||||
};
|
||||
if !error.is_connect_failure() {
|
||||
return Err(error);
|
||||
}
|
||||
let Some(backoff) = connect_retry_backoff(background_mode, connect_failures) else {
|
||||
return Err(error);
|
||||
};
|
||||
if let Some(deadline) = audit.external_call_deadline {
|
||||
// 睡完还得放得下 reserve + 一次完整调用,否则下一轮 maxQueueWait 必为 0,
|
||||
// 与其空睡不如立刻按当前错误返回。
|
||||
let next_round_floor = backoff
|
||||
.saturating_add(deadline_reserve)
|
||||
.saturating_add(Duration::from_millis(call_budget_ms));
|
||||
if deadline
|
||||
.checked_duration_since(Instant::now())
|
||||
.unwrap_or(Duration::ZERO)
|
||||
< next_round_floor
|
||||
{
|
||||
return Err(error);
|
||||
}
|
||||
}
|
||||
connect_failures += 1;
|
||||
tracing::warn!(
|
||||
background_mode = background_mode.as_str(),
|
||||
connect_failures,
|
||||
backoff_ms = duration_millis(backoff),
|
||||
error = %error.message,
|
||||
"bgfilter_worker_connect_retry"
|
||||
);
|
||||
bgfilter_metrics()
|
||||
.connect_retry_total
|
||||
.add(1, &[KeyValue::new("mode", background_mode.as_str())]);
|
||||
tokio::time::sleep(backoff).await;
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) async fn request_bgfilter_worker(
|
||||
state: &AppState,
|
||||
source_object_key: &str,
|
||||
@@ -2128,6 +2240,7 @@ pub(crate) async fn request_bgfilter_worker(
|
||||
status: StatusCode::BAD_GATEWAY,
|
||||
timeout: error.is_timeout(),
|
||||
transport: true,
|
||||
connect: error.is_connect(),
|
||||
})?;
|
||||
let status = response.status();
|
||||
if !status.is_success() {
|
||||
@@ -2175,6 +2288,7 @@ pub(crate) async fn request_bgfilter_worker(
|
||||
status: StatusCode::BAD_GATEWAY,
|
||||
timeout,
|
||||
transport: true,
|
||||
connect: false,
|
||||
});
|
||||
}
|
||||
Err(_) => {
|
||||
@@ -2227,6 +2341,7 @@ pub(crate) async fn request_bgfilter_worker(
|
||||
status: StatusCode::BAD_GATEWAY,
|
||||
timeout: false,
|
||||
transport: false,
|
||||
connect: false,
|
||||
})
|
||||
.map(|image| {
|
||||
tracing::info!(
|
||||
@@ -2345,6 +2460,7 @@ async fn read_worker_error(
|
||||
status: StatusCode::from_u16(status.as_u16()).unwrap_or(StatusCode::BAD_GATEWAY),
|
||||
timeout: code == "deadline_exceeded",
|
||||
transport: false,
|
||||
connect: false,
|
||||
})
|
||||
}
|
||||
|
||||
@@ -3152,6 +3268,10 @@ mod tests {
|
||||
.0;
|
||||
assert_eq!(parent_client.matches(".send()").count(), 1);
|
||||
assert!(!parent_client.contains("for attempt"));
|
||||
// 连接失败重试围绕唯一的 .send() 调用点循环:仅 TCP 从未建立的失败重入,
|
||||
// 收到任何 HTTP 响应都按原语义立刻返回。
|
||||
assert!(parent_client.contains("error.is_connect_failure()"));
|
||||
assert!(parent_client.contains("connect: error.is_connect()"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -3212,4 +3332,207 @@ mod tests {
|
||||
);
|
||||
assert_eq!(stable_error_code("unknown"), None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn connect_retry_backoff_is_bounded_per_mode() {
|
||||
// flat 只取前 2 项:额外延迟 ≤1.5s,不侵蚀 39s fallback 预留。
|
||||
assert_eq!(
|
||||
connect_retry_backoff(BgfilterBackgroundMode::Flat, 0),
|
||||
Some(Duration::from_millis(500))
|
||||
);
|
||||
assert_eq!(
|
||||
connect_retry_backoff(BgfilterBackgroundMode::Flat, 1),
|
||||
Some(Duration::from_millis(1_000))
|
||||
);
|
||||
assert_eq!(connect_retry_backoff(BgfilterBackgroundMode::Flat, 2), None);
|
||||
// complex 用完整退避序列(总额 22.5s),覆盖 RestartSec=5s + 进程启动窗口;
|
||||
// 序列有界,父无绝对 deadline 时也不会无限重试。
|
||||
let total: u64 = (0..BGFILTER_CONNECT_RETRY_BACKOFF_MS.len())
|
||||
.map(|attempt| {
|
||||
connect_retry_backoff(BgfilterBackgroundMode::Complex, attempt)
|
||||
.expect("complex 应可用完整退避序列")
|
||||
.as_millis() as u64
|
||||
})
|
||||
.sum();
|
||||
assert_eq!(total, 22_500);
|
||||
assert_eq!(
|
||||
connect_retry_backoff(
|
||||
BgfilterBackgroundMode::Complex,
|
||||
BGFILTER_CONNECT_RETRY_BACKOFF_MS.len()
|
||||
),
|
||||
None
|
||||
);
|
||||
}
|
||||
|
||||
fn parent_client_state(port: u16) -> AppState {
|
||||
let mut config = AppConfig::default();
|
||||
config.bgfilter_internal_token = Some("parent-test-token".to_string());
|
||||
config.bgfilter_worker_base_url = format!("http://127.0.0.1:{port}");
|
||||
// Windows 对已关闭的 loopback 端口不回 RST 而是挂到 connect timeout;
|
||||
// 调小超时让「连接失败」在两个平台上都以可控节奏出现。
|
||||
config.bgfilter_worker_connect_timeout_ms = 100;
|
||||
AppState::new(config).expect("test state should build")
|
||||
}
|
||||
|
||||
fn parent_audit(deadline: Option<Instant>) -> ExternalApiAuditContext {
|
||||
ExternalApiAuditContext {
|
||||
user_id: None,
|
||||
profile_id: None,
|
||||
request_id: None,
|
||||
external_call_deadline: deadline,
|
||||
}
|
||||
}
|
||||
|
||||
async fn reserved_loopback_port() -> u16 {
|
||||
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
|
||||
.await
|
||||
.expect("bind ephemeral port");
|
||||
let port = listener.local_addr().expect("local addr").port();
|
||||
drop(listener);
|
||||
port
|
||||
}
|
||||
|
||||
fn find_subslice(haystack: &[u8], needle: &[u8]) -> Option<usize> {
|
||||
haystack
|
||||
.windows(needle.len())
|
||||
.position(|window| window == needle)
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn connect_retry_exhausts_flat_quota_then_returns_connect_error() {
|
||||
let port = reserved_loopback_port().await;
|
||||
let state = parent_client_state(port);
|
||||
let audit = parent_audit(None);
|
||||
let started = Instant::now();
|
||||
let error = match request_bgfilter_worker_with_connect_retry(
|
||||
&state,
|
||||
"editor/test-object.png",
|
||||
BgfilterBackgroundMode::Flat,
|
||||
Some("#00FF00"),
|
||||
"test-model",
|
||||
false,
|
||||
Duration::ZERO,
|
||||
&audit,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(_) => panic!("端口关闭时应失败"),
|
||||
Err(error) => error,
|
||||
};
|
||||
assert!(error.is_connect_failure());
|
||||
// Linux 即时拒绝 → internal_error;Windows 挂到 connect timeout → deadline_exceeded。
|
||||
// 重试判定只看 is_connect_failure,code 允许两种平台形态。
|
||||
assert!(matches!(
|
||||
error.code(),
|
||||
"internal_error" | "deadline_exceeded"
|
||||
));
|
||||
// flat 配额 2 次重试对应 500ms + 1000ms 退避:证明确实退避过而不是立即放弃。
|
||||
assert!(started.elapsed() >= Duration::from_millis(1_400));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn connect_retry_crosses_worker_restart_window_and_stops_on_http_response() {
|
||||
use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _};
|
||||
|
||||
let port = reserved_loopback_port().await;
|
||||
let state = parent_client_state(port);
|
||||
let audit = parent_audit(None);
|
||||
let server = tokio::spawn(async move {
|
||||
// 模拟 worker 重启窗口:先保持端口关闭,约 700ms 后才开始监听。
|
||||
tokio::time::sleep(Duration::from_millis(700)).await;
|
||||
let listener = tokio::net::TcpListener::bind(("127.0.0.1", port))
|
||||
.await
|
||||
.expect("rebind test port");
|
||||
let (mut socket, _) = listener.accept().await.expect("accept parent retry");
|
||||
let mut request = Vec::new();
|
||||
let mut buffer = [0u8; 4096];
|
||||
loop {
|
||||
let read = socket.read(&mut buffer).await.expect("read request");
|
||||
if read == 0 {
|
||||
break;
|
||||
}
|
||||
request.extend_from_slice(&buffer[..read]);
|
||||
if let Some(headers_end) = find_subslice(&request, b"\r\n\r\n") {
|
||||
let headers = String::from_utf8_lossy(&request[..headers_end]);
|
||||
let content_length = headers
|
||||
.lines()
|
||||
.find_map(|line| {
|
||||
let (name, value) = line.split_once(':')?;
|
||||
if name.eq_ignore_ascii_case("content-length") {
|
||||
value.trim().parse::<usize>().ok()
|
||||
} else {
|
||||
None
|
||||
}
|
||||
})
|
||||
.unwrap_or(0);
|
||||
if request.len() >= headers_end + 4 + content_length {
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
let body = br#"{"error":{"code":"overloaded","message":"test","retryable":true}}"#;
|
||||
let response = format!(
|
||||
"HTTP/1.1 429 Too Many Requests\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n",
|
||||
body.len()
|
||||
);
|
||||
socket
|
||||
.write_all(response.as_bytes())
|
||||
.await
|
||||
.expect("write headers");
|
||||
socket.write_all(body).await.expect("write body");
|
||||
socket.shutdown().await.ok();
|
||||
});
|
||||
let started = Instant::now();
|
||||
let error = match request_bgfilter_worker_with_connect_retry(
|
||||
&state,
|
||||
"editor/test-object.png",
|
||||
BgfilterBackgroundMode::Complex,
|
||||
None,
|
||||
"test-model",
|
||||
false,
|
||||
Duration::ZERO,
|
||||
&audit,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(_) => panic!("mock 返回 429 时应失败"),
|
||||
Err(error) => error,
|
||||
};
|
||||
server.await.expect("mock server should finish");
|
||||
// 收到 HTTP 响应(429 → overloaded)即停止重试:错误按原语义返回,不再是 connect 失败。
|
||||
assert_eq!(error.code(), "overloaded");
|
||||
assert!(!error.is_connect_failure());
|
||||
// 至少经历一次退避才可能跨过 700ms 的监听空窗。
|
||||
assert!(started.elapsed() >= Duration::from_millis(700));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn connect_retry_stops_without_sleeping_when_deadline_cannot_fit_next_round() {
|
||||
let port = reserved_loopback_port().await;
|
||||
let state = parent_client_state(port);
|
||||
let call_budget = Duration::from_millis(state.config.bgfilter_call_budget_ms());
|
||||
// 父剩余刚好放得下第一次调用(约 300ms 排队额度),放不下「退避 + 再一次完整调用」。
|
||||
let audit = parent_audit(Some(
|
||||
Instant::now() + call_budget + Duration::from_millis(300),
|
||||
));
|
||||
let started = Instant::now();
|
||||
let error = match request_bgfilter_worker_with_connect_retry(
|
||||
&state,
|
||||
"editor/test-object.png",
|
||||
BgfilterBackgroundMode::Complex,
|
||||
None,
|
||||
"test-model",
|
||||
false,
|
||||
Duration::ZERO,
|
||||
&audit,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(_) => panic!("端口关闭时应失败"),
|
||||
Err(error) => error,
|
||||
};
|
||||
assert!(error.is_connect_failure());
|
||||
// 第一轮连接失败(≤100ms connect timeout)后应立即返回,不做 500ms 空睡。
|
||||
assert!(started.elapsed() < Duration::from_millis(450));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -3386,20 +3386,16 @@ async fn request_editor_background_removal_image_with_bgfilter_worker(
|
||||
source_object_key: &str,
|
||||
audit: &crate::external_api_audit::ExternalApiAuditContext,
|
||||
) -> Result<EditorBackgroundRemovalImage, AppError> {
|
||||
// complex 没有 flat fallback,排队上限只预留 2s 传输窗。
|
||||
let max_queue_wait_ms = crate::bgfilter_worker::max_queue_wait_ms(
|
||||
state.config.bgfilter_call_budget_ms(),
|
||||
audit.external_call_deadline,
|
||||
EDITOR_BGFILTER_PARENT_TRANSPORT_WINDOW,
|
||||
);
|
||||
let removed = crate::bgfilter_worker::request_bgfilter_worker(
|
||||
// complex 没有 flat fallback,排队上限只预留 2s 传输窗;worker 重启窗口内的
|
||||
// 连接失败由 client 内部按预算有界重试,避免 max_attempts=1 的队列任务终态失败。
|
||||
let removed = crate::bgfilter_worker::request_bgfilter_worker_with_connect_retry(
|
||||
state,
|
||||
source_object_key,
|
||||
crate::bgfilter_worker::BgfilterBackgroundMode::Complex,
|
||||
None,
|
||||
EDITOR_BGFILTER_DEFAULT_SEG_MODEL,
|
||||
EDITOR_BGFILTER_CROSS_CHECK_DISABLED,
|
||||
max_queue_wait_ms,
|
||||
EDITOR_BGFILTER_PARENT_TRANSPORT_WINDOW,
|
||||
audit,
|
||||
)
|
||||
.await
|
||||
@@ -3431,20 +3427,16 @@ pub(crate) async fn remove_editor_generated_screen_background_with_bgfilter(
|
||||
) -> Result<EditorScreenBackgroundRemovalOutput, AppError> {
|
||||
// flat 的排队上限从父剩余预算里预扣「阿里云 + 本地 + 传输」的 fallback 时间,
|
||||
// 保证排队吃不完 fallback 的执行窗口;上限为 0 时 client 直接返回
|
||||
// deadline_exceeded,本函数随即进入 fallback 链。
|
||||
let max_queue_wait_ms = crate::bgfilter_worker::max_queue_wait_ms(
|
||||
state.config.bgfilter_call_budget_ms(),
|
||||
audit.external_call_deadline,
|
||||
editor_bgfilter_flat_deadline_reserve(state.config.aliyun_matting_request_timeout_ms),
|
||||
);
|
||||
match crate::bgfilter_worker::request_bgfilter_worker(
|
||||
// deadline_exceeded,本函数随即进入 fallback 链。worker 重启窗口内的连接
|
||||
// 失败由 client 先做小额重试(≤2 次),仍失败才降级,避免不必要的质量损失。
|
||||
match crate::bgfilter_worker::request_bgfilter_worker_with_connect_retry(
|
||||
state,
|
||||
source_object_key,
|
||||
crate::bgfilter_worker::BgfilterBackgroundMode::Flat,
|
||||
Some(screen_color.hex),
|
||||
seg_model,
|
||||
cross_check,
|
||||
max_queue_wait_ms,
|
||||
editor_bgfilter_flat_deadline_reserve(state.config.aliyun_matting_request_timeout_ms),
|
||||
audit,
|
||||
)
|
||||
.await
|
||||
@@ -10625,11 +10617,10 @@ mod tests {
|
||||
"pub(crate) async fn remove_editor_generated_screen_background_with_bgfilter(",
|
||||
"async fn fallback_editor_screen_background_removal",
|
||||
&[
|
||||
"crate::bgfilter_worker::max_queue_wait_ms",
|
||||
"state.config.bgfilter_call_budget_ms()",
|
||||
"editor_bgfilter_flat_deadline_reserve",
|
||||
"crate::bgfilter_worker::request_bgfilter_worker",
|
||||
"crate::bgfilter_worker::request_bgfilter_worker_with_connect_retry",
|
||||
"crate::bgfilter_worker::BgfilterBackgroundMode::Flat",
|
||||
"editor_bgfilter_flat_deadline_reserve",
|
||||
"state.config.aliyun_matting_request_timeout_ms",
|
||||
"error.allows_flat_fallback()",
|
||||
"fallback_editor_screen_background_removal",
|
||||
],
|
||||
@@ -10662,12 +10653,11 @@ mod tests {
|
||||
"async fn request_editor_background_removal_image_with_bgfilter_worker",
|
||||
"pub(crate) struct EditorScreenBackgroundRemovalOutput",
|
||||
&[
|
||||
"crate::bgfilter_worker::max_queue_wait_ms",
|
||||
"state.config.bgfilter_call_budget_ms()",
|
||||
"crate::bgfilter_worker::request_bgfilter_worker",
|
||||
"crate::bgfilter_worker::request_bgfilter_worker_with_connect_retry",
|
||||
"crate::bgfilter_worker::BgfilterBackgroundMode::Complex",
|
||||
"EDITOR_BGFILTER_DEFAULT_SEG_MODEL",
|
||||
"EDITOR_BGFILTER_CROSS_CHECK_DISABLED",
|
||||
"EDITOR_BGFILTER_PARENT_TRANSPORT_WINDOW",
|
||||
".map_err(crate::bgfilter_worker::BgfilterClientError::into_app_error)?",
|
||||
],
|
||||
);
|
||||
|
||||
Reference in New Issue
Block a user