#![allow(dead_code)]
use anyhow::{Context, Result};
use rusqlite::{params, Connection};
use std::fs;
use std::path::PathBuf;
use std::time::SystemTime;
const DB_FILENAME: &str = "mrapids.sqlite";
pub struct AnalyticsEngine {
conn: Connection,
db_path: PathBuf,
}
#[derive(Debug)]
pub struct DbStatus {
pub exists: bool,
pub path: PathBuf,
pub size_bytes: u64,
pub size_human: String,
pub last_modified: Option<SystemTime>,
pub last_modified_human: String,
pub engine_version: String,
pub table_count: usize,
}
#[derive(Debug, Clone)]
pub struct TableInfo {
pub name: String,
pub columns: Vec<ColumnInfo>,
}
#[derive(Debug, Clone)]
pub struct ColumnInfo {
pub name: String,
pub data_type: String,
pub nullable: bool,
pub default_value: Option<String>,
}
impl AnalyticsEngine {
pub fn open() -> Result<Self> {
let db_path = Self::get_db_path()?;
Self::open_at_path(db_path)
}
pub fn open_at_path(db_path: PathBuf) -> Result<Self> {
if let Some(parent) = db_path.parent() {
fs::create_dir_all(parent)
.with_context(|| format!("Failed to create directory: {:?}", parent))?;
}
let conn = Connection::open(&db_path)
.with_context(|| format!("Failed to open SQLite at {:?}", db_path))?;
let engine = Self { conn, db_path };
engine.init_schema()?;
Ok(engine)
}
pub fn open_in_memory() -> Result<Self> {
let conn = Connection::open_in_memory().context("Failed to open in-memory SQLite")?;
let engine = Self {
conn,
db_path: PathBuf::from(":memory:"),
};
engine.init_schema()?;
Ok(engine)
}
pub fn get_db_path() -> Result<PathBuf> {
let home = dirs::home_dir().context("Could not determine home directory")?;
Ok(home.join(".mrapids").join(DB_FILENAME))
}
fn init_schema(&self) -> Result<()> {
self.conn.execute_batch(
r#"
CREATE TABLE IF NOT EXISTS _mrapids_meta (
key TEXT PRIMARY KEY,
value TEXT,
updated_at TEXT DEFAULT CURRENT_TIMESTAMP
);
CREATE TABLE IF NOT EXISTS api_requests (
id INTEGER PRIMARY KEY AUTOINCREMENT,
timestamp TEXT DEFAULT CURRENT_TIMESTAMP,
spec_file TEXT,
operation_id TEXT,
method TEXT,
path TEXT,
status_code INTEGER,
response_time_ms REAL,
request_size_bytes INTEGER,
response_size_bytes INTEGER,
success INTEGER,
error_message TEXT
);
CREATE TABLE IF NOT EXISTS collection_runs (
id INTEGER PRIMARY KEY AUTOINCREMENT,
timestamp TEXT DEFAULT CURRENT_TIMESTAMP,
collection_name TEXT,
total_requests INTEGER,
passed INTEGER,
failed INTEGER,
skipped INTEGER,
duration_ms REAL
);
CREATE TABLE IF NOT EXISTS runs (
id INTEGER PRIMARY KEY AUTOINCREMENT,
run_id TEXT UNIQUE NOT NULL,
timestamp TEXT DEFAULT CURRENT_TIMESTAMP,
spec_path TEXT,
environment TEXT,
total_requests INTEGER DEFAULT 0,
successful INTEGER DEFAULT 0,
failed INTEGER DEFAULT 0,
duration_ms REAL,
status TEXT DEFAULT 'running',
metadata TEXT
);
CREATE TABLE IF NOT EXISTS requests (
id INTEGER PRIMARY KEY AUTOINCREMENT,
run_id TEXT NOT NULL,
request_id TEXT UNIQUE NOT NULL,
timestamp TEXT DEFAULT CURRENT_TIMESTAMP,
operation_id TEXT,
endpoint TEXT NOT NULL,
method TEXT NOT NULL,
url TEXT,
headers TEXT,
query_params TEXT,
path_params TEXT,
payload TEXT,
payload_size_bytes INTEGER,
secrets_present INTEGER,
secret_fingerprints TEXT,
FOREIGN KEY (run_id) REFERENCES runs(run_id)
);
CREATE TABLE IF NOT EXISTS responses (
id INTEGER PRIMARY KEY AUTOINCREMENT,
request_id TEXT NOT NULL,
timestamp TEXT DEFAULT CURRENT_TIMESTAMP,
status_code INTEGER NOT NULL,
status_text TEXT,
headers TEXT,
body TEXT,
body_size_bytes INTEGER,
duration_ms REAL NOT NULL,
success INTEGER,
error_message TEXT,
FOREIGN KEY (request_id) REFERENCES requests(request_id)
);
CREATE TABLE IF NOT EXISTS comparisons (
id INTEGER PRIMARY KEY AUTOINCREMENT,
comparison_id TEXT UNIQUE NOT NULL,
left_run_id TEXT NOT NULL,
right_run_id TEXT NOT NULL,
created_at TEXT DEFAULT CURRENT_TIMESTAMP,
status TEXT DEFAULT 'pending',
total_diffs INTEGER DEFAULT 0,
summary TEXT,
metadata TEXT
);
CREATE TABLE IF NOT EXISTS comparison_diffs (
id INTEGER PRIMARY KEY AUTOINCREMENT,
comparison_id TEXT NOT NULL,
request_key TEXT NOT NULL,
diff_type TEXT NOT NULL,
left_value TEXT,
right_value TEXT,
field_path TEXT,
severity TEXT DEFAULT 'info',
description TEXT,
FOREIGN KEY (comparison_id) REFERENCES comparisons(comparison_id)
);
CREATE INDEX IF NOT EXISTS idx_runs_timestamp ON runs(timestamp);
CREATE INDEX IF NOT EXISTS idx_requests_run_id ON requests(run_id);
CREATE INDEX IF NOT EXISTS idx_requests_operation ON requests(operation_id);
CREATE INDEX IF NOT EXISTS idx_responses_request_id ON responses(request_id);
CREATE INDEX IF NOT EXISTS idx_responses_status ON responses(status_code);
CREATE INDEX IF NOT EXISTS idx_comparisons_created ON comparisons(created_at);
CREATE INDEX IF NOT EXISTS idx_comparisons_left_run ON comparisons(left_run_id);
CREATE INDEX IF NOT EXISTS idx_comparisons_right_run ON comparisons(right_run_id);
CREATE INDEX IF NOT EXISTS idx_diffs_comparison ON comparison_diffs(comparison_id);
CREATE INDEX IF NOT EXISTS idx_diffs_type ON comparison_diffs(diff_type);
CREATE TABLE IF NOT EXISTS mcp_decisions (
decision_id TEXT PRIMARY KEY,
session_id TEXT,
agent_id TEXT,
timestamp TEXT NOT NULL,
action_type TEXT NOT NULL,
operation_id TEXT,
method TEXT,
outcome TEXT NOT NULL,
policy_rule TEXT,
policy_reason TEXT,
claim_token_id TEXT,
preview_token_id TEXT,
environment TEXT,
duration_ms REAL,
metadata_json TEXT
);
CREATE INDEX IF NOT EXISTS idx_decisions_session ON mcp_decisions(session_id);
CREATE INDEX IF NOT EXISTS idx_decisions_agent ON mcp_decisions(agent_id);
"#,
)?;
Ok(())
}
pub fn get_status(&self) -> Result<DbStatus> {
let exists = self.db_path.exists() || self.db_path.to_string_lossy() == ":memory:";
let (size_bytes, last_modified) = if exists && self.db_path.to_string_lossy() != ":memory:"
{
let metadata = fs::metadata(&self.db_path)?;
(metadata.len(), metadata.modified().ok())
} else {
(0, None)
};
let size_human = format_size(size_bytes);
let last_modified_human = last_modified
.map(|t| {
let datetime: chrono::DateTime<chrono::Local> = t.into();
datetime.format("%Y-%m-%d %H:%M:%S").to_string()
})
.unwrap_or_else(|| "N/A".to_string());
let engine_version: String = self
.conn
.query_row("SELECT sqlite_version()", [], |row| row.get(0))
.unwrap_or_else(|_| "unknown".to_string());
let table_count: usize = self
.conn
.query_row(
"SELECT COUNT(*) FROM sqlite_master WHERE type = 'table'",
[],
|row| row.get(0),
)
.unwrap_or(0);
Ok(DbStatus {
exists,
path: self.db_path.clone(),
size_bytes,
size_human,
last_modified,
last_modified_human,
engine_version,
table_count,
})
}
pub fn get_schema(&self) -> Result<Vec<TableInfo>> {
let mut tables = Vec::new();
let table_names: Vec<String> = {
let mut stmt = self.conn.prepare(
"SELECT name FROM sqlite_master WHERE type = 'table' AND name NOT LIKE 'sqlite_%' ORDER BY name"
)?;
let mut rows = stmt.query([])?;
let mut names = Vec::new();
while let Some(row) = rows.next()? {
let name: String = row.get(0)?;
names.push(name);
}
names
};
for table_name in table_names {
let columns = self.get_table_columns(&table_name)?;
tables.push(TableInfo {
name: table_name,
columns,
});
}
Ok(tables)
}
fn get_table_columns(&self, table_name: &str) -> Result<Vec<ColumnInfo>> {
let mut stmt = self
.conn
.prepare(&format!("PRAGMA table_info({})", table_name))?;
let mut rows = stmt.query([])?;
let mut columns = Vec::new();
while let Some(row) = rows.next()? {
let name: String = row.get(1)?;
let data_type: String = row.get(2)?;
let notnull: bool = row.get(3)?;
let default_value: Option<String> = row.get(4)?;
columns.push(ColumnInfo {
name,
data_type,
nullable: !notnull,
default_value,
});
}
Ok(columns)
}
pub fn query_json(&self, sql: &str) -> Result<serde_json::Value> {
let mut stmt = self.conn.prepare(sql)?;
let column_names: Vec<String> = (0..stmt.column_count())
.map(|i| stmt.column_name(i).unwrap_or("?").to_string())
.collect();
let mut results = Vec::new();
let mut rows = stmt.query([])?;
while let Some(row) = rows.next()? {
let mut obj = serde_json::Map::new();
for (i, name) in column_names.iter().enumerate() {
let value: rusqlite::types::Value = row.get(i)?;
obj.insert(name.clone(), sqlite_value_to_json(value));
}
results.push(serde_json::Value::Object(obj));
}
Ok(serde_json::Value::Array(results))
}
pub fn query_csv(&self, sql: &str, include_header: bool) -> Result<String> {
let mut stmt = self.conn.prepare(sql)?;
let column_names: Vec<String> = (0..stmt.column_count())
.map(|i| stmt.column_name(i).unwrap_or("?").to_string())
.collect();
let mut csv_lines = Vec::new();
if include_header {
csv_lines.push(column_names.join(","));
}
let mut rows = stmt.query([])?;
while let Some(row) = rows.next()? {
let mut values = Vec::new();
for i in 0..column_names.len() {
let value: rusqlite::types::Value = row.get(i)?;
let csv_value = sqlite_value_to_csv(value);
values.push(csv_value);
}
csv_lines.push(values.join(","));
}
Ok(csv_lines.join("\n"))
}
pub fn query_table(&self, sql: &str) -> Result<(Vec<String>, Vec<Vec<String>>)> {
let mut stmt = self.conn.prepare(sql)?;
let column_names: Vec<String> = (0..stmt.column_count())
.map(|i| stmt.column_name(i).unwrap_or("?").to_string())
.collect();
let mut result_rows: Vec<Vec<String>> = Vec::new();
let mut rows = stmt.query([])?;
while let Some(row) = rows.next()? {
let mut values = Vec::new();
for i in 0..column_names.len() {
let value: rusqlite::types::Value = row.get(i)?;
values.push(sqlite_value_to_string(value));
}
result_rows.push(values);
}
Ok((column_names, result_rows))
}
pub fn log_request(
&self,
spec_file: &str,
operation_id: &str,
method: &str,
path: &str,
status_code: i32,
response_time_ms: f64,
request_size: i32,
response_size: i32,
success: bool,
error_message: Option<&str>,
) -> Result<()> {
self.conn.execute(
r#"
INSERT INTO api_requests
(spec_file, operation_id, method, path, status_code, response_time_ms,
request_size_bytes, response_size_bytes, success, error_message)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
"#,
params![
spec_file,
operation_id,
method,
path,
status_code,
response_time_ms,
request_size,
response_size,
success,
error_message
],
)?;
Ok(())
}
pub fn get_request_stats(&self) -> Result<serde_json::Value> {
self.query_json(
r#"
SELECT
COUNT(*) as total_requests,
SUM(CASE WHEN success THEN 1 ELSE 0 END) as successful,
SUM(CASE WHEN NOT success THEN 1 ELSE 0 END) as failed,
ROUND(AVG(response_time_ms), 2) as avg_response_time_ms,
ROUND(MIN(response_time_ms), 2) as min_response_time_ms,
ROUND(MAX(response_time_ms), 2) as max_response_time_ms,
COUNT(DISTINCT spec_file) as unique_specs,
COUNT(DISTINCT operation_id) as unique_operations
FROM api_requests
"#,
)
}
pub fn generate_run_id() -> String {
use uuid::Uuid;
Uuid::new_v4().to_string()[..8].to_string() }
pub fn generate_request_id() -> String {
use uuid::Uuid;
format!("req_{}", &Uuid::new_v4().to_string()[..8])
}
pub fn create_run(
&self,
run_id: &str,
spec_path: Option<&str>,
environment: Option<&serde_json::Value>,
) -> Result<()> {
let env_json = environment
.map(|e| serde_json::to_string(e).unwrap_or_default())
.unwrap_or_else(|| "{}".to_string());
self.conn.execute(
r#"
INSERT INTO runs (run_id, spec_path, environment, status)
VALUES (?, ?, ?, 'running')
"#,
params![run_id, spec_path, env_json],
)?;
Ok(())
}
pub fn log_request_v2(
&self,
run_id: &str,
request_id: &str,
operation_id: Option<&str>,
endpoint: &str,
method: &str,
url: Option<&str>,
headers: Option<&serde_json::Value>,
query_params: Option<&serde_json::Value>,
path_params: Option<&serde_json::Value>,
payload: Option<&str>,
) -> Result<()> {
use crate::utils::redaction::{compute_secret_fingerprints, sanitize_headers_for_storage};
let (sanitized_headers, secrets_present, fingerprints) = match headers {
Some(h) => {
let fingerprints = compute_secret_fingerprints(h);
let sanitized = sanitize_headers_for_storage(h);
let has_secrets = fingerprints.as_object().map_or(false, |m| !m.is_empty());
(Some(sanitized), has_secrets, Some(fingerprints))
}
None => (None, false, None),
};
let headers_json = sanitized_headers
.as_ref()
.map(|h| serde_json::to_string(h).unwrap_or_default());
let fingerprints_json = fingerprints
.as_ref()
.map(|f| serde_json::to_string(f).unwrap_or_default());
let query_json = query_params.map(|q| serde_json::to_string(q).unwrap_or_default());
let path_json = path_params.map(|p| serde_json::to_string(p).unwrap_or_default());
let payload_size = payload.map(|p| p.len() as i32);
self.conn.execute(
r#"
INSERT INTO requests
(run_id, request_id, operation_id, endpoint, method, url, headers, query_params, path_params, payload, payload_size_bytes, secrets_present, secret_fingerprints)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
"#,
params![
run_id,
request_id,
operation_id,
endpoint,
method,
url,
headers_json,
query_json,
path_json,
payload,
payload_size,
secrets_present,
fingerprints_json
],
)?;
Ok(())
}
pub fn log_response(
&self,
request_id: &str,
status_code: i32,
status_text: Option<&str>,
headers: Option<&serde_json::Value>,
body: Option<&str>,
duration_ms: f64,
success: bool,
error_message: Option<&str>,
) -> Result<()> {
use crate::utils::redaction::{sanitize_body_for_storage, sanitize_headers_for_storage};
let sanitized_headers = headers.map(|h| sanitize_headers_for_storage(h));
let headers_json = sanitized_headers
.as_ref()
.map(|h| serde_json::to_string(h).unwrap_or_default());
let body_size = body.map(|b| b.len() as i32);
let body_to_store = body.map(|b| {
let truncated = if b.len() > 1_000_000 {
format!("{}...[truncated, {} bytes total]", &b[..1000], b.len())
} else {
b.to_string()
};
sanitize_body_for_storage(&truncated)
});
self.conn.execute(
r#"
INSERT INTO responses
(request_id, status_code, status_text, headers, body, body_size_bytes, duration_ms, success, error_message)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
"#,
params![
request_id,
status_code,
status_text,
headers_json,
body_to_store,
body_size,
duration_ms,
success,
error_message
],
)?;
Ok(())
}
pub fn complete_run(
&self,
run_id: &str,
total_requests: i32,
successful: i32,
failed: i32,
duration_ms: f64,
status: &str,
) -> Result<()> {
self.conn.execute(
r#"
UPDATE runs
SET total_requests = ?, successful = ?, failed = ?, duration_ms = ?, status = ?
WHERE run_id = ?
"#,
params![
total_requests,
successful,
failed,
duration_ms,
status,
run_id
],
)?;
Ok(())
}
pub fn get_runs(&self, limit: usize) -> Result<serde_json::Value> {
self.query_json(&format!(
r#"
SELECT
run_id,
strftime('%Y-%m-%d %H:%M:%S', timestamp) as timestamp,
spec_path,
total_requests,
successful,
failed,
ROUND(duration_ms, 2) as duration_ms,
status
FROM runs
ORDER BY timestamp DESC
LIMIT {}
"#,
limit
))
}
pub fn get_run_requests(&self, run_id: &str) -> Result<serde_json::Value> {
self.query_json(&format!(
r#"
SELECT
r.request_id,
r.operation_id,
r.method,
r.endpoint,
resp.status_code,
ROUND(resp.duration_ms, 2) as duration_ms,
resp.success
FROM requests r
LEFT JOIN responses resp ON r.request_id = resp.request_id
WHERE r.run_id = '{}'
ORDER BY r.timestamp
"#,
run_id
))
}
pub fn get_request_detail(&self, request_id: &str) -> Result<serde_json::Value> {
self.query_json(&format!(
r#"
SELECT
r.request_id,
r.run_id,
r.operation_id,
r.method,
r.endpoint,
r.url,
r.headers as request_headers,
r.query_params,
r.path_params,
r.payload,
resp.status_code,
resp.status_text,
resp.headers as response_headers,
resp.body,
ROUND(resp.duration_ms, 2) as duration_ms,
resp.success,
resp.error_message
FROM requests r
LEFT JOIN responses resp ON r.request_id = resp.request_id
WHERE r.request_id = '{}'
"#,
request_id
))
}
pub fn validate(&self) -> Result<bool> {
self.conn.execute(
"INSERT OR REPLACE INTO _mrapids_meta (key, value) VALUES ('_test', 'validation')",
[],
)?;
let result: String = self.conn.query_row(
"SELECT value FROM _mrapids_meta WHERE key = '_test'",
[],
|row| row.get(0),
)?;
self.conn
.execute("DELETE FROM _mrapids_meta WHERE key = '_test'", [])?;
Ok(result == "validation")
}
pub fn connection(&self) -> &Connection {
&self.conn
}
pub fn log_decision(
&self,
decision_id: &str,
session_id: Option<&str>,
agent_id: Option<&str>,
action_type: &str,
operation_id: Option<&str>,
method: Option<&str>,
outcome: &str,
policy_rule: Option<&str>,
policy_reason: Option<&str>,
claim_token_id: Option<&str>,
preview_token_id: Option<&str>,
environment: Option<&str>,
duration_ms: Option<f64>,
metadata: Option<&serde_json::Value>,
) -> Result<()> {
let metadata_json = metadata.map(|v| serde_json::to_string(v).unwrap_or_default());
self.conn.execute(
r#"
INSERT INTO mcp_decisions
(decision_id, session_id, agent_id, timestamp, action_type, operation_id, method,
outcome, policy_rule, policy_reason, claim_token_id, preview_token_id,
environment, duration_ms, metadata_json)
VALUES (?, ?, ?, CURRENT_TIMESTAMP, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
"#,
params![
decision_id,
session_id,
agent_id,
action_type,
operation_id,
method,
outcome,
policy_rule,
policy_reason,
claim_token_id,
preview_token_id,
environment,
duration_ms,
metadata_json
],
)?;
Ok(())
}
pub fn query_decisions(
&self,
session_id: Option<&str>,
agent_id: Option<&str>,
limit: usize,
) -> Result<Vec<serde_json::Value>> {
let mut sql = String::from(
"SELECT decision_id, session_id, agent_id, timestamp, action_type, \
operation_id, method, outcome, policy_rule, policy_reason, \
claim_token_id, preview_token_id, environment, duration_ms, metadata_json \
FROM mcp_decisions WHERE 1=1",
);
let mut param_values: Vec<Box<dyn rusqlite::types::ToSql>> = Vec::new();
if let Some(sid) = session_id {
sql.push_str(" AND session_id = ?");
param_values.push(Box::new(sid.to_string()));
}
if let Some(aid) = agent_id {
sql.push_str(" AND agent_id = ?");
param_values.push(Box::new(aid.to_string()));
}
sql.push_str(" ORDER BY timestamp DESC LIMIT ?");
param_values.push(Box::new(limit as i64));
let params_ref: Vec<&dyn rusqlite::types::ToSql> =
param_values.iter().map(|p| p.as_ref()).collect();
let mut stmt = self.conn.prepare(&sql)?;
let column_names: Vec<String> = (0..stmt.column_count())
.map(|i| stmt.column_name(i).unwrap_or("?").to_string())
.collect();
let mut results = Vec::new();
let mut rows = stmt.query(params_ref.as_slice())?;
while let Some(row) = rows.next()? {
let mut obj = serde_json::Map::new();
for (i, name) in column_names.iter().enumerate() {
let value: rusqlite::types::Value = row.get(i)?;
obj.insert(name.clone(), sqlite_value_to_json(value));
}
results.push(serde_json::Value::Object(obj));
}
Ok(results)
}
pub fn generate_comparison_id() -> String {
use uuid::Uuid;
format!("cmp_{}", &Uuid::new_v4().to_string()[..8])
}
pub fn create_comparison(
&self,
comparison_id: &str,
left_run_id: &str,
right_run_id: &str,
) -> Result<()> {
self.conn.execute(
r#"
INSERT INTO comparisons (comparison_id, left_run_id, right_run_id, status)
VALUES (?, ?, ?, 'running')
"#,
params![comparison_id, left_run_id, right_run_id],
)?;
Ok(())
}
pub fn store_diff(
&self,
comparison_id: &str,
request_key: &str,
diff_type: &str,
left_value: Option<&serde_json::Value>,
right_value: Option<&serde_json::Value>,
field_path: Option<&str>,
severity: &str,
description: Option<&str>,
) -> Result<()> {
let left_json = left_value.map(|v| v.to_string());
let right_json = right_value.map(|v| v.to_string());
self.conn.execute(
r#"
INSERT INTO comparison_diffs (comparison_id, request_key, diff_type, left_value, right_value, field_path, severity, description)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
"#,
params![
comparison_id,
request_key,
diff_type,
left_json,
right_json,
field_path,
severity,
description
],
)?;
Ok(())
}
pub fn complete_comparison(
&self,
comparison_id: &str,
total_diffs: i32,
summary: Option<&serde_json::Value>,
) -> Result<()> {
let summary_json = summary.map(|v| v.to_string());
self.conn.execute(
r#"
UPDATE comparisons
SET status = 'completed', total_diffs = ?, summary = ?
WHERE comparison_id = ?
"#,
params![total_diffs, summary_json, comparison_id],
)?;
Ok(())
}
pub fn get_comparison(&self, comparison_id: &str) -> Result<serde_json::Value> {
self.query_json(&format!(
r#"
SELECT
comparison_id,
left_run_id,
right_run_id,
strftime('%Y-%m-%d %H:%M:%S', created_at) as created_at,
status,
total_diffs,
summary
FROM comparisons
WHERE comparison_id = '{}'
"#,
comparison_id
))
}
pub fn get_comparison_diffs(&self, comparison_id: &str) -> Result<serde_json::Value> {
self.query_json(&format!(
r#"
SELECT
request_key,
diff_type,
left_value,
right_value,
field_path,
severity,
description
FROM comparison_diffs
WHERE comparison_id = '{}'
ORDER BY severity DESC, request_key
"#,
comparison_id
))
}
pub fn get_comparisons(&self, limit: usize) -> Result<serde_json::Value> {
self.query_json(&format!(
r#"
SELECT
comparison_id,
left_run_id,
right_run_id,
strftime('%Y-%m-%d %H:%M:%S', created_at) as created_at,
status,
total_diffs
FROM comparisons
ORDER BY created_at DESC
LIMIT {}
"#,
limit
))
}
pub fn run_exists(&self, run_id: &str) -> Result<bool> {
let count: i32 = self.conn.query_row(
"SELECT COUNT(*) FROM runs WHERE run_id = ?",
params![run_id],
|row| row.get(0),
)?;
Ok(count > 0)
}
pub fn get_run_requests_for_comparison(
&self,
run_id: &str,
) -> Result<Vec<(String, String, serde_json::Value)>> {
let mut stmt = self.conn.prepare(
r#"
SELECT
r.endpoint,
r.method,
r.request_id,
r.url,
r.headers,
r.query_params,
r.payload,
resp.status_code,
resp.headers as response_headers,
resp.body,
resp.duration_ms,
resp.success
FROM requests r
LEFT JOIN responses resp ON r.request_id = resp.request_id
WHERE r.run_id = ?
ORDER BY r.endpoint, r.method
"#,
)?;
let mut results = Vec::new();
let mut rows = stmt.query(params![run_id])?;
while let Some(row) = rows.next()? {
let endpoint: String = row.get(0)?;
let method: String = row.get(1)?;
let request_id: Option<String> = row.get(2).ok();
let url: Option<String> = row.get(3).ok();
let headers: Option<String> = row.get(4).ok();
let query_params: Option<String> = row.get(5).ok();
let payload: Option<String> = row.get(6).ok();
let status_code: Option<i32> = row.get(7).ok();
let response_headers: Option<String> = row.get(8).ok();
let body: Option<String> = row.get(9).ok();
let duration_ms: Option<f64> = row.get(10).ok();
let success: Option<bool> = row.get(11).ok();
let data = serde_json::json!({
"request_id": request_id,
"url": url,
"headers": headers.and_then(|h| serde_json::from_str::<serde_json::Value>(&h).ok()),
"query_params": query_params.and_then(|q| serde_json::from_str::<serde_json::Value>(&q).ok()),
"payload": payload.and_then(|p| serde_json::from_str::<serde_json::Value>(&p).ok()),
"status_code": status_code,
"response_headers": response_headers.and_then(|h| serde_json::from_str::<serde_json::Value>(&h).ok()),
"response_body": body,
"duration_ms": duration_ms,
"success": success
});
results.push((endpoint, method, data));
}
Ok(results)
}
}
fn sqlite_value_to_json(value: rusqlite::types::Value) -> serde_json::Value {
match value {
rusqlite::types::Value::Null => serde_json::Value::Null,
rusqlite::types::Value::Integer(i) => serde_json::json!(i),
rusqlite::types::Value::Real(f) => serde_json::json!(f),
rusqlite::types::Value::Text(s) => serde_json::Value::String(s),
rusqlite::types::Value::Blob(b) => {
serde_json::Value::String(format!("<blob {} bytes>", b.len()))
}
}
}
fn sqlite_value_to_csv(value: rusqlite::types::Value) -> String {
match value {
rusqlite::types::Value::Null => String::new(),
rusqlite::types::Value::Integer(i) => i.to_string(),
rusqlite::types::Value::Real(f) => f.to_string(),
rusqlite::types::Value::Text(s) => {
if s.contains(',') || s.contains('"') || s.contains('\n') {
format!("\"{}\"", s.replace('"', "\"\""))
} else {
s
}
}
rusqlite::types::Value::Blob(b) => format!("<blob {} bytes>", b.len()),
}
}
fn sqlite_value_to_string(value: rusqlite::types::Value) -> String {
match value {
rusqlite::types::Value::Null => "NULL".to_string(),
rusqlite::types::Value::Integer(i) => i.to_string(),
rusqlite::types::Value::Real(f) => format!("{:.2}", f),
rusqlite::types::Value::Text(s) => s,
rusqlite::types::Value::Blob(b) => format!("<blob {} bytes>", b.len()),
}
}
fn format_size(bytes: u64) -> String {
const KB: u64 = 1024;
const MB: u64 = KB * 1024;
const GB: u64 = MB * 1024;
if bytes >= GB {
format!("{:.2} GB", bytes as f64 / GB as f64)
} else if bytes >= MB {
format!("{:.2} MB", bytes as f64 / MB as f64)
} else if bytes >= KB {
format!("{:.2} KB", bytes as f64 / KB as f64)
} else {
format!("{} bytes", bytes)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_open_in_memory() {
let engine = AnalyticsEngine::open_in_memory().unwrap();
let status = engine.get_status().unwrap();
assert!(status.exists);
assert!(status.table_count >= 6); }
#[test]
fn test_get_schema() {
let engine = AnalyticsEngine::open_in_memory().unwrap();
let schema = engine.get_schema().unwrap();
assert!(schema.len() >= 6);
let runs_table = schema.iter().find(|t| t.name == "runs");
assert!(runs_table.is_some());
let runs = runs_table.unwrap();
assert!(runs.columns.iter().any(|c| c.name == "run_id"));
assert!(runs.columns.iter().any(|c| c.name == "spec_path"));
let requests_table = schema.iter().find(|t| t.name == "requests");
assert!(requests_table.is_some());
let requests = requests_table.unwrap();
assert!(requests.columns.iter().any(|c| c.name == "request_id"));
assert!(requests.columns.iter().any(|c| c.name == "endpoint"));
let responses_table = schema.iter().find(|t| t.name == "responses");
assert!(responses_table.is_some());
let responses = responses_table.unwrap();
assert!(responses.columns.iter().any(|c| c.name == "status_code"));
assert!(responses.columns.iter().any(|c| c.name == "duration_ms"));
}
#[test]
fn test_validate() {
let engine = AnalyticsEngine::open_in_memory().unwrap();
assert!(engine.validate().unwrap());
}
#[test]
fn test_log_request() {
let engine = AnalyticsEngine::open_in_memory().unwrap();
engine
.log_request(
"petstore.yaml",
"getPets",
"GET",
"/pets",
200,
150.5,
0,
1024,
true,
None,
)
.unwrap();
let stats = engine.get_request_stats().unwrap();
let stats_arr = stats.as_array().unwrap();
assert_eq!(stats_arr.len(), 1);
assert_eq!(stats_arr[0]["total_requests"], 1);
}
#[test]
fn test_query_json() {
let engine = AnalyticsEngine::open_in_memory().unwrap();
let result = engine
.query_json("SELECT 1 as num, 'hello' as msg")
.unwrap();
let arr = result.as_array().unwrap();
assert_eq!(arr.len(), 1);
assert_eq!(arr[0]["num"], 1);
assert_eq!(arr[0]["msg"], "hello");
}
#[test]
fn test_log_and_query_decision() {
let engine = AnalyticsEngine::open_in_memory().unwrap();
let metadata = serde_json::json!({"query": "find pets", "result_count": 3});
engine
.log_decision(
"dec_abc123",
Some("session_xyz"),
Some("agent_test01"),
"api_find",
Some("listPets"),
Some("GET"),
"allowed",
None,
None,
None,
None,
Some("development"),
Some(12.5),
Some(&metadata),
)
.unwrap();
let results = engine.query_decisions(None, None, 10).unwrap();
assert_eq!(results.len(), 1);
let row = &results[0];
assert_eq!(row["decision_id"], "dec_abc123");
assert_eq!(row["session_id"], "session_xyz");
assert_eq!(row["agent_id"], "agent_test01");
assert_eq!(row["action_type"], "api_find");
assert_eq!(row["operation_id"], "listPets");
assert_eq!(row["method"], "GET");
assert_eq!(row["outcome"], "allowed");
assert_eq!(row["environment"], "development");
assert_eq!(row["duration_ms"], 12.5);
let stored_meta: serde_json::Value =
serde_json::from_str(row["metadata_json"].as_str().unwrap()).unwrap();
assert_eq!(stored_meta["result_count"], 3);
}
#[test]
fn test_query_decisions_by_session() {
let engine = AnalyticsEngine::open_in_memory().unwrap();
engine
.log_decision(
"dec_s1_a",
Some("session_A"),
Some("agent_1"),
"api_find",
None,
None,
"allowed",
None,
None,
None,
None,
None,
None,
None,
)
.unwrap();
engine
.log_decision(
"dec_s1_b",
Some("session_A"),
Some("agent_1"),
"api_claim",
Some("getPet"),
Some("GET"),
"allowed",
None,
None,
Some("claim_tok1"),
None,
None,
None,
None,
)
.unwrap();
engine
.log_decision(
"dec_s2_a",
Some("session_B"),
Some("agent_2"),
"api_run",
Some("deletePet"),
Some("DELETE"),
"denied",
Some("no_delete"),
Some("Delete not allowed in dev"),
None,
None,
None,
None,
None,
)
.unwrap();
let all = engine.query_decisions(None, None, 100).unwrap();
assert_eq!(all.len(), 3);
let session_a = engine
.query_decisions(Some("session_A"), None, 100)
.unwrap();
assert_eq!(session_a.len(), 2);
for r in &session_a {
assert_eq!(r["session_id"], "session_A");
}
let session_b = engine
.query_decisions(Some("session_B"), None, 100)
.unwrap();
assert_eq!(session_b.len(), 1);
assert_eq!(session_b[0]["outcome"], "denied");
assert_eq!(session_b[0]["policy_rule"], "no_delete");
let agent_1 = engine.query_decisions(None, Some("agent_1"), 100).unwrap();
assert_eq!(agent_1.len(), 2);
let filtered = engine
.query_decisions(Some("session_B"), Some("agent_2"), 100)
.unwrap();
assert_eq!(filtered.len(), 1);
let limited = engine.query_decisions(None, None, 2).unwrap();
assert_eq!(limited.len(), 2);
}
}