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#[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;