use crate::db;
use crate::util::UnwrapPoison;
use crate::util::json;
use anyhow::Context;
use db::{Row, Value, params};
use futures_util::FutureExt;
use serde::{Deserialize, Serialize};
use std::io;
use std::panic::AssertUnwindSafe;
use std::path::Path;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use tokio::sync::OnceCell;
use tokio::sync::mpsc::{UnboundedReceiver, UnboundedSender};
use tracing::level_filters::LevelFilter;
use tracing_subscriber::fmt::MakeWriter;
use tracing_subscriber::{EnvFilter, Layer, fmt, layer::SubscriberExt};
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct LogEntry {
pub timestamp: String,
pub level: String,
pub target: String,
pub message: String,
#[serde(default)]
pub fields: serde_json::Value,
#[serde(default)]
pub agent_id: String,
#[serde(default)]
pub agent_role: String,
#[serde(default)]
pub workspace: String,
}
crate::columns! {
LOGS_COLUMNS [LOGS] {
TIMESTAMP => "timestamp",
LEVEL => "level",
TARGET => "target",
MESSAGE => "message",
FIELDS => "fields",
AGENT_ID => "agent_id",
AGENT_ROLE => "agent_role",
WORKSPACE => "workspace",
}
}
#[derive(Clone, Debug)]
pub struct LogStore {
pub(crate) conn: crate::db::Connection,
}
pub static LOG_STORE: OnceCell<LogStore> = OnceCell::const_new();
#[derive(Debug)]
pub(crate) struct GrepTelemetryRow<'a> {
pub command: &'a str,
pub served: bool,
pub reason: &'a str,
pub recursive: bool,
pub piped: bool,
pub operand_count: usize,
pub flags: &'a str,
pub mode: &'a str,
pub workspace: &'a str,
pub grep_count: usize,
pub served_count: usize,
pub skipped_count: usize,
pub duration_ms: Option<i64>,
pub exit_code: Option<i32>,
}
pub(crate) struct StrayEvent<'a> {
pub source: &'a str,
pub message: &'a str,
pub recipe: &'a str,
pub scope: &'a str,
pub session: &'a str,
pub workspace: &'a str,
pub duration_ms: Option<i64>,
}
pub(crate) const STRAY_SOURCE_SHELL: &str = "shell";
pub(crate) const STRAY_SOURCE_CHROME_TOOL: &str = "chrome-tool";
impl LogStore {
pub(crate) async fn open(root: &Path) -> anyhow::Result<Self> {
let conn = crate::db::open_store(root, "logs", "").await?;
crate::db::migrations::run_migrations(&conn, crate::db::migrations::TargetDb::Logs).await?;
Ok(Self { conn })
}
pub(crate) async fn insert_batch(&self, entries: &[LogEntry]) -> anyhow::Result<()> {
if entries.is_empty() {
return Ok(());
}
let tx = self
.conn
.begin_tx()
.await
.context("Failed to begin log insert transaction")?;
for entry in entries {
tx.execute(
"INSERT INTO logs (timestamp, level, target, message, fields, agent_id, agent_role, workspace) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
params![
entry.timestamp.clone(),
entry.level.clone(),
entry.target.clone(),
entry.message.clone(),
serde_json::to_string(&entry.fields)
.expect("log entry fields serialization failed; this should not happen"),
entry.agent_id.clone(),
entry.agent_role.clone(),
entry.workspace.clone(),
],
)
.await
.context("Failed to insert log entry in batch")?;
}
tx.commit()
.await
.context("Failed to commit log insert transaction")?;
Ok(())
}
pub async fn delete_older_than(&self, level: &str, cutoff: &str) -> anyhow::Result<u64> {
let n = self
.conn
.execute(
"DELETE FROM logs WHERE level = ?1 AND timestamp < ?2",
params![level, cutoff],
)
.await
.context("Failed to delete old log entries")?;
Ok(n)
}
pub(crate) async fn clear_logs(&self, level_filter: Option<&str>) -> anyhow::Result<u64> {
let filters = LogQuery {
level: level_filter.map(str::to_owned),
..LogQuery::default()
};
let (where_sql, values) = build_where_clause(&filters);
self.conn
.execute(&format!("DELETE FROM logs {where_sql}"), values)
.await
.context("Failed to clear log entries")
}
pub(crate) async fn has_reason(&self, message: &str, reason: &str) -> anyhow::Result<bool> {
fn reason_of(fields: &str) -> Option<String> {
serde_json::from_str::<serde_json::Value>(fields)
.ok()?
.get("reason")
.and_then(serde_json::Value::as_str)
.map(str::to_owned)
}
let rows = self
.conn
.query(
"SELECT fields FROM logs WHERE message = ?1",
params![message],
)
.await
.context("Failed to read the log rows carrying a message")?;
for row in &rows {
if reason_of(&row.get::<String>(0)?).as_deref() == Some(reason) {
return Ok(true);
}
}
Ok(false)
}
pub async fn query(&self, filters: &LogQuery) -> anyhow::Result<(Vec<LogEntry>, usize)> {
let (where_sql, values) = build_where_clause(filters);
let count_sql = format!("SELECT COUNT(*) FROM logs {where_sql}");
let total = self
.conn
.query_row(&count_sql, values.clone(), |row| row.get::<i64>(0))
.await
.map(|n| usize::try_from(n).unwrap_or(0))?;
if total == 0 {
return Ok((vec![], 0));
}
let limit: i64 = i64::try_from(filters.limit.unwrap_or(100).min(1000))
.expect("log query limit overflowed i64; limit must be <= i64::MAX");
let offset: i64 = i64::try_from(filters.offset.unwrap_or(0))
.expect("log query offset overflowed i64; offset must be <= i64::MAX");
let mut data_values = values;
data_values.push(Value::Integer(limit));
data_values.push(Value::Integer(offset));
let data_sql = format!(
"SELECT {LOGS_COLUMNS} FROM logs {where_sql} ORDER BY id DESC LIMIT ? OFFSET ?",
);
let rows = self
.conn
.query(&data_sql, data_values)
.await
.context("Data query failed")?;
let mut entries = Vec::new();
for row in rows {
entries.push(log_entry_from_row(&row)?);
}
Ok((entries, total))
}
pub(crate) async fn record_grep_telemetry(
&self,
row: GrepTelemetryRow<'_>,
) -> anyhow::Result<()> {
self.conn
.execute(
"INSERT INTO grep_telemetry \
(recorded_at, command, served, reason, recursive, piped, operand_count, flags, \
mode, workspace, grep_count, served_count, skipped_count, duration_ms, exit_code) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15)",
params![
crate::db::now(),
row.command,
i64::from(row.served),
row.reason,
i64::from(row.recursive),
i64::from(row.piped),
i64::try_from(row.operand_count)?,
row.flags,
row.mode,
row.workspace,
i64::try_from(row.grep_count)?,
i64::try_from(row.served_count)?,
i64::try_from(row.skipped_count)?,
row.duration_ms,
row.exit_code.map(i64::from),
],
)
.await?;
Ok(())
}
}
pub(crate) async fn record_stray_event(event: StrayEvent<'_>) {
let Some(store) = LOG_STORE.get() else {
return;
};
let (agent_id, agent_role) = crate::agent::tool_record_attribution();
let written = write_warn_row(
store,
LogEntry {
target: event.source.to_string(),
message: event.message.to_string(),
fields: serde_json::json!({
"detail": event.recipe,
"scope": event.scope,
"session": event.session,
"duration_ms": event.duration_ms,
}),
workspace: event.workspace.to_string(),
agent_id,
agent_role,
..LogEntry::default()
},
)
.await;
if let Err(error) = written {
tracing::debug!(
%error,
target = %event.source,
"could not persist the leftover-process record"
);
}
}
pub(crate) enum IssueWrite {
Written,
AlreadyRecorded,
NotRecorded,
}
async fn write_warn_row(store: &LogStore, entry: LogEntry) -> Result<(), String> {
let entry = LogEntry {
timestamp: crate::db::now(),
level: "WARN".to_string(),
..entry
};
store
.insert_batch(&[entry])
.await
.map_err(|e| e.to_string())
}
pub(crate) async fn record_issue_once(
message: &str,
reason: &str,
target: &str,
fields: serde_json::Value,
) -> IssueWrite {
let Some(store) = LOG_STORE.get() else {
return IssueWrite::NotRecorded;
};
match store.has_reason(message, reason).await {
Ok(true) => return IssueWrite::AlreadyRecorded,
Ok(false) => {}
Err(e) => tracing::warn!(
error = %e,
"could not read whether this issue was already recorded — recording it again"
),
}
let mut fields = fields;
if let Some(object) = fields.as_object_mut() {
object.insert(
"reason".to_string(),
serde_json::Value::String(reason.to_string()),
);
}
match write_warn_row(
store,
LogEntry {
target: target.to_string(),
message: message.to_string(),
fields,
..LogEntry::default()
},
)
.await
{
Ok(()) => IssueWrite::Written,
Err(error) => {
tracing::warn!(%error, "could not record an issue for the Issues view");
IssueWrite::NotRecorded
}
}
}
#[derive(Debug, Clone, Default)]
pub struct LogQuery {
pub level: Option<String>,
pub target: Option<String>,
pub search: Option<String>,
pub since: Option<String>,
pub limit: Option<usize>,
pub offset: Option<usize>,
}
fn build_where_clause(filters: &LogQuery) -> (String, Vec<Value>) {
let mut conditions: Vec<String> = Vec::new();
let mut values: Vec<Value> = Vec::new();
if let Some(ref levels_str) = filters.level
&& !levels_str.is_empty()
{
let levels: Vec<Value> = levels_str
.split(',')
.map(str::trim)
.filter(|s| !s.is_empty())
.map(|s| Value::Text(s.to_string()))
.collect();
if !levels.is_empty() {
conditions.push(format!(
"level IN ({})",
db::sql_in_placeholders(levels.len()),
));
values.extend(levels);
}
}
if let Some(ref target) = filters.target {
conditions.push("target LIKE ?".into());
values.push(Value::Text(format!("{target}%")));
}
if let Some(ref search) = filters.search
&& !search.is_empty()
{
let val = Value::Text(format!("%{search}%"));
conditions.push("(target LIKE ? OR message LIKE ?)".into());
values.push(val.clone());
values.push(val);
}
if let Some(ref since) = filters.since {
conditions.push("timestamp >= ?".into());
values.push(Value::Text(since.clone()));
}
if conditions.is_empty() {
(String::new(), values)
} else {
(format!("WHERE {}", conditions.join(" AND ")), values)
}
}
fn log_entry_from_row(row: &Row) -> anyhow::Result<LogEntry> {
let timestamp = row.get::<String>(COL_LOGS_TIMESTAMP)?;
let level = row.get::<String>(COL_LOGS_LEVEL)?;
let target = row.get::<String>(COL_LOGS_TARGET)?;
let message = row.get::<String>(COL_LOGS_MESSAGE)?;
let fields_str = row.get::<String>(COL_LOGS_FIELDS)?;
let fields: serde_json::Value =
serde_json::from_str(&fields_str).unwrap_or(serde_json::Value::Null);
let agent_id = row.get::<String>(COL_LOGS_AGENT_ID)?;
let agent_role = row.get::<String>(COL_LOGS_AGENT_ROLE)?;
let workspace = row.get::<String>(COL_LOGS_WORKSPACE)?;
Ok(LogEntry {
timestamp,
level,
target,
message,
fields,
agent_id,
agent_role,
workspace,
})
}
pub(crate) const DEFAULT_LOG_FILTER: &str =
"info,turso_core=warn,tantivy=warn,fff=off,pdf_extract=error";
pub(crate) fn log_layers<W>(
writer: W,
env_filter: EnvFilter,
) -> impl tracing::Subscriber + Send + Sync + 'static
where
W: for<'a> fmt::MakeWriter<'a> + Send + Sync + 'static,
{
tracing_subscriber::registry()
.with(
fmt::Layer::new()
.json()
.with_writer(writer)
.with_ansi(false)
.with_filter(env_filter),
)
.with(
crate::db::checkpoint_cause::CauseCaptureLayer
.with_filter(crate::db::checkpoint_cause::filter()),
)
}
fn install_log_bridge(level: LevelFilter) -> anyhow::Result<()> {
use tracing_log::AsLog;
tracing_log::LogTracer::builder()
.with_max_level(level.as_log())
.init()
.map_err(|e| anyhow::anyhow!("failed to install the log→tracing bridge: {e}"))
}
fn log_bridge_level(env_filter: &EnvFilter) -> LevelFilter {
env_filter.max_level_hint().unwrap_or(LevelFilter::TRACE)
}
pub async fn init_tracing(
storage_root: &Path,
) -> anyhow::Result<(Arc<LogStore>, tokio::sync::broadcast::Sender<String>)> {
let store = match LogStore::open(storage_root).await {
Ok(store) => store,
Err(e) => {
crate::boot::clear_boot_diagnostics();
return Err(e);
}
};
LOG_STORE
.set(store.clone())
.map_err(|_| anyhow::anyhow!("LOG_STORE already initialized"))?;
let log_store = Arc::new(store);
let (log_tx, log_rx) = tokio::sync::mpsc::unbounded_channel();
let (broadcast_tx, _) = tokio::sync::broadcast::channel(256);
spawn_log_writer(Arc::clone(&log_store), log_rx, broadcast_tx.clone());
let env_filter =
EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new(DEFAULT_LOG_FILTER));
let bridge_level = log_bridge_level(&env_filter);
tracing::subscriber::set_global_default(log_layers(make_log_writer(log_tx), env_filter))
.expect("Unable to install the global tracing subscriber");
install_log_bridge(bridge_level)?;
crate::boot::mark_tracing_initialized();
crate::boot::replay_boot_diagnostics();
Ok((log_store, broadcast_tx))
}
const fn make_log_writer(tx: UnboundedSender<String>) -> LogWriter {
LogWriter { tx }
}
#[derive(Clone)]
struct LogWriter {
tx: UnboundedSender<String>,
}
impl io::Write for LogWriter {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
let line = String::from_utf8_lossy(buf).to_string();
let _ = self.tx.send(line);
Ok(buf.len())
}
fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}
impl MakeWriter<'_> for LogWriter {
type Writer = Self;
fn make_writer(&self) -> Self::Writer {
self.clone()
}
}
const LOG_BATCH_MAX: usize = 50;
const LOG_FLUSH_INTERVAL: std::time::Duration = std::time::Duration::from_millis(500);
fn spawn_log_writer(
store: Arc<LogStore>,
rx: UnboundedReceiver<String>,
broadcast: tokio::sync::broadcast::Sender<String>,
) {
spawn_log_writer_with_interval(store, rx, broadcast, LOG_FLUSH_INTERVAL);
}
fn spawn_log_writer_with_interval(
store: Arc<LogStore>,
mut rx: UnboundedReceiver<String>,
broadcast: tokio::sync::broadcast::Sender<String>,
flush_interval: std::time::Duration,
) {
tokio::spawn(async move {
let mut batch: Vec<LogEntry> = Vec::new();
let mut flush_timer = tokio::time::interval(flush_interval);
flush_timer.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
flush_timer.tick().await;
loop {
tokio::select! {
maybe_line = rx.recv() => {
let Some(line) = maybe_line else {
if !log_writer_stopped() {
absorb_flush(&store, &mut batch).await;
}
break;
};
let trimmed = line.trim();
if trimmed.is_empty() {
continue;
}
let Some(entry) = parse_tracing_json(trimmed) else {
continue;
};
let _ = broadcast.send(serde_json::to_string(&entry).expect(
"log entry broadcast serialization failed; this should not happen",
));
if log_writer_stopped() {
continue;
}
batch.push(entry);
if batch.len() >= LOG_BATCH_MAX {
absorb_flush(&store, &mut batch).await;
}
}
_ = flush_timer.tick() => {
if !batch.is_empty() && !log_writer_stopped() {
absorb_flush(&store, &mut batch).await;
}
}
}
}
});
}
async fn absorb_flush(store: &LogStore, batch: &mut Vec<LogEntry>) {
let result = AssertUnwindSafe(flush_log_batch(store, batch))
.catch_unwind()
.await;
match result {
Ok(()) => reset_log_writer_panic_state(),
Err(payload) => {
batch.clear();
let message = format!(
"log writer storage panic: {}",
crate::util::panic_message(&*payload)
);
let consecutive = record_log_writer_panic(&message);
if log_writer_stopped() {
crate::boot::timestamped_stderr(&format!(
"log store writer stopped after {consecutive} consecutive storage \
panics: {message}"
));
} else {
tokio::time::sleep(log_writer_panic_backoff(consecutive)).await;
}
}
}
}
async fn flush_log_batch(store: &LogStore, batch: &mut Vec<LogEntry>) {
if batch.is_empty() {
return;
}
let mut last_error: Option<anyhow::Error> = None;
for attempt in 0..LOG_INSERT_MAX_ATTEMPTS {
match store.insert_batch(batch).await {
Ok(()) => {
batch.clear();
return;
}
Err(e) => {
last_error = Some(e);
if attempt + 1 < LOG_INSERT_MAX_ATTEMPTS {
tokio::time::sleep(LOG_INSERT_RETRY_BACKOFF).await;
}
}
}
}
record_log_write_failure(last_error);
batch.clear();
}
const LOG_INSERT_MAX_ATTEMPTS: usize = 3;
const LOG_INSERT_RETRY_BACKOFF: std::time::Duration = std::time::Duration::from_millis(250);
const LOG_WRITE_STDERR_WARN_INTERVAL_MS: u64 = 60_000;
const LOG_WRITER_MAX_CONSECUTIVE_PANICS: u32 = 5;
const LOG_WRITER_PANIC_BACKOFF_MS: u64 = 500;
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct LogWriterPanicState {
consecutive_panics: u32,
pub writer_stopped: bool,
}
impl LogWriterPanicState {
#[must_use]
fn record_panic(&mut self) -> u32 {
self.consecutive_panics += 1;
if self.consecutive_panics >= LOG_WRITER_MAX_CONSECUTIVE_PANICS {
self.writer_stopped = true;
}
self.consecutive_panics
}
fn reset(&mut self) {
self.consecutive_panics = 0;
}
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub(crate) struct LogWriteErrorInfo {
pub count: u64,
pub last_timestamp: Option<String>,
pub last_message: Option<String>,
pub panic_state: LogWriterPanicState,
}
static LOG_WRITE_LAST_ERROR: std::sync::Mutex<LogWriteErrorInfo> =
std::sync::Mutex::new(LogWriteErrorInfo {
count: 0,
last_timestamp: None,
last_message: None,
panic_state: LogWriterPanicState {
consecutive_panics: 0,
writer_stopped: false,
},
});
static LOG_WRITE_LAST_STDERR_WARN_MS: AtomicU64 = AtomicU64::new(0);
#[must_use]
pub(crate) fn log_write_error_info() -> LogWriteErrorInfo {
LOG_WRITE_LAST_ERROR.lock().unwrap_poison().clone()
}
fn record_log_write_failure(error: Option<anyhow::Error>) {
let message = error.map_or_else(
|| "unknown log insert failure".to_string(),
|e| format!("{e:#}"),
);
record_log_write_failure_impl(&message, LogFailureKind::Insert);
}
fn record_log_writer_panic(message: &str) -> u32 {
record_log_write_failure_impl(message, LogFailureKind::WriterPanic)
}
#[derive(Clone, Copy)]
enum LogFailureKind {
Insert,
WriterPanic,
}
impl LogFailureKind {
fn label(self) -> &'static str {
match self {
Self::Insert => "insert failure",
Self::WriterPanic => "writer panic",
}
}
fn records_panic(self) -> bool {
matches!(self, Self::WriterPanic)
}
}
fn record_log_write_failure_impl(message: &str, kind: LogFailureKind) -> u32 {
let (count, consecutive) = {
let mut guard = LOG_WRITE_LAST_ERROR.lock().unwrap_poison();
guard.count += 1;
guard.last_timestamp = Some(db::now());
guard.last_message = Some(message.to_string());
let consecutive = if kind.records_panic() {
guard.panic_state.record_panic()
} else {
guard.panic_state.consecutive_panics
};
(guard.count, consecutive)
};
emit_stderr_warning(count, message, kind.label());
consecutive
}
fn reset_log_writer_panic_state() {
let mut guard = LOG_WRITE_LAST_ERROR.lock().unwrap_poison();
guard.panic_state.reset();
}
fn log_writer_stopped() -> bool {
LOG_WRITE_LAST_ERROR
.lock()
.unwrap_poison()
.panic_state
.writer_stopped
}
fn log_writer_panic_backoff(consecutive: u32) -> std::time::Duration {
let shift = consecutive.saturating_sub(1).min(6);
let ms = LOG_WRITER_PANIC_BACKOFF_MS.saturating_mul(1 << shift);
std::time::Duration::from_millis(ms.min(30_000))
}
fn emit_stderr_warning(count: u64, message: &str, kind: &str) {
let now_ms = crate::util::unix_millis();
let last_warn_ms = LOG_WRITE_LAST_STDERR_WARN_MS.load(Ordering::SeqCst);
if now_ms.saturating_sub(last_warn_ms) >= LOG_WRITE_STDERR_WARN_INTERVAL_MS {
LOG_WRITE_LAST_STDERR_WARN_MS.store(now_ms, Ordering::SeqCst);
crate::boot::timestamped_stderr(&format!("log store {kind} #{count}: {message}"));
}
}
fn get_str_or_empty(val: &serde_json::Value, key: &str) -> String {
json::get_opt_str(val, key).unwrap_or("").to_string()
}
fn parse_tracing_json(line: &str) -> Option<LogEntry> {
let val: serde_json::Value = serde_json::from_str(line).ok()?;
let timestamp = get_str_or_empty(&val, "timestamp");
let level = get_str_or_empty(&val, "level");
let target = get_str_or_empty(&val, "target");
let mut fields = val
.get("fields")
.cloned()
.unwrap_or(serde_json::Value::Null);
let message = get_str_or_empty(&fields, "message");
if let Some(obj) = fields.as_object_mut() {
obj.remove("message");
}
let fields = if fields.as_object().is_some_and(serde_json::Map::is_empty) {
serde_json::Value::Null
} else {
fields
};
let (agent_id, agent_role, workspace) = extract_agent_from_span(&val);
Some(LogEntry {
timestamp,
level,
target,
message,
fields,
agent_id,
agent_role,
workspace,
})
}
fn extract_agent_fields(span: &serde_json::Value) -> (String, String, String) {
(
get_str_or_empty(span, "agent_id"),
get_str_or_empty(span, "role"),
get_str_or_empty(span, "workspace"),
)
}
fn extract_agent_from_span(val: &serde_json::Value) -> (String, String, String) {
let mut agent_id = String::new();
let mut role = String::new();
let mut workspace = String::new();
for candidate in std::iter::once(val.get("fields"))
.chain(std::iter::once(val.get("span")))
.chain(std::iter::once(
val.get("spans")
.and_then(|v| v.as_array())
.and_then(|a| a.last()),
))
.flatten()
{
let (id, r, ws) = extract_agent_fields(candidate);
if agent_id.is_empty() {
agent_id = id;
}
if role.is_empty() {
role = r;
}
if workspace.is_empty() {
workspace = ws;
}
if !agent_id.is_empty() && !role.is_empty() && !workspace.is_empty() {
break;
}
}
(agent_id, role, workspace)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_parse_tracing_json_full() {
let line = r#"{"timestamp":"2025-05-06T12:34:56.000000Z","level":"INFO","target":"mahbot::orchestrator","span":{"name":"agent","agent_id":"00000000-0000-0000-0000-000000000000","role":"lead","workspace":"/some/workspace"},"fields":{"message":"Hello world","key":"value"}}"#;
let entry = parse_tracing_json(line).unwrap();
assert_eq!(entry.timestamp, "2025-05-06T12:34:56.000000Z");
assert_eq!(entry.level, "INFO");
assert_eq!(entry.target, "mahbot::orchestrator");
assert_eq!(entry.message, "Hello world");
assert_eq!(entry.fields, serde_json::json!({"key": "value"}));
assert_eq!(entry.agent_id, "00000000-0000-0000-0000-000000000000");
assert_eq!(entry.agent_role, "lead");
assert_eq!(entry.workspace, "/some/workspace");
}
#[test]
fn test_parse_tracing_json_no_fields() {
let line = r#"{"timestamp":"2025-05-06T12:34:56.000000Z","level":"WARN","target":"test","fields":{"message":"warning"}}"#;
let entry = parse_tracing_json(line).unwrap();
assert_eq!(entry.message, "warning");
assert_eq!(entry.fields, serde_json::Value::Null);
assert_eq!(entry.agent_id, "");
assert_eq!(entry.agent_role, "");
assert_eq!(entry.workspace, "");
}
#[test]
fn test_parse_tracing_json_lenient() {
let entry = parse_tracing_json(r#"{"incomplete": true}"#).unwrap();
assert_eq!(entry.timestamp, "");
assert_eq!(entry.level, "");
assert_eq!(entry.target, "");
assert_eq!(entry.message, "");
assert_eq!(entry.fields, serde_json::Value::Null);
assert_eq!(entry.agent_id, "");
assert_eq!(entry.agent_role, "");
assert_eq!(entry.workspace, "");
}
#[test]
fn test_parse_tracing_json_agent_attribution() {
let cases = [
(
"span only",
r#"{"timestamp":"...","level":"INFO","target":"test","span":{"name":"agent","agent_id":"abc-123","role":"analyst"},"fields":{"message":"researching"}}"#,
"abc-123",
"analyst",
"",
),
(
"spans array",
r#"{"timestamp":"...","level":"INFO","target":"test","spans":[{"name":"parent"},{"name":"agent","agent_id":"xyz-456","role":"coder","workspace":"/ws"}],"fields":{"message":"writing code"}}"#,
"xyz-456",
"coder",
"/ws",
),
(
"event fields without span",
r#"{"timestamp":"...","level":"ERROR","target":"mahbot::agent","fields":{"message":"Agent failed","agent_id":"ticket_123_engineer","role":"engineer","workspace":"my-ws","classification":"transport"}}"#,
"ticket_123_engineer",
"engineer",
"my-ws",
),
(
"event beats inherited span",
r#"{"timestamp":"...","level":"ERROR","target":"mahbot::agent","span":{"name":"agent","agent_id":"caller_42","role":"engineer","workspace":"parent-ws"},"fields":{"message":"Agent failed","agent_id":"analyze_ws_1_2_analyst","role":"analyst","workspace":"my-ws","classification":"runtime"}}"#,
"analyze_ws_1_2_analyst",
"analyst",
"my-ws",
),
(
"workspace-only event keeps span agent",
r#"{"timestamp":"...","level":"WARN","target":"mahbot::tools::edit","span":{"name":"agent","agent_id":"ticket_7_engineer","role":"engineer","workspace":"my-ws"},"fields":{"message":"Search index capacity exhausted","workspace":"my-ws","path":"src/a.rs"}}"#,
"ticket_7_engineer",
"engineer",
"my-ws",
),
(
"agent_id-only event merges span role/workspace",
r#"{"timestamp":"...","level":"WARN","target":"mahbot::agent","span":{"name":"agent","agent_id":"ticket_7_engineer","role":"engineer","workspace":"my-ws"},"fields":{"message":"Failed to persist incoming messages to session DB","agent_id":"ticket_7_engineer","error":"io"}}"#,
"ticket_7_engineer",
"engineer",
"my-ws",
),
];
for (name, line, id, role, ws) in cases {
let entry = parse_tracing_json(line).unwrap();
assert_eq!(entry.agent_id, id, "{name}: agent_id");
assert_eq!(entry.agent_role, role, "{name}: agent_role");
assert_eq!(entry.workspace, ws, "{name}: workspace");
}
}
async fn test_store() -> (Arc<LogStore>, tempfile::TempDir) {
let (store, dir) = crate::open_test_store!(LogStore, "log");
(Arc::new(store), dir)
}
async fn seed_entries(store: &LogStore, entries: &[LogEntry]) {
store.insert_batch(entries).await.unwrap();
}
#[tokio::test]
async fn clear_logs_scopes_match_the_tabs() {
let (store, _dir) = test_store().await;
seed_entries(
&store,
&[
LogEntry {
level: "INFO".into(),
..Default::default()
},
LogEntry {
level: "WARN".into(),
..Default::default()
},
LogEntry {
level: "ERROR".into(),
..Default::default()
},
],
)
.await;
assert_eq!(store.clear_logs(Some("ERROR,WARN")).await.unwrap(), 2);
let (entries, total) = store.query(&LogQuery::default()).await.unwrap();
assert_eq!(total, 1);
assert_eq!(entries[0].level, "INFO");
assert_eq!(store.clear_logs(None).await.unwrap(), 1);
assert_eq!(store.query(&LogQuery::default()).await.unwrap().1, 0);
}
#[tokio::test]
async fn a_reason_already_recorded_is_found_and_a_new_one_is_not() {
let (store, _dir) = test_store().await;
store
.insert_batch(&[
LogEntry {
message: "the update stopped".into(),
fields: serde_json::json!({ "reason": "first" }),
..Default::default()
},
LogEntry {
message: "the update stopped".into(),
fields: serde_json::json!({ "reason": "second" }),
..Default::default()
},
LogEntry {
message: "the update stopped".into(),
..Default::default()
},
LogEntry {
message: "another fact".into(),
fields: serde_json::json!({ "reason": "first" }),
..Default::default()
},
])
.await
.unwrap();
for (message, reason, found) in [
("the update stopped", "first", true),
("the update stopped", "second", true),
("the update stopped", "third", false),
("another fact", "first", true),
("another fact", "second", false),
("never recorded", "first", false),
] {
assert_eq!(
store.has_reason(message, reason).await.unwrap(),
found,
"{message:?} + {reason:?}"
);
}
}
async fn seeded_closed_store() -> tempfile::TempDir {
let tmp = tempfile::TempDir::new().expect("temp dir for test");
{
let store = LogStore::open(tmp.path())
.await
.expect("open healthy store");
seed_entries(
&store,
&[LogEntry {
timestamp: "2025-01-01T00:00:00Z".into(),
level: "INFO".into(),
target: "test".into(),
message: "pre-corruption".into(),
..Default::default()
}],
)
.await;
store
.conn
.checkpoint_ungated()
.await
.expect("checkpoint so pages land in the main DB file");
}
tmp
}
#[tokio::test]
async fn test_spawn_log_writer_writes_to_store() {
let (store, _dir) = test_store().await;
let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
let (broadcast_tx, _) = tokio::sync::broadcast::channel(256);
spawn_log_writer(store.clone(), rx, broadcast_tx);
tx.send(
r#"{"timestamp":"2025-01-01T00:00:00Z","level":"INFO","target":"test","fields":{"message":"hi"}}"#
.to_string(),
)
.unwrap();
tx.send(
r#"{"timestamp":"2025-01-01T00:00:01Z","level":"ERROR","target":"test","fields":{"message":"oh no","err":"boom"}}"#
.to_string(),
)
.unwrap();
drop(tx);
wait_for_total(&store, 2, std::time::Duration::from_secs(2)).await;
let (entries, total) = store.query(&LogQuery::default()).await.unwrap();
assert_eq!(total, 2);
assert_eq!(entries.len(), 2);
assert_eq!(entries[0].message, "oh no");
assert_eq!(entries[1].message, "hi");
}
async fn wait_for_total(store: &LogStore, expected: usize, timeout: std::time::Duration) {
let deadline = std::time::Instant::now() + timeout;
loop {
let (_, total) = store.query(&LogQuery::default()).await.unwrap();
if total == expected {
return;
}
assert!(
std::time::Instant::now() < deadline,
"timed out waiting for total == {expected}, got {total}"
);
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
}
}
#[tokio::test]
async fn test_flush_log_batch_records_failure_on_surface() {
let (store, _dir) = test_store().await;
let baseline = log_write_error_info().count;
let entry = LogEntry {
timestamp: "2025-01-01T00:00:00Z".to_string(),
level: "INFO".to_string(),
target: "test".to_string(),
message: "should not persist".to_string(),
..Default::default()
};
store
.conn
.execute("DROP TABLE logs", ())
.await
.expect("drop logs table for failure test");
let mut batch = vec![entry];
flush_log_batch(&store, &mut batch).await;
assert!(batch.is_empty(), "failed batch must still be cleared");
let info = log_write_error_info();
assert!(
info.count > baseline,
"failure count must advance: baseline {baseline}, now {}",
info.count
);
assert!(
info.last_message.is_some(),
"last-error message must be recorded"
);
}
#[tokio::test]
async fn test_flush_log_batch_persists_without_failure() {
let (store, _dir) = test_store().await;
let baseline = log_write_error_info().count;
let mut batch = vec![LogEntry {
timestamp: "2025-01-01T00:00:01Z".to_string(),
level: "INFO".to_string(),
target: "test".to_string(),
message: "persisted".to_string(),
..Default::default()
}];
flush_log_batch(&store, &mut batch).await;
assert!(batch.is_empty());
assert_eq!(
log_write_error_info().count,
baseline,
"no failure recorded"
);
let (entries, total) = store.query(&LogQuery::default()).await.unwrap();
assert_eq!(total, 1);
assert_eq!(entries[0].message, "persisted");
}
#[tokio::test]
async fn test_log_writer_batches_and_timer_flushes() {
let (store, _dir) = test_store().await;
let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
let (broadcast_tx, _) = tokio::sync::broadcast::channel(256);
spawn_log_writer_with_interval(
store.clone(),
rx,
broadcast_tx,
std::time::Duration::from_mins(1),
);
tx.send(
r#"{"timestamp":"2025-01-01T00:00:02Z","level":"INFO","target":"test","fields":{"message":"timer flush"}}"#
.to_string(),
)
.unwrap();
for i in 0..LOG_BATCH_MAX {
tx.send(
format!(
r#"{{"timestamp":"2025-01-01T00:00:03Z","level":"INFO","target":"test","fields":{{"message":"batch {i}"}}}}"#
),
)
.unwrap();
}
wait_for_total(&store, LOG_BATCH_MAX, std::time::Duration::from_secs(10)).await;
let (store2, _dir2) = test_store().await;
let (tx2, rx2) = tokio::sync::mpsc::unbounded_channel();
let (broadcast_tx2, _) = tokio::sync::broadcast::channel(256);
spawn_log_writer_with_interval(
store2.clone(),
rx2,
broadcast_tx2,
std::time::Duration::from_millis(50),
);
tx2.send(
r#"{"timestamp":"2025-01-01T00:00:04Z","level":"INFO","target":"test","fields":{"message":"timer fired"}}"#
.to_string(),
)
.unwrap();
wait_for_total(&store2, 1, std::time::Duration::from_secs(10)).await;
drop(tx);
drop(tx2);
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
}
#[tokio::test]
async fn test_like_search_substring() {
let (store, _dir) = test_store().await;
let entries = vec![
LogEntry {
timestamp: "2025-01-01T00:00:00Z".into(),
level: "INFO".into(),
target: "module_a".into(),
message: "processing request".into(),
..Default::default()
},
LogEntry {
timestamp: "2025-01-01T00:00:01Z".into(),
level: "ERROR".into(),
target: "module_b".into(),
message: "failed to process".into(),
..Default::default()
},
LogEntry {
timestamp: "2025-01-01T00:00:02Z".into(),
level: "INFO".into(),
target: "module_c".into(),
message: "started".into(),
..Default::default()
},
];
seed_entries(&store, &entries).await;
let (results, total) = store
.query(&LogQuery {
search: Some("proc".into()),
..Default::default()
})
.await
.unwrap();
assert_eq!(total, 2, "substring 'proc' should match both entries");
assert_eq!(results.len(), 2);
let (results, total) = store
.query(&LogQuery {
search: Some("request".into()),
..Default::default()
})
.await
.unwrap();
assert_eq!(total, 1);
assert_eq!(results[0].message, "processing request");
let (_results, total) = store
.query(&LogQuery {
search: Some("module".into()),
..Default::default()
})
.await
.unwrap();
assert_eq!(total, 3, "all targets contain 'module'");
}
#[tokio::test]
async fn test_like_search_combined_filters() {
let (store, _dir) = test_store().await;
let entries = vec![
LogEntry {
timestamp: "2025-01-01T00:00:00Z".into(),
level: "INFO".into(),
target: "mahbot::orchestrator".into(),
message: "processing request".into(),
..Default::default()
},
LogEntry {
timestamp: "2025-01-01T00:00:01Z".into(),
level: "ERROR".into(),
target: "mahbot::tools".into(),
message: "failed to process".into(),
fields: serde_json::json!({"code": 1}),
..Default::default()
},
LogEntry {
timestamp: "2025-01-01T00:00:02Z".into(),
level: "INFO".into(),
target: "mahbot::api".into(),
message: "started".into(),
..Default::default()
},
];
seed_entries(&store, &entries).await;
let (results, total) = store
.query(&LogQuery {
level: Some("ERROR".into()),
search: Some("process".into()),
..Default::default()
})
.await
.unwrap();
assert_eq!(total, 1, "only ERROR log matching 'process'");
assert_eq!(results[0].message, "failed to process");
let (_results, total) = store
.query(&LogQuery {
target: Some("mahbot::tools".into()),
search: Some("process".into()),
..Default::default()
})
.await
.unwrap();
assert_eq!(total, 1, "only tools target entry matching 'process'");
let (_results, total) = store
.query(&LogQuery {
since: Some("2025-01-01T00:00:01Z".into()),
search: Some("process".into()),
..Default::default()
})
.await
.unwrap();
assert_eq!(total, 1, "only entry after timestamp matching 'process'");
}
#[tokio::test]
async fn test_like_search_with_special_chars() {
let (store, _dir) = test_store().await;
let entries = vec![
LogEntry {
timestamp: "2025-01-01T00:00:00Z".into(),
level: "INFO".into(),
target: "module_a".into(),
message: "processing `Hello ${name}` template".into(),
..Default::default()
},
LogEntry {
timestamp: "2025-01-01T00:00:01Z".into(),
level: "ERROR".into(),
target: "module_b".into(),
message: "normal log entry".into(),
..Default::default()
},
];
seed_entries(&store, &entries).await;
let (results, total) = store
.query(&LogQuery {
search: Some("template".into()),
..Default::default()
})
.await
.unwrap();
assert_eq!(total, 1, "LIKE should match partial word in message");
assert!(
results[0].message.contains("template"),
"should match the correct entry"
);
let (_results, total) = store
.query(&LogQuery {
search: None,
..Default::default()
})
.await
.unwrap();
assert_eq!(total, 2, "no search filter should return all entries");
}
#[test]
fn test_log_writer_panic_state_machine() {
let mut state = LogWriterPanicState::default();
assert!(!state.writer_stopped);
for i in 1..=LOG_WRITER_MAX_CONSECUTIVE_PANICS {
let _ = state.record_panic();
assert_eq!(state.consecutive_panics, i);
}
assert!(
state.writer_stopped,
"writer must stop after the consecutive-panic bound"
);
state.reset();
assert_eq!(state.consecutive_panics, 0);
assert!(
state.writer_stopped,
"terminal stopped state is sticky across reset"
);
}
#[tokio::test]
async fn test_log_store_open_refuses_corrupt_store() {
let tmp = seeded_closed_store().await;
let root = tmp.path();
let db_path = db::store_db_path(root, "logs");
let bytes = std::fs::read(&db_path).expect("read db file");
assert!(bytes.len() > 8192, "test needs a multi-page db file");
let mut corrupted = bytes.clone();
corrupted[4096..8192].fill(0);
std::fs::write(&db_path, &corrupted).expect("corrupt db file");
let err = LogStore::open(root)
.await
.expect_err("a corrupt store must be refused, not recreated");
assert!(
err.downcast_ref::<crate::db::StoreRefusal>().is_some(),
"expected a StoreRefusal, got: {err:#}"
);
assert_eq!(
std::fs::read(&db_path).expect("read db file"),
corrupted,
"the refused store's main file must be unchanged"
);
crate::db::test_support::assert_not_quarantined(&root.join("db"), "a refused logs store");
}
#[test]
fn dependency_log_records_reach_the_log_store() {
const MESSAGE: &str = "a dependency record through the standard logging interface";
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
let env_filter = EnvFilter::new(DEFAULT_LOG_FILTER);
let bridge_level = log_bridge_level(&env_filter);
let _guard = tracing::subscriber::set_default(log_layers(make_log_writer(tx), env_filter));
install_log_bridge(bridge_level).expect("install the log→tracing bridge");
tracing_log::log::warn!("{MESSAGE}");
let mut written = String::new();
while let Ok(line) = rx.try_recv() {
written.push_str(&line);
}
let record = written
.lines()
.find(|line| line.contains(MESSAGE))
.unwrap_or_else(|| {
panic!(
"a record written through the standard `log` interface must reach the log \
layer: {written}"
)
});
let entry = parse_tracing_json(record).expect("the log store must parse the written line");
assert_eq!(
entry.message, MESSAGE,
"the store row must carry the record's message",
);
assert_eq!(
entry.target,
module_path!(),
"the record's real target must survive the bridge, not become the literal `log`",
);
}
#[test]
fn the_default_filter_holds_the_pdf_dependency_to_error() {
const TARGET: &str = "pdf_extract";
const DROPPED: &str = "a warning the dependency emits once per page";
const KEPT: &str = "the dependency's own error";
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
let _guard = tracing::subscriber::set_default(log_layers(
make_log_writer(tx),
EnvFilter::new(DEFAULT_LOG_FILTER),
));
tracing::warn!(target: TARGET, "{DROPPED}");
tracing::error!(target: TARGET, "{KEPT}");
let mut written = String::new();
while let Ok(line) = rx.try_recv() {
written.push_str(&line);
}
assert!(
!written.contains(DROPPED),
"the dependency's warnings must not reach the log layer: {written}"
);
assert!(
written.contains(KEPT),
"the dependency's errors must still reach the log layer: {written}"
);
}
#[test]
fn the_default_filter_silences_the_search_library() {
const BURST_TARGET: &str = "fff_search::stable_vec";
const BURST: &str = "StableVec: capacity exhausted — dropping item to prevent reallocation";
const GREP_TARGET: &str = "fff_search::grep";
const GREP: &str = "a grep diagnostic from the file-search library";
const SHARED_TARGET: &str = "fff_search::shared";
const SHARED: &str = "a watcher diagnostic from the file-search library";
const FFF_GREP_TARGET: &str = "fff_grep::matcher";
const FFF_GREP: &str = "a record from the family's grep crate";
const PARSER_TARGET: &str = "fff_query_parser::parse";
const PARSER: &str = "a record from the family's query parser";
const DEBOUNCER_TARGET: &str = "fff_notify_debouncer_full";
const DEBOUNCER: &str = "a record the family logs over the `log` interface";
const NOTICE_TARGET: &str = "mahbot::tools::edit";
const NOTICE: &str =
"Search index capacity exhausted after file write — background rescan needed";
const ENGINE_TARGET: &str = "mahbot::search_engine";
const ENGINE: &str = "Search engine created — background scan started";
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
let _guard = tracing::subscriber::set_default(log_layers(
make_log_writer(tx),
EnvFilter::new(DEFAULT_LOG_FILTER),
));
tracing::error!(target: BURST_TARGET, len = 1025, capacity = 1025, "{BURST}");
tracing::error!(target: GREP_TARGET, "{GREP}");
tracing::warn!(target: SHARED_TARGET, "{SHARED}");
tracing::info!(target: FFF_GREP_TARGET, "{FFF_GREP}");
tracing::error!(target: PARSER_TARGET, "{PARSER}");
tracing::info!(target: DEBOUNCER_TARGET, "{DEBOUNCER}");
tracing::warn!(target: NOTICE_TARGET, "{NOTICE}");
tracing::info!(target: ENGINE_TARGET, "{ENGINE}");
let mut written = String::new();
while let Ok(line) = rx.try_recv() {
written.push_str(&line);
}
for (target, message) in [
(BURST_TARGET, BURST),
(GREP_TARGET, GREP),
(SHARED_TARGET, SHARED),
(FFF_GREP_TARGET, FFF_GREP),
(PARSER_TARGET, PARSER),
(DEBOUNCER_TARGET, DEBOUNCER),
] {
assert!(
!written.contains(message),
"a record from {target} reached the log layer: {written}"
);
}
assert!(
written.contains(NOTICE),
"the product's own notice about the search index must still reach the log layer: \
{written}"
);
assert!(
written.contains(ENGINE),
"the product's own search-engine record must still reach the log layer: {written}"
);
}
}