Skip to main content

pitchfork_cli/log_store/
mod.rs

1use crate::Result;
2use crate::daemon_id::DaemonId;
3use chrono::{DateTime, Local};
4
5/// A single log entry.
6#[derive(Debug, Clone)]
7pub struct LogEntry {
8    pub id: i64,
9    pub daemon_id: String,
10    pub timestamp: DateTime<Local>,
11    pub message: String,
12}
13
14/// A filter applied to the message text of log entries.
15#[derive(Debug, Clone)]
16pub enum MessageFilter {
17    /// Case-insensitive substring match using SQLite LIKE.
18    Contains {
19        pattern: String,
20        case_sensitive: bool,
21    },
22    /// Regular expression match using SQLite REGEXP.
23    Regex { pattern: String },
24}
25
26impl MessageFilter {
27    #[allow(dead_code)]
28    pub fn contains(pattern: impl Into<String>) -> Self {
29        Self::Contains {
30            pattern: pattern.into(),
31            case_sensitive: false,
32        }
33    }
34
35    #[allow(dead_code)]
36    pub fn contains_case_sensitive(pattern: impl Into<String>) -> Self {
37        Self::Contains {
38            pattern: pattern.into(),
39            case_sensitive: true,
40        }
41    }
42
43    #[allow(dead_code)]
44    pub fn regex(pattern: impl Into<String>) -> Self {
45        Self::Regex {
46            pattern: pattern.into(),
47        }
48    }
49}
50
51/// Options for querying logs.
52#[derive(Debug, Clone, Default)]
53pub struct LogQuery {
54    pub daemon_ids: Vec<String>,
55    pub from: Option<DateTime<Local>>,
56    pub to: Option<DateTime<Local>>,
57    pub limit: Option<usize>,
58    pub order_desc: bool,
59    pub after_id: Option<i64>,
60    /// Filters applied to the message text. Multiple filters are combined with OR.
61    pub message_filters: Vec<MessageFilter>,
62}
63
64/// Escape special LIKE wildcard characters so user-supplied substrings are matched literally.
65pub fn escape_like_pattern(pattern: &str) -> String {
66    pattern
67        .replace('\\', "\\\\")
68        .replace('%', "\\%")
69        .replace('_', "\\_")
70}
71
72/// Parsed retention policy.
73#[derive(Debug, Clone, Copy)]
74pub struct RetentionPolicy {
75    /// Maximum age of entries to keep.
76    pub age: Option<chrono::Duration>,
77    /// Maximum number of entries to keep.
78    pub count: Option<u64>,
79}
80
81impl RetentionPolicy {
82    #[allow(dead_code)]
83    pub fn is_none(&self) -> bool {
84        self.age.is_none() && self.count.is_none()
85    }
86}
87
88/// Hook invoked before log entries are pruned by retention.
89#[derive(Debug, Clone)]
90pub struct ArchiveHook {
91    /// Shell command to run. It should read JSONL from stdin.
92    pub command: String,
93    /// Maximum number of log entries to pass to a single hook invocation.
94    pub batch_size: usize,
95}
96
97impl ArchiveHook {
98    /// Returns true if the hook is configured with a non-empty command.
99    pub fn is_enabled(&self) -> bool {
100        !self.command.trim().is_empty()
101    }
102}
103
104/// Unified interface for log storage and retrieval.
105pub trait LogStore: Send + Sync {
106    /// Append a single log line.
107    fn append(&self, daemon_id: &DaemonId, message: &str) -> Result<()>;
108
109    /// Append multiple log lines in a single transaction.
110    ///
111    /// The default implementation falls back to calling `append` for each line.
112    fn append_batch(&self, daemon_id: &DaemonId, messages: &[String]) -> Result<()> {
113        for msg in messages {
114            self.append(daemon_id, msg)?;
115        }
116        Ok(())
117    }
118
119    /// Query logs according to the given options.
120    fn query(&self, opts: &LogQuery) -> Result<Vec<LogEntry>>;
121
122    /// Read new log entries for a daemon.
123    /// When `after_id` is Some(id), returns only entries with row id > id.
124    /// When `after_id` is None, returns all entries for the daemon.
125    fn tail(&self, daemon_id: &DaemonId, after_id: Option<i64>) -> Result<Vec<LogEntry>>;
126
127    /// Clear all logs for the given daemon(s).
128    fn clear(&self, daemon_ids: &[DaemonId]) -> Result<()>;
129
130    /// Apply a retention policy (age-based and/or count-based pruning).
131    ///
132    /// By default this applies to all daemons; `excluded_daemon_ids` can be
133    /// used to skip daemons that have their own per-daemon overrides, so the
134    /// global policy does not accidentally prune entries those daemons intend
135    /// to keep.
136    fn apply_retention(
137        &self,
138        policy: &RetentionPolicy,
139        excluded_daemon_ids: &[DaemonId],
140        archive_hook: Option<&ArchiveHook>,
141    ) -> Result<u64> {
142        let _ = (policy, excluded_daemon_ids, archive_hook);
143        Ok(0)
144    }
145
146    /// Apply a retention policy to a specific daemon's logs.
147    fn apply_retention_for_daemon(
148        &self,
149        daemon_id: &DaemonId,
150        policy: &RetentionPolicy,
151        archive_hook: Option<&ArchiveHook>,
152    ) -> Result<u64> {
153        let _ = (daemon_id, policy, archive_hook);
154        Ok(0)
155    }
156
157    /// Return the highest row id for the given daemon, or `None` if no
158    /// entries exist. Used by tailing loops to advance the cursor past
159    /// non-matching rows without re-scanning them on every poll.
160    fn last_id(&self, daemon_id: &DaemonId) -> Result<Option<i64>> {
161        let entries = self.query(&LogQuery {
162            daemon_ids: vec![daemon_id.qualified()],
163            from: None,
164            to: None,
165            limit: Some(1),
166            order_desc: true,
167            after_id: None,
168            message_filters: Vec::new(),
169        })?;
170        Ok(entries.first().map(|e| e.id))
171    }
172
173    /// List all daemon IDs that have log entries.
174    fn list_daemon_ids(&self) -> Result<Vec<String>>;
175
176    /// Return the generation number for the daemon's last clear operation.
177    /// Each call to `clear` bumps the generation, so SSE streams can detect
178    /// when logs have been wiped and refresh their display.
179    fn last_clear_generation(&self, daemon_id: &DaemonId) -> Result<Option<u64>> {
180        let _ = daemon_id;
181        Ok(None)
182    }
183}
184
185pub mod sqlite;