use crate::global_store;
use crate::turso::{self, Connection};
use anyhow::{Context, Result};
use std::path::Path;
global_store! {
pub static STATS_STORE: StatsStore,
constructor = StatsStore::open,
}
#[derive(Clone, Debug)]
pub struct StatsStore {
pub(crate) conn: Connection,
}
const SCHEMA: &str = "\
CREATE TABLE IF NOT EXISTS tool_usage (
id INTEGER PRIMARY KEY AUTOINCREMENT,
agent_id TEXT NOT NULL,
role TEXT NOT NULL,
tool_name TEXT NOT NULL,
call_count INTEGER NOT NULL,
errors TEXT NOT NULL DEFAULT '[]',
workspace TEXT NOT NULL DEFAULT '',
recorded_at TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_tool_usage_agent_id ON tool_usage(agent_id);
CREATE INDEX IF NOT EXISTS idx_tool_usage_role ON tool_usage(role);
CREATE INDEX IF NOT EXISTS idx_tool_usage_recorded_at ON tool_usage(recorded_at);
CREATE INDEX IF NOT EXISTS idx_tool_usage_workspace ON tool_usage(workspace);";
#[derive(Debug, Clone)]
pub struct ToolErrorEntry {
pub tool_name: String,
pub role: String,
pub error: String,
pub workspace: String,
pub recorded_at: String,
}
#[derive(Debug, Clone, Default)]
pub struct ToolErrorQuery {
pub role_filter: Option<String>,
pub workspace_filter: Option<String>,
pub search: Option<String>,
}
impl StatsStore {
pub async fn open(root: &Path) -> Result<Self> {
let db_path = root.join("db/stats.db");
let conn = turso::open_with_schema(&db_path, SCHEMA).await?;
let user_version: i64 = conn
.query_row("PRAGMA user_version", turso::params![], |row| {
row.get::<Option<i64>>(0)
})
.await
.context("Failed to read PRAGMA user_version")?
.unwrap_or(0);
if user_version < 1 {
let _ = conn
.execute(
"ALTER TABLE tool_usage ADD COLUMN workspace TEXT NOT NULL DEFAULT ''",
turso::params![],
)
.await;
conn.execute("PRAGMA user_version = 1", turso::params![])
.await
.context("Failed to set PRAGMA user_version = 1")?;
}
Ok(Self { conn })
}
pub async fn query_tool_usage(&self, agent_id: &str, tool_name: &str) -> Result<Option<i64>> {
match self
.conn
.query_row(
"SELECT call_count FROM tool_usage \
WHERE agent_id = ?1 AND tool_name = ?2 \
ORDER BY id DESC LIMIT 1",
turso::params![agent_id, tool_name],
|row| row.get::<i64>(0),
)
.await
{
Ok(count) => Ok(Some(count)),
Err(::turso::Error::QueryReturnedNoRows) => Ok(None),
Err(e) => Err(e.into()),
}
}
#[must_use]
pub fn build_tool_error_filter(query: &ToolErrorQuery) -> (String, Vec<turso::Value>) {
let mut clauses = vec!["errors != '[]'".to_string()];
let mut params = Vec::new();
if let Some(ref role) = query.role_filter {
params.push(turso::Value::Text(role.clone()));
clauses.push("role = ?".to_string());
}
if let Some(ref workspace) = query.workspace_filter {
params.push(turso::Value::Text(workspace.clone()));
clauses.push("workspace = ?".to_string());
}
if let Some(ref search) = query.search {
params.push(turso::Value::Text(format!("%{search}%")));
clauses.push("json_each.value LIKE ?".to_string());
}
(clauses.join(" AND "), params)
}
pub async fn count_tool_errors(&self, query: &ToolErrorQuery) -> Result<usize> {
let (where_clause, params) = Self::build_tool_error_filter(query);
let sql = format!(
"SELECT COUNT(*) FROM tool_usage, json_each(tool_usage.errors) WHERE {where_clause}",
);
let rows = self
.conn
.query(&sql, turso::params_from_iter(params))
.await?;
let count = match rows.into_iter().next() {
Some(row) => match row.get_value(0)? {
turso::Value::Integer(n) => usize::try_from(n).unwrap_or(0),
_ => 0,
},
None => 0,
};
Ok(count)
}
pub async fn query_tool_errors(
&self,
query: &ToolErrorQuery,
limit: usize,
offset: usize,
) -> Result<(Vec<ToolErrorEntry>, usize)> {
let total = self.count_tool_errors(query).await?;
if total == 0 {
return Ok((vec![], 0));
}
let (where_clause, filter_params) = Self::build_tool_error_filter(query);
let limit_val = i64::try_from(limit).unwrap_or(50);
let offset_val = i64::try_from(offset).unwrap_or(0);
let sql = format!(
"SELECT tool_name, role, json_each.value AS error, workspace, recorded_at \
FROM tool_usage, json_each(tool_usage.errors) \
WHERE {where_clause} \
ORDER BY recorded_at DESC \
LIMIT ? OFFSET ?",
);
let mut all_params = filter_params;
all_params.push(turso::Value::Integer(limit_val));
all_params.push(turso::Value::Integer(offset_val));
let rows = self
.conn
.query(&sql, turso::params_from_iter(all_params))
.await?;
let mut entries = Vec::new();
for row in rows {
entries.push(ToolErrorEntry {
tool_name: turso::row_text(&row, 0)?,
role: turso::row_text(&row, 1)?,
error: turso::row_text(&row, 2)?,
workspace: turso::row_text(&row, 3)?,
recorded_at: turso::row_text(&row, 4)?,
});
}
Ok((entries, total))
}
pub async fn flush_batch(
&self,
agent_id: &str,
role: &str,
workspace: &str,
stats: &std::collections::HashMap<String, crate::ToolUsage>,
) -> Result<()> {
let recorded_at = turso::now();
for (tool_name, usage) in stats {
let errors_json = serde_json::to_string(&usage.errors).unwrap_or_default();
self.conn
.execute(
"INSERT INTO tool_usage (agent_id, role, tool_name, call_count, errors, workspace, recorded_at) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
turso::params![
agent_id,
role,
tool_name.clone(),
{
let raw_count = usage.call_count;
i64::try_from(raw_count).unwrap_or_else(|_| {
tracing::warn!(
agent_id = %agent_id,
role = %role,
tool_name = %tool_name,
call_count = raw_count,
"Tool call count overflowed i64, clamping to i64::MAX"
);
i64::MAX
})
},
errors_json,
workspace.to_string(),
recorded_at.clone(),
],
)
.await?;
}
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_build_tool_error_filter_no_filters() {
let query = ToolErrorQuery::default();
let (clause, params) = StatsStore::build_tool_error_filter(&query);
assert_eq!(clause, "errors != '[]'");
assert!(params.is_empty());
}
#[test]
fn test_build_tool_error_filter_role_only() {
let query = ToolErrorQuery {
role_filter: Some("Engineer".to_string()),
workspace_filter: None,
search: None,
};
let (clause, params) = StatsStore::build_tool_error_filter(&query);
assert_eq!(clause, "errors != '[]' AND role = ?");
assert_eq!(params.len(), 1);
assert_eq!(params[0], turso::Value::Text("Engineer".to_string()));
}
#[test]
fn test_build_tool_error_filter_search_only() {
let query = ToolErrorQuery {
role_filter: None,
workspace_filter: None,
search: Some("timeout".to_string()),
};
let (clause, params) = StatsStore::build_tool_error_filter(&query);
assert_eq!(clause, "errors != '[]' AND json_each.value LIKE ?");
assert_eq!(params.len(), 1);
assert_eq!(params[0], turso::Value::Text("%timeout%".to_string()));
}
#[test]
fn test_build_tool_error_filter_workspace_only() {
let query = ToolErrorQuery {
role_filter: None,
workspace_filter: Some("my-workspace".to_string()),
search: None,
};
let (clause, params) = StatsStore::build_tool_error_filter(&query);
assert_eq!(clause, "errors != '[]' AND workspace = ?");
assert_eq!(params.len(), 1);
assert_eq!(params[0], turso::Value::Text("my-workspace".to_string()));
}
#[test]
fn test_build_tool_error_filter_both() {
let query = ToolErrorQuery {
role_filter: Some("Analyst".to_string()),
workspace_filter: None,
search: Some("connection refused".to_string()),
};
let (clause, params) = StatsStore::build_tool_error_filter(&query);
assert_eq!(
clause,
"errors != '[]' AND role = ? AND json_each.value LIKE ?"
);
assert_eq!(params.len(), 2);
assert_eq!(params[0], turso::Value::Text("Analyst".to_string()));
assert_eq!(
params[1],
turso::Value::Text("%connection refused%".to_string())
);
}
#[test]
fn test_build_tool_error_filter_empty_strings() {
let query = ToolErrorQuery {
role_filter: Some(String::new()),
workspace_filter: None,
search: Some(String::new()),
};
let (clause, params) = StatsStore::build_tool_error_filter(&query);
assert_eq!(
clause,
"errors != '[]' AND role = ? AND json_each.value LIKE ?"
);
assert_eq!(params.len(), 2);
assert_eq!(params[0], turso::Value::Text(String::new()));
assert_eq!(params[1], turso::Value::Text("%%".to_string()));
}
}