后端:关停 flush 不再假承诺重试
Project CI / Backend tests (pull_request) Failing after 32s
Project CI / AI game creator shell Rust lane 2/2 (pull_request) Successful in 4m37s
Project CI / AI game creator shell Rust crates (pull_request) Successful in 2m46s
Project CI / Repository checks (pull_request) Failing after 34s
Project CI / AI game creator shell Rust lane 1/2 (pull_request) Successful in 5m21s
Project CI / Frontend tests (pull_request) Successful in 3m16s
Project CI / AI game creator shell web tests (pull_request) Successful in 3m3s
Project CI / Native shell tests (pull_request) Successful in 6m51s

- game_play_counter_worker.rs:flush_deltas 增加 requeue_on_build;worker 循环仍把 Build 失败放回下一轮,关停路径改为直接记丢失量,不再把增量放回一个不会再被 flush 的缓冲
This commit is contained in:
2026-10-03 19:28:39 +08:00
parent ea8734b700
commit 2122dc9469
@@ -28,7 +28,7 @@ pub(crate) fn spawn_game_play_counter_worker(state: AppState) {
let counter = state.game_play_counter();
counter.prune_expired(now);
if let Some(deltas) = counter.take_pending_if_due(now) {
flush_deltas(&state, deltas).await;
flush_deltas(&state, deltas, true).await;
}
}
});
@@ -40,16 +40,20 @@ pub(crate) async fn flush_game_play_counter_for_shutdown(state: &AppState) {
if deltas.is_empty() {
return;
}
flush_deltas(state, deltas).await;
flush_deltas(state, deltas, false).await;
}
async fn flush_deltas(state: &AppState, deltas: Vec<GamePlayCountDelta>) {
/// 把一批增量写库。
///
/// `requeue_on_build` 控制 `Build`(确定未发出)失败时是否放回:worker 循环为 `true`,等待下一轮;
/// 关停路径为 `false`——`finalize_shutdown` 之后不再重试,放回只会随进程退出丢失,应直接记丢失量。
async fn flush_deltas(state: &AppState, deltas: Vec<GamePlayCountDelta>, requeue_on_build: bool) {
let mut accepted = 0usize;
let mut dropped = 0u64;
for chunk in deltas.chunks(GAME_PLAY_COUNT_FLUSH_BATCH_SIZE) {
match write_batch(state, chunk).await {
Ok(()) => accepted += chunk.len(),
Err(SpacetimeClientError::Build(message)) => {
Err(SpacetimeClientError::Build(message)) if requeue_on_build => {
state.game_play_counter().requeue(chunk);
warn!(
error = %message,
@@ -57,6 +61,16 @@ async fn flush_deltas(state: &AppState, deltas: Vec<GamePlayCountDelta>) {
"游戏游玩计数写入连接未建立,已放回等待下一轮"
);
}
Err(SpacetimeClientError::Build(message)) => {
let lost = chunk.iter().map(|delta| delta.delta).sum::<u64>();
dropped = dropped.saturating_add(lost);
warn!(
error = %message,
games = chunk.len(),
lost,
"游戏游玩计数关停 flush 连接未建立,未落库增量丢弃"
);
}
Err(error) => {
let lost = chunk.iter().map(|delta| delta.delta).sum::<u64>();
dropped = dropped.saturating_add(lost);