From dbf3adeb137e7c43a904d6ba2596e72fe874552a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=8E=8B=E5=BE=B7=E5=AE=87?= Date: Sat, 3 Oct 2026 18:14:42 +0800 Subject: [PATCH] =?UTF-8?q?=E5=90=8E=E7=AB=AF=EF=BC=9A=E6=8E=A5=E5=85=A5?= =?UTF-8?q?=E6=B8=B8=E7=8E=A9=E8=AE=A1=E6=95=B0=E5=86=85=E5=AD=98=E7=BC=93?= =?UTF-8?q?=E5=86=B2=E4=B8=8E=20flush=20worker?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - game_play_counter.rs:新增纯内存计数器(30min 身份去重、IP+game 固定窗口限流、按时/按量取增量、requeue、过期清理)与 9 个单测 - game_play_counter_worker.rs:新增 flush worker,Build 失败放回重试、Timeout/ConnectDropped 丢弃并记录丢失量,关停前强制 flush - main.rs:HTTP 角色注册 worker,并在 finalize_shutdown 内按 outbox 超时强制落库 - state.rs/config.rs:AppState 持有计数器,新增 GENARRATIVE_GAME_PLAY_COUNTER_FLUSH_INTERVAL_MS(默认 5s) - request_context.rs:抽出 client_ip_from_headers 并补测试,runtime_profile.rs 改为复用 --- server-rs/crates/api-server/src/config.rs | 15 + .../api-server/src/game_play_counter.rs | 409 ++++++++++++++++++ .../src/game_play_counter_worker.rs | 95 ++++ server-rs/crates/api-server/src/main.rs | 23 + .../crates/api-server/src/request_context.rs | 40 +- .../crates/api-server/src/runtime_profile.rs | 16 +- server-rs/crates/api-server/src/state.rs | 15 + 7 files changed, 597 insertions(+), 16 deletions(-) create mode 100644 server-rs/crates/api-server/src/game_play_counter.rs create mode 100644 server-rs/crates/api-server/src/game_play_counter_worker.rs diff --git a/server-rs/crates/api-server/src/config.rs b/server-rs/crates/api-server/src/config.rs index 1b71b2cda..9cb9334c1 100644 --- a/server-rs/crates/api-server/src/config.rs +++ b/server-rs/crates/api-server/src/config.rs @@ -74,6 +74,8 @@ pub struct AppConfig { pub tracking_outbox_batch_size: usize, pub tracking_outbox_flush_interval: Duration, pub tracking_outbox_max_bytes: u64, + /// 游玩计数内存缓冲的 flush 间隔;写入是批量 procedure,间隔决定展示滞后上限。 + pub game_play_counter_flush_interval: Duration, pub wallet_refund_outbox_enabled: bool, pub wallet_refund_outbox_dir: PathBuf, pub wallet_refund_outbox_batch_size: usize, @@ -374,6 +376,7 @@ impl Default for AppConfig { tracking_outbox_batch_size: 500, tracking_outbox_flush_interval: Duration::from_millis(1_000), tracking_outbox_max_bytes: 256 * 1024 * 1024, + game_play_counter_flush_interval: Duration::from_millis(5_000), wallet_refund_outbox_enabled: true, wallet_refund_outbox_dir: PathBuf::from("server-rs/.data/wallet-refund-outbox"), wallet_refund_outbox_batch_size: 100, @@ -858,6 +861,11 @@ impl AppConfig { { config.tracking_outbox_max_bytes = max_bytes; } + if let Some(flush_interval_ms) = + read_first_positive_u64_env(&["GENARRATIVE_GAME_PLAY_COUNTER_FLUSH_INTERVAL_MS"]) + { + config.game_play_counter_flush_interval = Duration::from_millis(flush_interval_ms); + } if let Some(enabled) = read_first_bool_env(&["GENARRATIVE_WALLET_REFUND_OUTBOX_ENABLED"]) { config.wallet_refund_outbox_enabled = enabled; } @@ -2527,6 +2535,7 @@ mod tests { std::env::remove_var("GENARRATIVE_TRACKING_OUTBOX_BATCH_SIZE"); std::env::remove_var("GENARRATIVE_TRACKING_OUTBOX_FLUSH_INTERVAL_MS"); std::env::remove_var("GENARRATIVE_TRACKING_OUTBOX_MAX_BYTES"); + std::env::remove_var("GENARRATIVE_GAME_PLAY_COUNTER_FLUSH_INTERVAL_MS"); std::env::remove_var("GENARRATIVE_WALLET_REFUND_OUTBOX_ENABLED"); std::env::remove_var("GENARRATIVE_WALLET_REFUND_OUTBOX_DIR"); std::env::remove_var("GENARRATIVE_WALLET_REFUND_OUTBOX_BATCH_SIZE"); @@ -2546,6 +2555,7 @@ mod tests { std::env::set_var("GENARRATIVE_TRACKING_OUTBOX_BATCH_SIZE", "250"); std::env::set_var("GENARRATIVE_TRACKING_OUTBOX_FLUSH_INTERVAL_MS", "2000"); std::env::set_var("GENARRATIVE_TRACKING_OUTBOX_MAX_BYTES", "1048576"); + std::env::set_var("GENARRATIVE_GAME_PLAY_COUNTER_FLUSH_INTERVAL_MS", "4000"); std::env::set_var("GENARRATIVE_WALLET_REFUND_OUTBOX_ENABLED", "false"); std::env::set_var( "GENARRATIVE_WALLET_REFUND_OUTBOX_DIR", @@ -2577,6 +2587,10 @@ mod tests { std::time::Duration::from_millis(2_000) ); assert_eq!(config.tracking_outbox_max_bytes, 1_048_576); + assert_eq!( + config.game_play_counter_flush_interval, + std::time::Duration::from_millis(4_000) + ); assert!(!config.wallet_refund_outbox_enabled); assert_eq!( config.wallet_refund_outbox_dir, @@ -2601,6 +2615,7 @@ mod tests { std::env::remove_var("GENARRATIVE_TRACKING_OUTBOX_BATCH_SIZE"); std::env::remove_var("GENARRATIVE_TRACKING_OUTBOX_FLUSH_INTERVAL_MS"); std::env::remove_var("GENARRATIVE_TRACKING_OUTBOX_MAX_BYTES"); + std::env::remove_var("GENARRATIVE_GAME_PLAY_COUNTER_FLUSH_INTERVAL_MS"); std::env::remove_var("GENARRATIVE_WALLET_REFUND_OUTBOX_ENABLED"); std::env::remove_var("GENARRATIVE_WALLET_REFUND_OUTBOX_DIR"); std::env::remove_var("GENARRATIVE_WALLET_REFUND_OUTBOX_BATCH_SIZE"); diff --git a/server-rs/crates/api-server/src/game_play_counter.rs b/server-rs/crates/api-server/src/game_play_counter.rs new file mode 100644 index 000000000..60f131c7d --- /dev/null +++ b/server-rs/crates/api-server/src/game_play_counter.rs @@ -0,0 +1,409 @@ +//! 游戏游玩次数的进程内聚合缓冲。 +//! +//! 只在内存里累计「开始游戏」上报:按 `identity + game_id` 做去重窗口、按 `IP + game_id` 做固定 +//! 窗口限流,并按 flush 间隔或待落库游戏数上限决定何时把增量交给写库方。 +//! +//! 网络写入不在本模块内:`take_pending` / `requeue` 让调用方在锁外发起 procedure,锁只覆盖 +//! HashMap 操作。时间点全部由调用方传入 `now`,因此本模块不依赖运行时,可直接单测。 + +use std::{ + collections::HashMap, + sync::{Mutex, MutexGuard}, + time::{Duration, Instant}, +}; + +/// 默认 flush 间隔;与 `tracking_outbox` 的秒级节奏一致,够短以保证展示及时。 +const DEFAULT_FLUSH_INTERVAL: Duration = Duration::from_secs(5); +/// 同一 `identity + game_id` 的去重窗口。 +const DEFAULT_DEDUP_WINDOW: Duration = Duration::from_secs(30 * 60); +/// `IP + game_id` 固定窗口长度。 +const DEFAULT_RATE_WINDOW: Duration = Duration::from_secs(60); +/// 单个 `IP + game_id` 在每个固定窗口内允许的上报次数。 +const DEFAULT_RATE_LIMIT: u32 = 60; +/// 待落库游戏数上限;达到后下一次检查立即 flush,而不是等满一整个间隔。 +const DEFAULT_MAX_PENDING_GAMES: usize = 4096; + +/// 一次上报携带的最小信息。`identity` 由调用方决定:登录用户是 userId,匿名是 clientId, +/// 都拿不到时才回退 `IP + UA`。 +#[derive(Clone, Copy, Debug)] +pub struct GamePlayReport<'a> { + pub game_id: &'a str, + pub identity: &'a str, + pub client_ip: &'a str, +} + +/// 单次上报的判定结果。 +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum GamePlayOutcome { + /// 已计入待落库增量。 + Counted, + /// 落在去重窗口内,未计入。 + Deduped, + /// 超过 `IP + game_id` 固定窗口上限,未计入。 + RateLimited, +} + +/// 待落库的增量;同一 `game_id` 在一个批次内只出现一次。 +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct GamePlayCountDelta { + pub game_id: String, + pub delta: u64, +} + +/// 计数器参数。默认值覆盖决策口径,`flush_interval` 由 `AppConfig` 覆盖。 +#[derive(Clone, Copy, Debug)] +pub struct GamePlayCounterSettings { + pub flush_interval: Duration, + pub dedup_window: Duration, + pub rate_window: Duration, + pub rate_limit: u32, + pub max_pending_games: usize, +} + +impl Default for GamePlayCounterSettings { + fn default() -> Self { + Self { + flush_interval: DEFAULT_FLUSH_INTERVAL, + dedup_window: DEFAULT_DEDUP_WINDOW, + rate_window: DEFAULT_RATE_WINDOW, + rate_limit: DEFAULT_RATE_LIMIT, + max_pending_games: DEFAULT_MAX_PENDING_GAMES, + } + } +} + +/// 游玩计数内存缓冲。 +pub struct GamePlayCounter { + settings: GamePlayCounterSettings, + inner: Mutex, +} + +struct GamePlayCounterState { + /// `game_id -> 待落库增量`;只包含通过校验的上报,键空间由公开游戏目录界定。 + pending: HashMap, + /// `identity + game_id -> 最近一次计数时间`。 + seen: HashMap, + /// `IP + game_id -> 当前固定窗口`。 + rate: HashMap, + /// 最近一次真正取走增量的时间,用于判断是否到达 flush 间隔。 + last_flush_at: Instant, +} + +struct RateWindow { + started_at: Instant, + count: u32, +} + +impl GamePlayCounter { + pub fn new(settings: GamePlayCounterSettings, now: Instant) -> Self { + Self { + settings, + inner: Mutex::new(GamePlayCounterState { + pending: HashMap::new(), + seen: HashMap::new(), + rate: HashMap::new(), + last_flush_at: now, + }), + } + } + + pub fn flush_interval(&self) -> Duration { + self.settings.flush_interval + } + + /// 处理一次上报:先去重,再限流,最后累加增量。 + /// + /// 被去重命中的上报不消耗限流额度;限流只挡同一 IP 对同一游戏的超额上报。 + pub fn record(&self, report: GamePlayReport<'_>, now: Instant) -> GamePlayOutcome { + let dedup_key = format!("{}\u{1f}{}", report.identity, report.game_id); + let rate_key = format!("{}\u{1f}{}", report.client_ip, report.game_id); + let mut state = self.lock(); + + if let Some(seen_at) = state.seen.get(&dedup_key) + && now.saturating_duration_since(*seen_at) < self.settings.dedup_window + { + return GamePlayOutcome::Deduped; + } + + let rate = state.rate.entry(rate_key).or_insert_with(|| RateWindow { + started_at: now, + count: 0, + }); + if now.saturating_duration_since(rate.started_at) >= self.settings.rate_window { + rate.started_at = now; + rate.count = 0; + } + if rate.count >= self.settings.rate_limit { + return GamePlayOutcome::RateLimited; + } + rate.count = rate.count.saturating_add(1); + + state.seen.insert(dedup_key, now); + let pending = state.pending.entry(report.game_id.to_string()).or_insert(0); + *pending = pending.saturating_add(1); + GamePlayOutcome::Counted + } + + /// 到达 flush 间隔或待落库游戏数达到上限时取走全部增量;否则返回 `None`。 + pub fn take_pending_if_due(&self, now: Instant) -> Option> { + let mut state = self.lock(); + if state.pending.is_empty() { + return None; + } + let interval_elapsed = + now.saturating_duration_since(state.last_flush_at) >= self.settings.flush_interval; + let at_capacity = state.pending.len() >= self.settings.max_pending_games; + if !interval_elapsed && !at_capacity { + return None; + } + state.last_flush_at = now; + Some(drain_pending(&mut state.pending)) + } + + /// 无条件取走全部增量,供关停 flush 使用。 + pub fn take_pending(&self) -> Vec { + let mut state = self.lock(); + drain_pending(&mut state.pending) + } + + /// 写库失败时把增量放回,等待下一次 flush。 + pub fn requeue(&self, deltas: &[GamePlayCountDelta]) { + let mut state = self.lock(); + for delta in deltas { + let pending = state.pending.entry(delta.game_id.clone()).or_insert(0); + *pending = pending.saturating_add(delta.delta); + } + } + + /// 清掉过期的去重与限流条目,避免 map 无界增长。 + pub fn prune_expired(&self, now: Instant) { + let mut state = self.lock(); + let dedup_window = self.settings.dedup_window; + let rate_window = self.settings.rate_window; + state + .seen + .retain(|_, seen_at| now.saturating_duration_since(*seen_at) < dedup_window); + state + .rate + .retain(|_, window| now.saturating_duration_since(window.started_at) < rate_window); + } + + #[cfg(test)] + pub fn pending_game_count(&self) -> usize { + self.lock().pending.len() + } + + #[cfg(test)] + pub fn pending_total(&self) -> u64 { + self.lock().pending.values().copied().sum() + } + + fn lock(&self) -> MutexGuard<'_, GamePlayCounterState> { + // 计数是尽力而为的展示指标:锁中毒时继续用内部状态,不让一次 panic 永久关闭计数。 + self.inner + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + } +} + +fn drain_pending(pending: &mut HashMap) -> Vec { + let mut deltas = pending + .drain() + .filter_map(|(game_id, delta)| (delta > 0).then_some(GamePlayCountDelta { game_id, delta })) + .collect::>(); + // 稳定批次顺序,便于测试与日志比对。 + deltas.sort_by(|left, right| left.game_id.cmp(&right.game_id)); + deltas +} + +#[cfg(test)] +mod tests { + use super::*; + + fn settings() -> GamePlayCounterSettings { + GamePlayCounterSettings { + flush_interval: Duration::from_secs(5), + dedup_window: Duration::from_secs(30 * 60), + rate_window: Duration::from_secs(60), + rate_limit: 3, + max_pending_games: 4, + } + } + + fn report<'a>(game_id: &'a str, identity: &'a str, client_ip: &'a str) -> GamePlayReport<'a> { + GamePlayReport { + game_id, + identity, + client_ip, + } + } + + #[test] + fn repeated_report_within_dedup_window_is_deduped() { + let start = Instant::now(); + let counter = GamePlayCounter::new(settings(), start); + + assert_eq!( + counter.record(report("g1", "u1", "1.1.1.1"), start), + GamePlayOutcome::Counted + ); + assert_eq!( + counter.record( + report("g1", "u1", "1.1.1.1"), + start + Duration::from_secs(60) + ), + GamePlayOutcome::Deduped + ); + assert_eq!(counter.pending_total(), 1); + } + + #[test] + fn dedup_expires_after_window() { + let start = Instant::now(); + let counter = GamePlayCounter::new(settings(), start); + + counter.record(report("g1", "u1", "1.1.1.1"), start); + assert_eq!( + counter.record( + report("g1", "u1", "1.1.1.1"), + start + Duration::from_secs(30 * 60) + ), + GamePlayOutcome::Counted + ); + assert_eq!(counter.pending_total(), 2); + } + + #[test] + fn different_identities_count_separately() { + let start = Instant::now(); + let counter = GamePlayCounter::new(settings(), start); + + counter.record(report("g1", "u1", "1.1.1.1"), start); + counter.record(report("g1", "u2", "1.1.1.1"), start); + assert_eq!(counter.pending_total(), 2); + assert_eq!(counter.pending_game_count(), 1); + } + + #[test] + fn rate_limit_blocks_excess_and_resets_after_window() { + let start = Instant::now(); + let counter = GamePlayCounter::new(settings(), start); + + for index in 0..3 { + assert_eq!( + counter.record(report("g1", &format!("u{index}"), "1.1.1.1"), start), + GamePlayOutcome::Counted + ); + } + assert_eq!( + counter.record(report("g1", "u9", "1.1.1.1"), start), + GamePlayOutcome::RateLimited + ); + + assert_eq!( + counter.record( + report("g1", "u9", "1.1.1.1"), + start + Duration::from_secs(60) + ), + GamePlayOutcome::Counted + ); + } + + #[test] + fn take_pending_if_due_waits_for_interval_then_drains() { + let start = Instant::now(); + let counter = GamePlayCounter::new(settings(), start); + + counter.record(report("g1", "u1", "1.1.1.1"), start); + assert!( + counter + .take_pending_if_due(start + Duration::from_secs(4)) + .is_none() + ); + + let deltas = counter + .take_pending_if_due(start + Duration::from_secs(5)) + .expect("到达间隔后应当取走增量"); + assert_eq!( + deltas, + vec![GamePlayCountDelta { + game_id: "g1".to_string(), + delta: 1, + }] + ); + assert_eq!(counter.pending_total(), 0); + } + + #[test] + fn capacity_triggers_early_flush() { + let start = Instant::now(); + let counter = GamePlayCounter::new(settings(), start); + + for index in 0..4u64 { + counter.record( + report(&format!("g{index}"), "u1", "1.1.1.1"), + start + Duration::from_millis(index * 10), + ); + } + let deltas = counter + .take_pending_if_due(start + Duration::from_secs(1)) + .expect("达到待落库游戏数上限应当立即 flush"); + assert_eq!(deltas.len(), 4); + } + + #[test] + fn aggregates_per_game_and_sorts() { + let start = Instant::now(); + let counter = GamePlayCounter::new(settings(), start); + + counter.record(report("g2", "u1", "1.1.1.1"), start); + counter.record(report("g1", "u2", "1.1.1.1"), start); + counter.record(report("g2", "u2", "1.1.1.1"), start); + + let deltas = counter.take_pending(); + assert_eq!( + deltas, + vec![ + GamePlayCountDelta { + game_id: "g1".to_string(), + delta: 1, + }, + GamePlayCountDelta { + game_id: "g2".to_string(), + delta: 2, + }, + ] + ); + } + + #[test] + fn requeue_restores_deltas() { + let start = Instant::now(); + let counter = GamePlayCounter::new(settings(), start); + + counter.record(report("g1", "u1", "1.1.1.1"), start); + let deltas = counter.take_pending(); + assert_eq!(counter.pending_total(), 0); + + counter.requeue(&deltas); + assert_eq!(counter.pending_total(), 1); + } + + #[test] + fn prune_expired_drops_old_dedup_and_rate_entries() { + let start = Instant::now(); + let counter = GamePlayCounter::new(settings(), start); + + counter.record(report("g1", "u1", "1.1.1.1"), start); + counter.prune_expired(start + Duration::from_secs(30 * 60 + 1)); + + // 去重条目过期后同一身份还能重新计数。 + assert_eq!( + counter.record( + report("g1", "u1", "1.1.1.1"), + start + Duration::from_secs(30 * 60 + 2) + ), + GamePlayOutcome::Counted + ); + } +} diff --git a/server-rs/crates/api-server/src/game_play_counter_worker.rs b/server-rs/crates/api-server/src/game_play_counter_worker.rs new file mode 100644 index 000000000..6df0607e6 --- /dev/null +++ b/server-rs/crates/api-server/src/game_play_counter_worker.rs @@ -0,0 +1,95 @@ +//! 游玩计数的 flush worker:按间隔把内存增量批量写进 SpacetimeDB。 +//! +//! 失败语义与 ADR 一致:连接还没建起来(`Build`)属于确定未发出,交给计数器放回下一轮重试; +//! 其余错误(`Timeout` / `ConnectDropped` / `Procedure`)无法判断是否已提交,直接丢弃该批并 +//! 记录丢失量,避免系统性双计。 + +use std::time::{Duration, Instant}; + +use spacetime_client::{GameDistributionPlayCountIncrementRecordInput, SpacetimeClientError}; +use tokio::time::sleep; +use tracing::{info, warn}; + +use crate::{game_play_counter::GamePlayCountDelta, state::AppState}; + +/// 一次 procedure 最多携带多少条游戏增量;只影响帧大小,不影响累计结果。 +pub(crate) const GAME_PLAY_COUNT_FLUSH_BATCH_SIZE: usize = 500; + +/// worker 检查节拍:比默认 flush 间隔细,以便待落库游戏数达到上限时提前 flush。 +const GAME_PLAY_COUNTER_TICK: Duration = Duration::from_secs(1); + +/// 启动 flush worker;只在 HTTP 角色进程注册一次。 +pub(crate) fn spawn_game_play_counter_worker(state: AppState) { + let tick = GAME_PLAY_COUNTER_TICK.min(state.game_play_counter().flush_interval()); + tokio::spawn(async move { + loop { + sleep(tick).await; + let now = Instant::now(); + 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 一次;由 `finalize_shutdown` 在总超时内调用。 +pub(crate) async fn flush_game_play_counter_for_shutdown(state: &AppState) { + let deltas = state.game_play_counter().take_pending(); + if deltas.is_empty() { + return; + } + flush_deltas(state, deltas).await; +} + +async fn flush_deltas(state: &AppState, deltas: Vec) { + 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)) => { + state.game_play_counter().requeue(chunk); + warn!( + error = %message, + games = chunk.len(), + "游戏游玩计数写入连接未建立,已放回等待下一轮" + ); + } + Err(error) => { + let lost = chunk.iter().map(|delta| delta.delta).sum::(); + dropped = dropped.saturating_add(lost); + warn!( + error = %error, + games = chunk.len(), + lost, + "游戏游玩计数写入结果未知,丢弃本批以避免双计" + ); + } + } + } + if accepted > 0 { + info!(games = accepted, "游戏游玩计数已批量落库"); + } + if dropped > 0 { + warn!(dropped, "游戏游玩计数存在丢弃量"); + } +} + +async fn write_batch( + state: &AppState, + deltas: &[GamePlayCountDelta], +) -> Result<(), SpacetimeClientError> { + let increments = deltas + .iter() + .map(|delta| GameDistributionPlayCountIncrementRecordInput { + game_id: delta.game_id.clone(), + delta: delta.delta, + }) + .collect(); + state + .spacetime_client() + .increment_game_distribution_game_play_counts(increments) + .await +} diff --git a/server-rs/crates/api-server/src/main.rs b/server-rs/crates/api-server/src/main.rs index 8515721f5..4954845cc 100644 --- a/server-rs/crates/api-server/src/main.rs +++ b/server-rs/crates/api-server/src/main.rs @@ -48,6 +48,8 @@ mod external_generation_worker_controller; mod external_mcp; mod external_skill_api; mod frontend_runtime_config; +mod game_play_counter; +mod game_play_counter_worker; mod generated_image_assets; mod health; mod http_error; @@ -680,6 +682,26 @@ async fn finalize_shutdown(context: ShutdownContext) { } } } + + if let Some(state) = context.app_state.as_ref() { + info!(timeout_ms, "api-server 退出前 flush 游戏游玩计数内存缓冲"); + match timeout( + context.outbox_flush_timeout, + crate::game_play_counter_worker::flush_game_play_counter_for_shutdown(state), + ) + .await + { + Ok(()) => { + info!("api-server 退出前游戏游玩计数 flush 完成"); + } + Err(_) => { + warn!( + timeout_ms, + "api-server 退出前游戏游玩计数 flush 超时,未落库增量已丢弃" + ); + } + } + } } fn spawn_common_app_state_background_workers(state: &AppState) { @@ -695,6 +717,7 @@ fn spawn_common_app_state_background_workers(state: &AppState) { fn spawn_http_app_state_background_workers(state: &AppState, process_role: ProcessRole) { spawn_common_app_state_background_workers(state); crate::error_reports::spawn_cleanup_worker(state.clone()); + crate::game_play_counter_worker::spawn_game_play_counter_worker(state.clone()); if should_start_profile_recharge_expiration_listener(process_role) { spawn_profile_recharge_expiration_listener(state.clone()); spawn_profile_recharge_refund_reconciliation_worker(state.clone()); diff --git a/server-rs/crates/api-server/src/request_context.rs b/server-rs/crates/api-server/src/request_context.rs index 57a109d28..9cd65b410 100644 --- a/server-rs/crates/api-server/src/request_context.rs +++ b/server-rs/crates/api-server/src/request_context.rs @@ -2,7 +2,7 @@ use std::time::{Duration, Instant}; use axum::{ extract::Request, - http::{HeaderValue, Request as HttpRequest, header::HeaderName}, + http::{HeaderMap, HeaderValue, Request as HttpRequest, header::HeaderName}, middleware::Next, response::Response, }; @@ -107,6 +107,26 @@ pub async fn attach_request_context(mut request: Request, next: Next) -> Respons .await } +/// 从代理头解析客户端 IP:反代固定用 `x-forwarded-for` 的第一个地址,直连(本地开发) +/// 回退 `x-real-ip`,都拿不到时才兜底回环地址。 +pub fn client_ip_from_headers(headers: &HeaderMap) -> String { + headers + .get("x-forwarded-for") + .and_then(|value| value.to_str().ok()) + .and_then(|value| value.split(',').next()) + .map(str::trim) + .filter(|value| !value.is_empty()) + .or_else(|| { + headers + .get("x-real-ip") + .and_then(|value| value.to_str().ok()) + .map(str::trim) + .filter(|value| !value.is_empty()) + }) + .unwrap_or("127.0.0.1") + .to_string() +} + pub fn resolve_request_id(request: &HttpRequest) -> Option { request .extensions() @@ -148,4 +168,22 @@ mod tests { assert_eq!(context.external_call_deadline(), None); } + + #[test] + fn client_ip_prefers_first_forwarded_address() { + let mut headers = HeaderMap::new(); + headers.insert( + "x-forwarded-for", + HeaderValue::from_static("203.0.113.7, 10.0.0.1"), + ); + assert_eq!(client_ip_from_headers(&headers), "203.0.113.7"); + } + + #[test] + fn client_ip_falls_back_to_real_ip_then_loopback() { + let mut headers = HeaderMap::new(); + headers.insert("x-real-ip", HeaderValue::from_static("198.51.100.9")); + assert_eq!(client_ip_from_headers(&headers), "198.51.100.9"); + assert_eq!(client_ip_from_headers(&HeaderMap::new()), "127.0.0.1"); + } } diff --git a/server-rs/crates/api-server/src/runtime_profile.rs b/server-rs/crates/api-server/src/runtime_profile.rs index db27bb20f..5f8f96d78 100644 --- a/server-rs/crates/api-server/src/runtime_profile.rs +++ b/server-rs/crates/api-server/src/runtime_profile.rs @@ -1632,21 +1632,7 @@ fn is_wechat_recharge_payment_channel(payment_channel: &str) -> bool { } fn resolve_wechat_pay_client_ip(headers: &HeaderMap) -> String { - headers - .get("x-forwarded-for") - .and_then(|value| value.to_str().ok()) - .and_then(|value| value.split(',').next()) - .map(str::trim) - .filter(|value| !value.is_empty()) - .or_else(|| { - headers - .get("x-real-ip") - .and_then(|value| value.to_str().ok()) - .map(str::trim) - .filter(|value| !value.is_empty()) - }) - .unwrap_or("127.0.0.1") - .to_string() + crate::request_context::client_ip_from_headers(headers) } async fn resolve_wechat_identity_for_payment( diff --git a/server-rs/crates/api-server/src/state.rs b/server-rs/crates/api-server/src/state.rs index a7b970bfc..e89cb4829 100644 --- a/server-rs/crates/api-server/src/state.rs +++ b/server-rs/crates/api-server/src/state.rs @@ -45,6 +45,7 @@ use crate::editor_generation_config::{ EditorGenerationPricingConfig, EditorGenerationPricingError, EditorGenerationPricingStore, EditorGenerationPricingUnit, }; +use crate::game_play_counter::{GamePlayCounter, GamePlayCounterSettings}; use crate::tracking_outbox::TrackingOutbox; use crate::wallet_refund_outbox::{ProfileWalletRefundOutboxWorker, WalletRefundOutbox}; use crate::wechat::pay::{build_wechat_pay_config, map_wechat_pay_init_error}; @@ -312,6 +313,8 @@ pub struct AppStateInner { tracking_outbox: Option>, wallet_refund_outbox: Option>, profile_wallet_refund_outbox_worker: Arc, + /// 游玩计数的进程内聚合缓冲;写入由 `game_play_counter_worker` 负责。 + game_play_counter: GamePlayCounter, editor_generation_pricing_store: EditorGenerationPricingStore, llm_client: Option, vector_engine_llm_client: Option, @@ -617,6 +620,13 @@ impl AppState { WalletRefundOutbox::from_config(&config, spacetime_client.clone()); let profile_wallet_refund_outbox_worker = ProfileWalletRefundOutboxWorker::from_config(&config, spacetime_client.clone()); + let game_play_counter = GamePlayCounter::new( + GamePlayCounterSettings { + flush_interval: config.game_play_counter_flush_interval, + ..GamePlayCounterSettings::default() + }, + std::time::Instant::now(), + ); let editor_generation_pricing_store = EditorGenerationPricingStore::load( config.editor_generation_pricing_override_path.clone(), ) @@ -713,6 +723,7 @@ impl AppState { tracking_outbox, wallet_refund_outbox, profile_wallet_refund_outbox_worker, + game_play_counter, editor_generation_pricing_store, llm_client, vector_engine_llm_client, @@ -1691,6 +1702,10 @@ impl AppState { self.profile_wallet_refund_outbox_worker.clone() } + pub fn game_play_counter(&self) -> &GamePlayCounter { + &self.game_play_counter + } + pub fn llm_client(&self) -> Option<&LlmClient> { self.llm_client.as_ref() }