Skip to main content

pitchfork_cli/log_store/
mod.rs

1use crate::Result;
2use crate::daemon_id::DaemonId;
3use crate::log_parse::ParsedLog;
4use chrono::{DateTime, Local};
5
6/// A single log entry.
7#[derive(Debug, Clone)]
8pub struct LogEntry {
9    pub id: i64,
10    pub daemon_id: String,
11    pub timestamp: DateTime<Local>,
12    pub message: String,
13    /// Normalized log level (`error`/`warn`/`info`/`debug`/`trace`), or `None`
14    /// for unstructured log lines.
15    pub level: Option<String>,
16    /// Extracted human-readable message (from `msg`/`message`/`event`/...).
17    pub msg: Option<String>,
18    /// Logger name (from `logger`/`name`/`component`/...).
19    pub logger: Option<String>,
20    /// The full parsed JSON object as a string, for `json_extract` queries.
21    /// `None` for plain-text lines that were not parsed.
22    pub fields_json: Option<String>,
23}
24
25/// A filter applied to the message text of log entries.
26#[derive(Debug, Clone)]
27pub enum MessageFilter {
28    /// Case-insensitive substring match using SQLite LIKE.
29    Contains {
30        pattern: String,
31        case_sensitive: bool,
32    },
33    /// Regular expression match using SQLite REGEXP.
34    Regex { pattern: String },
35}
36
37/// A filter applied to structured log fields.
38#[derive(Debug, Clone)]
39pub enum FieldFilter {
40    /// Match entries at or above a severity threshold.
41    ///
42    /// `--level warn` keeps `warn` and `error` (higher severity).
43    /// Severity order (low→high): trace < debug < info < warn < error.
44    LevelMin(String),
45    /// Match entries where `json_extract(fields_json, '$.key') = value`.
46    FieldEq { key: String, value: String },
47    /// Match entries where the `logger` column contains the substring.
48    LoggerContains(String),
49}
50
51/// Levels at or above the given threshold, ordered low→high.
52pub fn levels_at_or_above(min: &str) -> Vec<&'static str> {
53    match min {
54        "error" => vec!["error"],
55        "warn" => vec!["warn", "error"],
56        "info" => vec!["info", "warn", "error"],
57        "debug" => vec!["debug", "info", "warn", "error"],
58        "trace" => vec!["trace", "debug", "info", "warn", "error"],
59        _ => vec![],
60    }
61}
62
63impl MessageFilter {
64    #[allow(dead_code)]
65    pub fn contains(pattern: impl Into<String>) -> Self {
66        Self::Contains {
67            pattern: pattern.into(),
68            case_sensitive: false,
69        }
70    }
71
72    #[allow(dead_code)]
73    pub fn contains_case_sensitive(pattern: impl Into<String>) -> Self {
74        Self::Contains {
75            pattern: pattern.into(),
76            case_sensitive: true,
77        }
78    }
79
80    #[allow(dead_code)]
81    pub fn regex(pattern: impl Into<String>) -> Self {
82        Self::Regex {
83            pattern: pattern.into(),
84        }
85    }
86}
87
88/// Options for querying logs.
89#[derive(Debug, Clone, Default)]
90pub struct LogQuery {
91    pub daemon_ids: Vec<String>,
92    pub from: Option<DateTime<Local>>,
93    pub to: Option<DateTime<Local>>,
94    pub limit: Option<usize>,
95    pub order_desc: bool,
96    pub after_id: Option<i64>,
97    /// When set, only return entries with id < before_id (used for backward
98    /// pagination — loading older history on scroll-up).
99    pub before_id: Option<i64>,
100    /// Filters applied to the message text. Multiple filters are combined with OR.
101    pub message_filters: Vec<MessageFilter>,
102    /// Filters applied to structured fields. Multiple filters are combined with AND.
103    pub field_filters: Vec<FieldFilter>,
104    /// Whether to SELECT the structured columns (level, msg, logger, fields_json).
105    /// When false, those fields are NULL in the result, avoiding unnecessary
106    /// string allocations for callers that only need the raw message.
107    pub include_structured: bool,
108}
109
110/// Escape special LIKE wildcard characters so user-supplied substrings are matched literally.
111pub fn escape_like_pattern(pattern: &str) -> String {
112    pattern
113        .replace('\\', "\\\\")
114        .replace('%', "\\%")
115        .replace('_', "\\_")
116}
117
118/// Parsed retention policy.
119#[derive(Debug, Clone, Copy)]
120pub struct RetentionPolicy {
121    /// Maximum age of entries to keep.
122    pub age: Option<chrono::Duration>,
123    /// Maximum number of entries to keep.
124    pub count: Option<u64>,
125}
126
127impl RetentionPolicy {
128    #[allow(dead_code)]
129    pub fn is_none(&self) -> bool {
130        self.age.is_none() && self.count.is_none()
131    }
132}
133
134/// Hook invoked before log entries are pruned by retention.
135#[derive(Debug, Clone)]
136pub struct ArchiveHook {
137    /// Shell command to run. It should read JSONL from stdin.
138    pub command: String,
139    /// Maximum number of log entries to pass to a single hook invocation.
140    pub batch_size: usize,
141}
142
143impl ArchiveHook {
144    /// Returns true if the hook is configured with a non-empty command.
145    pub fn is_enabled(&self) -> bool {
146        !self.command.trim().is_empty()
147    }
148}
149
150/// Unified interface for log storage and retrieval.
151pub trait LogStore: Send + Sync {
152    /// Append a single log line (unstructured).
153    fn append(&self, daemon_id: &DaemonId, message: &str) -> Result<()>;
154
155    /// Append a single parsed log line with structured fields.
156    ///
157    /// The default implementation discards structured fields and falls back
158    /// to `append`.
159    fn append_structured(&self, daemon_id: &DaemonId, parsed: &ParsedLog) -> Result<()> {
160        self.append(daemon_id, &parsed.message)
161    }
162
163    /// Append multiple parsed log lines in a single transaction.
164    ///
165    /// The default implementation calls `append_structured` for each line.
166    fn append_structured_batch(&self, daemon_id: &DaemonId, entries: &[ParsedLog]) -> Result<()> {
167        for entry in entries {
168            self.append_structured(daemon_id, entry)?;
169        }
170        Ok(())
171    }
172
173    /// Append multiple log lines in a single transaction.
174    ///
175    /// The default implementation falls back to calling `append` for each line.
176    #[allow(dead_code)]
177    fn append_batch(&self, daemon_id: &DaemonId, messages: &[String]) -> Result<()> {
178        for msg in messages {
179            self.append(daemon_id, msg)?;
180        }
181        Ok(())
182    }
183
184    /// Query logs according to the given options.
185    fn query(&self, opts: &LogQuery) -> Result<Vec<LogEntry>>;
186
187    /// Read new log entries for a daemon.
188    /// When `after_id` is Some(id), returns only entries with row id > id.
189    /// When `after_id` is None, returns all entries for the daemon.
190    fn tail(&self, daemon_id: &DaemonId, after_id: Option<i64>) -> Result<Vec<LogEntry>>;
191
192    /// Clear all logs for the given daemon(s).
193    fn clear(&self, daemon_ids: &[DaemonId]) -> Result<()>;
194
195    /// Apply a retention policy (age-based and/or count-based pruning).
196    ///
197    /// By default this applies to all daemons; `excluded_daemon_ids` can be
198    /// used to skip daemons that have their own per-daemon overrides, so the
199    /// global policy does not accidentally prune entries those daemons intend
200    /// to keep.
201    fn apply_retention(
202        &self,
203        policy: &RetentionPolicy,
204        excluded_daemon_ids: &[DaemonId],
205        archive_hook: Option<&ArchiveHook>,
206    ) -> Result<u64> {
207        let _ = (policy, excluded_daemon_ids, archive_hook);
208        Ok(0)
209    }
210
211    /// Apply a retention policy to a specific daemon's logs.
212    fn apply_retention_for_daemon(
213        &self,
214        daemon_id: &DaemonId,
215        policy: &RetentionPolicy,
216        archive_hook: Option<&ArchiveHook>,
217    ) -> Result<u64> {
218        let _ = (daemon_id, policy, archive_hook);
219        Ok(0)
220    }
221
222    /// Return the highest row id for the given daemon, or `None` if no
223    /// entries exist. Used by tailing loops to advance the cursor past
224    /// non-matching rows without re-scanning them on every poll.
225    fn last_id(&self, daemon_id: &DaemonId) -> Result<Option<i64>> {
226        let entries = self.query(&LogQuery {
227            daemon_ids: vec![daemon_id.qualified()],
228            from: None,
229            to: None,
230            limit: Some(1),
231            order_desc: true,
232            after_id: None,
233            before_id: None,
234            message_filters: Vec::new(),
235            field_filters: Vec::new(),
236            include_structured: false,
237        })?;
238        Ok(entries.first().map(|e| e.id))
239    }
240
241    /// List all daemon IDs that have log entries.
242    fn list_daemon_ids(&self) -> Result<Vec<String>>;
243
244    /// Return the generation number for the daemon's last clear operation.
245    /// Each call to `clear` bumps the generation, so SSE streams can detect
246    /// when logs have been wiped and refresh their display.
247    fn last_clear_generation(&self, daemon_id: &DaemonId) -> Result<Option<u64>> {
248        let _ = daemon_id;
249        Ok(None)
250    }
251
252    /// Query logs and the current clear generation atomically.
253    ///
254    /// This acquires a single connection lock and wraps both reads in one
255    /// transaction so that a concurrent `clear` cannot pair stale history
256    /// with a new generation (which would evade clear detection).
257    fn query_with_generation(
258        &self,
259        opts: &LogQuery,
260        daemon_id: &DaemonId,
261    ) -> Result<(Vec<LogEntry>, Option<u64>)> {
262        let entries = self.query(opts)?;
263        let generation = self.last_clear_generation(daemon_id)?;
264        Ok((entries, generation))
265    }
266}
267
268pub mod sqlite;