Skip to main content

aft/db/
compression_events.rs

1use std::collections::HashMap;
2#[cfg(test)]
3use std::sync::atomic::{AtomicUsize, Ordering};
4
5use parking_lot::Mutex;
6use rusqlite::{params, Connection};
7
8pub struct CompressionEventRow<'a> {
9    pub harness: &'a str,
10    pub session_id: Option<&'a str>,
11    pub project_key: &'a str,
12    pub tool: &'a str,
13    pub task_id: Option<&'a str>,
14    pub command: Option<&'a str>,
15    pub compressor: &'a str,
16    pub original_bytes: i64,
17    pub compressed_bytes: i64,
18    pub original_tokens: u32,
19    pub compressed_tokens: u32,
20    pub created_at: i64,
21}
22
23#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, serde::Serialize)]
24pub struct CompressionAggregate {
25    pub events: u64,
26    pub original_tokens: u64,
27    pub compressed_tokens: u64,
28}
29
30impl CompressionAggregate {
31    pub fn savings_tokens(&self) -> u64 {
32        self.original_tokens.saturating_sub(self.compressed_tokens)
33    }
34
35    fn add_event(&mut self, row: &CompressionEventRow<'_>) {
36        self.events = self.events.saturating_add(1);
37        self.original_tokens = self
38            .original_tokens
39            .saturating_add(u64::from(row.original_tokens));
40        self.compressed_tokens = self
41            .compressed_tokens
42            .saturating_add(u64::from(row.compressed_tokens));
43    }
44}
45
46#[derive(Debug, Clone, PartialEq, Eq, Hash)]
47struct ProjectAggregateKey {
48    harness: String,
49    project_key: String,
50}
51
52impl ProjectAggregateKey {
53    fn new(harness: &str, project_key: &str) -> Self {
54        Self {
55            harness: harness.to_string(),
56            project_key: project_key.to_string(),
57        }
58    }
59}
60
61#[derive(Debug, Clone, PartialEq, Eq, Hash)]
62struct SessionAggregateKey {
63    project: ProjectAggregateKey,
64    session_id: String,
65}
66
67impl SessionAggregateKey {
68    fn new(harness: &str, project_key: &str, session_id: &str) -> Self {
69        Self {
70            project: ProjectAggregateKey::new(harness, project_key),
71            session_id: session_id.to_string(),
72        }
73    }
74}
75
76#[derive(Debug, Clone, Copy)]
77struct CachedAggregate {
78    aggregate: CompressionAggregate,
79    watermark: i64,
80}
81
82#[derive(Debug, Default)]
83struct CompressionAggregateCacheInner {
84    connection_identity: Option<usize>,
85    projects: HashMap<ProjectAggregateKey, CachedAggregate>,
86    sessions: HashMap<SessionAggregateKey, CachedAggregate>,
87}
88
89/// Process-local compression totals backed by the durable event table.
90///
91/// Status reads validate entries with the table's maximum row id, an indexed
92/// lookup that detects writes from other AFT processes. Full aggregate scans run
93/// only for a cold or stale key. Successful local inserts advance warm entries
94/// directly while the caller still owns the database connection mutex.
95#[derive(Debug, Default)]
96pub struct CompressionAggregateCache {
97    inner: Mutex<CompressionAggregateCacheInner>,
98    #[cfg(test)]
99    aggregate_scan_count: AtomicUsize,
100}
101
102impl CompressionAggregateCache {
103    pub fn aggregates_for_session(
104        &self,
105        conn: &Connection,
106        harness: &str,
107        project_key: &str,
108        session_id: &str,
109    ) -> rusqlite::Result<(CompressionAggregate, CompressionAggregate)> {
110        let watermark = compression_event_watermark(conn)?;
111        let project_key = ProjectAggregateKey::new(harness, project_key);
112        let session_key = SessionAggregateKey::new(harness, &project_key.project_key, session_id);
113        let mut inner = self.inner.lock();
114        reset_for_connection_change(&mut inner, conn);
115
116        let project = match inner.projects.get(&project_key) {
117            Some(cached) if cached.watermark == watermark => cached.aggregate,
118            _ => {
119                self.note_aggregate_scan();
120                let aggregate = aggregate_for_project(conn, harness, &project_key.project_key)?;
121                inner.projects.insert(
122                    project_key.clone(),
123                    CachedAggregate {
124                        aggregate,
125                        watermark,
126                    },
127                );
128                aggregate
129            }
130        };
131
132        let session = match inner.sessions.get(&session_key) {
133            Some(cached) if cached.watermark == watermark => cached.aggregate,
134            _ => {
135                self.note_aggregate_scan();
136                let aggregate =
137                    aggregate_for_session(conn, harness, &project_key.project_key, session_id)?;
138                inner.sessions.insert(
139                    session_key,
140                    CachedAggregate {
141                        aggregate,
142                        watermark,
143                    },
144                );
145                aggregate
146            }
147        };
148
149        Ok((project, session))
150    }
151
152    /// Apply a row that was inserted successfully on `conn`.
153    ///
154    /// A warm entry is advanced only when its watermark matches the row that
155    /// immediately preceded `inserted_row_id`. If another process wrote first,
156    /// the entry remains stale and the next status read rebuilds it from SQL.
157    pub fn record_successful_insert(
158        &self,
159        conn: &Connection,
160        row: &CompressionEventRow<'_>,
161        inserted_row_id: i64,
162    ) {
163        let previous_watermark = compression_event_watermark_before(conn, inserted_row_id);
164        let project_key = ProjectAggregateKey::new(row.harness, row.project_key);
165        let session_key = row
166            .session_id
167            .map(|session_id| SessionAggregateKey::new(row.harness, row.project_key, session_id));
168        let mut inner = self.inner.lock();
169        reset_for_connection_change(&mut inner, conn);
170
171        let Ok(previous_watermark) = previous_watermark else {
172            *inner = CompressionAggregateCacheInner {
173                connection_identity: inner.connection_identity,
174                ..CompressionAggregateCacheInner::default()
175            };
176            return;
177        };
178
179        for (key, cached) in &mut inner.projects {
180            if cached.watermark != previous_watermark {
181                continue;
182            }
183            if key == &project_key {
184                cached.aggregate.add_event(row);
185            }
186            cached.watermark = inserted_row_id;
187        }
188        for (key, cached) in &mut inner.sessions {
189            if cached.watermark != previous_watermark {
190                continue;
191            }
192            if session_key.as_ref() == Some(key) {
193                cached.aggregate.add_event(row);
194            }
195            cached.watermark = inserted_row_id;
196        }
197    }
198
199    pub fn clear(&self) {
200        *self.inner.lock() = CompressionAggregateCacheInner::default();
201    }
202
203    #[cfg(test)]
204    fn aggregate_scan_count_for_test(&self) -> usize {
205        self.aggregate_scan_count.load(Ordering::Relaxed)
206    }
207
208    #[cfg(test)]
209    fn note_aggregate_scan(&self) {
210        self.aggregate_scan_count.fetch_add(1, Ordering::Relaxed);
211    }
212
213    #[cfg(not(test))]
214    fn note_aggregate_scan(&self) {}
215}
216
217/// Insert one event and return its row id. Duplicate identities are ignored and
218/// return `None`, allowing in-process aggregates to advance only for durable rows.
219pub fn insert_compression_event(
220    conn: &Connection,
221    row: &CompressionEventRow<'_>,
222) -> rusqlite::Result<Option<i64>> {
223    let inserted = conn.execute(
224        r#"
225        INSERT OR IGNORE INTO compression_events (
226            harness, session_id, project_key, tool, task_id, command, compressor,
227            original_bytes, compressed_bytes, original_tokens, compressed_tokens, created_at
228        )
229        VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)
230        "#,
231        params![
232            row.harness,
233            row.session_id,
234            row.project_key,
235            row.tool,
236            row.task_id,
237            row.command,
238            row.compressor,
239            row.original_bytes,
240            row.compressed_bytes,
241            row.original_tokens,
242            row.compressed_tokens,
243            row.created_at,
244        ],
245    )?;
246    Ok((inserted > 0).then(|| conn.last_insert_rowid()))
247}
248
249pub fn aggregate_for_project(
250    conn: &Connection,
251    harness: &str,
252    project_key: &str,
253) -> rusqlite::Result<CompressionAggregate> {
254    conn.query_row(
255        r#"
256        SELECT SUM(events), SUM(original), SUM(compressed) FROM (
257            SELECT COUNT(*) AS events,
258                   COALESCE(SUM(original_tokens), 0) AS original,
259                   COALESCE(SUM(compressed_tokens), 0) AS compressed
260            FROM compression_events WHERE harness = ?1 AND project_key = ?2
261            UNION ALL
262            SELECT events, original_tokens, compressed_tokens
263            FROM compression_event_rollups WHERE harness = ?1 AND project_key = ?2
264        )
265        "#,
266        params![harness, project_key],
267        |row| {
268            Ok(CompressionAggregate {
269                events: row.get::<_, i64>(0)? as u64,
270                original_tokens: row.get::<_, i64>(1)? as u64,
271                compressed_tokens: row.get::<_, i64>(2)? as u64,
272            })
273        },
274    )
275}
276
277pub fn aggregate_for_session(
278    conn: &Connection,
279    harness: &str,
280    project_key: &str,
281    session_id: &str,
282) -> rusqlite::Result<CompressionAggregate> {
283    conn.query_row(
284        r#"
285        SELECT SUM(events), SUM(original), SUM(compressed) FROM (
286            SELECT COUNT(*) AS events,
287                   COALESCE(SUM(original_tokens), 0) AS original,
288                   COALESCE(SUM(compressed_tokens), 0) AS compressed
289            FROM compression_events
290            WHERE harness = ?1 AND project_key = ?2 AND session_id = ?3
291            UNION ALL
292            SELECT events, original_tokens, compressed_tokens
293            FROM compression_event_rollups
294            WHERE harness = ?1 AND project_key = ?2 AND session_is_null = 0 AND session_id = ?3
295        )
296        "#,
297        params![harness, project_key, session_id],
298        |row| {
299            Ok(CompressionAggregate {
300                events: row.get::<_, i64>(0)? as u64,
301                original_tokens: row.get::<_, i64>(1)? as u64,
302                compressed_tokens: row.get::<_, i64>(2)? as u64,
303            })
304        },
305    )
306}
307
308fn reset_for_connection_change(inner: &mut CompressionAggregateCacheInner, conn: &Connection) {
309    let identity = conn as *const Connection as usize;
310    if inner.connection_identity != Some(identity) {
311        *inner = CompressionAggregateCacheInner {
312            connection_identity: Some(identity),
313            ..CompressionAggregateCacheInner::default()
314        };
315    }
316}
317
318fn compression_event_watermark(conn: &Connection) -> rusqlite::Result<i64> {
319    conn.query_row(
320        "SELECT COALESCE(MAX(id), 0) FROM compression_events",
321        [],
322        |row| row.get(0),
323    )
324}
325
326fn compression_event_watermark_before(
327    conn: &Connection,
328    inserted_row_id: i64,
329) -> rusqlite::Result<i64> {
330    conn.query_row(
331        "SELECT COALESCE(MAX(id), 0) FROM compression_events WHERE id < ?1",
332        [inserted_row_id],
333        |row| row.get(0),
334    )
335}
336
337/// Raw history is kept for thirty days; lifetime counters survive in rollups.
338pub const RETENTION_AGE_MS: i64 = 30 * 24 * 60 * 60 * 1000;
339const RETENTION_BATCH: i64 = 500;
340
341const RETENTION_CANDIDATES: &str = "
342    SELECT id, created_at, harness, project_key, session_id, original_tokens, compressed_tokens,
343           task_id IS NOT NULL AND EXISTS (
344               SELECT 1 FROM bash_tasks b
345               WHERE b.harness = e.harness AND b.session_id IS e.session_id AND b.task_id = e.task_id
346                 AND b.status NOT IN ('completed', 'failed', 'killed', 'timed_out', 'fate_unknown')
347           )
348    FROM compression_events e
349    WHERE (created_at, id) > (?1, ?2) AND created_at < ?3
350    ORDER BY created_at, id LIMIT ?4";
351
352/// Fold and remove at most 500 old events in one atomic transaction.
353///
354/// The cursor bounds rows examined as well as rows deleted, so long-running
355/// tasks cannot make every sweep rescan the same protected history. It wraps
356/// after the last old row. The highest event ID stays raw to preserve the warm
357/// aggregate cache's insertion watermark. Live tasks retain their identities
358/// because they can still emit compression events; completed history outside
359/// the retention window no longer participates in duplicate suppression.
360pub fn prune_compression_events(conn: &mut Connection, now_ms: i64) -> rusqlite::Result<usize> {
361    use rusqlite::{OptionalExtension, TransactionBehavior};
362    let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
363    let (created_at, event_id) = tx
364        .query_row(
365            "SELECT created_at, event_id FROM compression_retention_cursor WHERE singleton = 1",
366            [],
367            |row| Ok((row.get::<_, i64>(0)?, row.get::<_, i64>(1)?)),
368        )
369        .optional()?
370        .unwrap_or((i64::MIN, 0));
371    let max_id = compression_event_watermark(&tx)?;
372    let candidates = tx
373        .prepare(RETENTION_CANDIDATES)?
374        .query_map(
375            params![
376                created_at,
377                event_id,
378                now_ms.saturating_sub(RETENTION_AGE_MS),
379                RETENTION_BATCH
380            ],
381            |row| {
382                Ok((
383                    row.get::<_, i64>(0)?,
384                    row.get::<_, i64>(1)?,
385                    row.get::<_, String>(2)?,
386                    row.get::<_, String>(3)?,
387                    row.get::<_, Option<String>>(4)?,
388                    row.get::<_, i64>(5)?,
389                    row.get::<_, i64>(6)?,
390                    row.get::<_, bool>(7)?,
391                ))
392            },
393        )?
394        .collect::<rusqlite::Result<Vec<_>>>()?;
395    let mut folded: HashMap<(String, String, Option<String>), (i64, i64, i64)> = HashMap::new();
396    let mut deleted = 0;
397    for (id, _, harness, project, session, original, compressed, live) in &candidates {
398        if *live || *id == max_id {
399            continue;
400        }
401        let totals = folded
402            .entry((harness.clone(), project.clone(), session.clone()))
403            .or_default();
404        totals.0 += 1;
405        totals.1 += original;
406        totals.2 += compressed;
407        deleted += tx.execute("DELETE FROM compression_events WHERE id = ?1", [id])?;
408    }
409    for ((harness, project, session), (events, original, compressed)) in folded {
410        tx.execute(
411            "INSERT INTO compression_event_rollups
412             (harness, project_key, session_is_null, session_id, events, original_tokens, compressed_tokens)
413             VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
414             ON CONFLICT(harness, project_key, session_is_null, session_id) DO UPDATE SET
415               events = events + excluded.events,
416               original_tokens = original_tokens + excluded.original_tokens,
417               compressed_tokens = compressed_tokens + excluded.compressed_tokens",
418            params![harness, project, session.is_none(), session.unwrap_or_default(), events, original, compressed],
419        )?;
420    }
421    let (next_created, next_id) = candidates
422        .last()
423        .map(|row| (row.1, row.0))
424        .unwrap_or((i64::MIN, 0));
425    tx.execute(
426        "INSERT INTO compression_retention_cursor VALUES (1, ?1, ?2)
427         ON CONFLICT(singleton) DO UPDATE SET created_at = excluded.created_at, event_id = excluded.event_id",
428        params![next_created, next_id],
429    )?;
430    tx.commit()?;
431    Ok(deleted)
432}
433
434/// Schedule bounded retention away from the daemon and standalone request loops.
435/// A process can have only one pass in flight and attempts at most once a minute.
436pub fn maybe_spawn_retention(
437    db: Option<std::sync::Arc<std::sync::Mutex<crate::db::TrackedConnection>>>,
438) {
439    use std::sync::{
440        atomic::{AtomicBool, Ordering},
441        OnceLock,
442    };
443    use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
444    static IN_FLIGHT: AtomicBool = AtomicBool::new(false);
445    static LAST: OnceLock<Mutex<Option<Instant>>> = OnceLock::new();
446    let Some(db) = db else {
447        return;
448    };
449    let mut last = LAST.get_or_init(|| Mutex::new(None)).lock();
450    if last.is_some_and(|value| value.elapsed() < Duration::from_secs(60))
451        || IN_FLIGHT.swap(true, Ordering::AcqRel)
452    {
453        return;
454    }
455    *last = Some(Instant::now());
456    if let Err(error) = std::thread::Builder::new()
457        .name("aft-compression-retention".into())
458        .spawn(move || {
459            if let Ok(mut conn) = db.try_lock() {
460                let now = SystemTime::now()
461                    .duration_since(UNIX_EPOCH)
462                    .unwrap_or_default()
463                    .as_millis();
464                match prune_compression_events(&mut conn, i64::try_from(now).unwrap_or(i64::MAX)) {
465                    Ok(0) => {}
466                    Ok(rows) => {
467                        crate::slog_info!("compression retention: folded {} raw events", rows)
468                    }
469                    Err(error) => crate::slog_warn!("compression retention failed: {}", error),
470                }
471            }
472            IN_FLIGHT.store(false, Ordering::Release);
473        })
474    {
475        IN_FLIGHT.store(false, Ordering::Release);
476        crate::slog_warn!("compression retention worker failed: {}", error);
477    }
478}
479
480#[cfg(test)]
481mod tests {
482    use super::*;
483    use tempfile::tempdir;
484
485    #[test]
486    fn retention_preserves_lifetime_totals_and_live_task_identity() {
487        let dir = tempdir().unwrap();
488        let mut conn = crate::db::open(&dir.path().join("aft.db")).unwrap();
489        let now = RETENTION_AGE_MS + 100;
490        for index in 0..510 {
491            let task = format!("task-{index}");
492            let mut event = row(
493                if index % 2 == 0 {
494                    "project-a"
495                } else {
496                    "project-b"
497                },
498                &task,
499                100,
500                40,
501                1,
502            );
503            event.session_id = match index % 3 {
504                0 => None,
505                1 => Some(""),
506                _ => Some("session-1"),
507            };
508            insert_compression_event(&conn, &event).unwrap();
509        }
510        conn.execute("INSERT INTO bash_tasks (harness, session_id, task_id, project_key, command, cwd, status, started_at)
511            VALUES ('opencode', 'session-1', 'task-2', 'project-a', 'sleep', '.', 'running', 1)", []).unwrap();
512        let recent = row("project-a", "recent", 17, 9, 100);
513        insert_compression_event(&conn, &recent).unwrap();
514        let cache = CompressionAggregateCache::default();
515        let before = ["project-a", "project-b"].map(|project| {
516            (
517                aggregate_for_project(&conn, "opencode", project).unwrap(),
518                aggregate_for_session(&conn, "opencode", project, "session-1").unwrap(),
519                aggregate_for_session(&conn, "opencode", project, "").unwrap(),
520            )
521        });
522        let warm = cache
523            .aggregates_for_session(&conn, "opencode", "project-a", "session-1")
524            .unwrap();
525        assert_eq!(prune_compression_events(&mut conn, now).unwrap(), 499);
526        assert_eq!(prune_compression_events(&mut conn, now).unwrap(), 10);
527        assert_eq!(prune_compression_events(&mut conn, now).unwrap(), 0);
528        for (index, project) in ["project-a", "project-b"].iter().enumerate() {
529            assert_eq!(
530                aggregate_for_project(&conn, "opencode", project).unwrap(),
531                before[index].0
532            );
533            assert_eq!(
534                aggregate_for_session(&conn, "opencode", project, "session-1").unwrap(),
535                before[index].1
536            );
537            assert_eq!(
538                aggregate_for_session(&conn, "opencode", project, "").unwrap(),
539                before[index].2
540            );
541        }
542        assert_eq!(
543            cache
544                .aggregates_for_session(&conn, "opencode", "project-a", "session-1")
545                .unwrap(),
546            warm
547        );
548        assert!(insert_compression_event(&conn, &recent).unwrap().is_none());
549        assert!(
550            insert_compression_event(&conn, &row("project-a", "task-2", 100, 40, 1))
551                .unwrap()
552                .is_none()
553        );
554        conn.execute("UPDATE bash_tasks SET status = 'completed'", [])
555            .unwrap();
556        assert_eq!(prune_compression_events(&mut conn, now).unwrap(), 1);
557        assert_eq!(
558            aggregate_for_project(&conn, "opencode", "project-a").unwrap(),
559            before[0].0
560        );
561        let next = row("project-a", "next", 20, 10, now);
562        let id = insert_compression_event(&conn, &next).unwrap().unwrap();
563        cache.record_successful_insert(&conn, &next, id);
564        let totals = cache
565            .aggregates_for_session(&conn, "opencode", "project-a", "session-1")
566            .unwrap();
567        assert_eq!(
568            totals.0,
569            aggregate_for_project(&conn, "opencode", "project-a").unwrap()
570        );
571        assert_eq!(
572            totals.1,
573            aggregate_for_session(&conn, "opencode", "project-a", "session-1").unwrap()
574        );
575    }
576
577    #[test]
578    fn retention_rollup_failure_rolls_back_raw_deletes_and_cursor() {
579        let dir = tempdir().unwrap();
580        let mut conn = crate::db::open(&dir.path().join("aft.db")).unwrap();
581        insert_compression_event(&conn, &row("project-a", "old", 100, 40, 1)).unwrap();
582        insert_compression_event(&conn, &row("project-a", "watermark", 100, 40, 1)).unwrap();
583        conn.execute_batch("CREATE TRIGGER reject_fold BEFORE INSERT ON compression_event_rollups BEGIN SELECT RAISE(ABORT, 'fold failure'); END;").unwrap();
584        assert!(prune_compression_events(&mut conn, RETENTION_AGE_MS + 100).is_err());
585        assert_eq!(
586            conn.query_row("SELECT count(*) FROM compression_events", [], |r| r
587                .get::<_, i64>(0))
588                .unwrap(),
589            2
590        );
591        assert_eq!(
592            conn.query_row(
593                "SELECT count(*) FROM compression_retention_cursor",
594                [],
595                |r| r.get::<_, i64>(0)
596            )
597            .unwrap(),
598            0
599        );
600    }
601
602    #[test]
603    fn retention_selection_is_indexed_and_keeps_the_watermark() {
604        let dir = tempdir().unwrap();
605        let mut conn = crate::db::open(&dir.path().join("aft.db")).unwrap();
606        insert_compression_event(&conn, &row("project-a", "watermark", 100, 40, 1)).unwrap();
607        let plan = conn
608            .prepare(&format!("EXPLAIN QUERY PLAN {RETENTION_CANDIDATES}"))
609            .unwrap()
610            .query_map(
611                params![i64::MIN, 0, RETENTION_AGE_MS, RETENTION_BATCH],
612                |r| r.get::<_, String>(3),
613            )
614            .unwrap()
615            .collect::<rusqlite::Result<Vec<_>>>()
616            .unwrap()
617            .join("\n");
618        assert!(plan.contains("idx_compression_created"), "{plan}");
619        assert!(!plan.contains("TEMP B-TREE"), "{plan}");
620        assert_eq!(
621            prune_compression_events(&mut conn, RETENTION_AGE_MS + 100).unwrap(),
622            0
623        );
624        assert_eq!(compression_event_watermark(&conn).unwrap(), 1);
625    }
626
627    #[test]
628    fn duplicate_identity_is_ignored_without_cross_project_suppression() {
629        let dir = tempdir().expect("tempdir");
630        let conn = crate::db::open(&dir.path().join("aft.db")).expect("open db");
631
632        assert!(
633            insert_compression_event(&conn, &row("project-a", "task-1", 100, 40, 1))
634                .expect("insert first")
635                .is_some()
636        );
637        assert!(
638            insert_compression_event(&conn, &row("project-a", "task-1", 900, 10, 2))
639                .expect("ignore duplicate")
640                .is_none()
641        );
642        assert!(
643            insert_compression_event(&conn, &row("project-b", "task-1", 200, 80, 3))
644                .expect("insert same task id for other project")
645                .is_some()
646        );
647
648        let project_a = aggregate_for_project(&conn, "opencode", "project-a").unwrap();
649        assert_eq!(project_a.events, 1);
650        assert_eq!(project_a.original_tokens, 100);
651        assert_eq!(project_a.compressed_tokens, 40);
652
653        let project_b = aggregate_for_project(&conn, "opencode", "project-b").unwrap();
654        assert_eq!(project_b.events, 1);
655        assert_eq!(project_b.original_tokens, 200);
656        assert_eq!(project_b.compressed_tokens, 80);
657    }
658
659    #[test]
660    fn cached_aggregates_match_sql_after_generated_inserts_and_duplicates() {
661        let dir = tempdir().expect("tempdir");
662        let conn = crate::db::open(&dir.path().join("aft.db")).expect("open db");
663        let cache = CompressionAggregateCache::default();
664        let (project, session) = cache
665            .aggregates_for_session(&conn, "opencode", "project-a", "session-1")
666            .expect("warm cache");
667        assert_eq!(project, CompressionAggregate::default());
668        assert_eq!(session, CompressionAggregate::default());
669        cache
670            .aggregates_for_session(&conn, "opencode", "project-a", "session-2")
671            .expect("warm sibling session");
672        cache
673            .aggregates_for_session(&conn, "opencode", "project-b", "session-1")
674            .expect("warm sibling project");
675        assert_eq!(cache.aggregate_scan_count_for_test(), 5);
676
677        let mut previous_task = String::new();
678        for index in 0..64u32 {
679            let task_id = if index % 5 == 4 {
680                previous_task.clone()
681            } else {
682                let task_id = format!("task-{index}");
683                previous_task = task_id.clone();
684                task_id
685            };
686            let row = row(
687                "project-a",
688                &task_id,
689                100 + index,
690                40 + (index % 17),
691                i64::from(index),
692            );
693            if let Some(row_id) = insert_compression_event(&conn, &row).expect("insert event") {
694                cache.record_successful_insert(&conn, &row, row_id);
695            }
696
697            for (project_key, session_id) in [
698                ("project-a", "session-1"),
699                ("project-a", "session-2"),
700                ("project-b", "session-1"),
701            ] {
702                let cached = cache
703                    .aggregates_for_session(&conn, "opencode", project_key, session_id)
704                    .expect("read cache");
705                let scanned = (
706                    aggregate_for_project(&conn, "opencode", project_key).expect("scan project"),
707                    aggregate_for_session(&conn, "opencode", project_key, session_id)
708                        .expect("scan session"),
709                );
710                assert_eq!(cached, scanned, "aggregate mismatch after step {index}");
711            }
712            assert_eq!(
713                cache.aggregate_scan_count_for_test(),
714                5,
715                "local inserts must advance warm entries without rescanning"
716            );
717        }
718    }
719
720    fn row<'a>(
721        project_key: &'a str,
722        task_id: &'a str,
723        original_tokens: u32,
724        compressed_tokens: u32,
725        created_at: i64,
726    ) -> CompressionEventRow<'a> {
727        CompressionEventRow {
728            harness: "opencode",
729            session_id: Some("session-1"),
730            project_key,
731            tool: "bash",
732            task_id: Some(task_id),
733            command: Some("echo ok"),
734            compressor: "zstd",
735            original_bytes: i64::from(original_tokens) * 4,
736            compressed_bytes: i64::from(compressed_tokens) * 4,
737            original_tokens,
738            compressed_tokens,
739            created_at,
740        }
741    }
742}