Files
Genarrative/server-rs/crates/spacetime-module/src/agc_analytics.rs
T
lhk229 1e186369c9
Project CI / AI game creator shell Rust crates (push) Successful in 1m24s
Project CI / AI game creator shell Rust smoke (push) Successful in 1m56s
Project CI / AI game creator shell Rust lane 1/2 (push) Has been cancelled
Project CI / Frontend tests (push) Has been cancelled
Project CI / Backend tests (push) Has been cancelled
Project CI / Repository checks (push) Has been cancelled
Project CI / AI game creator shell web tests (push) Has been cancelled
Project CI / AI game creator shell Rust lane 2/2 (push) Has been cancelled
Project CI / Native shell tests (push) Has been cancelled
客户端埋点设置 (#446)
Reviewed-on: https://git.genarrative.world/git/GenarrativeAI/Genarrative/pulls/446
Co-authored-by: Linghong <ink29535@proton.me>
Co-committed-by: Linghong <ink29535@proton.me>
2026-09-23 00:09:12 +08:00

315 lines
12 KiB
Rust

use crate::*;
use module_runtime::agc_analytics::{
AgcTrackingCursor, agc_time_micros, agc_tracking_filter_key, validate_agc_analytics_batch,
validate_agc_tracking_query,
};
use shared_contracts::admin::{
AdminAgcTrackingEventEntry, AdminAgcTrackingEventListPayload, AdminAgcTrackingEventListQuery,
};
use shared_contracts::agc_analytics::{AgcAnalyticsAcknowledgement, AgcAnalyticsBatch, Event};
use std::collections::BTreeMap;
#[spacetimedb::table(
accessor = agc_tracking_event,
index(accessor = by_agc_tracking_received, btree(columns = [received_at])),
index(accessor = by_agc_tracking_user_received, btree(columns = [user_id, received_at])),
index(accessor = by_agc_tracking_name_received, btree(columns = [event_name, received_at]))
)]
#[derive(Clone, Debug, PartialEq)]
pub struct AgcTrackingEvent {
#[primary_key]
pub event_id: String,
pub schema_version: u32,
pub event_name: String,
pub event_time: Timestamp,
pub user_id: String,
pub editor_session_id: String,
pub project_id: Option<String>,
pub creative_task_id: Option<String>,
pub agent_run_id: Option<String>,
pub agent_turn_id: Option<String>,
pub status: Option<String>,
pub error_code: Option<String>,
pub source: String,
pub client_version: String,
pub properties_json: String,
pub batch_id: String,
pub received_at: Timestamp,
}
fn event_row(
event: &Event,
batch_id: &str,
received_at: Timestamp,
) -> Result<AgcTrackingEvent, String> {
let enum_string = |v: serde_json::Value| v.as_str().unwrap_or_default().to_owned();
Ok(AgcTrackingEvent {
event_id: event.event_id.clone(),
schema_version: event.schema_version,
event_name: event.event_name.clone(),
event_time: Timestamp::from_micros_since_unix_epoch(agc_time_micros(&event.event_time)?),
user_id: event.user_id.clone().ok_or("invalid_batch")?,
editor_session_id: event.editor_session_id.clone(),
project_id: event.project_id.clone(),
creative_task_id: event.creative_task_id.clone(),
agent_run_id: event.agent_run_id.clone(),
agent_turn_id: event.agent_turn_id.clone(),
status: event
.status
.map(|v| enum_string(serde_json::to_value(v).unwrap())),
error_code: event
.error_code
.map(|v| enum_string(serde_json::to_value(v).unwrap())),
source: enum_string(serde_json::to_value(event.source).unwrap()),
client_version: event.client_version.clone(),
properties_json: serde_json::to_string(&event.properties)
.map_err(|_| "invalid_properties")?,
batch_id: batch_id.to_owned(),
received_at,
})
}
// 接收批次和时间不属于原始事件内容,重传保留首次值。
fn same_event(existing: &AgcTrackingEvent, incoming: &AgcTrackingEvent) -> bool {
let mut incoming = incoming.clone();
incoming.batch_id = existing.batch_id.clone();
incoming.received_at = existing.received_at;
let same_properties = serde_json::from_str::<serde_json::Value>(&existing.properties_json).ok()
== serde_json::from_str::<serde_json::Value>(&incoming.properties_json).ok();
incoming.properties_json = existing.properties_json.clone();
same_properties && existing == &incoming
}
#[spacetimedb::procedure]
pub fn upload_agc_analytics_batch(
ctx: &mut ProcedureContext,
payload_json: String,
) -> Result<String, String> {
if payload_json.len() > 2 * 1024 * 1024 {
return Err("invalid_batch".into());
}
let batch: AgcAnalyticsBatch =
serde_json::from_str(&payload_json).map_err(|_| "invalid_batch")?;
validate_agc_analytics_batch(&batch)?;
let caller = ctx.sender();
ctx.try_with_tx(|tx| {
crate::editor_project_storage::require_editor_generation_runtime_service_identity(
tx, caller,
)?;
// 同一事务内任一冲突返回 Err,之前插入的行一起回滚。
for event in &batch.events {
let incoming = event_row(event, &batch.batch_id, tx.timestamp)?;
if let Some(existing) = tx.db.agc_tracking_event().event_id().find(&event.event_id) {
if !same_event(&existing, &incoming) {
return Err("agc_event_conflict".into());
}
} else {
tx.db.agc_tracking_event().insert(incoming);
}
}
serde_json::to_string(&AgcAnalyticsAcknowledgement {
acknowledged_batch_ids: vec![batch.batch_id.clone()],
event_count: batch.events.len() as u32,
})
.map_err(|_| "acknowledgement_failed".into())
})
}
fn matches_query(
row: &AgcTrackingEvent,
query: &AdminAgcTrackingEventListQuery,
start: Option<i64>,
end: Option<i64>,
) -> bool {
query.user_id.as_ref().is_none_or(|v| v == &row.user_id)
&& query
.project_id
.as_ref()
.is_none_or(|v| Some(v) == row.project_id.as_ref())
&& query
.creative_task_id
.as_ref()
.is_none_or(|v| Some(v) == row.creative_task_id.as_ref())
&& query
.agent_run_id
.as_ref()
.is_none_or(|v| Some(v) == row.agent_run_id.as_ref())
&& query
.event_name
.as_ref()
.is_none_or(|v| v == &row.event_name)
&& query
.client_version
.as_ref()
.is_none_or(|v| v == &row.client_version)
&& start.is_none_or(|v| row.event_time.to_micros_since_unix_epoch() >= v)
&& end.is_none_or(|v| row.event_time.to_micros_since_unix_epoch() < v)
}
fn entry(row: AgcTrackingEvent) -> Result<AdminAgcTrackingEventEntry, String> {
Ok(AdminAgcTrackingEventEntry {
event_id: row.event_id,
schema_version: row.schema_version,
event_name: row.event_name,
event_time: module_runtime::format_utc_micros(row.event_time.to_micros_since_unix_epoch()),
user_id: row.user_id,
editor_session_id: row.editor_session_id,
project_id: row.project_id,
creative_task_id: row.creative_task_id,
agent_run_id: row.agent_run_id,
agent_turn_id: row.agent_turn_id,
status: row.status,
error_code: row.error_code,
source: row.source,
client_version: row.client_version,
properties: serde_json::from_str(&row.properties_json)
.map_err(|_| "invalid_stored_properties")?,
batch_id: row.batch_id,
received_at: module_runtime::format_utc_micros(
row.received_at.to_micros_since_unix_epoch(),
),
})
}
#[spacetimedb::procedure]
pub fn list_agc_tracking_events(
ctx: &mut ProcedureContext,
query_json: String,
) -> Result<String, String> {
if query_json.len() > 16384 {
return Err("invalid_agc_query".into());
}
let query: AdminAgcTrackingEventListQuery =
serde_json::from_str(&query_json).map_err(|_| "invalid_agc_query")?;
let cursor = validate_agc_tracking_query(&query)?;
let start = query
.start_time
.as_deref()
.map(agc_time_micros)
.transpose()?;
let end = query.end_time.as_deref().map(agc_time_micros).transpose()?;
let caller = ctx.sender();
ctx.try_with_tx(|tx| {
crate::editor_project_storage::require_editor_generation_runtime_service_identity(
tx, caller,
)?;
let snapshot = cursor
.as_ref()
.map(|v| v.snapshot)
.unwrap_or(tx.timestamp.to_micros_since_unix_epoch());
let limit = query.limit.unwrap_or(50).clamp(1, 200) as usize;
let upper = Timestamp::from_micros_since_unix_epoch(snapshot);
let mut selected = BTreeMap::new();
// 在数据库事务中扫描完整索引候选,再保留最大的 limit+1 个键。
// 不依赖 SDK 迭代器顺序,不将任意 LIMIT 截断误作最新一页,内存只保留一页。
let mut consider = |row: AgcTrackingEvent| {
let key = (
row.event_time.to_micros_since_unix_epoch(),
row.event_id.clone(),
);
if row.received_at.to_micros_since_unix_epoch() > snapshot
|| cursor
.as_ref()
.is_some_and(|v| key >= (v.event_time, v.event_id.clone()))
|| !matches_query(&row, &query, start, end)
{
return;
}
selected.insert(key, row);
if selected.len() > limit + 1 {
selected.pop_first();
}
};
if let Some(user_id) = &query.user_id {
for row in tx
.db
.agc_tracking_event()
.by_agc_tracking_user_received()
.filter((user_id.as_str(), ..=upper))
{
consider(row);
}
} else if let Some(event_name) = &query.event_name {
for row in tx
.db
.agc_tracking_event()
.by_agc_tracking_name_received()
.filter((event_name.as_str(), ..=upper))
{
consider(row);
}
} else {
for row in tx
.db
.agc_tracking_event()
.by_agc_tracking_received()
.filter(..=upper)
{
consider(row);
}
}
let has_more = selected.len() > limit;
let mut rows: Vec<_> = selected.into_iter().rev().take(limit).collect();
let next_cursor = if has_more {
let ((event_time, event_id), _) = rows.last().expect("nonempty page");
Some(
serde_json::to_string(&AgcTrackingCursor {
snapshot,
event_time: *event_time,
event_id: event_id.clone(),
filter_key: agc_tracking_filter_key(&query),
})
.map_err(|_| "cursor_failed")?,
)
} else {
None
};
let entries = rows
.drain(..)
.map(|(_, row)| entry(row))
.collect::<Result<Vec<_>, _>>()?;
serde_json::to_string(&AdminAgcTrackingEventListPayload {
entries,
next_cursor,
})
.map_err(|_| "query_failed".into())
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn agc_replay_compares_business_fields_but_preserves_first_receipt() {
let event: Event = serde_json::from_value(serde_json::json!({
"schema_version":1,"event_id":"790a1275-a0a0-405c-8bfd-287201bef10a",
"event_name":"editor_session_start","event_time":"2026-09-21T12:00:00.000Z",
"user_id":"a","editor_session_id":"790a1275-a0a0-405c-8bfd-287201bef10a",
"project_id":null,"creative_task_id":null,"agent_run_id":null,"agent_turn_id":null,
"status":"success","error_code":null,"source":"editor","client_version":"1",
"properties":{"entry_source":"direct_launch","first_project_id":null}
}))
.unwrap();
let first = event_row(
&event,
"batch-1",
Timestamp::from_micros_since_unix_epoch(10),
)
.unwrap();
let mut replay = event_row(
&event,
"batch-2",
Timestamp::from_micros_since_unix_epoch(20),
)
.unwrap();
replay.properties_json =
r#"{"first_project_id":null,"entry_source":"direct_launch"}"#.into();
assert!(same_event(&first, &replay));
replay.user_id = "b".into();
assert!(!same_event(&first, &replay));
replay.user_id = "a".into();
replay.client_version = "2".into();
assert!(!same_event(&first, &replay));
}
}