use crate::Result;
use crate::daemon_id::DaemonId;
use crate::log_parse::ParsedLog;
use crate::log_store::{
ArchiveHook, FieldFilter, LogEntry, LogQuery, LogStore, MessageFilter, escape_like_pattern,
};
use chrono::{DateTime, Local, TimeZone};
use log::error;
use miette::IntoDiagnostic;
use rusqlite::{Connection, OpenFlags, OptionalExtension, params};
use std::collections::HashSet;
use std::io::{BufRead, BufReader, Write};
use std::path::PathBuf;
use std::sync::Mutex;
fn text_to_json_literal(value: &str) -> String {
if value.eq_ignore_ascii_case("true") {
return "true".to_string();
}
if value.eq_ignore_ascii_case("false") {
return "false".to_string();
}
if value.eq_ignore_ascii_case("null") {
return "null".to_string();
}
if let Ok(serde_json::Value::Number(_)) = serde_json::from_str(value) {
return value.to_string();
}
serde_json::to_string(value).unwrap_or_else(|_| format!(r#""{value}""#))
}
fn add_regexp_function(conn: &Connection) -> Result<()> {
use std::cell::RefCell;
let cache: RefCell<lru::LruCache<String, regex::Regex>> =
RefCell::new(lru::LruCache::new(std::num::NonZeroUsize::new(32).unwrap()));
conn.create_scalar_function(
"regexp",
2,
rusqlite::functions::FunctionFlags::SQLITE_UTF8
| rusqlite::functions::FunctionFlags::SQLITE_DETERMINISTIC,
move |ctx| {
let pattern: String = ctx.get(0)?;
let text: String = ctx.get(1)?;
let mut cache = cache.borrow_mut();
let re = match cache.get(&pattern) {
Some(re) => re.clone(),
None => {
let re = regex::Regex::new(&pattern)
.map_err(|e| rusqlite::Error::UserFunctionError(e.to_string().into()))?;
cache.put(pattern.clone(), re.clone());
re
}
};
Ok(re.is_match(&text))
},
)
.into_diagnostic()
}
pub struct SqliteLogStore {
conn: Mutex<Connection>,
path: PathBuf,
}
const PARALLEL_QUERY_THRESHOLD: usize = 200_000;
const BUSY_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
fn enable_wal(conn: &Connection) {
const ATTEMPTS: usize = 5;
for attempt in 0..ATTEMPTS {
match conn.query_row("PRAGMA journal_mode", [], |row| row.get::<_, String>(0)) {
Ok(mode) if mode.eq_ignore_ascii_case("wal") => return,
Ok(_) => {}
Err(e) => {
debug!("could not read journal_mode: {e}");
return;
}
}
if conn.execute_batch("PRAGMA journal_mode = WAL;").is_ok() {
return;
}
if attempt + 1 < ATTEMPTS {
std::thread::sleep(std::time::Duration::from_millis(20));
}
}
debug!("log store is not in WAL mode; another connection may be switching it");
}
impl SqliteLogStore {
pub fn open(path: impl Into<PathBuf>) -> Result<Self> {
let path = path.into();
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent).into_diagnostic()?;
}
let conn = Connection::open(&path).into_diagnostic()?;
conn.busy_timeout(BUSY_TIMEOUT).into_diagnostic()?;
add_regexp_function(&conn)?;
enable_wal(&conn);
conn.execute_batch(
"PRAGMA synchronous = NORMAL;
PRAGMA mmap_size = 268435456;",
)
.into_diagnostic()?;
conn.execute(
"CREATE TABLE IF NOT EXISTS log_entries (
id INTEGER PRIMARY KEY AUTOINCREMENT,
daemon_id TEXT NOT NULL,
timestamp INTEGER NOT NULL,
message TEXT NOT NULL,
level TEXT,
msg TEXT,
logger TEXT,
fields_json TEXT
);",
[],
)
.into_diagnostic()?;
conn.execute(
"CREATE INDEX IF NOT EXISTS idx_daemon_ts ON log_entries(daemon_id, timestamp);",
[],
)
.into_diagnostic()?;
conn.execute(
"CREATE INDEX IF NOT EXISTS idx_daemon_id ON log_entries(daemon_id, id);",
[],
)
.into_diagnostic()?;
conn.execute(
"CREATE INDEX IF NOT EXISTS idx_timestamp ON log_entries(timestamp);",
[],
)
.into_diagnostic()?;
let existing_cols: Vec<String> = {
let mut stmt = conn
.prepare("PRAGMA table_info(log_entries)")
.into_diagnostic()?;
let rows = stmt
.query_map([], |row| row.get::<_, String>(1))
.into_diagnostic()?;
rows.filter_map(|r| r.ok()).collect()
};
for col in ["level", "msg", "logger", "fields_json"] {
if !existing_cols.iter().any(|c| c == col) {
conn.execute(
&format!("ALTER TABLE log_entries ADD COLUMN {col} TEXT"),
[],
)
.into_diagnostic()?;
}
}
conn.execute(
"CREATE INDEX IF NOT EXISTS idx_daemon_level_ts ON log_entries(daemon_id, level, timestamp);",
[],
)
.into_diagnostic()?;
conn.execute(
"CREATE TABLE IF NOT EXISTS log_clear_generations (
daemon_id TEXT PRIMARY KEY,
generation INTEGER NOT NULL DEFAULT 0
);",
[],
)
.into_diagnostic()?;
Ok(Self {
conn: Mutex::new(conn),
path,
})
}
fn row_to_entry(row: &rusqlite::Row) -> rusqlite::Result<LogEntry> {
let id: i64 = row.get(0)?;
let daemon_id: String = row.get(1)?;
let ts_millis: i64 = row.get(2)?;
let message: String = row.get(3)?;
let level: Option<String> = row.get(4)?;
let msg: Option<String> = row.get(5)?;
let logger: Option<String> = row.get(6)?;
let fields_json: Option<String> = row.get(7)?;
let timestamp = Local
.timestamp_millis_opt(ts_millis)
.single()
.unwrap_or_else(Local::now);
Ok(LogEntry {
id,
daemon_id,
timestamp,
message,
level,
msg,
logger,
fields_json,
})
}
fn archive_entries(
&self,
entries: &[LogEntry],
archive_hook: &ArchiveHook,
daemon_id: &DaemonId,
reason: &str,
) -> Result<()> {
use std::process::{Command, Stdio};
if entries.is_empty() {
return Ok(());
}
for chunk in entries.chunks(archive_hook.batch_size.max(1)) {
let mut child = Command::new("sh")
.arg("-c")
.arg(&archive_hook.command)
.stdin(Stdio::piped())
.stdout(Stdio::null())
.stderr(Stdio::piped())
.env("PITCHFORK_DAEMON_ID", daemon_id.qualified())
.env("PITCHFORK_ARCHIVE_REASON", reason)
.spawn()
.into_diagnostic()
.map_err(|e| miette::miette!("failed to spawn archive hook: {e}"))?;
let write_result = {
let stdin = child.stdin.take().expect("piped stdin should be available");
let mut stdin = std::io::BufWriter::new(stdin);
let mut result = Ok(());
for entry in chunk {
let line = serde_json::json!({
"id": entry.id,
"daemon_id": entry.daemon_id,
"timestamp": entry.timestamp.to_rfc3339(),
"message": entry.message,
});
if let Err(e) = writeln!(stdin, "{}", line) {
result = Err(miette::miette!(
"failed to write to archive hook stdin: {e}"
));
break;
}
}
if result.is_ok()
&& let Err(e) = stdin.flush()
{
result = Err(miette::miette!("failed to flush archive hook stdin: {e}"));
}
result
};
if let Err(e) = write_result {
let _ = child.kill();
let _ = child.wait();
return Err(e);
}
let output = child.wait_with_output().into_diagnostic()?;
if !output.status.success() {
let stderr = String::from_utf8_lossy(&output.stderr);
return Err(miette::miette!(
"archive hook failed with status {}: {stderr}",
output.status
));
}
}
Ok(())
}
fn delete_by_ids(&self, ids: &[i64]) -> Result<u64> {
const SQLITE_MAX_VARS: usize = 999;
let mut total = 0u64;
let conn = self.conn.lock().unwrap();
for chunk in ids.chunks(SQLITE_MAX_VARS) {
if chunk.is_empty() {
continue;
}
let placeholders: Vec<String> = (1..=chunk.len()).map(|i| format!("?{i}")).collect();
let sql = format!(
"DELETE FROM log_entries WHERE id IN ({})",
placeholders.join(", ")
);
total += conn
.execute(&sql, rusqlite::params_from_iter(chunk.iter()))
.into_diagnostic()? as u64;
}
Ok(total)
}
pub fn rotate_by_age(
&self,
daemon_id: &DaemonId,
max_age: chrono::Duration,
archive_hook: Option<&ArchiveHook>,
) -> Result<u64> {
let cutoff = (Local::now() - max_age).timestamp_millis();
let hook = archive_hook.filter(|h| h.is_enabled());
if let Some(hook) = hook {
let mut total_deleted = 0u64;
loop {
let entries: Vec<LogEntry> = {
let conn = self.conn.lock().unwrap();
let mut stmt = conn
.prepare(
"SELECT id, daemon_id, timestamp, message, level, msg, logger, fields_json FROM log_entries
WHERE daemon_id = ?1 AND timestamp < ?2
ORDER BY timestamp ASC, id ASC
LIMIT ?3",
)
.into_diagnostic()?;
stmt.query_map(
params![daemon_id.qualified(), cutoff, hook.batch_size as i64],
Self::row_to_entry,
)
.into_diagnostic()?
.collect::<rusqlite::Result<Vec<_>>>()
.into_diagnostic()?
};
if entries.is_empty() {
break;
}
let batch_ids: Vec<i64> = entries.iter().map(|e| e.id).collect();
self.archive_entries(&entries, hook, daemon_id, "age")?;
let deleted = self.delete_by_ids(&batch_ids)?;
total_deleted += deleted;
}
Ok(total_deleted)
} else {
let conn = self.conn.lock().unwrap();
let rows = conn
.execute(
"DELETE FROM log_entries WHERE daemon_id = ?1 AND timestamp < ?2",
params![daemon_id.qualified(), cutoff],
)
.into_diagnostic()?;
Ok(rows as u64)
}
}
pub fn rotate_by_count(
&self,
daemon_id: &DaemonId,
max_count: u64,
archive_hook: Option<&ArchiveHook>,
) -> Result<u64> {
let hook = archive_hook.filter(|h| h.is_enabled());
let to_delete: i64 = {
let conn = self.conn.lock().unwrap();
let count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM log_entries WHERE daemon_id = ?1",
[daemon_id.qualified()],
|row| row.get(0),
)
.into_diagnostic()?;
count.saturating_sub(max_count as i64)
};
if to_delete <= 0 {
return Ok(0);
}
if let Some(hook) = hook {
let mut total_deleted = 0u64;
let mut remaining = to_delete;
loop {
let batch_len = remaining.min(hook.batch_size as i64);
let entries: Vec<LogEntry> = {
let conn = self.conn.lock().unwrap();
let mut stmt = conn
.prepare(
"SELECT id, daemon_id, timestamp, message, level, msg, logger, fields_json FROM log_entries
WHERE daemon_id = ?1
ORDER BY timestamp ASC, id ASC
LIMIT ?2",
)
.into_diagnostic()?;
stmt.query_map(
params![daemon_id.qualified(), batch_len],
Self::row_to_entry,
)
.into_diagnostic()?
.collect::<rusqlite::Result<Vec<_>>>()
.into_diagnostic()?
};
if entries.is_empty() {
break;
}
let fetched = entries.len() as i64;
let batch_ids: Vec<i64> = entries.iter().map(|e| e.id).collect();
self.archive_entries(&entries, hook, daemon_id, "count")?;
let deleted = self.delete_by_ids(&batch_ids)?;
total_deleted += deleted;
remaining -= fetched;
}
Ok(total_deleted)
} else {
let conn = self.conn.lock().unwrap();
let rows = conn
.execute(
"DELETE FROM log_entries WHERE id IN (
SELECT id FROM log_entries WHERE daemon_id = ?1
ORDER BY timestamp ASC, id ASC LIMIT ?2
)",
params![daemon_id.qualified(), to_delete],
)
.into_diagnostic()?;
Ok(rows as u64)
}
}
pub fn migrate_daemon_text_logs(&self, daemon_id: &DaemonId) -> Result<u64> {
let text_path = daemon_id.log_path();
if !text_path.exists() {
return Ok(0);
}
let file = std::fs::File::open(&text_path).into_diagnostic()?;
let reader = BufReader::new(file);
let re = regex::Regex::new(r"^(\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}) ([\w./-]+) (.*)$")
.expect("invalid regex");
let mut current_timestamp: Option<DateTime<Local>> = None;
let mut current_message = String::new();
let mut entries = Vec::with_capacity(1000);
let mut total_migrated: u64 = 0;
for line in reader.lines() {
let line = line.into_diagnostic()?;
if let Some(caps) = re.captures(&line) {
if let Some(ts) = current_timestamp.take() {
entries.push((ts, std::mem::take(&mut current_message)));
}
let ts_str = caps.get(1).map(|m| m.as_str()).unwrap_or_default();
let msg = caps.get(3).map(|m| m.as_str()).unwrap_or_default();
if let Ok(naive) =
chrono::NaiveDateTime::parse_from_str(ts_str, "%Y-%m-%d %H:%M:%S")
{
current_timestamp = Local.from_local_datetime(&naive).single();
current_message = msg.to_string();
}
} else if current_timestamp.is_some() {
current_message.push('\n');
current_message.push_str(&line);
}
if entries.len() >= 1000 {
total_migrated += self.insert_batch(daemon_id, &entries)?;
entries.clear();
}
}
if let Some(ts) = current_timestamp {
entries.push((ts, std::mem::take(&mut current_message)));
}
if !entries.is_empty() {
total_migrated += self.insert_batch(daemon_id, &entries)?;
}
if total_migrated > 0
&& let Err(e) = std::fs::remove_file(&text_path)
{
log::warn!(
"failed to remove legacy log file after migration {}: {e}",
text_path.display()
);
}
Ok(total_migrated)
}
fn insert_batch(
&self,
daemon_id: &DaemonId,
entries: &[(DateTime<Local>, String)],
) -> Result<u64> {
let mut conn = self.conn.lock().unwrap();
let tx = conn.transaction().into_diagnostic()?;
let mut count = 0u64;
{
let mut stmt = tx
.prepare(
"INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
)
.into_diagnostic()?;
for (ts, msg) in entries {
stmt.execute(params![daemon_id.qualified(), ts.timestamp_millis(), msg])
.into_diagnostic()?;
count += 1;
}
}
tx.commit().into_diagnostic()?;
Ok(count)
}
fn build_query_sql(
opts: &LogQuery,
id_range: Option<(i64, i64)>,
) -> (String, Vec<Box<dyn rusqlite::ToSql>>) {
let mut conditions = Vec::new();
let mut query_params: Vec<Box<dyn rusqlite::ToSql>> = Vec::new();
if !opts.daemon_ids.is_empty() {
let placeholders: Vec<String> = (1..=opts.daemon_ids.len())
.map(|i| format!("?{}", i))
.collect();
conditions.push(format!("daemon_id IN ({})", placeholders.join(", ")));
for id in &opts.daemon_ids {
query_params.push(Box::new(id.clone()));
}
}
if let Some(from) = opts.from {
conditions.push(format!("timestamp >= ?{}", query_params.len() + 1));
query_params.push(Box::new(from.timestamp_millis()));
}
if let Some(to) = opts.to {
conditions.push(format!("timestamp <= ?{}", query_params.len() + 1));
query_params.push(Box::new(to.timestamp_millis()));
}
if let Some(after_id) = opts.after_id {
conditions.push(format!("id > ?{}", query_params.len() + 1));
query_params.push(Box::new(after_id));
}
if let Some(before_id) = opts.before_id {
conditions.push(format!("id < ?{}", query_params.len() + 1));
query_params.push(Box::new(before_id));
}
if let Some((start, end)) = id_range {
conditions.push(format!("id > ?{}", query_params.len() + 1));
query_params.push(Box::new(start));
conditions.push(format!("id <= ?{}", query_params.len() + 1));
query_params.push(Box::new(end));
}
let mut message_conditions = Vec::new();
for filter in &opts.message_filters {
match filter {
MessageFilter::Contains {
pattern,
case_sensitive,
} => {
let param_index = query_params.len() + 1;
if *case_sensitive {
message_conditions
.push(format!("INSTR(message, ?{idx}) > 0", idx = param_index));
query_params.push(Box::new(pattern.clone()));
} else {
let escaped = escape_like_pattern(pattern);
let param = format!("%{}%", escaped);
message_conditions.push(format!(
"LOWER(message) LIKE LOWER(?{idx}) ESCAPE '\\'",
idx = param_index
));
query_params.push(Box::new(param));
}
}
MessageFilter::Regex { pattern } => {
let param_index = query_params.len() + 1;
message_conditions.push(format!("message REGEXP ?{param_index}"));
query_params.push(Box::new(pattern.clone()));
}
}
}
if !message_conditions.is_empty() {
conditions.push(format!("({})", message_conditions.join(" OR ")));
}
for filter in &opts.field_filters {
match filter {
FieldFilter::LevelMin(level) => {
let matching = crate::log_store::levels_at_or_above(level);
if matching.is_empty() {
conditions.push("0".to_string());
} else {
let placeholders = matching
.iter()
.map(|l| {
let idx = query_params.len() + 1;
query_params.push(Box::new((*l).to_string()));
format!("?{idx}")
})
.collect::<Vec<_>>()
.join(", ");
conditions.push(format!("level IN ({placeholders})"));
}
}
FieldFilter::FieldEq { key, value } => {
let key_idx = query_params.len() + 1;
let val_idx = query_params.len() + 2;
let json_literal = text_to_json_literal(value);
let raw_idx = query_params.len() + 3;
conditions.push(format!(
"EXISTS (SELECT 1 FROM json_each(fields_json) \
WHERE json_each.key = ?{key_idx} \
AND (json_each.value IS json_extract(?{val_idx}, '$') \
OR json_each.value IS ?{raw_idx}))"
));
query_params.push(Box::new(key.clone()));
query_params.push(Box::new(json_literal));
query_params.push(Box::new(value.clone()));
}
FieldFilter::LoggerContains(pattern) => {
let param_index = query_params.len() + 1;
conditions.push(format!("logger LIKE ?{param_index} ESCAPE '\\'"));
let escaped = crate::log_store::escape_like_pattern(pattern);
query_params.push(Box::new(format!("%{escaped}%")));
}
}
}
let where_clause = if conditions.is_empty() {
String::new()
} else {
format!("WHERE {}", conditions.join(" AND "))
};
let order = if opts.order_desc { "DESC" } else { "ASC" };
let limit_clause = opts
.limit
.map(|n| format!("LIMIT {}", n))
.unwrap_or_default();
let columns = if opts.include_structured {
"id, daemon_id, timestamp, message, level, msg, logger, fields_json"
} else {
"id, daemon_id, timestamp, message, NULL, NULL, NULL, NULL"
};
let sql = format!(
"SELECT {columns} FROM log_entries {} ORDER BY timestamp {}, id {} {}",
where_clause, order, order, limit_clause
);
(sql, query_params)
}
fn should_parallelize(opts: &LogQuery) -> bool {
if opts.daemon_ids.len() != 1 {
return false;
}
if opts.after_id.is_some() {
return false;
}
let limit = opts.limit.unwrap_or(usize::MAX);
if limit < PARALLEL_QUERY_THRESHOLD {
return false;
}
std::thread::available_parallelism()
.map(|n| n.get() >= 2)
.unwrap_or(false)
}
fn query_parallel(&self, opts: &LogQuery) -> Result<Vec<LogEntry>> {
let max_threads = 2;
let num_threads = std::thread::available_parallelism()
.map(|n| n.get().min(max_threads))
.unwrap_or(1);
let max_id: Option<i64> = {
let conn = self.conn.lock().unwrap();
conn.query_row(
"SELECT MAX(id) FROM log_entries WHERE daemon_id = ?1",
params![&opts.daemon_ids[0]],
|row| row.get(0),
)
.ok()
};
let Some(max_id) = max_id else {
return Ok(Vec::new());
};
if max_id == 0 {
return Ok(Vec::new());
}
let shard_size = (max_id as usize).div_ceil(num_threads);
let path = self.path.clone();
let opts = opts.clone();
let needs_regexp = opts
.message_filters
.iter()
.any(|f| matches!(f, MessageFilter::Regex { .. }));
let shards: Vec<Result<Vec<LogEntry>>> = std::thread::scope(|s| {
(0..num_threads)
.map(|i| {
let start = (i * shard_size) as i64;
let end = if i == num_threads - 1 {
max_id
} else {
((i + 1) * shard_size) as i64
};
let opts = opts.clone();
let path = path.clone();
s.spawn(move || -> Result<Vec<LogEntry>> {
let conn = Connection::open_with_flags(
&path,
OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_NO_MUTEX,
)
.into_diagnostic()?;
conn.execute_batch(
"PRAGMA mmap_size = 268435456;
PRAGMA query_only = 1;",
)
.into_diagnostic()?;
if needs_regexp {
add_regexp_function(&conn)?;
}
let (sql, query_params) = Self::build_query_sql(&opts, Some((start, end)));
Self::execute_built_query(&conn, &sql, &query_params)
})
})
.map(|h| h.join().unwrap())
.collect()
});
let mut merged = Vec::new();
if opts.order_desc {
for shard in shards.into_iter().rev() {
merged.extend(shard?);
}
} else {
for shard in shards {
merged.extend(shard?);
}
}
if let Some(limit) = opts.limit
&& merged.len() > limit
{
merged.truncate(limit);
}
Ok(merged)
}
fn execute_built_query(
conn: &Connection,
sql: &str,
query_params: &[Box<dyn rusqlite::ToSql>],
) -> Result<Vec<LogEntry>> {
let mut stmt = conn.prepare(sql).into_diagnostic()?;
let params_ref: Vec<&dyn rusqlite::ToSql> =
query_params.iter().map(|p| p.as_ref()).collect();
let rows = stmt
.query_map(params_ref.as_slice(), Self::row_to_entry)
.into_diagnostic()?;
let mut entries = Vec::new();
for row in rows {
entries.push(row.into_diagnostic()?);
}
Ok(entries)
}
pub fn distinct_loggers(&self, daemon_id: &str) -> Result<Vec<String>> {
let conn = self.conn.lock().unwrap();
let mut stmt = conn
.prepare(
"SELECT DISTINCT logger FROM ( \
SELECT logger FROM log_entries \
WHERE daemon_id = ?1 AND logger IS NOT NULL \
ORDER BY id DESC LIMIT 5000 \
) ORDER BY logger",
)
.into_diagnostic()?;
let rows = stmt
.query_map(params![daemon_id], |row| row.get::<_, String>(0))
.into_diagnostic()?;
let mut loggers = Vec::new();
for row in rows {
loggers.push(row.into_diagnostic()?);
}
Ok(loggers)
}
pub fn distinct_field_keys(&self, daemon_id: &str) -> Result<Vec<String>> {
let conn = self.conn.lock().unwrap();
let mut stmt = conn
.prepare(
"SELECT DISTINCT je.key \
FROM ( \
SELECT fields_json FROM log_entries \
WHERE daemon_id = ?1 AND fields_json IS NOT NULL \
ORDER BY id DESC LIMIT 5000 \
), json_each(fields_json) AS je \
ORDER BY je.key",
)
.into_diagnostic()?;
let rows = stmt
.query_map(params![daemon_id], |row| row.get::<_, String>(0))
.into_diagnostic()?;
let mut keys = Vec::new();
for row in rows {
keys.push(row.into_diagnostic()?);
}
Ok(keys)
}
}
impl LogStore for SqliteLogStore {
fn append(&self, daemon_id: &DaemonId, message: &str) -> Result<()> {
let id = daemon_id.qualified();
let msg = message.to_string();
let conn = self.conn.lock().unwrap();
let ts = Local::now().timestamp_millis();
let _ = conn
.execute(
"INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
params![id, ts, msg],
)
.into_diagnostic()?;
Ok(())
}
fn append_batch(&self, daemon_id: &DaemonId, messages: &[String]) -> Result<()> {
if messages.is_empty() {
return Ok(());
}
let id = daemon_id.qualified();
let mut conn = self.conn.lock().unwrap();
let tx = conn
.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
.into_diagnostic()?;
let base_ts = Local::now().timestamp_millis();
{
let mut stmt = tx
.prepare(
"INSERT INTO log_entries (daemon_id, timestamp, message) VALUES (?1, ?2, ?3)",
)
.into_diagnostic()?;
for msg in messages.iter() {
let ts = base_ts;
stmt.execute(params![id, ts, msg]).into_diagnostic()?;
}
}
tx.commit().into_diagnostic()?;
Ok(())
}
fn append_structured(&self, daemon_id: &DaemonId, parsed: &ParsedLog) -> Result<()> {
let id = daemon_id.qualified();
let conn = self.conn.lock().unwrap();
let ts = Local::now().timestamp_millis();
let _ = conn
.execute(
"INSERT INTO log_entries (daemon_id, timestamp, message, level, msg, logger, fields_json) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
params![id, ts, parsed.message, parsed.level, parsed.msg, parsed.logger, parsed.fields_json],
)
.into_diagnostic()?;
Ok(())
}
fn append_structured_batch(&self, daemon_id: &DaemonId, entries: &[ParsedLog]) -> Result<()> {
if entries.is_empty() {
return Ok(());
}
let id = daemon_id.qualified();
let mut conn = self.conn.lock().unwrap();
let tx = conn
.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
.into_diagnostic()?;
let base_ts = Local::now().timestamp_millis();
{
let mut stmt = tx
.prepare(
"INSERT INTO log_entries (daemon_id, timestamp, message, level, msg, logger, fields_json) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
)
.into_diagnostic()?;
for entry in entries.iter() {
let ts = base_ts;
stmt.execute(params![
id,
ts,
entry.message,
entry.level,
entry.msg,
entry.logger,
entry.fields_json
])
.into_diagnostic()?;
}
}
tx.commit().into_diagnostic()?;
Ok(())
}
fn query(&self, opts: &LogQuery) -> Result<Vec<LogEntry>> {
if Self::should_parallelize(opts)
&& self.path.as_os_str() != ":memory:"
&& let Ok(entries) = self.query_parallel(opts)
{
return Ok(entries);
}
let conn = self.conn.lock().unwrap();
let (sql, query_params) = Self::build_query_sql(opts, None);
Self::execute_built_query(&conn, &sql, &query_params)
}
fn tail(&self, daemon_id: &DaemonId, after_id: Option<i64>) -> Result<Vec<LogEntry>> {
self.query(&LogQuery {
daemon_ids: vec![daemon_id.qualified()],
from: None,
to: None,
limit: None,
order_desc: false,
after_id,
before_id: None,
message_filters: Vec::new(),
field_filters: Vec::new(),
include_structured: false,
})
}
fn clear(&self, daemon_ids: &[DaemonId]) -> Result<()> {
let mut conn = self.conn.lock().unwrap();
let tx = conn.transaction().into_diagnostic()?;
for id in daemon_ids {
tx.execute(
"DELETE FROM log_entries WHERE daemon_id = ?1",
params![id.qualified()],
)
.into_diagnostic()?;
tx.execute(
"INSERT INTO log_clear_generations (daemon_id, generation)
VALUES (?1, 1)
ON CONFLICT(daemon_id) DO UPDATE SET generation = generation + 1",
params![id.qualified()],
)
.into_diagnostic()?;
}
tx.commit().into_diagnostic()?;
Ok(())
}
fn last_id(&self, daemon_id: &DaemonId) -> Result<Option<i64>> {
let conn = self.conn.lock().unwrap();
let id: Option<i64> = conn
.query_row(
"SELECT MAX(id) FROM log_entries WHERE daemon_id = ?1",
params![daemon_id.qualified()],
|row| row.get(0),
)
.into_diagnostic()?;
Ok(id)
}
fn list_daemon_ids(&self) -> Result<Vec<String>> {
let conn = self.conn.lock().unwrap();
let mut stmt = conn
.prepare("SELECT DISTINCT daemon_id FROM log_entries")
.into_diagnostic()?;
let ids = stmt
.query_map([], |row| {
let id: String = row.get(0)?;
Ok(id)
})
.into_diagnostic()?
.filter_map(|r| r.ok())
.collect();
Ok(ids)
}
fn apply_retention(
&self,
policy: &super::RetentionPolicy,
excluded_daemon_ids: &[DaemonId],
archive_hook: Option<&ArchiveHook>,
) -> Result<u64> {
let daemon_ids = self.list_daemon_ids()?;
let excluded: HashSet<String> = excluded_daemon_ids.iter().map(|d| d.qualified()).collect();
let mut total = 0u64;
for id_str in daemon_ids {
if excluded.contains(&id_str) {
continue;
}
let id = DaemonId::parse(&id_str).unwrap_or_else(|_| {
DaemonId::try_new("global", &id_str).unwrap_or_else(|_| DaemonId::pitchfork())
});
if let Some(dur) = policy.age {
total += self.rotate_by_age(&id, dur, archive_hook)?;
}
if let Some(n) = policy.count {
total += self.rotate_by_count(&id, n, archive_hook)?;
}
}
Ok(total)
}
fn apply_retention_for_daemon(
&self,
daemon_id: &DaemonId,
policy: &super::RetentionPolicy,
archive_hook: Option<&ArchiveHook>,
) -> Result<u64> {
let mut total = 0u64;
if let Some(dur) = policy.age {
total += self.rotate_by_age(daemon_id, dur, archive_hook)?;
}
if let Some(n) = policy.count {
total += self.rotate_by_count(daemon_id, n, archive_hook)?;
}
Ok(total)
}
fn last_clear_generation(&self, daemon_id: &DaemonId) -> Result<Option<u64>> {
let conn = self.conn.lock().unwrap();
let mut stmt = conn
.prepare("SELECT generation FROM log_clear_generations WHERE daemon_id = ?1")
.into_diagnostic()?;
let generation: Option<i64> = stmt
.query_row(params![daemon_id.qualified()], |row| row.get(0))
.optional()
.into_diagnostic()?;
generation
.map(|generation| {
u64::try_from(generation)
.map_err(|_| miette::miette!("log clear generation cannot be negative"))
})
.transpose()
}
fn query_with_generation(
&self,
opts: &LogQuery,
daemon_id: &DaemonId,
) -> Result<(Vec<LogEntry>, Option<u64>)> {
let conn = self.conn.lock().unwrap();
conn.execute_batch("BEGIN").into_diagnostic()?;
let result = (|| {
let (sql, query_params) = Self::build_query_sql(opts, None);
let entries = Self::execute_built_query(&conn, &sql, &query_params)?;
let generation: Option<i64> = conn
.query_row(
"SELECT generation FROM log_clear_generations WHERE daemon_id = ?1",
params![daemon_id.qualified()],
|row| row.get(0),
)
.optional()
.into_diagnostic()?;
let generation = generation
.map(|g| {
u64::try_from(g)
.map_err(|_| miette::miette!("log clear generation cannot be negative"))
})
.transpose()?;
Ok((entries, generation))
})();
let _ = conn.execute_batch("ROLLBACK");
result
}
}
use once_cell::sync::Lazy;
use std::sync::Arc;
pub static LOG_STORE: Lazy<Arc<SqliteLogStore>> = Lazy::new(|| {
let path = crate::env::PITCHFORK_LOGS_DIR.join("logs.db");
let mut is_fallback = false;
let store = Arc::new(SqliteLogStore::open(&path).unwrap_or_else(|e| {
error!(
"failed to open log store at {}: {e}. Falling back to in-memory store; logs will not persist across restarts.",
path.display()
);
is_fallback = true;
SqliteLogStore::open(":memory:").expect("in-memory SQLite should always open")
}));
if !is_fallback {
if let Err(e) = auto_migrate_legacy_logs(&store) {
warn!("legacy log auto-migration failed: {e}");
}
} else {
warn!(
"skipping legacy log auto-migration because log store is in-memory (no durable destination)"
);
}
store
});
fn auto_migrate_legacy_logs(store: &SqliteLogStore) -> Result<()> {
let logs_dir = &*crate::env::PITCHFORK_LOGS_DIR;
if !logs_dir.exists() {
return Ok(());
}
let Ok(entries) = std::fs::read_dir(logs_dir) else {
return Ok(());
};
let mut total_migrated = 0u64;
let mut migrated_ids = Vec::new();
for entry in entries.flatten() {
let path = entry.path();
if !path.is_dir() {
continue;
}
let file_name = path
.file_name()
.map_or(String::new(), |n| n.to_string_lossy().to_string());
if file_name == "pitchfork" {
continue;
}
if !file_name.contains("--") {
continue;
}
let log_file = path.join(format!("{file_name}.log"));
if !log_file.exists() {
continue;
}
let daemon_id = match DaemonId::from_safe_path(&file_name) {
Ok(id) => id,
Err(_) => continue,
};
if daemon_id == DaemonId::pitchfork() {
continue;
}
match store.migrate_daemon_text_logs(&daemon_id) {
Ok(0) => {}
Ok(n) => {
total_migrated += n;
migrated_ids.push(daemon_id.qualified());
}
Err(e) => {
warn!(
"failed to migrate text logs for {}: {e}",
daemon_id.qualified()
);
}
}
}
if total_migrated > 0 {
warn!(
"auto-migrated {total_migrated} legacy log entries from {count} daemon(s): {ids}",
count = migrated_ids.len(),
ids = migrated_ids.join(", ")
);
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::text_to_json_literal;
use super::*;
use crate::log_store::LogStore;
#[test]
fn test_json_literal_booleans() {
assert_eq!(text_to_json_literal("true"), "true");
assert_eq!(text_to_json_literal("TRUE"), "true");
assert_eq!(text_to_json_literal("false"), "false");
assert_eq!(text_to_json_literal("False"), "false");
}
#[test]
fn test_json_literal_null() {
assert_eq!(text_to_json_literal("null"), "null");
assert_eq!(text_to_json_literal("NULL"), "null");
}
#[test]
fn test_json_literal_valid_numbers() {
assert_eq!(text_to_json_literal("42"), "42");
assert_eq!(text_to_json_literal("-1"), "-1");
assert_eq!(text_to_json_literal("0"), "0");
assert_eq!(text_to_json_literal("3.14"), "3.14");
assert_eq!(text_to_json_literal("1e10"), "1e10");
assert_eq!(text_to_json_literal("1.5e-3"), "1.5e-3");
}
#[test]
fn test_json_literal_rejects_invalid_json_numbers() {
assert_eq!(text_to_json_literal("+42"), r#""+42""#);
assert_eq!(text_to_json_literal("1."), r#""1.""#);
assert_eq!(text_to_json_literal(".5"), r#"".5""#);
assert_eq!(text_to_json_literal("inf"), r#""inf""#);
assert_eq!(text_to_json_literal("nan"), r#""nan""#);
assert_eq!(text_to_json_literal("infinity"), r#""infinity""#);
}
#[test]
fn test_json_literal_strings() {
assert_eq!(text_to_json_literal("hello"), r#""hello""#);
assert_eq!(text_to_json_literal("req_1"), r#""req_1""#);
assert_eq!(text_to_json_literal(r#"a"b"#), r#""a\"b""#);
}
#[test]
fn write_waits_for_a_contended_lock() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("logs.db");
let store = SqliteLogStore::open(&path).unwrap();
let blocker = Connection::open(&path).unwrap();
blocker.busy_timeout(BUSY_TIMEOUT).unwrap();
blocker.execute_batch("BEGIN IMMEDIATE").unwrap();
let releaser = std::thread::spawn(move || {
std::thread::sleep(std::time::Duration::from_millis(300));
blocker.execute_batch("COMMIT").unwrap();
});
let id = DaemonId::try_new("test", "blocked").unwrap();
let entries = vec![crate::log_parse::parse("held-lock-line", "text")];
store
.append_structured_batch(&id, &entries)
.expect("a contended write must wait for the lock, not fail");
releaser.join().unwrap();
let found = store
.query(&LogQuery {
daemon_ids: vec![id.qualified()],
..Default::default()
})
.unwrap()
.len();
assert_eq!(found, 1, "the batch written under contention was lost");
}
#[test]
fn concurrent_writers_do_not_lose_batches() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("logs.db");
const WRITERS: usize = 4;
const BATCHES: usize = 15;
const PER_BATCH: usize = 10;
let barrier = std::sync::Barrier::new(WRITERS);
std::thread::scope(|scope| {
for writer in 0..WRITERS {
let path = path.clone();
let barrier = &barrier;
scope.spawn(move || {
barrier.wait();
let store = SqliteLogStore::open(&path).unwrap();
let id = DaemonId::try_new("test", format!("w{writer}")).unwrap();
for batch in 0..BATCHES {
let entries: Vec<ParsedLog> = (0..PER_BATCH)
.map(|i| {
crate::log_parse::parse(&format!("w{writer}-{batch}-{i}"), "text")
})
.collect();
store
.append_structured_batch(&id, &entries)
.expect("concurrent batch write must not fail");
}
});
}
});
let store = SqliteLogStore::open(&path).unwrap();
for writer in 0..WRITERS {
let id = DaemonId::try_new("test", format!("w{writer}")).unwrap();
let found = store
.query(&LogQuery {
daemon_ids: vec![id.qualified()],
..Default::default()
})
.unwrap()
.len();
assert_eq!(
found,
BATCHES * PER_BATCH,
"writer {writer} lost entries: got {found}"
);
}
}
#[test]
fn batch_rows_share_single_timestamp() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("logs.db");
let store = SqliteLogStore::open(&path).unwrap();
let id = DaemonId::try_new("test", "burst").unwrap();
let entries: Vec<ParsedLog> = (0..50)
.map(|i| crate::log_parse::parse(&format!("line-{i}"), "text"))
.collect();
store.append_structured_batch(&id, &entries).unwrap();
let conn = store.conn.lock().unwrap();
let timestamps: Vec<i64> = {
let mut stmt = conn
.prepare("SELECT timestamp FROM log_entries ORDER BY id")
.unwrap();
let rows = stmt.query_map([], |row| row.get::<_, i64>(0)).unwrap();
rows.map(|r| r.unwrap()).collect()
};
assert_eq!(timestamps.len(), 50);
let distinct: std::collections::HashSet<i64> = timestamps.into_iter().collect();
assert_eq!(
distinct.len(),
1,
"all rows in a batch must share one timestamp so (timestamp, id) ordering matches insertion order"
);
}
#[test]
fn burst_batches_preserve_insertion_order() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("logs.db");
let store = SqliteLogStore::open(&path).unwrap();
let id = DaemonId::try_new("test", "burst").unwrap();
const BATCHES: usize = 3;
const PER_BATCH: usize = 300;
for batch in 0..BATCHES {
let entries: Vec<ParsedLog> = (0..PER_BATCH)
.map(|i| crate::log_parse::parse(&format!("{}", batch * PER_BATCH + i), "text"))
.collect();
store.append_structured_batch(&id, &entries).unwrap();
}
let entries = store
.query(&LogQuery {
daemon_ids: vec![id.qualified()],
..Default::default()
})
.unwrap();
assert_eq!(entries.len(), BATCHES * PER_BATCH);
for (n, entry) in entries.iter().enumerate() {
assert_eq!(
entry.message,
n.to_string(),
"log line {n} read back out of order (got '{}')",
entry.message
);
}
}
}