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}
48
49/// Levels at or above the given threshold, ordered low→high.
50pub fn levels_at_or_above(min: &str) -> Vec<&'static str> {
51    match min {
52        "error" => vec!["error"],
53        "warn" => vec!["warn", "error"],
54        "info" => vec!["info", "warn", "error"],
55        "debug" => vec!["debug", "info", "warn", "error"],
56        "trace" => vec!["trace", "debug", "info", "warn", "error"],
57        _ => vec![],
58    }
59}
60
61impl MessageFilter {
62    #[allow(dead_code)]
63    pub fn contains(pattern: impl Into<String>) -> Self {
64        Self::Contains {
65            pattern: pattern.into(),
66            case_sensitive: false,
67        }
68    }
69
70    #[allow(dead_code)]
71    pub fn contains_case_sensitive(pattern: impl Into<String>) -> Self {
72        Self::Contains {
73            pattern: pattern.into(),
74            case_sensitive: true,
75        }
76    }
77
78    #[allow(dead_code)]
79    pub fn regex(pattern: impl Into<String>) -> Self {
80        Self::Regex {
81            pattern: pattern.into(),
82        }
83    }
84}
85
86/// Options for querying logs.
87#[derive(Debug, Clone, Default)]
88pub struct LogQuery {
89    pub daemon_ids: Vec<String>,
90    pub from: Option<DateTime<Local>>,
91    pub to: Option<DateTime<Local>>,
92    pub limit: Option<usize>,
93    pub order_desc: bool,
94    pub after_id: Option<i64>,
95    /// Filters applied to the message text. Multiple filters are combined with OR.
96    pub message_filters: Vec<MessageFilter>,
97    /// Filters applied to structured fields. Multiple filters are combined with AND.
98    pub field_filters: Vec<FieldFilter>,
99    /// Whether to SELECT the structured columns (level, msg, logger, fields_json).
100    /// When false, those fields are NULL in the result, avoiding unnecessary
101    /// string allocations for callers that only need the raw message.
102    pub include_structured: bool,
103}
104
105/// Escape special LIKE wildcard characters so user-supplied substrings are matched literally.
106pub fn escape_like_pattern(pattern: &str) -> String {
107    pattern
108        .replace('\\', "\\\\")
109        .replace('%', "\\%")
110        .replace('_', "\\_")
111}
112
113/// Parsed retention policy.
114#[derive(Debug, Clone, Copy)]
115pub struct RetentionPolicy {
116    /// Maximum age of entries to keep.
117    pub age: Option<chrono::Duration>,
118    /// Maximum number of entries to keep.
119    pub count: Option<u64>,
120}
121
122impl RetentionPolicy {
123    #[allow(dead_code)]
124    pub fn is_none(&self) -> bool {
125        self.age.is_none() && self.count.is_none()
126    }
127}
128
129/// Hook invoked before log entries are pruned by retention.
130#[derive(Debug, Clone)]
131pub struct ArchiveHook {
132    /// Shell command to run. It should read JSONL from stdin.
133    pub command: String,
134    /// Maximum number of log entries to pass to a single hook invocation.
135    pub batch_size: usize,
136}
137
138impl ArchiveHook {
139    /// Returns true if the hook is configured with a non-empty command.
140    pub fn is_enabled(&self) -> bool {
141        !self.command.trim().is_empty()
142    }
143}
144
145/// Unified interface for log storage and retrieval.
146pub trait LogStore: Send + Sync {
147    /// Append a single log line (unstructured).
148    fn append(&self, daemon_id: &DaemonId, message: &str) -> Result<()>;
149
150    /// Append a single parsed log line with structured fields.
151    ///
152    /// The default implementation discards structured fields and falls back
153    /// to `append`.
154    fn append_structured(&self, daemon_id: &DaemonId, parsed: &ParsedLog) -> Result<()> {
155        self.append(daemon_id, &parsed.message)
156    }
157
158    /// Append multiple parsed log lines in a single transaction.
159    ///
160    /// The default implementation calls `append_structured` for each line.
161    fn append_structured_batch(&self, daemon_id: &DaemonId, entries: &[ParsedLog]) -> Result<()> {
162        for entry in entries {
163            self.append_structured(daemon_id, entry)?;
164        }
165        Ok(())
166    }
167
168    /// Append multiple log lines in a single transaction.
169    ///
170    /// The default implementation falls back to calling `append` for each line.
171    #[allow(dead_code)]
172    fn append_batch(&self, daemon_id: &DaemonId, messages: &[String]) -> Result<()> {
173        for msg in messages {
174            self.append(daemon_id, msg)?;
175        }
176        Ok(())
177    }
178
179    /// Query logs according to the given options.
180    fn query(&self, opts: &LogQuery) -> Result<Vec<LogEntry>>;
181
182    /// Read new log entries for a daemon.
183    /// When `after_id` is Some(id), returns only entries with row id > id.
184    /// When `after_id` is None, returns all entries for the daemon.
185    fn tail(&self, daemon_id: &DaemonId, after_id: Option<i64>) -> Result<Vec<LogEntry>>;
186
187    /// Clear all logs for the given daemon(s).
188    fn clear(&self, daemon_ids: &[DaemonId]) -> Result<()>;
189
190    /// Apply a retention policy (age-based and/or count-based pruning).
191    ///
192    /// By default this applies to all daemons; `excluded_daemon_ids` can be
193    /// used to skip daemons that have their own per-daemon overrides, so the
194    /// global policy does not accidentally prune entries those daemons intend
195    /// to keep.
196    fn apply_retention(
197        &self,
198        policy: &RetentionPolicy,
199        excluded_daemon_ids: &[DaemonId],
200        archive_hook: Option<&ArchiveHook>,
201    ) -> Result<u64> {
202        let _ = (policy, excluded_daemon_ids, archive_hook);
203        Ok(0)
204    }
205
206    /// Apply a retention policy to a specific daemon's logs.
207    fn apply_retention_for_daemon(
208        &self,
209        daemon_id: &DaemonId,
210        policy: &RetentionPolicy,
211        archive_hook: Option<&ArchiveHook>,
212    ) -> Result<u64> {
213        let _ = (daemon_id, policy, archive_hook);
214        Ok(0)
215    }
216
217    /// Return the highest row id for the given daemon, or `None` if no
218    /// entries exist. Used by tailing loops to advance the cursor past
219    /// non-matching rows without re-scanning them on every poll.
220    fn last_id(&self, daemon_id: &DaemonId) -> Result<Option<i64>> {
221        let entries = self.query(&LogQuery {
222            daemon_ids: vec![daemon_id.qualified()],
223            from: None,
224            to: None,
225            limit: Some(1),
226            order_desc: true,
227            after_id: None,
228            message_filters: Vec::new(),
229            field_filters: Vec::new(),
230            include_structured: false,
231        })?;
232        Ok(entries.first().map(|e| e.id))
233    }
234
235    /// List all daemon IDs that have log entries.
236    fn list_daemon_ids(&self) -> Result<Vec<String>>;
237
238    /// Return the generation number for the daemon's last clear operation.
239    /// Each call to `clear` bumps the generation, so SSE streams can detect
240    /// when logs have been wiped and refresh their display.
241    fn last_clear_generation(&self, daemon_id: &DaemonId) -> Result<Option<u64>> {
242        let _ = daemon_id;
243        Ok(None)
244    }
245}
246
247pub mod sqlite;