后端:接入游玩计数内存缓冲与 flush worker
- 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 改为复用
This commit is contained in:
@@ -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");
|
||||
|
||||
@@ -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<GamePlayCounterState>,
|
||||
}
|
||||
|
||||
struct GamePlayCounterState {
|
||||
/// `game_id -> 待落库增量`;只包含通过校验的上报,键空间由公开游戏目录界定。
|
||||
pending: HashMap<String, u64>,
|
||||
/// `identity + game_id -> 最近一次计数时间`。
|
||||
seen: HashMap<String, Instant>,
|
||||
/// `IP + game_id -> 当前固定窗口`。
|
||||
rate: HashMap<String, RateWindow>,
|
||||
/// 最近一次真正取走增量的时间,用于判断是否到达 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<Vec<GamePlayCountDelta>> {
|
||||
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<GamePlayCountDelta> {
|
||||
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<String, u64>) -> Vec<GamePlayCountDelta> {
|
||||
let mut deltas = pending
|
||||
.drain()
|
||||
.filter_map(|(game_id, delta)| (delta > 0).then_some(GamePlayCountDelta { game_id, delta }))
|
||||
.collect::<Vec<_>>();
|
||||
// 稳定批次顺序,便于测试与日志比对。
|
||||
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
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -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<GamePlayCountDelta>) {
|
||||
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::<u64>();
|
||||
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
|
||||
}
|
||||
@@ -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());
|
||||
|
||||
@@ -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<B>(request: &HttpRequest<B>) -> Option<String> {
|
||||
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");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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<Arc<TrackingOutbox>>,
|
||||
wallet_refund_outbox: Option<Arc<WalletRefundOutbox>>,
|
||||
profile_wallet_refund_outbox_worker: Arc<ProfileWalletRefundOutboxWorker>,
|
||||
/// 游玩计数的进程内聚合缓冲;写入由 `game_play_counter_worker` 负责。
|
||||
game_play_counter: GamePlayCounter,
|
||||
editor_generation_pricing_store: EditorGenerationPricingStore,
|
||||
llm_client: Option<LlmClient>,
|
||||
vector_engine_llm_client: Option<LlmClient>,
|
||||
@@ -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()
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user