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;