diff --git a/server-rs/crates/spacetime-module/src/external_generation.rs b/server-rs/crates/spacetime-module/src/external_generation.rs index ab09c27c3..91f0758ba 100644 --- a/server-rs/crates/spacetime-module/src/external_generation.rs +++ b/server-rs/crates/spacetime-module/src/external_generation.rs @@ -1,5 +1,6 @@ use crate::*; use std::cmp::Ordering; +use std::collections::HashMap; use std::ops::RangeFrom; const EXTERNAL_GENERATION_STATUS_PENDING: &str = "pending"; @@ -885,14 +886,23 @@ fn claim_external_generation_jobs_tx( }); let mut claimed = Vec::new(); + let mut concurrency = ExternalGenerationOwnerConcurrencyCache::default(); for mut row in candidates.into_iter().take(limit) { if external_generation_job_has_exhausted_attempts(&row) { + let owner_user_id = row.owner_user_id.clone(); finalize_external_generation_job_after_lease_exhaustion(ctx, row, claim_time)?; + // 终结可能把 running 行改成 failed,账号计数缓存随之失效。 + concurrency.invalidate_running(&owner_user_id); continue; } // 会员并发上限:只约束 `pending → running` 的认领。达上限的任务留在 `pending` 天然排队; // 账号自己那条已过期的 `running` 回收不受限,否则超限账号的卡死任务永远无法回收。 - if !external_generation_owner_has_concurrency_capacity(ctx, &row, claim_time) { + if !external_generation_owner_has_concurrency_capacity( + ctx, + &mut concurrency, + &row, + claim_time, + ) { continue; } let next_attempt = row.attempt.saturating_add(1); @@ -921,33 +931,77 @@ fn claim_external_generation_jobs_tx( Some(worker_id.clone()), claim_time, ); + // 本次认领把一行推进到 running,账号计数缓存失效,后续候选行重新现查。 + concurrency.invalidate_running(&row.owner_user_id); claimed.push(map_external_generation_job_row(row)); } Ok(claimed) } +/// 单次 claim 事务内按账号缓存有效并发上限与 running 行数。 +/// +/// 上限在一次事务内是常量;running 计数在账号的任一行被认领或终结后立即失效、下次重查, +/// 因此判定结果与逐候选行现查一致,只是把同一账号的重复查询折叠为一次。 +#[derive(Default)] +struct ExternalGenerationOwnerConcurrencyCache { + limits: HashMap, + running: HashMap, +} + +impl ExternalGenerationOwnerConcurrencyCache { + fn limit(&mut self, ctx: &ReducerContext, owner_user_id: &str) -> u32 { + if let Some(limit) = self.limits.get(owner_user_id) { + return *limit; + } + let limit = crate::effective_profile_concurrent_job_limit(ctx, owner_user_id); + self.limits.insert(owner_user_id.to_string(), limit); + limit + } + + fn running_count(&mut self, ctx: &ReducerContext, owner_user_id: &str) -> u32 { + if let Some(count) = self.running.get(owner_user_id) { + return *count; + } + let count = count_running_external_generation_jobs_for_owner(ctx, owner_user_id); + self.running.insert(owner_user_id.to_string(), count); + count + } + + /// 账号任一 `running` 行发生状态变化(认领 / 终结)后调用,保证下次读取重新现查。 + fn invalidate_running(&mut self, owner_user_id: &str) { + self.running.remove(owner_user_id); + } + + fn allows(&mut self, ctx: &ReducerContext, owner_user_id: &str, is_recycling: bool) -> bool { + let limit = self.limit(ctx, owner_user_id); + if is_unlimited_concurrency(limit) { + return true; + } + // 回收豁免不需要计数:自身过期 running 已计入在飞数,放行不会推高它。 + let running = if is_recycling { + 0 + } else { + self.running_count(ctx, owner_user_id) + }; + external_generation_claim_within_concurrency_limit(is_recycling, running, limit) + } +} + /// 认领前的并发上限判定。 /// /// 计数只算 `status = running`(含 lease 过期待回收的 `expired_running`):`pending`(含延时重试) /// 不占名额,失败退回 `pending` 会自然释放。`128` 哨兵由 [`is_unlimited_concurrency`] 解释。 fn external_generation_owner_has_concurrency_capacity( ctx: &ReducerContext, + cache: &mut ExternalGenerationOwnerConcurrencyCache, row: &ExternalGenerationJob, now: Timestamp, ) -> bool { - let limit = crate::effective_profile_concurrent_job_limit(ctx, &row.owner_user_id); - if is_unlimited_concurrency(limit) { - return true; - } // 回收豁免:候选本身已是该账号过期的 `running`,它已计入在飞数;放行回收不会推高在飞任务数。 let is_recycling = row.status == EXTERNAL_GENERATION_STATUS_RUNNING && is_external_generation_job_claimable(row, now); - external_generation_claim_within_concurrency_limit( - is_recycling, - count_running_external_generation_jobs_for_owner(ctx, &row.owner_user_id), - limit, - ) + cache.allows(ctx, &row.owner_user_id, is_recycling) } /// 并发上限判定(纯函数):`running_count >= limit` 时只放行对账号自身过期 `running` 的回收。 @@ -3936,6 +3990,18 @@ mod tests { )); } + #[test] + fn owner_concurrency_cache_only_invalidates_running_count() { + let mut cache = ExternalGenerationOwnerConcurrencyCache::default(); + cache.running.insert("owner-a".to_string(), 2); + cache.limits.insert("owner-a".to_string(), 3); + // 认领 / 终结后 running 计数必须重查,否则会把上一行的变更带到下一候选行。 + cache.invalidate_running("owner-a"); + assert!(cache.running.get("owner-a").is_none()); + // 上限在一次事务内是常量,不受 running 计数失效影响。 + assert_eq!(cache.limits.get("owner-a"), Some(&3)); + } + fn micros(value: i64) -> Timestamp { Timestamp::from_micros_since_unix_epoch(value) }