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;