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#[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 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
217pub 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
337pub 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
352pub 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
434pub 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}