use rusqlite::Connection;
use std::path::Path;
use thiserror::Error;
use time::format_description::well_known::Rfc3339;
use time::OffsetDateTime;
use uuid::Uuid;
use crate::envelope::CheckResult;
pub const SCHEMA_VERSION: u32 = 2;
pub const FAILOPEN_RULE: &str = "pushkin.failopen.malformed_input";
pub const ESCALATION_RULE: &str = "pushkin.escalation";
const MIGRATIONS: &[&str] = &[
"
CREATE TABLE IF NOT EXISTS schema_meta (
version INTEGER NOT NULL
);
CREATE TABLE IF NOT EXISTS events (
id INTEGER PRIMARY KEY AUTOINCREMENT,
session TEXT NOT NULL,
seq INTEGER NOT NULL,
ts TEXT NOT NULL, -- UTC ISO-8601 (RFC 3339)
decision TEXT NOT NULL,
rule TEXT,
file TEXT,
payload TEXT NOT NULL,
UNIQUE (session, seq)
);
CREATE TRIGGER IF NOT EXISTS events_no_update
BEFORE UPDATE ON events
BEGIN SELECT RAISE(ABORT, 'events are append-only'); END;
CREATE TRIGGER IF NOT EXISTS events_no_delete
BEFORE DELETE ON events
BEGIN SELECT RAISE(ABORT, 'events are append-only'); END;
",
"
CREATE TABLE IF NOT EXISTS delivered_slices (
session TEXT NOT NULL,
cwd TEXT NOT NULL,
slice_key TEXT NOT NULL,
delivered_at_emission INTEGER NOT NULL,
PRIMARY KEY (session, cwd, slice_key)
);
CREATE TABLE IF NOT EXISTS delivery_counters (
session TEXT NOT NULL,
cwd TEXT NOT NULL,
emissions INTEGER NOT NULL,
PRIMARY KEY (session, cwd)
);
",
];
#[derive(Debug, Error)]
pub enum EventLogError {
#[error("event log storage error: {0}")]
Storage(#[from] rusqlite::Error),
#[error("event serialization error: {0}")]
Serialize(#[from] serde_json::Error),
#[error("timestamp formatting error: {0}")]
Timestamp(#[from] time::error::Format),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SessionId(String);
impl SessionId {
#[must_use]
pub fn from_name(name: &str) -> Self {
Self(name.to_owned())
}
#[must_use]
pub fn as_str(&self) -> &str {
&self.0
}
}
#[derive(Debug)]
pub struct StatsSummary {
pub blocks_by_rule: Vec<(String, u64)>,
pub compression_events: u64,
pub compression_saved_chars: u64,
pub nudge_arms: Vec<(String, u64)>,
pub failopen_events: u64,
}
#[derive(Debug)]
pub struct Telemetry<'a> {
pub rule: &'a str,
pub payload: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct GateEvent {
pub session: String,
pub seq: u64,
pub ts: String,
}
pub struct EventLog {
conn: Connection,
}
impl EventLog {
pub fn open(path: impl AsRef<Path>) -> Result<Self, EventLogError> {
let conn = Connection::open(path)?;
apply_migrations(&conn)?;
Ok(Self { conn })
}
pub fn begin_session(&self) -> Result<SessionId, EventLogError> {
Ok(SessionId(Uuid::new_v4().to_string()))
}
pub fn append(
&self,
session: &SessionId,
result: &CheckResult,
) -> Result<GateEvent, EventLogError> {
let rule = result.violations.first().map(|v| v.rule.clone());
let file = result.violations.first().map(|v| v.file.clone());
self.append_row(
session,
result.decision.as_str(),
rule,
file,
serde_json::to_string(result)?,
)
}
pub fn append_failopen(
&self,
session: &SessionId,
detail: &str,
) -> Result<GateEvent, EventLogError> {
let payload = serde_json::json!({ "failopen": true, "detail": detail }).to_string();
self.append_row(
session,
"allow",
Some(FAILOPEN_RULE.to_owned()),
None,
payload,
)
}
pub fn append_escalation(
&self,
session: &SessionId,
detail: &str,
) -> Result<GateEvent, EventLogError> {
let payload = serde_json::json!({ "escalation": true, "detail": detail }).to_string();
self.append_row(
session,
"block",
Some(ESCALATION_RULE.to_owned()),
None,
payload,
)
}
pub fn append_telemetry(
&self,
session: &SessionId,
telemetry: Telemetry<'_>,
) -> Result<GateEvent, EventLogError> {
self.append_row(
session,
"telemetry",
Some(telemetry.rule.to_owned()),
None,
telemetry.payload,
)
}
pub fn attempts(
&self,
session: &SessionId,
rule: &str,
file: &str,
) -> Result<u64, EventLogError> {
let count: u64 = self.conn.query_row(
"SELECT COUNT(*) FROM events
WHERE session = ?1 AND rule IN (?2, ?3) AND file = ?4",
(
session.as_str(),
rule,
crate::legacy::legacy_rule_id(rule).as_ref(),
file,
),
|row| row.get(0),
)?;
Ok(count)
}
pub fn denied_since(&self, rule: &str, file: &str, since: &str) -> Result<bool, EventLogError> {
let count: u64 = self.conn.query_row(
"SELECT COUNT(*) FROM events
WHERE decision = 'block' AND rule IN (?1, ?2)
AND file = ?3 AND substr(ts, 1, 19) >= substr(?4, 1, 19)",
(
rule,
crate::legacy::legacy_rule_id(rule).as_ref(),
file,
since,
),
|row| row.get(0),
)?;
Ok(count > 0)
}
pub fn block_count(&self, session: &SessionId) -> Result<u64, EventLogError> {
let count: u64 = self.conn.query_row(
"SELECT COUNT(*) FROM events WHERE session = ?1 AND decision = 'block'",
[session.as_str()],
|row| row.get(0),
)?;
Ok(count)
}
pub fn rule_count(&self, session: &SessionId, rule: &str) -> Result<u64, EventLogError> {
let count: u64 = self.conn.query_row(
"SELECT COUNT(*) FROM events WHERE session = ?1 AND rule IN (?2, ?3)",
(
session.as_str(),
rule,
crate::legacy::legacy_rule_id(rule).as_ref(),
),
|row| row.get(0),
)?;
Ok(count)
}
pub fn failopen_count(&self, session: &SessionId) -> Result<u64, EventLogError> {
let count: u64 = self.conn.query_row(
"SELECT COUNT(*) FROM events WHERE session = ?1 AND rule IN (?2, ?3)",
(
session.as_str(),
FAILOPEN_RULE,
crate::legacy::legacy_rule_id(FAILOPEN_RULE).as_ref(),
),
|row| row.get(0),
)?;
Ok(count)
}
pub fn stats(&self) -> Result<StatsSummary, EventLogError> {
let mut blocks = self.conn.prepare(
"SELECT rule, COUNT(*) FROM events
WHERE decision = 'block' AND rule IS NOT NULL
GROUP BY rule ORDER BY COUNT(*) DESC, rule",
)?;
let raw_blocks = blocks
.query_map([], |row| Ok((row.get(0)?, row.get(1)?)))?
.collect::<Result<Vec<(String, u64)>, _>>()?;
let blocks_by_rule = group_modern(raw_blocks);
let (compression_events, compression_saved_chars): (u64, u64) = self.conn.query_row(
"SELECT COUNT(*), COALESCE(SUM(json_extract(payload, '$.saved_chars')), 0)
FROM events WHERE rule IN (?1, ?2)",
(
"pushkin.compression",
crate::legacy::legacy_rule_id("pushkin.compression").as_ref(),
),
|row| Ok((row.get(0)?, row.get(1)?)),
)?;
let mut arms = self.conn.prepare(
"SELECT COALESCE(json_extract(payload, '$.arm'), 'unknown'), COUNT(*)
FROM events WHERE rule IN (?1, ?2) GROUP BY 1 ORDER BY 1",
)?;
let nudge_arms = arms
.query_map(
(
"pushkin.nudge",
crate::legacy::legacy_rule_id("pushkin.nudge").as_ref(),
),
|row| Ok((row.get(0)?, row.get(1)?)),
)?
.collect::<Result<Vec<(String, u64)>, _>>()?;
let failopen_events: u64 = self.conn.query_row(
"SELECT COUNT(*) FROM events WHERE rule IN (?1, ?2)",
(
FAILOPEN_RULE,
crate::legacy::legacy_rule_id(FAILOPEN_RULE).as_ref(),
),
|row| row.get(0),
)?;
Ok(StatsSummary {
blocks_by_rule,
compression_events,
compression_saved_chars,
nudge_arms,
failopen_events,
})
}
pub fn decision_counts(&self) -> Result<(u64, u64), EventLogError> {
let row = self.conn.query_row(
"SELECT
COUNT(*) FILTER (WHERE decision IN ('allow', 'block')),
COUNT(*) FILTER (WHERE decision = 'block')
FROM events WHERE rule IS NULL OR rule NOT IN (?1, ?2)",
(
ESCALATION_RULE,
crate::legacy::legacy_rule_id(ESCALATION_RULE).as_ref(),
),
|row| Ok((row.get(0)?, row.get(1)?)),
)?;
Ok(row)
}
pub fn latest_session_escalated(&self) -> Result<bool, EventLogError> {
let escalated: u64 = self.conn.query_row(
"SELECT COUNT(*) FROM events
WHERE rule IN (?1, ?2) AND session =
(SELECT session FROM events ORDER BY id DESC LIMIT 1)",
(
ESCALATION_RULE,
crate::legacy::legacy_rule_id(ESCALATION_RULE).as_ref(),
),
|row| row.get(0),
)?;
Ok(escalated > 0)
}
pub fn mean_attempts_to_compliance(&self) -> Result<Option<f64>, EventLogError> {
let mut statement = self.conn.prepare(
"SELECT session, decision FROM events
WHERE decision IN ('allow', 'block')
AND (rule IS NULL OR rule NOT IN (?1, ?2))
ORDER BY session, seq",
)?;
let rows = statement
.query_map(
(
ESCALATION_RULE,
crate::legacy::legacy_rule_id(ESCALATION_RULE).as_ref(),
),
|row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
)?
.collect::<Result<Vec<_>, _>>()?;
let mut recoveries: Vec<u64> = Vec::new();
let mut current_session = String::new();
let mut open_blocks: u64 = 0;
for (session, decision) in rows {
if session != current_session {
current_session = session;
open_blocks = 0;
}
match decision.as_str() {
"block" => open_blocks += 1,
"allow" if open_blocks > 0 => {
recoveries.push(open_blocks + 1);
open_blocks = 0;
}
_ => {}
}
}
if recoveries.is_empty() {
return Ok(None);
}
#[allow(clippy::cast_precision_loss)]
let mean = recoveries.iter().sum::<u64>() as f64 / recoveries.len() as f64;
Ok(Some(mean))
}
pub fn schema_version(&self) -> Result<u32, EventLogError> {
let version: u32 = self
.conn
.query_row("SELECT version FROM schema_meta", [], |row| row.get(0))?;
Ok(version)
}
fn append_row(
&self,
session: &SessionId,
decision: &str,
rule: Option<String>,
file: Option<String>,
payload: String,
) -> Result<GateEvent, EventLogError> {
let next_seq: u64 = self.conn.query_row(
"SELECT COALESCE(MAX(seq), 0) + 1 FROM events WHERE session = ?1",
[session.as_str()],
|row| row.get(0),
)?;
let ts = OffsetDateTime::now_utc().format(&Rfc3339)?;
self.conn.execute(
"INSERT INTO events (session, seq, ts, decision, rule, file, payload)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
(
session.as_str(),
next_seq,
&ts,
decision,
rule,
file,
payload,
),
)?;
Ok(GateEvent {
session: session.as_str().to_owned(),
seq: next_seq,
ts,
})
}
}
fn group_modern(rows: Vec<(String, u64)>) -> Vec<(String, u64)> {
let mut merged: std::collections::BTreeMap<String, u64> = std::collections::BTreeMap::new();
for (rule, count) in rows {
*merged
.entry(crate::legacy::modern_rule_id(&rule).into_owned())
.or_insert(0) += count;
}
let mut grouped: Vec<(String, u64)> = merged.into_iter().collect();
grouped.sort_by(|a, b| b.1.cmp(&a.1).then_with(|| a.0.cmp(&b.0)));
grouped
}
pub(crate) fn apply_migrations(conn: &Connection) -> Result<(), EventLogError> {
let current: u32 = conn
.query_row("SELECT version FROM schema_meta", [], |row| row.get(0))
.unwrap_or(0);
for (index, step) in MIGRATIONS.iter().enumerate() {
let step_version = u32::try_from(index).unwrap_or(u32::MAX).saturating_add(1);
if step_version > current {
conn.execute_batch(step)?;
}
}
if current == 0 {
conn.execute(
"INSERT INTO schema_meta (version) VALUES (?1)",
[SCHEMA_VERSION],
)?;
} else if current < SCHEMA_VERSION {
conn.execute("UPDATE schema_meta SET version = ?1", [SCHEMA_VERSION])?;
}
Ok(())
}