75ad087e2c
- main.rs:finalize_shutdown 删除游戏游玩计数强制 flush,关停不再为它等待网络,内存里剩余增量随进程结束丢弃 - game_play_counter_worker.rs:删除 flush_game_play_counter_for_shutdown,flush_deltas 去掉 requeue_on_build 参数(关停路径已不存在),Build 仍放回等待下一轮 - game_play_counter.rs:take_pending 标记为仅测试使用,模块文档改指 take_pending_if_due
474 lines
16 KiB
Rust
474 lines
16 KiB
Rust
//! 游戏游玩次数的进程内聚合缓冲。
|
||
//!
|
||
//! 只在内存里累计「开始游戏」上报:按 `identity + game_id` 做去重窗口、按 `IP + game_id` 做固定
|
||
//! 窗口限流,并按 flush 间隔或待落库游戏数上限决定何时把增量交给写库方。
|
||
//!
|
||
//! 网络写入不在本模块内:`take_pending_if_due` / `requeue` 让调用方在锁外发起 procedure,锁只覆盖
|
||
//! HashMap 操作。时间点全部由调用方传入 `now`,因此本模块不依赖运行时,可直接单测。
|
||
|
||
use std::{
|
||
collections::HashMap,
|
||
sync::{Mutex, MutexGuard},
|
||
time::{Duration, Instant},
|
||
};
|
||
|
||
use tracing::warn;
|
||
|
||
/// 默认 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, String), Instant>,
|
||
/// `(IP, game_id) -> 当前固定窗口`。
|
||
rate: HashMap<(String, 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 {
|
||
// 键用元组而不是拼接字符串:clientId 来自 JSON,可能包含任意字节(含 U+001F),
|
||
// 拼接会产生本不存在的键碰撞。
|
||
let dedup_key = (report.identity.to_string(), report.game_id.to_string());
|
||
let rate_key = (report.client_ip.to_string(), report.game_id.to_string());
|
||
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
|
||
}
|
||
|
||
/// 限流预检:只读,不消耗额度、不改任何状态。
|
||
///
|
||
/// 公开上报接口在昂贵的公开可见性查询之前先用它挡掉明显超限的请求;真正计数时
|
||
/// `record` 仍会重新判定一次,所以这里只用于省一次远端查询,不承担正确性。
|
||
pub fn is_rate_limited(&self, game_id: &str, client_ip: &str, now: Instant) -> bool {
|
||
let key = (client_ip.to_string(), game_id.to_string());
|
||
let state = self.lock();
|
||
state.rate.get(&key).is_some_and(|window| {
|
||
now.saturating_duration_since(window.started_at) < self.settings.rate_window
|
||
&& window.count >= self.settings.rate_limit
|
||
})
|
||
}
|
||
|
||
/// 到达 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;
|
||
let deltas = drain_pending(&mut state.pending);
|
||
drop(state);
|
||
Some(sort_deltas(deltas))
|
||
}
|
||
|
||
/// 无条件取走全部增量;仅供测试使用,生产路径走 `take_pending_if_due`。
|
||
#[cfg(test)]
|
||
pub fn take_pending(&self) -> Vec<GamePlayCountDelta> {
|
||
let mut state = self.lock();
|
||
let deltas = drain_pending(&mut state.pending);
|
||
drop(state);
|
||
sort_deltas(deltas)
|
||
}
|
||
|
||
/// 写库失败时把增量放回,等待下一次 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 永久关闭计数。
|
||
// 但中毒意味着上一次 panic 可能留下部分更新的状态,必须留下可关联的日志。
|
||
self.inner.lock().unwrap_or_else(|poisoned| {
|
||
warn!("游戏游玩计数锁已中毒,继续使用内部状态");
|
||
poisoned.into_inner()
|
||
})
|
||
}
|
||
}
|
||
|
||
fn drain_pending(pending: &mut HashMap<String, u64>) -> Vec<GamePlayCountDelta> {
|
||
pending
|
||
.drain()
|
||
.filter_map(|(game_id, delta)| (delta > 0).then_some(GamePlayCountDelta { game_id, delta }))
|
||
.collect()
|
||
}
|
||
|
||
/// 稳定批次顺序,便于测试与日志比对;调用方已释放计数锁,排序不阻塞并发 `record`。
|
||
fn sort_deltas(mut deltas: Vec<GamePlayCountDelta>) -> Vec<GamePlayCountDelta> {
|
||
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 rate_limit_precheck_is_read_only_and_window_scoped() {
|
||
let start = Instant::now();
|
||
let counter = GamePlayCounter::new(settings(), start);
|
||
|
||
assert!(!counter.is_rate_limited("g1", "1.1.1.1", start));
|
||
for index in 0..3 {
|
||
counter.record(report("g1", &format!("u{index}"), "1.1.1.1"), start);
|
||
}
|
||
|
||
// 预检只读:连续调用不消耗额度,也不改变判定。
|
||
assert!(counter.is_rate_limited("g1", "1.1.1.1", start));
|
||
assert!(counter.is_rate_limited("g1", "1.1.1.1", start));
|
||
// 另一个 IP、另一个游戏都不受影响。
|
||
assert!(!counter.is_rate_limited("g1", "2.2.2.2", start));
|
||
assert!(!counter.is_rate_limited("g2", "1.1.1.1", start));
|
||
// 固定窗口结束后恢复。
|
||
assert!(!counter.is_rate_limited("g1", "1.1.1.1", start + Duration::from_secs(60)));
|
||
}
|
||
|
||
#[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
|
||
);
|
||
}
|
||
|
||
#[test]
|
||
fn separator_like_bytes_in_key_parts_do_not_collide() {
|
||
let start = Instant::now();
|
||
let counter = GamePlayCounter::new(settings(), start);
|
||
|
||
// 旧实现用 U+001F 拼接 identity/game_id,下面两个不同元组会被拼成同一个键。
|
||
assert_eq!(
|
||
counter.record(report("b\u{1f}c", "a", "1.1.1.1"), start),
|
||
GamePlayOutcome::Counted
|
||
);
|
||
assert_eq!(
|
||
counter.record(report("c", "a\u{1f}b", "1.1.1.1"), start),
|
||
GamePlayOutcome::Counted
|
||
);
|
||
assert_eq!(counter.pending_total(), 2);
|
||
}
|
||
}
|