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;
340const BASH_TASK_MUTATION_BATCH: usize = 250;
343
344#[derive(Debug, Clone, Copy, PartialEq, Eq)]
345pub struct RetentionTick {
346 pub bash_tasks: crate::db::bash_tasks::TerminalRowsPrune,
347 pub compression_events_removed: usize,
348}
349
350#[derive(Debug, Clone, Copy, PartialEq, Eq)]
351pub struct RetentionPhaseTimings {
352 pub selection_lock_micros: u128,
353 pub stat_micros: u128,
354 pub task_delete_micros: u128,
355 pub event_prune_micros: u128,
356 pub commit_micros: u128,
357 pub mutation_lock_micros: u128,
358}
359
360impl RetentionPhaseTimings {
361 pub fn total_lock_micros(self) -> u128 {
362 self.selection_lock_micros
363 .saturating_add(self.mutation_lock_micros)
364 }
365}
366
367#[derive(Debug, Clone, Copy, PartialEq, Eq)]
368pub struct RetentionPass {
369 pub tick: RetentionTick,
370 pub timings: RetentionPhaseTimings,
371}
372
373const RETENTION_CANDIDATES: &str = "
374 SELECT id, created_at, harness, project_key, session_id, original_tokens, compressed_tokens,
375 task_id IS NOT NULL AND EXISTS (
376 SELECT 1 FROM bash_tasks b
377 WHERE b.harness = e.harness AND b.session_id IS e.session_id AND b.task_id = e.task_id
378 AND b.status NOT IN ('completed', 'failed', 'killed', 'timed_out', 'fate_unknown')
379 )
380 FROM compression_events e
381 WHERE (created_at, id) > (?1, ?2) AND created_at < ?3
382 ORDER BY created_at, id LIMIT ?4";
383
384pub fn prune_compression_events(conn: &mut Connection, now_ms: i64) -> rusqlite::Result<usize> {
393 use rusqlite::TransactionBehavior;
394 let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
395 let deleted = prune_compression_events_in_transaction(&tx, now_ms)?;
396 tx.commit()?;
397 Ok(deleted)
398}
399
400fn prune_compression_events_in_transaction(
401 conn: &Connection,
402 now_ms: i64,
403) -> rusqlite::Result<usize> {
404 use rusqlite::OptionalExtension;
405 let (created_at, event_id) = conn
406 .query_row(
407 "SELECT created_at, event_id FROM compression_retention_cursor WHERE singleton = 1",
408 [],
409 |row| Ok((row.get::<_, i64>(0)?, row.get::<_, i64>(1)?)),
410 )
411 .optional()?
412 .unwrap_or((i64::MIN, 0));
413 let max_id = compression_event_watermark(conn)?;
414 let candidates = conn
415 .prepare(RETENTION_CANDIDATES)?
416 .query_map(
417 params![
418 created_at,
419 event_id,
420 now_ms.saturating_sub(RETENTION_AGE_MS),
421 RETENTION_BATCH
422 ],
423 |row| {
424 Ok((
425 row.get::<_, i64>(0)?,
426 row.get::<_, i64>(1)?,
427 row.get::<_, String>(2)?,
428 row.get::<_, String>(3)?,
429 row.get::<_, Option<String>>(4)?,
430 row.get::<_, i64>(5)?,
431 row.get::<_, i64>(6)?,
432 row.get::<_, bool>(7)?,
433 ))
434 },
435 )?
436 .collect::<rusqlite::Result<Vec<_>>>()?;
437 let mut folded: HashMap<(String, String, Option<String>), (i64, i64, i64)> = HashMap::new();
438 let mut deleted = 0;
439 for (id, _, harness, project, session, original, compressed, live) in &candidates {
440 if *live || *id == max_id {
441 continue;
442 }
443 let totals = folded
444 .entry((harness.clone(), project.clone(), session.clone()))
445 .or_default();
446 totals.0 += 1;
447 totals.1 += original;
448 totals.2 += compressed;
449 deleted += conn.execute("DELETE FROM compression_events WHERE id = ?1", [id])?;
450 }
451 for ((harness, project, session), (events, original, compressed)) in folded {
452 conn.execute(
453 "INSERT INTO compression_event_rollups
454 (harness, project_key, session_is_null, session_id, events, original_tokens, compressed_tokens)
455 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
456 ON CONFLICT(harness, project_key, session_is_null, session_id) DO UPDATE SET
457 events = events + excluded.events,
458 original_tokens = original_tokens + excluded.original_tokens,
459 compressed_tokens = compressed_tokens + excluded.compressed_tokens",
460 params![harness, project, session.is_none(), session.unwrap_or_default(), events, original, compressed],
461 )?;
462 }
463 let (next_created, next_id) = candidates
464 .last()
465 .map(|row| (row.1, row.0))
466 .unwrap_or((i64::MIN, 0));
467 conn.execute(
468 "INSERT INTO compression_retention_cursor VALUES (1, ?1, ?2)
469 ON CONFLICT(singleton) DO UPDATE SET created_at = excluded.created_at, event_id = excluded.event_id",
470 params![next_created, next_id],
471 )?;
472 Ok(deleted)
473}
474
475pub fn prune_retention_tick(conn: &mut Connection, now_ms: i64) -> rusqlite::Result<RetentionTick> {
476 let plan = crate::db::bash_tasks::select_terminal_prune_candidates(conn, now_ms, 500)?;
477 let prepared = crate::db::bash_tasks::prepare_terminal_prune(plan, |_| false);
478 apply_prepared_retention_tick(conn, now_ms, prepared)
479}
480
481struct RetentionMutation {
482 tick: RetentionTick,
483 task_delete_micros: u128,
484 event_prune_micros: u128,
485 commit_micros: u128,
486}
487
488fn apply_prepared_retention_tick(
489 conn: &mut Connection,
490 now_ms: i64,
491 prepared: crate::db::bash_tasks::PreparedTerminalPrune,
492) -> rusqlite::Result<RetentionTick> {
493 apply_prepared_retention_tick_timed(conn, now_ms, prepared).map(|mutation| mutation.tick)
494}
495
496fn apply_prepared_retention_tick_timed(
497 conn: &mut Connection,
498 now_ms: i64,
499 prepared: crate::db::bash_tasks::PreparedTerminalPrune,
500) -> rusqlite::Result<RetentionMutation> {
501 use rusqlite::TransactionBehavior;
502 use std::time::Instant;
503
504 let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
505 let task_delete_started = Instant::now();
508 let bash_tasks = crate::db::bash_tasks::delete_prepared_terminal_rows(&tx, prepared)?;
509 let task_delete_micros = task_delete_started.elapsed().as_micros();
510 let event_prune_started = Instant::now();
511 let compression_events_removed = prune_compression_events_in_transaction(&tx, now_ms)?;
512 let event_prune_micros = event_prune_started.elapsed().as_micros();
513 let commit_started = Instant::now();
514 tx.commit()?;
515 let commit_micros = commit_started.elapsed().as_micros();
516 Ok(RetentionMutation {
517 tick: RetentionTick {
518 bash_tasks,
519 compression_events_removed,
520 },
521 task_delete_micros,
522 event_prune_micros,
523 commit_micros,
524 })
525}
526
527pub fn prune_retention_once(
528 db: &std::sync::Arc<std::sync::Mutex<crate::db::TrackedConnection>>,
529 now_ms: i64,
530 registries: Option<&[crate::bash_background::BgTaskRegistry]>,
531) -> Result<Option<RetentionPass>, String> {
532 prune_retention_once_observed(db, now_ms, registries, || {})
533}
534
535fn prune_retention_once_observed(
536 db: &std::sync::Arc<std::sync::Mutex<crate::db::TrackedConnection>>,
537 now_ms: i64,
538 registries: Option<&[crate::bash_background::BgTaskRegistry]>,
539 observe_stat_phase: impl FnOnce(),
540) -> Result<Option<RetentionPass>, String> {
541 use std::sync::TryLockError;
542 use std::time::Instant;
543
544 let conn = match db.try_lock() {
545 Ok(conn) => conn,
546 Err(TryLockError::WouldBlock) => return Ok(None),
547 Err(TryLockError::Poisoned(_)) => {
548 return Err("retention database mutex poisoned".to_string())
549 }
550 };
551 let selection_started = Instant::now();
552 let plan = crate::db::bash_tasks::select_terminal_prune_candidates(&conn, now_ms, 500)
553 .map_err(|error| error.to_string())?;
554 drop(conn);
555 let selection_lock_micros = selection_started.elapsed().as_micros();
556
557 let stat_started = Instant::now();
558 let mut prepared = crate::db::bash_tasks::prepare_terminal_prune_observed(
559 plan,
560 |task_id| {
561 registries.is_none_or(|registries| {
562 registries
563 .iter()
564 .any(|registry| registry.active_watch_count(task_id) > 0)
565 })
566 },
567 observe_stat_phase,
568 );
569 crate::db::bash_tasks::cap_prepared_terminal_rows(&mut prepared, BASH_TASK_MUTATION_BATCH);
570 let stat_micros = stat_started.elapsed().as_micros();
571
572 let mut conn = match db.try_lock() {
573 Ok(conn) => conn,
574 Err(TryLockError::WouldBlock) => return Ok(None),
575 Err(TryLockError::Poisoned(_)) => {
576 return Err("retention database mutex poisoned".to_string())
577 }
578 };
579 let mutation_started = Instant::now();
580 let mutation = apply_prepared_retention_tick_timed(&mut conn, now_ms, prepared)
581 .map_err(|error| error.to_string())?;
582 drop(conn);
583 let mutation_lock_micros = mutation_started.elapsed().as_micros();
584
585 Ok(Some(RetentionPass {
586 tick: mutation.tick,
587 timings: RetentionPhaseTimings {
588 selection_lock_micros,
589 stat_micros,
590 task_delete_micros: mutation.task_delete_micros,
591 event_prune_micros: mutation.event_prune_micros,
592 commit_micros: mutation.commit_micros,
593 mutation_lock_micros,
594 },
595 }))
596}
597
598pub fn maybe_spawn_retention(
601 db: Option<std::sync::Arc<std::sync::Mutex<crate::db::TrackedConnection>>>,
602 registries: Option<Vec<crate::bash_background::BgTaskRegistry>>,
603) {
604 use std::sync::{
605 atomic::{AtomicBool, Ordering},
606 OnceLock,
607 };
608 use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
609 static IN_FLIGHT: AtomicBool = AtomicBool::new(false);
610 static LAST: OnceLock<Mutex<Option<Instant>>> = OnceLock::new();
611 let Some(db) = db else {
612 return;
613 };
614 let mut last = LAST.get_or_init(|| Mutex::new(None)).lock();
615 if last.is_some_and(|value| value.elapsed() < Duration::from_secs(60))
616 || IN_FLIGHT.swap(true, Ordering::AcqRel)
617 {
618 return;
619 }
620 *last = Some(Instant::now());
621 if let Err(error) = std::thread::Builder::new()
622 .name("aft-retention".into())
623 .spawn(move || {
624 let now = SystemTime::now()
625 .duration_since(UNIX_EPOCH)
626 .unwrap_or_default()
627 .as_millis();
628 match prune_retention_once(
629 &db,
630 i64::try_from(now).unwrap_or(i64::MAX),
631 registries.as_deref(),
632 ) {
633 Ok(Some(pass)) => {
634 crate::slog_info!(
635 "bash task retention: removed={} remaining_candidates={}",
636 pass.tick.bash_tasks.removed,
637 pass.tick.bash_tasks.remaining_candidates
638 );
639 if pass.tick.compression_events_removed > 0 {
640 crate::slog_info!(
641 "compression retention: folded {} raw events",
642 pass.tick.compression_events_removed
643 );
644 }
645 }
646 Ok(None) => {}
647 Err(error) => crate::slog_warn!("retention failed: {}", error),
648 }
649 IN_FLIGHT.store(false, Ordering::Release);
650 })
651 {
652 IN_FLIGHT.store(false, Ordering::Release);
653 crate::slog_warn!("compression retention worker failed: {}", error);
654 }
655}
656
657#[cfg(test)]
658mod tests {
659 use super::*;
660 use tempfile::tempdir;
661
662 #[test]
663 fn retention_releases_database_mutex_before_layout_stats() {
664 let dir = tempdir().unwrap();
665 let conn = crate::db::open(&dir.path().join("aft.db")).unwrap();
666 conn.execute(
667 "INSERT INTO bash_tasks (
668 harness, session_id, task_id, project_key, command, cwd, status,
669 started_at, completed_at, completion_delivered
670 ) VALUES ('opencode', 'session', 'bash-0000000000000001', 'project',
671 'true', '.', 'completed', 1, 1, 1)",
672 [],
673 )
674 .unwrap();
675 let db = std::sync::Arc::new(std::sync::Mutex::new(conn));
676 let observed = std::sync::atomic::AtomicBool::new(false);
677
678 let pass = prune_retention_once_observed(
679 &db,
680 crate::db::bash_tasks::TERMINAL_ROW_RETENTION_AGE_MS + 100,
681 Some(&[]),
682 || {
683 let _guard = db
684 .try_lock()
685 .expect("database mutex held during layout stat phase");
686 observed.store(true, Ordering::SeqCst);
687 },
688 )
689 .unwrap()
690 .unwrap();
691
692 assert!(observed.load(Ordering::SeqCst));
693 assert_eq!(pass.tick.bash_tasks.removed, 1);
694 }
695
696 #[test]
697 fn retention_preserves_lifetime_totals_and_live_task_identity() {
698 let dir = tempdir().unwrap();
699 let mut conn = crate::db::open(&dir.path().join("aft.db")).unwrap();
700 let now = RETENTION_AGE_MS + 100;
701 for index in 0..510 {
702 let task = format!("task-{index}");
703 let mut event = row(
704 if index % 2 == 0 {
705 "project-a"
706 } else {
707 "project-b"
708 },
709 &task,
710 100,
711 40,
712 1,
713 );
714 event.session_id = match index % 3 {
715 0 => None,
716 1 => Some(""),
717 _ => Some("session-1"),
718 };
719 insert_compression_event(&conn, &event).unwrap();
720 }
721 conn.execute("INSERT INTO bash_tasks (harness, session_id, task_id, project_key, command, cwd, status, started_at)
722 VALUES ('opencode', 'session-1', 'task-2', 'project-a', 'sleep', '.', 'running', 1)", []).unwrap();
723 let recent = row("project-a", "recent", 17, 9, 100);
724 insert_compression_event(&conn, &recent).unwrap();
725 let cache = CompressionAggregateCache::default();
726 let before = ["project-a", "project-b"].map(|project| {
727 (
728 aggregate_for_project(&conn, "opencode", project).unwrap(),
729 aggregate_for_session(&conn, "opencode", project, "session-1").unwrap(),
730 aggregate_for_session(&conn, "opencode", project, "").unwrap(),
731 )
732 });
733 let warm = cache
734 .aggregates_for_session(&conn, "opencode", "project-a", "session-1")
735 .unwrap();
736 assert_eq!(prune_compression_events(&mut conn, now).unwrap(), 499);
737 assert_eq!(prune_compression_events(&mut conn, now).unwrap(), 10);
738 assert_eq!(prune_compression_events(&mut conn, now).unwrap(), 0);
739 for (index, project) in ["project-a", "project-b"].iter().enumerate() {
740 assert_eq!(
741 aggregate_for_project(&conn, "opencode", project).unwrap(),
742 before[index].0
743 );
744 assert_eq!(
745 aggregate_for_session(&conn, "opencode", project, "session-1").unwrap(),
746 before[index].1
747 );
748 assert_eq!(
749 aggregate_for_session(&conn, "opencode", project, "").unwrap(),
750 before[index].2
751 );
752 }
753 assert_eq!(
754 cache
755 .aggregates_for_session(&conn, "opencode", "project-a", "session-1")
756 .unwrap(),
757 warm
758 );
759 assert!(insert_compression_event(&conn, &recent).unwrap().is_none());
760 assert!(
761 insert_compression_event(&conn, &row("project-a", "task-2", 100, 40, 1))
762 .unwrap()
763 .is_none()
764 );
765 conn.execute("UPDATE bash_tasks SET status = 'completed'", [])
766 .unwrap();
767 assert_eq!(prune_compression_events(&mut conn, now).unwrap(), 1);
768 assert_eq!(
769 aggregate_for_project(&conn, "opencode", "project-a").unwrap(),
770 before[0].0
771 );
772 let next = row("project-a", "next", 20, 10, now);
773 let id = insert_compression_event(&conn, &next).unwrap().unwrap();
774 cache.record_successful_insert(&conn, &next, id);
775 let totals = cache
776 .aggregates_for_session(&conn, "opencode", "project-a", "session-1")
777 .unwrap();
778 assert_eq!(
779 totals.0,
780 aggregate_for_project(&conn, "opencode", "project-a").unwrap()
781 );
782 assert_eq!(
783 totals.1,
784 aggregate_for_session(&conn, "opencode", "project-a", "session-1").unwrap()
785 );
786 }
787
788 #[test]
789 fn retention_rollup_failure_rolls_back_raw_deletes_and_cursor() {
790 let dir = tempdir().unwrap();
791 let mut conn = crate::db::open(&dir.path().join("aft.db")).unwrap();
792 insert_compression_event(&conn, &row("project-a", "old", 100, 40, 1)).unwrap();
793 insert_compression_event(&conn, &row("project-a", "watermark", 100, 40, 1)).unwrap();
794 conn.execute_batch("CREATE TRIGGER reject_fold BEFORE INSERT ON compression_event_rollups BEGIN SELECT RAISE(ABORT, 'fold failure'); END;").unwrap();
795 assert!(prune_compression_events(&mut conn, RETENTION_AGE_MS + 100).is_err());
796 assert_eq!(
797 conn.query_row("SELECT count(*) FROM compression_events", [], |r| r
798 .get::<_, i64>(0))
799 .unwrap(),
800 2
801 );
802 assert_eq!(
803 conn.query_row(
804 "SELECT count(*) FROM compression_retention_cursor",
805 [],
806 |r| r.get::<_, i64>(0)
807 )
808 .unwrap(),
809 0
810 );
811 }
812
813 #[test]
814 fn retention_selection_is_indexed_and_keeps_the_watermark() {
815 let dir = tempdir().unwrap();
816 let mut conn = crate::db::open(&dir.path().join("aft.db")).unwrap();
817 insert_compression_event(&conn, &row("project-a", "watermark", 100, 40, 1)).unwrap();
818 let plan = conn
819 .prepare(&format!("EXPLAIN QUERY PLAN {RETENTION_CANDIDATES}"))
820 .unwrap()
821 .query_map(
822 params![i64::MIN, 0, RETENTION_AGE_MS, RETENTION_BATCH],
823 |r| r.get::<_, String>(3),
824 )
825 .unwrap()
826 .collect::<rusqlite::Result<Vec<_>>>()
827 .unwrap()
828 .join("\n");
829 assert!(plan.contains("idx_compression_created"), "{plan}");
830 assert!(!plan.contains("TEMP B-TREE"), "{plan}");
831 assert_eq!(
832 prune_compression_events(&mut conn, RETENTION_AGE_MS + 100).unwrap(),
833 0
834 );
835 assert_eq!(compression_event_watermark(&conn).unwrap(), 1);
836 }
837
838 #[test]
839 fn duplicate_identity_is_ignored_without_cross_project_suppression() {
840 let dir = tempdir().expect("tempdir");
841 let conn = crate::db::open(&dir.path().join("aft.db")).expect("open db");
842
843 assert!(
844 insert_compression_event(&conn, &row("project-a", "task-1", 100, 40, 1))
845 .expect("insert first")
846 .is_some()
847 );
848 assert!(
849 insert_compression_event(&conn, &row("project-a", "task-1", 900, 10, 2))
850 .expect("ignore duplicate")
851 .is_none()
852 );
853 assert!(
854 insert_compression_event(&conn, &row("project-b", "task-1", 200, 80, 3))
855 .expect("insert same task id for other project")
856 .is_some()
857 );
858
859 let project_a = aggregate_for_project(&conn, "opencode", "project-a").unwrap();
860 assert_eq!(project_a.events, 1);
861 assert_eq!(project_a.original_tokens, 100);
862 assert_eq!(project_a.compressed_tokens, 40);
863
864 let project_b = aggregate_for_project(&conn, "opencode", "project-b").unwrap();
865 assert_eq!(project_b.events, 1);
866 assert_eq!(project_b.original_tokens, 200);
867 assert_eq!(project_b.compressed_tokens, 80);
868 }
869
870 #[test]
871 fn cached_aggregates_match_sql_after_generated_inserts_and_duplicates() {
872 let dir = tempdir().expect("tempdir");
873 let conn = crate::db::open(&dir.path().join("aft.db")).expect("open db");
874 let cache = CompressionAggregateCache::default();
875 let (project, session) = cache
876 .aggregates_for_session(&conn, "opencode", "project-a", "session-1")
877 .expect("warm cache");
878 assert_eq!(project, CompressionAggregate::default());
879 assert_eq!(session, CompressionAggregate::default());
880 cache
881 .aggregates_for_session(&conn, "opencode", "project-a", "session-2")
882 .expect("warm sibling session");
883 cache
884 .aggregates_for_session(&conn, "opencode", "project-b", "session-1")
885 .expect("warm sibling project");
886 assert_eq!(cache.aggregate_scan_count_for_test(), 5);
887
888 let mut previous_task = String::new();
889 for index in 0..64u32 {
890 let task_id = if index % 5 == 4 {
891 previous_task.clone()
892 } else {
893 let task_id = format!("task-{index}");
894 previous_task = task_id.clone();
895 task_id
896 };
897 let row = row(
898 "project-a",
899 &task_id,
900 100 + index,
901 40 + (index % 17),
902 i64::from(index),
903 );
904 if let Some(row_id) = insert_compression_event(&conn, &row).expect("insert event") {
905 cache.record_successful_insert(&conn, &row, row_id);
906 }
907
908 for (project_key, session_id) in [
909 ("project-a", "session-1"),
910 ("project-a", "session-2"),
911 ("project-b", "session-1"),
912 ] {
913 let cached = cache
914 .aggregates_for_session(&conn, "opencode", project_key, session_id)
915 .expect("read cache");
916 let scanned = (
917 aggregate_for_project(&conn, "opencode", project_key).expect("scan project"),
918 aggregate_for_session(&conn, "opencode", project_key, session_id)
919 .expect("scan session"),
920 );
921 assert_eq!(cached, scanned, "aggregate mismatch after step {index}");
922 }
923 assert_eq!(
924 cache.aggregate_scan_count_for_test(),
925 5,
926 "local inserts must advance warm entries without rescanning"
927 );
928 }
929 }
930
931 fn row<'a>(
932 project_key: &'a str,
933 task_id: &'a str,
934 original_tokens: u32,
935 compressed_tokens: u32,
936 created_at: i64,
937 ) -> CompressionEventRow<'a> {
938 CompressionEventRow {
939 harness: "opencode",
940 session_id: Some("session-1"),
941 project_key,
942 tool: "bash",
943 task_id: Some(task_id),
944 command: Some("echo ok"),
945 compressor: "zstd",
946 original_bytes: i64::from(original_tokens) * 4,
947 compressed_bytes: i64::from(compressed_tokens) * 4,
948 original_tokens,
949 compressed_tokens,
950 created_at,
951 }
952 }
953}