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