Merge remote-tracking branch 'origin/master' into feat/game-works-management
Project CI / AI game creator shell Rust crates (pull_request) Successful in 2m32s
Project CI / AI game creator shell Rust lane 1/2 (pull_request) Successful in 5m5s
Project CI / AI game creator shell Rust lane 2/2 (pull_request) Successful in 3m33s
Project CI / Frontend tests (pull_request) Successful in 3m39s
Project CI / Backend tests (pull_request) Successful in 6m49s
Project CI / AI game creator shell web tests (pull_request) Successful in 2m39s
Project CI / Repository checks (pull_request) Successful in 6m43s
Project CI / Native shell tests (pull_request) Successful in 7m39s

# Conflicts:
#	docs/project-memory/shared-memory/decision-log.md
#	docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md
#	server-rs/crates/api-server/src/modules/game_distribution.rs
#	server-rs/crates/spacetime-client/src/active.rs
#	server-rs/crates/spacetime-module/src/game_distribution.rs
#	src/components/game-distribution/GameDistributionPages.test.tsx
This commit is contained in:
2026-10-04 15:00:23 +08:00
46 changed files with 1806 additions and 130 deletions
+15
View File
@@ -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,473 @@
//! 游戏游玩次数的进程内聚合缓冲。
//!
//! 只在内存里累计「开始游戏」上报:按 `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);
}
}
@@ -0,0 +1,112 @@
//! 游玩计数的 flush worker:按间隔把内存增量批量写进 SpacetimeDB。
//!
//! 失败语义与 ADR 一致:连接还没建起来(`Build`)属于确定未发出,交给计数器放回下一轮重试;
//! 其余错误(`Timeout` / `ConnectDropped` / `Procedure`)无法判断是否已提交,直接丢弃该批并
//! 记录丢失量,避免系统性双计。
//!
//! 任一分片失败都会终止本次 flush 的后续分片:连接不通时剩余分片只会重复同样的失败,逐个重试
//! 会把 worker 卡在多次连接超时上;剩余增量按“直接丢弃”处理,尽快回到 tick。
//!
//! 进程关停不做强制 flush:内存里剩下的增量随进程结束丢弃,关停路径不为它等待网络。
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;
}
}
});
}
/// 把一批增量写库。
///
/// `Build`(确定未发出)放回等待下一轮;其余错误无法判断是否已提交,直接丢弃并记丢失量。
async fn flush_deltas(state: &AppState, deltas: Vec<GamePlayCountDelta>) {
let mut accepted = 0usize;
let mut dropped = 0u64;
let mut start = 0usize;
while start < deltas.len() {
let end = (start + GAME_PLAY_COUNT_FLUSH_BATCH_SIZE).min(deltas.len());
let chunk = &deltas[start..end];
match write_batch(state, chunk).await {
Ok(()) => {
accepted += chunk.len();
start = end;
}
// 连接没建起来(确定未发出):本批放回,但后面的分片会重复同样的失败,
// 直接丢弃剩余,避免每个分片各等一次连接超时把 worker 卡住。
Err(SpacetimeClientError::Build(message)) => {
state.game_play_counter().requeue(chunk);
let lost = total_delta(&deltas[end..]);
dropped = dropped.saturating_add(lost);
warn!(
error = %message,
games = chunk.len(),
lost,
"游戏游玩计数写入连接未建立:本批放回等待下一轮,剩余增量直接丢弃"
);
break;
}
// 结果未知:连本批一起丢弃剩余,避免双计,也不再把 worker 卡在逐个重连上。
Err(error) => {
let lost = total_delta(&deltas[start..]);
dropped = dropped.saturating_add(lost);
warn!(
error = %error,
games = deltas.len() - start,
lost,
"游戏游玩计数写入失败,剩余增量直接丢弃"
);
break;
}
}
}
if accepted > 0 {
info!(games = accepted, "游戏游玩计数已批量落库");
}
if dropped > 0 {
warn!(dropped, "游戏游玩计数存在丢弃量");
}
}
fn total_delta(deltas: &[GamePlayCountDelta]) -> u64 {
deltas.iter().map(|delta| delta.delta).sum()
}
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
}
+3
View File
@@ -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;
@@ -700,6 +702,7 @@ fn spawn_http_app_state_background_workers(state: &AppState, process_role: Proce
crate::payment_webhook::PaymentWebhookWorker::new(state.spacetime_client().clone())
.spawn_worker();
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());
@@ -1,7 +1,7 @@
use std::{
collections::{BTreeMap, HashMap, VecDeque},
sync::{Arc, Mutex, OnceLock},
time::{SystemTime, UNIX_EPOCH},
time::{Instant, SystemTime, UNIX_EPOCH},
};
use axum::{
@@ -67,10 +67,11 @@ use uuid::Uuid;
use crate::{
admin::{AuthenticatedAdmin, require_admin_auth},
api_response::json_success_body,
auth::{AuthenticatedAccessToken, require_bearer_auth},
auth::{AuthenticatedAccessToken, optional_access_token_from_headers, require_bearer_auth},
game_play_counter::{GamePlayOutcome, GamePlayReport},
http_error::AppError,
platform_errors::{map_llm_error, map_oss_error},
request_context::RequestContext,
request_context::{RequestContext, client_ip_from_headers},
state::AppState,
tracking::{TrackingEventDraft, record_tracking_event_after_success},
};
@@ -385,6 +386,10 @@ pub fn router(state: AppState) -> Router<AppState> {
let public_games = Router::new()
.route("/api/game-distribution/games", get(list_games))
.route("/api/game-distribution/games/{game_id}", get(get_game))
.route(
"/api/game-distribution/games/{game_id}/plays",
post(record_game_play),
)
.route_layer(middleware::from_fn(add_no_store_response_headers));
Router::new()
@@ -1037,6 +1042,112 @@ async fn get_game(
Ok(json_success_body(Some(&ctx), public_game_payload(game)))
}
/// 一次游玩上报的请求体;只有匿名身份需要 `clientId`,登录身份由 bearer 决定。
#[derive(Debug, Default, Deserialize)]
#[serde(rename_all = "camelCase")]
struct RecordGamePlayRequest {
#[serde(default)]
client_id: Option<String>,
}
/// 记录一次「开始游戏」。
///
/// 公开端点:登录用户按 `userId` 去重,匿名按 `clientId`(缺失时回退 `IP + UA`)去重;
/// 命中 30 分钟去重窗口或超过 `IP + game` 限流时不增加计数。计数只进内存缓冲,
/// 立即返回 `recorded`,任何失败都不影响游玩本身。
async fn record_game_play(
State(state): State<AppState>,
Extension(ctx): Extension<RequestContext>,
Path(game_id): Path<String>,
headers: HeaderMap,
body: Bytes,
) -> Result<Json<Value>, AppError> {
let game_id = game_id.trim().to_string();
if game_id.is_empty() {
return Err(AppError::from_status(StatusCode::NOT_FOUND));
}
let client_ip = client_ip_from_headers(&headers);
// 先在内存里挡掉明显超限的请求,避免它们也去打一次 SpacetimeDB;真正计数时 record 会再判一次。
if state
.game_play_counter()
.is_rate_limited(&game_id, &client_ip, Instant::now())
{
return Err(AppError::from_status(StatusCode::TOO_MANY_REQUESTS));
}
// 非公开 / 已下架 / 已暂停的游戏不计数,按不存在返回。
let is_public = state
.spacetime_client()
.get_public_game_distribution_game(game_id.clone())
.await
.map_err(map_spacetime_error)?
.is_some();
if !is_public {
return Err(AppError::from_status(StatusCode::NOT_FOUND));
}
let user_agent = user_agent_tag(&headers);
let authenticated = optional_access_token_from_headers(
&state,
format!("/api/game-distribution/games/{game_id}/plays"),
headers,
ctx.request_id().to_string(),
)
.await
.unwrap_or_else(|error| {
// 可选 bearer:无效 token 按匿名处理,绝不能因为它挡掉一次真实游玩。
debug!(error = %error, "游戏游玩计数忽略无效 bearer,按匿名计数");
None
});
let identity = authenticated
.as_ref()
.map(|token| format!("user:{}", token.claims().user_id()))
.or_else(|| request_client_id(&body).map(|client_id| format!("client:{client_id}")))
.unwrap_or_else(|| format!("ip:{client_ip}|ua:{user_agent}"));
let outcome = state.game_play_counter().record(
GamePlayReport {
game_id: &game_id,
identity: &identity,
client_ip: &client_ip,
},
Instant::now(),
);
if outcome == GamePlayOutcome::RateLimited {
return Err(AppError::from_status(StatusCode::TOO_MANY_REQUESTS));
}
Ok(json_success_body(
Some(&ctx),
json!({ "recorded": outcome == GamePlayOutcome::Counted }),
))
}
fn request_client_id(body: &Bytes) -> Option<String> {
if body.is_empty() {
return None;
}
let request = serde_json::from_slice::<RecordGamePlayRequest>(body).ok()?;
request
.client_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.map(|value| value.chars().take(128).collect())
}
fn user_agent_tag(headers: &HeaderMap) -> String {
headers
.get(header::USER_AGENT)
.and_then(|value| value.to_str().ok())
.map(str::trim)
.filter(|value| !value.is_empty())
.unwrap_or("unknown")
.chars()
.take(64)
.collect()
}
/// 作者自有游戏列表:只返回当前认证主体名下的游戏与最近版本状态。
async fn list_my_games(
State(state): State<AppState>,
@@ -4876,4 +4987,53 @@ mod tests {
assert_eq!(delete_audit.metadata["title"], "已删除作品");
assert_eq!(delete_audit.metadata["expectedPublicationRevision"], 5);
}
#[tokio::test]
async fn game_play_route_is_public_and_no_store_without_spacetime_connection() {
use axum::http::Request;
use tower::ServiceExt;
let app =
crate::app::build_router(AppState::new(crate::config::AppConfig::default()).unwrap());
let response = app
.oneshot(
Request::builder()
.method("POST")
.uri("/api/game-distribution/games/game_1/plays")
.header(header::CONTENT_TYPE, "application/json")
.body(Body::from(r#"{"clientId":"client-1"}"#))
.unwrap(),
)
.await
.unwrap();
// 未连接 SpacetimeDB 时公开可见性读取失败,但路由可达且不需要登录;只有确认公开后才计数。
assert_eq!(response.status(), StatusCode::BAD_GATEWAY);
assert_eq!(response.headers()[header::CACHE_CONTROL], "no-store");
}
#[test]
fn play_request_client_id_trims_limits_and_rejects_blank() {
assert_eq!(request_client_id(&Bytes::from_static(b"")), None);
assert_eq!(request_client_id(&Bytes::from_static(b"not json")), None);
assert_eq!(request_client_id(&Bytes::from_static(b"{}")), None);
assert_eq!(
request_client_id(&Bytes::from_static(br#"{"clientId":" abc "}"#)),
Some("abc".to_string())
);
assert_eq!(
request_client_id(&Bytes::from_static(br#"{"clientId":" "}"#)),
None
);
let long = "x".repeat(200);
let body = Bytes::from(format!(r#"{{"clientId":"{long}"}}"#));
assert_eq!(request_client_id(&body).unwrap().chars().count(), 128);
}
#[test]
fn play_report_user_agent_is_bounded_and_falls_back() {
assert_eq!(user_agent_tag(&HeaderMap::new()), "unknown");
let mut headers = HeaderMap::new();
headers.insert(header::USER_AGENT, " test-agent ".parse().unwrap());
assert_eq!(user_agent_tag(&headers), "test-agent");
}
}
@@ -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,29 @@ pub async fn attach_request_context(mut request: Request, next: Next) -> Respons
.await
}
/// 从代理头解析客户端 IP。
///
/// 优先 `x-real-ip`:nginx 用 `$remote_addr` 覆盖写入,是真实 TCP 对端,调用方无法伪造。
/// `x-forwarded-for` 只作回退,并取**最后一段**——nginx 用 `$proxy_add_x_forwarded_for` 会把真实
/// 对端追加在末尾,前面几段是调用方自带的、可伪造。两者都拿不到时才兜底回环地址。
pub fn client_ip_from_headers(headers: &HeaderMap) -> String {
headers
.get("x-real-ip")
.and_then(|value| value.to_str().ok())
.map(str::trim)
.filter(|value| !value.is_empty())
.or_else(|| {
headers
.get("x-forwarded-for")
.and_then(|value| value.to_str().ok())
.and_then(|value| value.rsplit(',').next())
.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 +171,35 @@ mod tests {
assert_eq!(context.external_call_deadline(), None);
}
#[test]
fn client_ip_prefers_real_ip_over_forwarded_for() {
let mut headers = HeaderMap::new();
headers.insert(
"x-forwarded-for",
HeaderValue::from_static("203.0.113.7, 10.0.0.1"),
);
headers.insert("x-real-ip", HeaderValue::from_static("198.51.100.9"));
assert_eq!(client_ip_from_headers(&headers), "198.51.100.9");
}
#[test]
fn client_ip_forwarded_for_fallback_uses_last_address() {
// nginx 把真实对端追加在末尾,前面是调用方可伪造的值,只能取最后一段。
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), "10.0.0.1");
}
#[test]
fn client_ip_ignores_blank_real_ip_and_falls_back_then_loopback() {
let mut headers = HeaderMap::new();
headers.insert("x-real-ip", HeaderValue::from_static(" "));
headers.insert("x-forwarded-for", HeaderValue::from_static("10.0.0.1"));
assert_eq!(client_ip_from_headers(&headers), "10.0.0.1");
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(
+15
View File
@@ -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};
@@ -320,6 +321,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>,
@@ -625,6 +628,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(),
)
@@ -722,6 +732,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,
@@ -1729,6 +1740,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()
}