Skip to main content

remem/db/query/
stats.rs

1use anyhow::Result;
2use rusqlite::{params, Connection, OptionalExtension};
3
4use crate::db::{
5    AiUsageBreakdown, AiUsageSourceTotals, AiUsageTotals, DailyAiUsage, FailureLifecycleStats,
6    WeeklyAiUsage,
7};
8
9pub use super::legacy_surfaces::LegacySurfaceStats;
10use super::shared::collect_rows;
11
12#[derive(Debug, Clone, PartialEq, Eq)]
13pub struct SystemStats {
14    pub active_memories: i64,
15    pub active_observations: i64,
16    pub total_observations: i64,
17    pub session_summaries: i64,
18    pub raw_messages: i64,
19    pub raw_ingest_failures: i64,
20    pub raw_ingest_parse_errors: i64,
21    pub raw_ingest_insert_errors: i64,
22    pub latest_raw_ingest_failure_epoch: Option<i64>,
23    pub latest_raw_ingest_failure_kind: Option<String>,
24    pub latest_raw_ingest_failure_path: Option<String>,
25    pub latest_raw_ingest_failure_message: Option<String>,
26    pub captured_events: i64,
27    pub latest_captured_event_epoch: Option<i64>,
28    pub latest_capture_activity_epoch: Option<i64>,
29    pub capture_drop_events: i64,
30    pub actionable_capture_drops: i64,
31    pub unrecovered_capture_spills: i64,
32    pub latest_capture_drop_epoch: Option<i64>,
33    pub latest_capture_drop_reason: Option<String>,
34    pub latest_capture_drop_detail: Option<String>,
35    pub pending_extraction_tasks: i64,
36    pub processing_extraction_tasks: i64,
37    pub expired_processing_extraction_tasks: i64,
38    pub failed_extraction_tasks: i64,
39    pub retryable_extraction_replay_ranges: i64,
40    pub active_extraction_replay_ranges: i64,
41    pub quarantined_extraction_replay_ranges: i64,
42    pub oldest_pending_extraction_epoch: Option<i64>,
43    pub pending_memory_candidates: i64,
44    pub total_memory_candidates: i64,
45    pub promoted_memory_candidates: i64,
46    pub pending_review_memory_candidates: i64,
47    pub pending_graph_candidates: i64,
48    pub pending_observations: i64,
49    pub ready_pending_observations: i64,
50    pub delayed_pending_observations: i64,
51    pub processing_pending_observations: i64,
52    pub expired_processing_pending_observations: i64,
53    pub failed_pending_observations: i64,
54    pub oldest_ready_pending_epoch: Option<i64>,
55    pub pending_jobs: i64,
56    pub processing_jobs: i64,
57    pub failed_jobs: i64,
58    pub stuck_jobs: i64,
59    pub failure_lifecycle: FailureLifecycleStats,
60    pub worker_daemon_healthy: bool,
61    pub worker_heartbeat_owner: Option<String>,
62    pub worker_heartbeat_age_secs: Option<i64>,
63    pub legacy_surfaces: Vec<LegacySurfaceStats>,
64}
65
66#[derive(Debug, Clone, PartialEq, Eq)]
67pub struct MemoryFactsStats {
68    pub table_exists: bool,
69    pub total: i64,
70    pub retrieval_eligible: i64,
71    pub active_memories: i64,
72    pub captured_events: i64,
73}
74
75#[derive(Debug, Clone, PartialEq, Eq)]
76pub struct DailyActivityStats {
77    pub memories: i64,
78    pub observations: i64,
79}
80
81#[derive(Debug, Clone, PartialEq, Eq)]
82pub struct ProjectCount {
83    pub project: String,
84    pub count: i64,
85}
86
87#[derive(Debug, Clone, PartialEq, Eq)]
88pub struct CandidatePromotionStat {
89    pub source_kind: String,
90    pub review_status: String,
91    pub block_reason: Option<String>,
92    pub total: i64,
93    pub last_7_days: i64,
94}
95
96pub fn query_system_stats(conn: &Connection) -> Result<SystemStats> {
97    let now = chrono::Utc::now().timestamp();
98    let raw_ingest = query_raw_ingest_failure_stats(conn)?;
99    let capture_drop = crate::db::query_capture_drop_stats(conn)?;
100    let captured_events =
101        conn.query_row("SELECT COUNT(*) FROM captured_events", [], |row| row.get(0))?;
102    let latest_captured_event_epoch = conn.query_row(
103        "SELECT MAX(inserted_at_epoch) FROM captured_events",
104        [],
105        |row| row.get(0),
106    )?;
107    let latest_raw_message_epoch = conn.query_row(
108        "SELECT MAX(created_at_epoch) FROM raw_messages",
109        [],
110        |row| row.get::<_, Option<i64>>(0),
111    )?;
112    let latest_capture_activity_epoch = [
113        latest_captured_event_epoch,
114        capture_drop.latest_epoch,
115        latest_raw_message_epoch,
116    ]
117    .into_iter()
118    .flatten()
119    .max();
120    let replay_ranges = query_extraction_replay_range_stats(conn)?;
121    let failure_lifecycle = crate::db::query_failure_lifecycle_stats(conn, now)?;
122    let worker_heartbeat = crate::db::worker::latest_daemon_worker_heartbeat(conn)?;
123    let healthy_worker_heartbeat = crate::db::worker::healthy_daemon_worker_heartbeat(
124        conn,
125        crate::db::worker::WORKER_HEARTBEAT_HEALTH_SECS,
126    )?;
127    let worker_heartbeat_age_secs = worker_heartbeat
128        .as_ref()
129        .map(|heartbeat| now.saturating_sub(heartbeat.updated_at_epoch));
130    let worker_daemon_healthy = healthy_worker_heartbeat.is_some();
131    let legacy_surfaces = super::legacy_surfaces::query_legacy_surface_stats(conn)?;
132    Ok(SystemStats {
133        active_memories: query_current_active_memory_count(conn)?,
134        active_observations: conn.query_row(
135            "SELECT COUNT(*) FROM observations WHERE status = 'active'",
136            [],
137            |row| row.get(0),
138        )?,
139        total_observations: conn.query_row("SELECT COUNT(*) FROM observations", [], |row| {
140            row.get(0)
141        })?,
142        session_summaries: conn.query_row("SELECT COUNT(*) FROM session_summaries", [], |row| {
143            row.get(0)
144        })?,
145        raw_messages: conn.query_row("SELECT COUNT(*) FROM raw_messages", [], |row| row.get(0))?,
146        raw_ingest_failures: raw_ingest.failures,
147        raw_ingest_parse_errors: raw_ingest.parse_errors,
148        raw_ingest_insert_errors: raw_ingest.insert_errors,
149        latest_raw_ingest_failure_epoch: raw_ingest.latest_epoch,
150        latest_raw_ingest_failure_kind: raw_ingest.latest_kind,
151        latest_raw_ingest_failure_path: raw_ingest.latest_path,
152        latest_raw_ingest_failure_message: raw_ingest.latest_message,
153        captured_events,
154        latest_captured_event_epoch,
155        latest_capture_activity_epoch,
156        capture_drop_events: capture_drop.total,
157        actionable_capture_drops: capture_drop.actionable,
158        unrecovered_capture_spills: capture_drop.unrecovered_spills,
159        latest_capture_drop_epoch: capture_drop.latest_epoch,
160        latest_capture_drop_reason: capture_drop.latest_reason,
161        latest_capture_drop_detail: capture_drop.latest_detail,
162        pending_extraction_tasks: conn.query_row(
163            "SELECT COUNT(*) FROM extraction_tasks WHERE status = 'pending'",
164            [],
165            |row| row.get(0),
166        )?,
167        processing_extraction_tasks: conn.query_row(
168            "SELECT COUNT(*) FROM extraction_tasks WHERE status = 'processing'",
169            [],
170            |row| row.get(0),
171        )?,
172        expired_processing_extraction_tasks: conn.query_row(
173            "SELECT COUNT(*)
174             FROM extraction_tasks
175             WHERE status = 'processing'
176               AND lease_expires_epoch IS NOT NULL
177               AND lease_expires_epoch < ?1",
178            params![now],
179            |row| row.get(0),
180        )?,
181        failed_extraction_tasks: failure_lifecycle.extraction_task.actionable_total,
182        retryable_extraction_replay_ranges: replay_ranges.retryable,
183        active_extraction_replay_ranges: replay_ranges.active,
184        quarantined_extraction_replay_ranges: replay_ranges.quarantined,
185        oldest_pending_extraction_epoch: conn.query_row(
186            "SELECT MIN(created_at_epoch) FROM extraction_tasks WHERE status = 'pending'",
187            [],
188            |row| row.get(0),
189        )?,
190        pending_memory_candidates: conn.query_row(
191            "SELECT COUNT(*) FROM memory_candidates WHERE review_status = 'pending_review'",
192            [],
193            |row| row.get(0),
194        )?,
195        total_memory_candidates: conn.query_row("SELECT COUNT(*) FROM memory_candidates", [], |row| {
196            row.get(0)
197        })?,
198        promoted_memory_candidates: conn.query_row(
199            "SELECT COUNT(*) FROM memory_candidates
200             WHERE review_status IN ('auto_promoted', 'approved', 'edited')",
201            [],
202            |row| row.get(0),
203        )?,
204        pending_review_memory_candidates: conn.query_row(
205            "SELECT COUNT(*) FROM memory_candidates WHERE review_status = 'pending_review'",
206            [],
207            |row| row.get(0),
208        )?,
209        pending_graph_candidates: query_pending_graph_candidates(conn)?,
210        pending_observations: conn.query_row(
211            "SELECT COUNT(*) FROM pending_observations WHERE status = 'pending'",
212            [],
213            |row| row.get(0),
214        )?,
215        ready_pending_observations: conn.query_row(
216            "SELECT COUNT(*) FROM pending_observations
217             WHERE status = 'pending'
218               AND (next_retry_epoch IS NULL OR next_retry_epoch <= ?1)
219               AND (lease_owner IS NULL OR lease_expires_epoch IS NULL OR lease_expires_epoch < ?1)",
220            params![now],
221            |row| row.get(0),
222        )?,
223        delayed_pending_observations: conn.query_row(
224            "SELECT COUNT(*) FROM pending_observations
225             WHERE status = 'pending'
226               AND next_retry_epoch IS NOT NULL
227               AND next_retry_epoch > ?1",
228            params![now],
229            |row| row.get(0),
230        )?,
231        processing_pending_observations: conn.query_row(
232            "SELECT COUNT(*) FROM pending_observations WHERE status = 'processing'",
233            [],
234            |row| row.get(0),
235        )?,
236        expired_processing_pending_observations: conn.query_row(
237            "SELECT COUNT(*) FROM pending_observations
238             WHERE status = 'processing'
239               AND lease_expires_epoch IS NOT NULL
240               AND lease_expires_epoch < ?1",
241            params![now],
242            |row| row.get(0),
243        )?,
244        failed_pending_observations: failure_lifecycle.pending_observation.actionable_total,
245        oldest_ready_pending_epoch: conn.query_row(
246            "SELECT MIN(created_at_epoch) FROM pending_observations
247             WHERE status = 'pending'
248               AND (next_retry_epoch IS NULL OR next_retry_epoch <= ?1)
249               AND (lease_owner IS NULL OR lease_expires_epoch IS NULL OR lease_expires_epoch < ?1)",
250            params![now],
251            |row| row.get(0),
252        )?,
253        pending_jobs: conn.query_row("SELECT COUNT(*) FROM jobs WHERE state = 'pending'", [], |row| {
254            row.get(0)
255        })?,
256        processing_jobs: conn.query_row(
257            "SELECT COUNT(*) FROM jobs WHERE state = 'processing'",
258            [],
259            |row| row.get(0),
260        )?,
261        failed_jobs: failure_lifecycle.job.actionable_total,
262        stuck_jobs: conn.query_row(
263            "SELECT COUNT(*) FROM jobs WHERE state = 'processing' \
264             AND lease_expires_epoch < strftime('%s', 'now')",
265            [],
266            |row| row.get(0),
267        )?,
268        worker_daemon_healthy,
269        worker_heartbeat_owner: worker_heartbeat.map(|heartbeat| heartbeat.owner),
270        worker_heartbeat_age_secs,
271        failure_lifecycle,
272        legacy_surfaces,
273    })
274}
275
276pub fn query_memory_facts_stats(conn: &Connection) -> Result<MemoryFactsStats> {
277    let memory_facts_exists = table_exists(conn, "memory_facts")?;
278    let memories_exists = table_exists(conn, "memories")?;
279    let captured_events_exists = table_exists(conn, "captured_events")?;
280    let total = if memory_facts_exists {
281        conn.query_row("SELECT COUNT(*) FROM memory_facts", [], |row| row.get(0))?
282    } else {
283        0
284    };
285    let retrieval_eligible = if memory_facts_exists && memories_exists {
286        let fact_current_filter = crate::memory::facts::current_fact_filter_sql(
287            "f",
288            crate::memory::facts::invalidated_at_epoch_available(conn)?,
289        );
290        let sql = format!(
291            "SELECT COUNT(*)
292             FROM memory_facts f
293             JOIN memories m ON m.id = f.source_memory_id
294             WHERE {fact_current_filter}
295               AND f.valid_from_epoch IS NOT NULL
296               AND m.status = 'active'
297               AND (
298                 m.expires_at_epoch IS NULL
299                 OR m.expires_at_epoch > CAST(strftime('%s', 'now') AS INTEGER)
300               )"
301        );
302        conn.query_row(&sql, [], |row| row.get(0))?
303    } else {
304        0
305    };
306    let active_memories = if memories_exists {
307        query_current_active_memory_count(conn)?
308    } else {
309        0
310    };
311    let captured_events = if captured_events_exists {
312        conn.query_row("SELECT COUNT(*) FROM captured_events", [], |row| row.get(0))?
313    } else {
314        0
315    };
316    Ok(MemoryFactsStats {
317        table_exists: memory_facts_exists,
318        total,
319        retrieval_eligible,
320        active_memories,
321        captured_events,
322    })
323}
324
325#[derive(Debug, Clone, Default)]
326struct RawIngestFailureStats {
327    failures: i64,
328    parse_errors: i64,
329    insert_errors: i64,
330    latest_epoch: Option<i64>,
331    latest_kind: Option<String>,
332    latest_path: Option<String>,
333    latest_message: Option<String>,
334}
335
336#[derive(Debug, Clone, Default)]
337struct ExtractionReplayRangeStats {
338    retryable: i64,
339    active: i64,
340    quarantined: i64,
341}
342
343fn query_extraction_replay_range_stats(conn: &Connection) -> Result<ExtractionReplayRangeStats> {
344    if !table_exists(conn, "extraction_replay_ranges")? {
345        return Ok(ExtractionReplayRangeStats::default());
346    }
347    let archived_filter = if column_exists(conn, "extraction_replay_ranges", "archived_at_epoch")? {
348        "AND r.archived_at_epoch IS NULL"
349    } else {
350        ""
351    };
352    let archived_filter_bare =
353        if column_exists(conn, "extraction_replay_ranges", "archived_at_epoch")? {
354            "AND archived_at_epoch IS NULL"
355        } else {
356            ""
357        };
358
359    Ok(ExtractionReplayRangeStats {
360        retryable: conn.query_row(
361            &format!(
362                "SELECT COUNT(*)
363             FROM extraction_replay_ranges r
364             WHERE r.status IN ('pending', 'failed')
365               {archived_filter}
366               AND NOT EXISTS (
367                 SELECT 1
368                 FROM extraction_tasks t
369                 WHERE t.replay_range_id = r.id
370                   AND t.status IN ('pending', 'processing')
371               )"
372            ),
373            [],
374            |row| row.get(0),
375        )?,
376        active: conn.query_row(
377            &format!(
378                "SELECT COUNT(*)
379             FROM extraction_replay_ranges r
380             WHERE r.status = 'requeued'
381                {archived_filter}
382                OR (r.status IN ('pending', 'failed') AND EXISTS (
383                 SELECT 1
384                 FROM extraction_tasks t
385                 WHERE t.replay_range_id = r.id
386                   AND t.status IN ('pending', 'processing')
387               ) {archived_filter})"
388            ),
389            [],
390            |row| row.get(0),
391        )?,
392        quarantined: conn.query_row(
393            &format!(
394                "SELECT COUNT(*)
395             FROM extraction_replay_ranges
396             WHERE status = 'quarantined'
397               {archived_filter_bare}"
398            ),
399            [],
400            |row| row.get(0),
401        )?,
402    })
403}
404
405fn query_raw_ingest_failure_stats(conn: &Connection) -> Result<RawIngestFailureStats> {
406    if !table_exists(conn, "raw_ingest_failures")? {
407        return Ok(RawIngestFailureStats::default());
408    }
409
410    let (failures, parse_errors, insert_errors) = conn.query_row(
411        "SELECT COUNT(*), COALESCE(SUM(parse_errors), 0), COALESCE(SUM(insert_errors), 0)
412         FROM raw_ingest_failures",
413        [],
414        |row| {
415            Ok((
416                row.get::<_, i64>(0)?,
417                row.get::<_, i64>(1)?,
418                row.get::<_, i64>(2)?,
419            ))
420        },
421    )?;
422    let latest = conn
423        .query_row(
424            "SELECT created_at_epoch, error_kind, transcript_path, error_message
425             FROM raw_ingest_failures
426             ORDER BY created_at_epoch DESC, id DESC
427             LIMIT 1",
428            [],
429            |row| {
430                Ok((
431                    row.get::<_, i64>(0)?,
432                    row.get::<_, String>(1)?,
433                    row.get::<_, Option<String>>(2)?,
434                    row.get::<_, String>(3)?,
435                ))
436            },
437        )
438        .optional()?;
439
440    let (latest_epoch, latest_kind, latest_path, latest_message) = match latest {
441        Some((epoch, kind, path, message)) => (Some(epoch), Some(kind), path, Some(message)),
442        None => (None, None, None, None),
443    };
444
445    Ok(RawIngestFailureStats {
446        failures,
447        parse_errors,
448        insert_errors,
449        latest_epoch,
450        latest_kind,
451        latest_path,
452        latest_message,
453    })
454}
455
456fn query_current_active_memory_count(conn: &Connection) -> Result<i64> {
457    Ok(conn.query_row(
458        "SELECT COUNT(*)
459         FROM memories
460         WHERE status = 'active'
461           AND (
462             expires_at_epoch IS NULL
463             OR expires_at_epoch > CAST(strftime('%s', 'now') AS INTEGER)
464           )",
465        [],
466        |row| row.get(0),
467    )?)
468}
469
470fn table_exists(conn: &Connection, table: &str) -> Result<bool> {
471    Ok(conn
472        .query_row(
473            "SELECT 1 FROM sqlite_master WHERE type='table' AND name=?1",
474            [table],
475            |_| Ok(()),
476        )
477        .optional()?
478        .is_some())
479}
480
481fn column_exists(conn: &Connection, table: &str, column: &str) -> Result<bool> {
482    if !table_exists(conn, table)? {
483        return Ok(false);
484    }
485    let mut stmt = conn.prepare(&format!("PRAGMA table_info({})", quote_identifier(table)))?;
486    let mut rows = stmt.query([])?;
487    while let Some(row) = rows.next()? {
488        let name: String = row.get(1)?;
489        if name == column {
490            return Ok(true);
491        }
492    }
493    Ok(false)
494}
495
496fn quote_identifier(identifier: &str) -> String {
497    format!("\"{}\"", identifier.replace('"', "\"\""))
498}
499
500fn query_pending_graph_candidates(conn: &Connection) -> Result<i64> {
501    if !table_exists(conn, "graph_candidates")? {
502        return Ok(0);
503    }
504    conn.query_row(
505		"SELECT COUNT(*) FROM graph_candidates WHERE review_status IN ('pending_review', 'deferred')",
506		[],
507		|row| row.get(0),
508	)
509    .map_err(Into::into)
510}
511
512pub fn query_daily_activity_stats(
513    conn: &Connection,
514    since_epoch: i64,
515) -> Result<DailyActivityStats> {
516    Ok(DailyActivityStats {
517        memories: conn.query_row(
518            "SELECT COUNT(*) FROM memories WHERE created_at_epoch >= ?1",
519            params![since_epoch],
520            |row| row.get(0),
521        )?,
522        observations: conn.query_row(
523            "SELECT COUNT(*) FROM observations WHERE created_at_epoch >= ?1",
524            params![since_epoch],
525            |row| row.get(0),
526        )?,
527    })
528}
529
530pub fn query_candidate_promotion_stats(
531    conn: &Connection,
532    now_epoch: i64,
533) -> Result<Vec<CandidatePromotionStat>> {
534    let week_ago = now_epoch - 7 * 24 * 3600;
535    let mut stmt = conn.prepare(
536        "SELECT source_kind,
537                review_status,
538                auto_promote_block_reason,
539                COUNT(*) AS total,
540                SUM(CASE WHEN created_at_epoch >= ?1 THEN 1 ELSE 0 END) AS last_7_days
541         FROM memory_candidates
542         GROUP BY source_kind, review_status, auto_promote_block_reason
543         ORDER BY total DESC, source_kind ASC, review_status ASC, auto_promote_block_reason ASC",
544    )?;
545    let rows = stmt.query_map(params![week_ago], |row| {
546        Ok(CandidatePromotionStat {
547            source_kind: row.get(0)?,
548            review_status: row.get(1)?,
549            block_reason: row.get(2)?,
550            total: row.get(3)?,
551            last_7_days: row.get(4)?,
552        })
553    })?;
554    collect_rows(rows)
555}
556
557pub fn query_top_projects(conn: &Connection, limit: i64) -> Result<Vec<ProjectCount>> {
558    let mut stmt = conn.prepare(
559        "SELECT project, COUNT(*) as cnt
560         FROM memories
561         WHERE status = 'active'
562           AND (
563             expires_at_epoch IS NULL
564             OR expires_at_epoch > CAST(strftime('%s', 'now') AS INTEGER)
565           )
566         GROUP BY project ORDER BY cnt DESC, project ASC LIMIT ?1",
567    )?;
568    let rows = stmt.query_map(params![limit], |row| {
569        Ok(ProjectCount {
570            project: row.get(0)?,
571            count: row.get(1)?,
572        })
573    })?;
574    collect_rows(rows)
575}
576
577pub fn query_ai_usage_totals(
578    conn: &Connection,
579    since_epoch: Option<i64>,
580    project: Option<&str>,
581) -> Result<AiUsageTotals> {
582    conn.query_row(
583        "SELECT COUNT(*),
584                COALESCE(SUM(input_tokens), 0),
585                COALESCE(SUM(output_tokens), 0),
586                COALESCE(SUM(reasoning_tokens), 0),
587                COALESCE(SUM(cache_creation_tokens), 0),
588                COALESCE(SUM(cache_read_tokens), 0),
589                COALESCE(SUM(total_tokens), 0),
590                COALESCE(SUM(estimated_cost_usd), 0.0)
591         FROM ai_usage_events
592         WHERE (?1 IS NULL OR created_at_epoch >= ?1)
593           AND (?2 IS NULL OR project = ?2)",
594        params![since_epoch, project],
595        |row| {
596            Ok(AiUsageTotals {
597                calls: row.get(0)?,
598                input_tokens: row.get(1)?,
599                output_tokens: row.get(2)?,
600                reasoning_tokens: row.get(3)?,
601                cache_creation_tokens: row.get(4)?,
602                cache_read_tokens: row.get(5)?,
603                total_tokens: row.get(6)?,
604                estimated_cost_usd: row.get(7)?,
605            })
606        },
607    )
608    .map_err(Into::into)
609}
610
611pub fn query_ai_usage_source_totals(
612    conn: &Connection,
613    since_epoch: Option<i64>,
614    project: Option<&str>,
615) -> Result<Vec<AiUsageSourceTotals>> {
616    let mut stmt = conn.prepare(
617        "SELECT usage_source,
618                pricing_source,
619                COUNT(*),
620                COALESCE(SUM(total_tokens), 0),
621                COALESCE(SUM(estimated_cost_usd), 0.0)
622         FROM ai_usage_events
623         WHERE (?1 IS NULL OR created_at_epoch >= ?1)
624           AND (?2 IS NULL OR project = ?2)
625         GROUP BY usage_source, pricing_source
626         ORDER BY SUM(total_tokens) DESC",
627    )?;
628    let rows = stmt.query_map(params![since_epoch, project], |row| {
629        Ok(AiUsageSourceTotals {
630            usage_source: row.get(0)?,
631            pricing_source: row.get(1)?,
632            calls: row.get(2)?,
633            total_tokens: row.get(3)?,
634            estimated_cost_usd: row.get(4)?,
635        })
636    })?;
637    collect_rows(rows)
638}
639
640pub fn query_ai_usage_breakdown(
641    conn: &Connection,
642    since_epoch: Option<i64>,
643    project: Option<&str>,
644    limit: i64,
645) -> Result<Vec<AiUsageBreakdown>> {
646    if limit <= 0 {
647        return Ok(Vec::new());
648    }
649
650    let mut stmt = conn.prepare(
651        "SELECT project,
652                executor,
653                usage_source,
654                pricing_source,
655                COUNT(*),
656                COALESCE(SUM(total_tokens), 0),
657                COALESCE(SUM(estimated_cost_usd), 0.0)
658         FROM ai_usage_events
659         WHERE (?1 IS NULL OR created_at_epoch >= ?1)
660           AND (?2 IS NULL OR project = ?2)
661         GROUP BY project, executor, usage_source, pricing_source
662         ORDER BY SUM(estimated_cost_usd) DESC,
663                  SUM(total_tokens) DESC,
664                  COUNT(*) DESC,
665                  project ASC,
666                  executor ASC
667         LIMIT ?3",
668    )?;
669    let rows = stmt.query_map(params![since_epoch, project, limit], |row| {
670        Ok(AiUsageBreakdown {
671            project: row.get(0)?,
672            executor: row.get(1)?,
673            usage_source: row.get(2)?,
674            pricing_source: row.get(3)?,
675            calls: row.get(4)?,
676            total_tokens: row.get(5)?,
677            estimated_cost_usd: row.get(6)?,
678        })
679    })?;
680    collect_rows(rows)
681}
682
683pub fn query_daily_ai_usage(
684    conn: &Connection,
685    since_epoch: i64,
686    project: Option<&str>,
687    limit: i64,
688) -> Result<Vec<DailyAiUsage>> {
689    let mut stmt = conn.prepare(
690        "SELECT strftime('%Y-%m-%d', created_at_epoch, 'unixepoch') AS day,
691                COUNT(*),
692                COALESCE(SUM(input_tokens), 0),
693                COALESCE(SUM(output_tokens), 0),
694                COALESCE(SUM(reasoning_tokens), 0),
695                COALESCE(SUM(cache_creation_tokens), 0),
696                COALESCE(SUM(cache_read_tokens), 0),
697                COALESCE(SUM(total_tokens), 0),
698                COALESCE(SUM(estimated_cost_usd), 0.0)
699         FROM ai_usage_events
700         WHERE created_at_epoch >= ?1
701           AND (?2 IS NULL OR project = ?2)
702         GROUP BY day
703         ORDER BY day DESC
704         LIMIT ?3",
705    )?;
706    let rows = stmt.query_map(params![since_epoch, project, limit], |row| {
707        Ok(DailyAiUsage {
708            day: row.get(0)?,
709            calls: row.get(1)?,
710            input_tokens: row.get(2)?,
711            output_tokens: row.get(3)?,
712            reasoning_tokens: row.get(4)?,
713            cache_creation_tokens: row.get(5)?,
714            cache_read_tokens: row.get(6)?,
715            total_tokens: row.get(7)?,
716            estimated_cost_usd: row.get(8)?,
717        })
718    })?;
719    collect_rows(rows)
720}
721
722pub fn query_weekly_ai_usage(
723    conn: &Connection,
724    since_epoch: i64,
725    project: Option<&str>,
726    limit: i64,
727) -> Result<Vec<WeeklyAiUsage>> {
728    let mut stmt = conn.prepare(
729        "SELECT strftime('%Y-W%W', created_at_epoch, 'unixepoch') AS week,
730                COUNT(*),
731                COALESCE(SUM(input_tokens), 0),
732                COALESCE(SUM(output_tokens), 0),
733                COALESCE(SUM(reasoning_tokens), 0),
734                COALESCE(SUM(cache_creation_tokens), 0),
735                COALESCE(SUM(cache_read_tokens), 0),
736                COALESCE(SUM(total_tokens), 0),
737                COALESCE(SUM(estimated_cost_usd), 0.0)
738         FROM ai_usage_events
739         WHERE created_at_epoch >= ?1
740           AND (?2 IS NULL OR project = ?2)
741         GROUP BY week
742         ORDER BY week DESC
743         LIMIT ?3",
744    )?;
745    let rows = stmt.query_map(params![since_epoch, project, limit], |row| {
746        Ok(WeeklyAiUsage {
747            week: row.get(0)?,
748            calls: row.get(1)?,
749            input_tokens: row.get(2)?,
750            output_tokens: row.get(3)?,
751            reasoning_tokens: row.get(4)?,
752            cache_creation_tokens: row.get(5)?,
753            cache_read_tokens: row.get(6)?,
754            total_tokens: row.get(7)?,
755            estimated_cost_usd: row.get(8)?,
756        })
757    })?;
758    collect_rows(rows)
759}
760
761#[cfg(test)]
762mod tests;