use crate::models::LogLine;
use sqlx::{SqlitePool, sqlite::SqliteConnectOptions};
use std::path::Path;
use std::str::FromStr;
pub async fn init_db(db_path: &str) -> Result<SqlitePool, sqlx::Error> {
let connection_string = format!("sqlite://{}", db_path);
let options = SqliteConnectOptions::from_str(&connection_string)?
.create_if_missing(true)
.journal_mode(sqlx::sqlite::SqliteJournalMode::Wal)
.synchronous(sqlx::sqlite::SqliteSynchronous::Normal);
let pool = SqlitePool::connect_with(options).await?;
// Enable auto_vacuum if not already enabled, so that deleting rows reclaims space.
// Note: AUTO_VACUUM must be set before tables are created.
let auto_vacuum_status: (String,) = sqlx::query_as("PRAGMA auto_vacuum")
.fetch_one(&pool)
.await
.unwrap_or(("0".to_string(),));
if auto_vacuum_status.0 == "0" || auto_vacuum_status.0 == "NONE" {
// Need to set auto_vacuum and then VACUUM to apply it
sqlx::query("PRAGMA auto_vacuum = FULL")
.execute(&pool)
.await?;
sqlx::query("VACUUM").execute(&pool).await?;
}
// Create logs table
sqlx::query(
"CREATE TABLE IF NOT EXISTS logs (
id INTEGER PRIMARY KEY AUTOINCREMENT,
timestamp INTEGER NOT NULL,
severity TEXT NOT NULL,
message TEXT NOT NULL,
service TEXT NOT NULL,
metadata TEXT NOT NULL
);",
)
.execute(&pool)
.await?;
// Create indexes for efficient querying
sqlx::query("CREATE INDEX IF NOT EXISTS idx_logs_timestamp ON logs(timestamp);")
.execute(&pool)
.await?;
sqlx::query("CREATE INDEX IF NOT EXISTS idx_logs_severity ON logs(severity);")
.execute(&pool)
.await?;
sqlx::query("CREATE INDEX IF NOT EXISTS idx_logs_service ON logs(service);")
.execute(&pool)
.await?;
Ok(pool)
}
pub async fn insert_logs(pool: &SqlitePool, logs: &[LogLine]) -> Result<(), sqlx::Error> {
if logs.is_empty() {
return Ok(());
}
let mut tx = pool.begin().await?;
for log in logs {
let metadata_str =
serde_json::to_string(&log.metadata).unwrap_or_else(|_| "{}".to_string());
let timestamp_ms = log.timestamp.timestamp_millis();
sqlx::query(
"INSERT INTO logs (timestamp, severity, message, service, metadata)
VALUES (?, ?, ?, ?, ?)",
)
.bind(timestamp_ms)
.bind(&log.severity)
.bind(&log.message)
.bind(&log.service)
.bind(metadata_str)
.execute(&mut *tx)
.await?;
}
tx.commit().await?;
Ok(())
}
pub async fn enforce_size_limit(
pool: &SqlitePool,
db_path: &str,
limit_bytes: u64,
) -> Result<(), sqlx::Error> {
let path = Path::new(db_path);
if !path.exists() {
return Ok(());
}
let mut file_size = match std::fs::metadata(path) {
Ok(meta) => meta.len(),
Err(_) => return Ok(()),
};
// Include the size of the -wal file if it exists, since WAL mode writes to it first.
let wal_path = format!("{}-wal", db_path);
if let Ok(meta) = std::fs::metadata(&wal_path) {
file_size += meta.len();
}
// Include the size of the -shm file if it exists
let shm_path = format!("{}-shm", db_path);
if let Ok(meta) = std::fs::metadata(&shm_path) {
file_size += meta.len();
}
if file_size <= limit_bytes {
return Ok(());
}
// Keep deleting old logs and vacuuming until we are under the limit
while file_size > limit_bytes {
let row_count: (i64,) = sqlx::query_as("SELECT COUNT(*) FROM logs")
.fetch_one(pool)
.await?;
if row_count.0 == 0 {
break;
}
// Delete the oldest 15% of records or at least 1000 rows
let delete_count = std::cmp::max(1000, row_count.0 * 15 / 100);
sqlx::query(
"DELETE FROM logs WHERE id IN (SELECT id FROM logs ORDER BY timestamp ASC LIMIT ?)",
)
.bind(delete_count)
.execute(pool)
.await?;
// Run incremental/full vacuum to reclaim filesystem space.
// Since auto_vacuum = FULL is set, SQLite reclaims pages automatically, but sometimes
// running VACUUM explicitly or letting database shrink is necessary to immediately update file size.
sqlx::query("VACUUM").execute(pool).await?;
// Recheck file size
file_size = match std::fs::metadata(path) {
Ok(meta) => meta.len(),
Err(_) => break,
};
if let Ok(meta) = std::fs::metadata(&wal_path) {
file_size += meta.len();
}
if let Ok(meta) = std::fs::metadata(&shm_path) {
file_size += meta.len();
}
if delete_count >= row_count.0 {
break;
}
}
Ok(())
}