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// Probe up to 500 filesystem rows without holding the shared database mutex.
341// Capping the write phase at 250 keeps total connection lock holds below 100 ms.
342const 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
384/// Fold and remove at most 500 old events in one atomic transaction.
385///
386/// The cursor bounds rows examined as well as rows deleted, so long-running
387/// tasks cannot make every sweep rescan the same protected history. It wraps
388/// after the last old row. The highest event ID stays raw to preserve the warm
389/// aggregate cache's insertion watermark. Live tasks retain their identities
390/// because they can still emit compression events; completed history outside
391/// the retention window no longer participates in duplicate suppression.
392pub 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    // Task rows go first so the event pass in this transaction observes their
506    // final liveness, while the commit publishes both retention decisions at once.
507    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
598/// Schedule bounded retention away from the daemon and standalone request loops.
599/// A process can have only one pass in flight and attempts at most once a minute.
600pub 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}