From 3336d4663cb6a1a25377fe05b84c9d4c7b702e32 Mon Sep 17 00:00:00 2001 From: "Mars.P" Date: Tue, 19 May 2026 13:02:54 +0800 Subject: [PATCH] Serialise apply_delta to fix concurrent-ingest XP loss MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit # The race `StateManager::apply_delta` does a read-modify-write on `pet_state`: let prior = self.db.get_pet_state(pet_id).await?; // READ let xp_after = prior.total_xp + delta; self.db.upsert_pet_state(pet_id, xp_after, ...).await?; // WRITE With concurrent ingest the read+write isn't atomic across tasks: Task A: read pet_state.total_xp = 100 Task B: read pet_state.total_xp = 100 ← same snapshot Task A: write 100 + 10 = 110 Task B: write 100 + 5 = 105 ← overwrites A's 110 A's 10 XP is "lost" — the `xp_event` row is still written (different UUID, different dedup key), so `sum(xp_event)` > `pet_state.total_xp`. The pet's displayed level / stage is computed from the (stale) state, so the user sees ~8% fewer XP than the events database implies. Empirically reproduced: a 1.5-hour heavy Claude-Code Unicorn session ended with `sum(xp_event)=988` but `pet_state.total_xp=910` — 78 XP / 7.9% loss. # The fix Single tokio Mutex held across the whole read-modify-write inside `apply_delta` (and `rebuild`, which has the same shape). petpet runs at most one active pet at a time, so a global mutex is the right granularity — contention is bounded by the live ingest rate (handful of events per second worst-case). If multi-pet concurrent ingest ever lands, swap for a per-pet mutex map (HashMap>> behind a parent RwLock). # Regression coverage New test `apply_delta_serialises_concurrent_ingest`: 1. Spawn 200 concurrent `apply_delta(+1)` tasks against one pet 2. Await all, then read final `pet_state.total_xp` 3. Assert == 200 (pre-fix: would be < 200 with high probability) 309/309 lib tests pass; race regression stable across 5 reruns. # Why not SQL `BEGIN IMMEDIATE` transaction? Cleaner long-term, but would require wider plumbing changes (DbHandle exposes `execute` shapes, not multi-statement transaction blocks). The Mutex approach is mechanically equivalent in this codebase since there's only one DbHandle and one active pet — write contention is logically equivalent to mutex contention. Worth revisiting if we add multi-pet concurrent ingest OR an external SQL client that might bypass the engine's serialisation. --- src/xp/state.rs | 108 +++++++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 106 insertions(+), 2 deletions(-) diff --git a/src/xp/state.rs b/src/xp/state.rs index 5371017..f9e283f 100644 --- a/src/xp/state.rs +++ b/src/xp/state.rs @@ -53,15 +53,43 @@ impl AppliedDelta { pub struct StateManager { db: Arc, + /// Serialises the read-modify-write inside `apply_delta` / + /// `rebuild` so concurrent ingest doesn't lose XP via TOCTOU. + /// + /// Without this lock, two concurrent `ingest_usage` calls can: + /// 1. Task A reads pet_state.total_xp = 100 + /// 2. Task B reads pet_state.total_xp = 100 (same snapshot) + /// 3. Task A writes 100 + 10 = 110 + /// 4. Task B writes 100 + 5 = 105 ← overwrites A's 110 + /// Net: A's 10 XP is "lost from `pet_state`" — the xp_event row + /// is still written (different PK), so `sum(xp_event) > pet_state`. + /// + /// Empirically reproduced as ~8% XP loss in a 1.5-hour heavy + /// Claude-Code session (988 XP in xp_event vs 910 XP in + /// pet_state). + /// + /// petpet runs at most one active pet at a time, so a single + /// global mutex is fine — contention is bounded by the live + /// ingest rate (a handful of events per second worst-case). + /// If we ever support multi-pet concurrent ingest, swap this + /// for a per-pet mutex map. + apply_lock: tokio::sync::Mutex<()>, } impl StateManager { pub fn new(db: Arc) -> Self { - Self { db } + Self { + db, + apply_lock: tokio::sync::Mutex::new(()), + } } /// Apply an XP delta. The caller supplies the pet's full level /// curve + stages (loaded from pet.json snapshot). + /// + /// The whole read-modify-write of `pet_state` is held under + /// `apply_lock` so concurrent ingest cannot lose deltas — see + /// the struct comment for the race details. pub async fn apply_delta( &self, pet_id: &str, @@ -71,6 +99,7 @@ impl StateManager { occurred_at: DateTime, pet_age_days: u32, ) -> Result { + let _guard = self.apply_lock.lock().await; let prior = self.db.get_pet_state(pet_id).await?; let xp_before = prior.as_ref().map(|s| s.total_xp).unwrap_or(0); let xp_after = xp_before + delta; @@ -114,7 +143,8 @@ impl StateManager { }) } - /// Full recompute from xp_event history. + /// Full recompute from xp_event history. Takes the same lock as + /// `apply_delta` so a rebuild can't race with an in-flight ingest. pub async fn rebuild( &self, pet_id: &str, @@ -122,6 +152,7 @@ impl StateManager { stages: &[Stage], pet_age_days: u32, ) -> Result { + let _guard = self.apply_lock.lock().await; let sum = self.db.sum_xp_for_pet(pet_id).await?; let prior = self.db.get_pet_state(pet_id).await?; let xp_before = prior.as_ref().map(|s| s.total_xp).unwrap_or(0); @@ -298,4 +329,77 @@ mod tests { // L99 = max → None assert_eq!(lc.xp_for_next_level(99), None); } + + // ─── Concurrent ingest race regression ────────────────────────── + // + // Before this PR, two concurrent `apply_delta` calls could race on + // the pet_state read-modify-write — both read the same prior xp, + // both wrote prior+own_delta, the later write overwrote the + // earlier. Net: one delta's worth of XP was silently dropped from + // `pet_state` even though both `xp_event` rows landed. Empirically + // ~8% XP loss over a heavy ingest session. + // + // The fix is a single Mutex held across the read-modify-write. + // This test fires N concurrent applies of +1 XP each and asserts + // the final state matches N exactly — pre-fix, a few would be + // lost; post-fix, none. + #[tokio::test] + async fn apply_delta_serialises_concurrent_ingest() { + use crate::db::DbHandle; + + let dir = tempfile::tempdir().expect("tempdir"); + let db = DbHandle::open(&dir.path().join("test.db")) + .await + .expect("open db"); + let pet_id = "race-test-pet"; + db.insert_pet( + pet_id, + "Racer", + "test-template", + "/tmp/snap", + chrono::Utc::now(), + true, + "device", + ) + .await + .expect("insert pet"); + db.upsert_pet_state(pet_id, 0, 0, None).await.expect("init state"); + + let levels = tiered_levels_fixture(); + let stages: Vec = Vec::new(); + let mgr = Arc::new(StateManager::new(db.clone())); + let now = chrono::Utc::now(); + + // Fan out N concurrent apply_delta(+1) calls. + let n: usize = 200; + let mut handles = Vec::with_capacity(n); + for _ in 0..n { + let mgr = mgr.clone(); + let levels = levels.clone(); + let stages = stages.clone(); + let pet_id = pet_id.to_string(); + handles.push(tokio::spawn(async move { + mgr.apply_delta(&pet_id, &levels, &stages, 1, now, 0) + .await + .expect("apply_delta") + })); + } + for h in handles { + let _ = h.await.expect("task"); + } + + // After N applies of +1 each, pet_state.total_xp MUST be N. + // Pre-fix this test would fail with total_xp < N (some + // concurrent writes lost). + let final_state = db + .get_pet_state(pet_id) + .await + .expect("state") + .expect("row"); + assert_eq!( + final_state.total_xp, n as i64, + "concurrent apply_delta lost XP: expected {n}, got {}", + final_state.total_xp + ); + } }