feat(会员): 外部生成认领按账号并发上限过滤

- external_generation_job 追加 (owner_user_id, status) 组合索引,认领事务内现算 running 行数,不建计数表
- claim_external_generation_jobs_tx 认领前判定会员并发上限:达上限的 pending 保持排队
- 账号自身过期 running 的回收豁免上限,避免超限账号卡死任务无法回收
- profile.rs 新增 effective_profile_concurrent_job_limit:缺会员行按 Normal、缺目录行失败关闭到 1
- 抽取纯函数 external_generation_claim_within_concurrency_limit 并补单元测试
This commit is contained in:
2026-10-03 16:20:34 +08:00
parent 325466d9b4
commit 103280025c
2 changed files with 106 additions and 0 deletions
@@ -47,6 +47,11 @@ const INLINE_MEDIA_WARNING_REDACTED_MESSAGE: &str =
accessor = by_external_generation_job_owner_user_id,
btree(columns = [owner_user_id])
),
// 会员并发上限在认领事务里按账号现算 `status = running` 行数;不建计数表,`running` 行即真相。
index(
accessor = by_external_generation_job_owner_status,
btree(columns = [owner_user_id, status])
),
index(
accessor = by_external_generation_job_cursor,
btree(columns = [job_id, source_module])
@@ -885,6 +890,11 @@ fn claim_external_generation_jobs_tx(
finalize_external_generation_job_after_lease_exhaustion(ctx, row, claim_time)?;
continue;
}
// 会员并发上限:只约束 `pending → running` 的认领。达上限的任务留在 `pending` 天然排队;
// 账号自己那条已过期的 `running` 回收不受限,否则超限账号的卡死任务永远无法回收。
if !external_generation_owner_has_concurrency_capacity(ctx, &row, claim_time) {
continue;
}
let next_attempt = row.attempt.saturating_add(1);
let lease_token = build_external_generation_lease_token(
&row.job_id,
@@ -917,6 +927,53 @@ fn claim_external_generation_jobs_tx(
Ok(claimed)
}
/// 认领前的并发上限判定。
///
/// 计数只算 `status = running`(含 lease 过期待回收的 `expired_running`):`pending`(含延时重试)
/// 不占名额,失败退回 `pending` 会自然释放。`128` 哨兵由 [`is_unlimited_concurrency`] 解释。
fn external_generation_owner_has_concurrency_capacity(
ctx: &ReducerContext,
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,
)
}
/// 并发上限判定(纯函数):`running_count >= limit` 时只放行对账号自身过期 `running` 的回收。
fn external_generation_claim_within_concurrency_limit(
is_recycling: bool,
running_count: u32,
limit: u32,
) -> bool {
is_unlimited_concurrency(limit) || is_recycling || running_count < limit
}
/// 账号当前在飞的 `running` 行数;事务内的认领写入会立即反映到下一次计数。
fn count_running_external_generation_jobs_for_owner(
ctx: &ReducerContext,
owner_user_id: &str,
) -> u32 {
ctx.db
.external_generation_job()
.by_external_generation_job_owner_status()
.filter(&owner_user_id.to_string())
.filter(|row| row.status == EXTERNAL_GENERATION_STATUS_RUNNING)
.count()
.try_into()
.unwrap_or(u32::MAX)
}
fn finalize_external_generation_job_after_lease_exhaustion(
ctx: &ReducerContext,
row: ExternalGenerationJob,
@@ -3850,6 +3907,35 @@ mod tests {
}
}
#[test]
fn concurrency_limit_blocks_pending_at_limit_but_recycles_expired_running() {
// pending 在飞数达到上限后不再认领,任务保持 pending 排队。
assert!(external_generation_claim_within_concurrency_limit(
false, 0, 1
));
assert!(!external_generation_claim_within_concurrency_limit(
false, 1, 1
));
assert!(external_generation_claim_within_concurrency_limit(
false, 2, 3
));
assert!(!external_generation_claim_within_concurrency_limit(
false, 3, 3
));
// 回收账号自身过期的 running 不受上限限制,否则超限账号的卡死任务永远无法回收。
assert!(external_generation_claim_within_concurrency_limit(
true, 3, 1
));
// 128 哨兵表示不设上限。
assert!(external_generation_claim_within_concurrency_limit(
false,
10_000,
MEMBERSHIP_UNLIMITED_CONCURRENCY,
));
}
fn micros(value: i64) -> Timestamp {
Timestamp::from_micros_since_unix_epoch(value)
}
@@ -9480,6 +9480,26 @@ fn membership_plan_row(
ctx.db.profile_membership_plan().plan().find(&plan)
}
/// 账号当前有效的并发上限:`profile_membership.plan` → 目录行 `concurrent_job_limit`。
///
/// 缺会员行按 `Normal`(非会员);缺目录行**失败关闭**到 `Normal`(=1),
/// `128` 哨兵由 [`is_unlimited_concurrency`] 在调用方解释。
pub(crate) fn effective_profile_concurrent_job_limit(
ctx: &ReducerContext,
owner_user_id: &str,
) -> u32 {
let plan = ctx
.db
.profile_membership()
.user_id()
.find(&owner_user_id.to_string())
.map(|row| row.plan)
.unwrap_or(RuntimeProfileMembershipPlan::Normal);
membership_plan_row(ctx, plan)
.map(|row| row.concurrent_job_limit)
.unwrap_or(1)
}
fn membership_plan_period_points(ctx: &ReducerContext, plan: RuntimeProfileMembershipPlan) -> u64 {
membership_plan_row(ctx, plan)
.map(|row| row.period_points)