use crate::models::{Actor, ActorType, Event, EventType, Payload};
use anyhow::Result;
use chrono::{DateTime, Utc};
use rusqlite::{params, Connection};
use std::path::Path;
pub struct DbManager {
conn: Connection,
}
#[derive(Debug)]
pub struct SessionSummary {
pub session_id: String,
pub start_time: DateTime<Utc>,
pub end_time: DateTime<Utc>,
#[allow(dead_code)]
pub total_events: i64,
pub files_changed: usize,
pub agent_name: Option<String>,
pub exit_code: i32,
}
impl DbManager {
pub fn open<P: AsRef<Path>>(path: P) -> Result<Self> {
let conn = Connection::open(path)?;
let manager = Self { conn };
manager.init_schema()?;
Ok(manager)
}
fn init_schema(&self) -> Result<()> {
let schema = r#"
CREATE TABLE IF NOT EXISTS events (
id INTEGER PRIMARY KEY AUTOINCREMENT,
event_id TEXT UNIQUE NOT NULL,
timestamp DATETIME NOT NULL,
session_id TEXT NOT NULL,
event_type TEXT NOT NULL,
actor_type TEXT,
agent_name TEXT,
payload_json TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_session_id ON events(session_id);
CREATE INDEX IF NOT EXISTS idx_timestamp ON events(timestamp);
CREATE INDEX IF NOT EXISTS idx_event_type ON events(event_type);
CREATE INDEX IF NOT EXISTS idx_session_time ON events(session_id, timestamp);
"#;
self.conn.execute_batch(schema)?;
Ok(())
}
pub fn insert_event(&self, event: &Event) -> Result<()> {
let actor_type = match event.actor.r#type {
ActorType::Ai => "ai",
ActorType::Human => "human",
};
let agent_name = event.actor.agent_name.as_deref();
let payload_json = serde_json::to_string(&event.payload)?;
let event_type_str = match event.r#type {
EventType::SessionStart => "session_start",
EventType::SessionEnd => "session_end",
EventType::ProcessExec => "process_exec",
EventType::FileChange => "file_change",
EventType::GitCommit => "git_commit",
EventType::Error => "error",
};
self.conn.execute(
"INSERT INTO events (event_id, timestamp, session_id, event_type, actor_type, agent_name, payload_json)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
ON CONFLICT(event_id) DO UPDATE SET
timestamp = excluded.timestamp,
session_id = excluded.session_id,
event_type = excluded.event_type,
actor_type = excluded.actor_type,
agent_name = excluded.agent_name,
payload_json = excluded.payload_json",
params![
event.event_id,
event.timestamp.to_rfc3339(),
event.session_id,
event_type_str,
actor_type,
agent_name,
payload_json
],
)?;
Ok(())
}
fn row_to_event(&self, row: &rusqlite::Row) -> rusqlite::Result<Event> {
let event_id: String = row.get(0)?;
let timestamp_str: String = row.get(1)?;
let timestamp = timestamp_str
.parse::<DateTime<Utc>>()
.unwrap_or_else(|_| Utc::now());
let session_id: String = row.get(2)?;
let event_type_str: String = row.get(3)?;
let actor_type_str: Option<String> = row.get(4)?;
let agent_name: Option<String> = row.get(5)?;
let payload_json: String = row.get(6)?;
let event_type = match event_type_str.as_str() {
"session_start" => EventType::SessionStart,
"session_end" => EventType::SessionEnd,
"process_exec" => EventType::ProcessExec,
"file_change" => EventType::FileChange,
"git_commit" => EventType::GitCommit,
_ => EventType::Error,
};
let actor_type = match actor_type_str.as_deref() {
Some("ai") => ActorType::Ai,
Some("human") => ActorType::Human,
_ => ActorType::Human,
};
let payload: Payload = serde_json::from_str(&payload_json).unwrap_or(Payload::Error {
message: "Failed to parse payload".to_string(),
});
Ok(Event {
event_id,
timestamp,
session_id,
r#type: event_type,
actor: Actor {
r#type: actor_type,
agent_name,
},
payload,
})
}
pub fn get_events_by_session(&self, session_id: &str) -> Result<Vec<Event>> {
let mut stmt = self.conn.prepare(
"SELECT event_id, timestamp, session_id, event_type, actor_type, agent_name, payload_json
FROM events
WHERE session_id = ?1
ORDER BY timestamp ASC"
)?;
let events = stmt
.query_map([session_id], |row| self.row_to_event(row))?
.collect::<std::result::Result<Vec<_>, _>>()?;
Ok(events)
}
pub fn get_all_events(&self) -> Result<Vec<Event>> {
let mut stmt = self.conn.prepare(
"SELECT event_id, timestamp, session_id, event_type, actor_type, agent_name, payload_json
FROM events
ORDER BY timestamp DESC"
)?;
let events = stmt
.query_map([], |row| self.row_to_event(row))?
.collect::<std::result::Result<Vec<_>, _>>()?;
Ok(events)
}
pub fn get_sessions(&self) -> Result<Vec<SessionSummary>> {
let mut stmt = self.conn.prepare(
"SELECT
session_id,
MIN(timestamp) as start_time,
MAX(timestamp) as end_time,
COUNT(*) as total_events,
SUM(CASE WHEN event_type = 'file_change' THEN 1 ELSE 0 END) as files_changed,
MAX(agent_name) as agent_name,
MAX(CASE WHEN event_type = 'session_end' THEN json_extract(payload_json, '$.exit_code') END) as exit_code
FROM events
GROUP BY session_id
ORDER BY MIN(timestamp) DESC"
)?;
let sessions = stmt
.query_map([], |row| {
let start_time_str: String = row.get(1)?;
let end_time_str: String = row.get(2)?;
Ok(SessionSummary {
session_id: row.get(0)?,
start_time: start_time_str
.parse::<DateTime<Utc>>()
.unwrap_or_else(|_| Utc::now()),
end_time: end_time_str
.parse::<DateTime<Utc>>()
.unwrap_or_else(|_| Utc::now()),
total_events: row.get(3)?,
files_changed: row.get::<_, i64>(4).unwrap_or(0) as usize,
agent_name: row.get(5)?,
exit_code: row.get::<_, Option<i64>>(6)?.unwrap_or(0) as i32,
})
})?
.collect::<std::result::Result<Vec<_>, _>>()?;
Ok(sessions)
}
pub fn get_file_changes(&self, file_pattern: &str) -> Result<Vec<Event>> {
let mut stmt = self.conn.prepare(
"SELECT event_id, timestamp, session_id, event_type, actor_type, agent_name, payload_json
FROM events
WHERE event_type = 'file_change'
AND json_extract(payload_json, '$.file') LIKE ?1
ORDER BY timestamp ASC"
)?;
let events = stmt
.query_map([format!("%{}%", file_pattern)], |row| {
self.row_to_event(row)
})?
.collect::<std::result::Result<Vec<_>, _>>()?;
Ok(events)
}
pub fn export_to_jsonl<P: AsRef<Path>>(&self, output_path: P) -> Result<()> {
use std::fs::File;
use std::io::Write;
let mut file = File::create(output_path)?;
let events = self.get_all_events()?;
for event in events {
let json = serde_json::to_string(&event)?;
writeln!(file, "{}", json)?;
}
Ok(())
}
pub fn get_event_count(&self) -> Result<i64> {
let count: i64 = self
.conn
.query_row("SELECT COUNT(*) FROM events", [], |row| row.get(0))?;
Ok(count)
}
pub fn search_sessions(
&self,
since: Option<DateTime<Utc>>,
until: Option<DateTime<Utc>>,
agent: Option<&str>,
file_pattern: Option<&str>,
) -> Result<Vec<SessionSummary>> {
let mut params: Vec<String> = vec![];
let mut where_conditions = vec!["1=1"];
if let Some(since_dt) = since {
where_conditions.push("timestamp >= ?");
params.push(since_dt.to_rfc3339());
}
if let Some(until_dt) = until {
where_conditions.push("timestamp <= ?");
params.push(until_dt.to_rfc3339());
}
if let Some(agent_name) = agent {
where_conditions.push("agent_name = ?");
params.push(agent_name.to_string());
}
if let Some(pattern) = file_pattern {
where_conditions
.push("event_type = 'file_change' AND json_extract(payload_json, '$.file') LIKE ?");
params.push(format!("%{}%", pattern));
}
let sql = format!(
"WITH filtered_events AS (
SELECT * FROM events
WHERE {}
)
SELECT
session_id,
MIN(timestamp) as start_time,
MAX(timestamp) as end_time,
COUNT(*) as total_events,
SUM(CASE WHEN event_type = 'file_change' THEN 1 ELSE 0 END) as files_changed,
MAX(agent_name) as agent_name,
MAX(CASE WHEN event_type = 'session_end' THEN json_extract(payload_json, '$.exit_code') END) as exit_code
FROM filtered_events
GROUP BY session_id
ORDER BY MIN(timestamp) DESC",
where_conditions.join(" AND ")
);
let mut stmt = self.conn.prepare(&sql)?;
let sessions = stmt
.query_map(rusqlite::params_from_iter(params.iter()), |row| {
let start_time_str: String = row.get(1)?;
let end_time_str: String = row.get(2)?;
Ok(SessionSummary {
session_id: row.get(0)?,
start_time: start_time_str
.parse::<DateTime<Utc>>()
.unwrap_or_else(|_| Utc::now()),
end_time: end_time_str
.parse::<DateTime<Utc>>()
.unwrap_or_else(|_| Utc::now()),
total_events: row.get(3)?,
files_changed: row.get::<_, i64>(4).unwrap_or(0) as usize,
agent_name: row.get(5)?,
exit_code: row.get::<_, Option<i64>>(6)?.unwrap_or(0) as i32,
})
})?
.collect::<std::result::Result<Vec<_>, _>>()?;
Ok(sessions)
}
}