use rusqlite::{params, Connection, Result};
use std::time::UNIX_EPOCH;
const DECAY_THRESHOLD_DAYS: i64 = 30;
const ACCESS_COUNT_THRESHOLD: i64 = 10;
pub fn init_decay_table(conn: &Connection) -> Result<()> {
conn.execute(
"CREATE TABLE IF NOT EXISTS decay (
event_id TEXT PRIMARY KEY,
access_count INTEGER NOT NULL DEFAULT 0,
last_accessed INTEGER NOT NULL,
pinned BOOLEAN NOT NULL DEFAULT 0,
created_at INTEGER NOT NULL DEFAULT (unixepoch())
)",
[],
)?;
Ok(())
}
pub fn init_shadow_table(conn: &Connection) -> Result<()> {
conn.execute(
"CREATE TABLE IF NOT EXISTS shadow_state (
event_id TEXT PRIMARY KEY,
decay_score REAL NOT NULL,
flagged_at INTEGER NOT NULL DEFAULT (unixepoch()),
FOREIGN KEY (event_id) REFERENCES events(id) ON DELETE CASCADE
)",
[],
)?;
Ok(())
}
pub fn init_decay_tables(conn: &Connection) -> Result<()> {
init_decay_table(conn)?;
init_shadow_table(conn)?;
Ok(())
}
pub fn track_access(conn: &Connection, event_id: &str) -> Result<()> {
conn.execute(
"INSERT INTO decay (event_id, last_accessed)
VALUES (?1, unixepoch())
ON CONFLICT(event_id) DO UPDATE SET
access_count = decay.access_count + 1,
last_accessed = unixepoch(),
pinned = decay.pinned",
[event_id],
)?;
Ok(())
}
pub fn get_decay_score(conn: &Connection, event_id: &str) -> Result<f64> {
let score: Option<f64> = conn.query_row(
"SELECT
CASE
WHEN last_accessed = 0 THEN 0
WHEN CAST((unixepoch() - last_accessed) / 86400 AS INTEGER) < 1 THEN CAST(access_count AS REAL)
ELSE CAST(access_count AS REAL) /
CAST((unixepoch() - last_accessed) / 86400 AS REAL)
END
FROM decay
WHERE event_id = ?1",
[event_id],
|row| row.get(0),
)?;
Ok(score.unwrap_or(0.0))
}
pub fn is_flagged(conn: &Connection, event_id: &str) -> Result<bool> {
let flagged: bool = conn.query_row(
"SELECT EXISTS(
SELECT 1 FROM decay
WHERE event_id = ?1
AND access_count < ?2
AND (unixepoch() - last_accessed) > ?3 * 86400
AND pinned = 0
)",
params![event_id, ACCESS_COUNT_THRESHOLD, DECAY_THRESHOLD_DAYS],
|row| row.get(0),
)?;
Ok(flagged)
}
pub fn get_flagged_events(conn: &Connection) -> Result<Vec<String>> {
let mut ids = Vec::new();
let mut stmt = conn.prepare(
"SELECT e.id FROM events e
JOIN decay d ON e.id = d.event_id
LEFT JOIN shadow_state s ON e.id = s.event_id
WHERE d.access_count < ?
AND (unixepoch() - d.last_accessed) > ? * 86400
AND d.pinned = 0
AND s.event_id IS NULL
ORDER BY d.last_accessed ASC",
)?;
let rows = stmt.query_map(
params![ACCESS_COUNT_THRESHOLD, DECAY_THRESHOLD_DAYS],
|row| row.get(0),
)?;
for row in rows {
ids.push(row?);
}
Ok(ids)
}
pub fn move_to_shadow(conn: &Connection) -> Result<usize> {
let flagged_ids = get_flagged_events(conn)?;
if flagged_ids.is_empty() {
return Ok(0);
}
let mut scores: Vec<(String, f64)> = Vec::new();
for id in &flagged_ids {
if let Ok(score) = get_decay_score(conn, id) {
scores.push((id.clone(), score));
}
}
let tx = conn.unchecked_transaction()?;
let mut moved = 0;
for (id, score) in scores {
let now = UNIX_EPOCH.elapsed().unwrap().as_secs() as i64;
tx.execute(
"INSERT INTO shadow_state (event_id, decay_score, flagged_at)
VALUES (?1, ?2, ?3)
ON CONFLICT(event_id) DO UPDATE SET
decay_score = excluded.decay_score,
flagged_at = excluded.flagged_at",
params![id, score, now],
)?;
moved += 1;
}
tx.commit()?;
Ok(moved)
}
pub fn get_shadow_events(conn: &Connection) -> Result<Vec<ShadowEvent>> {
let mut events = Vec::new();
let mut stmt = conn.prepare(
"SELECT e.id, e.timestamp, e.source, e.content, e.meta, e.ingested_at, e.content_hash, s.decay_score, s.flagged_at
FROM shadow_state s
JOIN events e ON e.id = s.event_id
ORDER BY flagged_at DESC"
)?;
let rows = stmt.query_map([], |row| {
Ok(ShadowEvent {
id: row.get(0)?,
timestamp: row.get(1)?,
source: row.get(2)?,
content: row.get(3)?,
meta: row.get(4)?,
ingested_at: row.get(5)?,
content_hash: row.get(6)?,
decay_score: row.get(7)?,
flagged_at: row.get(8)?,
})
})?;
for row in rows {
events.push(row?);
}
Ok(events)
}
pub fn pin_event(conn: &Connection, event_id: &str) -> Result<()> {
conn.execute(
"UPDATE decay SET pinned = 1 WHERE event_id = ?1",
[event_id],
)?;
Ok(())
}
pub fn unpin_event(conn: &Connection, event_id: &str) -> Result<()> {
conn.execute(
"UPDATE decay SET pinned = 0 WHERE event_id = ?1",
[event_id],
)?;
Ok(())
}
pub fn restore_from_shadow(conn: &Connection, event_id: &str) -> Result<()> {
conn.execute("DELETE FROM shadow_state WHERE event_id = ?1", [event_id])?;
conn.execute(
"INSERT INTO decay (event_id, access_count, last_accessed, pinned)
VALUES (?1, 0, unixepoch(), 0)
ON CONFLICT(event_id) DO UPDATE SET
access_count = 0,
last_accessed = unixepoch(),
pinned = 0",
[event_id],
)?;
Ok(())
}
pub fn get_decay_stats(conn: &Connection) -> Result<DecayStats> {
let total_events: i64 = conn.query_row("SELECT COUNT(*) FROM decay", [], |row| row.get(0))?;
let total_access: i64 =
conn.query_row("SELECT SUM(access_count) FROM decay", [], |row| row.get(0))?;
let pinned_count: i64 =
conn.query_row("SELECT COUNT(*) FROM decay WHERE pinned = 1", [], |row| {
row.get(0)
})?;
let avg_access: f64 = if total_events > 0 {
total_access as f64 / total_events as f64
} else {
0.0
};
let flagged_count: i64 = conn.query_row(
"SELECT COUNT(*)
FROM decay d
LEFT JOIN shadow_state s ON d.event_id = s.event_id
WHERE d.access_count < ?1
AND (unixepoch() - d.last_accessed) > ?2 * 86400
AND d.pinned = 0
AND s.event_id IS NULL",
params![ACCESS_COUNT_THRESHOLD, DECAY_THRESHOLD_DAYS],
|row| row.get(0),
)?;
Ok(DecayStats {
total_events,
total_access,
pinned_count,
avg_access,
flagged_count,
})
}
#[derive(Debug, Clone)]
pub struct ShadowEvent {
pub id: String,
pub timestamp: i64,
pub source: String,
pub content: String,
pub meta: Option<String>,
pub ingested_at: i64,
pub content_hash: Option<String>,
pub decay_score: f64,
pub flagged_at: i64,
}
#[derive(Debug, Clone)]
pub struct DecayStats {
pub total_events: i64,
pub total_access: i64,
pub pinned_count: i64,
pub avg_access: f64,
pub flagged_count: i64,
}