Skip to main content

pitchfork_cli/
log_parse.rs

1//! Structured log line parsing.
2//!
3//! Parses daemon stdout/stderr lines into structured fields (level, msg,
4//! logger, fields_json) based on the configured `log_format`. Supports JSON
5//! and logfmt formats.
6
7use serde_json::{Map, Value};
8
9/// Result of parsing a single log line.
10///
11/// `message` always holds the original raw line text. The structured fields
12/// are `None` when the line could not be parsed (e.g. plain text or parse
13/// failure).
14#[derive(Debug, Clone, Default)]
15pub struct ParsedLog {
16    /// The original raw line text, always preserved.
17    pub message: String,
18    /// Normalized log level: `error` | `warn` | `info` | `debug` | `trace`.
19    pub level: Option<String>,
20    /// Extracted human-readable message (from `msg`/`message`/`event`/...).
21    pub msg: Option<String>,
22    /// Logger name (from `logger`/`name`/`component`/...).
23    pub logger: Option<String>,
24    /// The full parsed JSON object as a string, for `json_extract` queries.
25    /// `None` for plain-text or logfmt lines (logfmt fields are also stored
26    /// here as a JSON object string).
27    pub fields_json: Option<String>,
28}
29
30impl ParsedLog {
31    /// Create a plain-text ParsedLog with no structured fields.
32    fn plain(message: impl Into<String>) -> Self {
33        Self {
34            message: message.into(),
35            ..Default::default()
36        }
37    }
38}
39
40/// Maximum line length for structured parsing. Lines exceeding this are
41/// stored as plain text without field extraction, protecting against
42/// pathological inputs that would cause excessive memory or CPU use.
43const MAX_PARSE_LINE_LEN: usize = 65536;
44
45/// Parse a log line according to the given format string.
46///
47/// Format values: `"json"`, `"logfmt"`, `"text"` (or any other
48/// value is treated as text). Parse failures fall back to plain text.
49pub fn parse(line: &str, format: &str) -> ParsedLog {
50    if line.len() > MAX_PARSE_LINE_LEN {
51        return ParsedLog::plain(line);
52    }
53    match format {
54        "json" => parse_json(line).unwrap_or_else(|| ParsedLog::plain(line)),
55        "logfmt" => parse_logfmt(line).unwrap_or_else(|| ParsedLog::plain(line)),
56        _ => ParsedLog::plain(line),
57    }
58}
59
60// ---------------------------------------------------------------------------
61// JSON parsing
62// ---------------------------------------------------------------------------
63
64fn parse_json(line: &str) -> Option<ParsedLog> {
65    let value: Value = serde_json::from_str(line.trim()).ok()?;
66    let obj = value.as_object()?;
67
68    let level = extract_level(obj);
69    let msg = extract_msg(obj);
70    let logger = extract_logger(obj);
71
72    // Re-serialize as compact JSON for storage. Using the original object
73    // (not a filtered subset) so all fields are available for json_extract.
74    let fields_json = serde_json::to_string(&value).ok()?;
75
76    Some(ParsedLog {
77        message: line.to_string(),
78        level,
79        msg,
80        logger,
81        fields_json: Some(fields_json),
82    })
83}
84
85// ---------------------------------------------------------------------------
86// logfmt parsing
87// ---------------------------------------------------------------------------
88
89fn parse_logfmt(line: &str) -> Option<ParsedLog> {
90    let pairs = parse_logfmt_pairs(line)?;
91
92    // Build a JSON object from the key-value pairs.
93    let mut obj = Map::new();
94    for (key, value) in &pairs {
95        // Try to parse the value as a JSON type (number, bool, null);
96        // fall back to string.
97        let json_val = if value.is_empty() {
98            Value::Bool(true)
99        } else if let Ok(n) = value.parse::<i64>() {
100            Value::Number(n.into())
101        } else if let Ok(n) = value.parse::<f64>() {
102            serde_json::Number::from_f64(n)
103                .map(Value::Number)
104                .unwrap_or_else(|| Value::String(value.clone()))
105        } else if value.eq_ignore_ascii_case("true") {
106            Value::Bool(true)
107        } else if value.eq_ignore_ascii_case("false") {
108            Value::Bool(false)
109        } else if value.eq_ignore_ascii_case("null") {
110            Value::Null
111        } else {
112            Value::String(value.clone())
113        };
114        obj.insert(key.clone(), json_val);
115    }
116
117    let level = extract_level(&obj);
118    let msg = extract_msg(&obj);
119    let logger = extract_logger(&obj);
120    let value = Value::Object(obj);
121    let fields_json = serde_json::to_string(&value).ok()?;
122
123    Some(ParsedLog {
124        message: line.to_string(),
125        level,
126        msg,
127        logger,
128        fields_json: Some(fields_json),
129    })
130}
131
132/// Parse a logfmt line into (key, value) pairs.
133///
134/// Grammar (simplified from kr/logfmt):
135/// ```text
136/// pair = key '=' value | key '=' | key
137/// key  = ident
138/// value = ident | '"...' '"'
139/// ```
140///
141/// Returns `None` if the line doesn't look like logfmt (no `=` found, or
142/// parsing yields zero pairs).
143fn parse_logfmt_pairs(line: &str) -> Option<Vec<(String, String)>> {
144    let bytes = line.as_bytes();
145    let mut pairs = Vec::new();
146    let mut i = 0;
147
148    while i < bytes.len() {
149        // Skip whitespace and garbage between pairs.
150        while i < bytes.len() && bytes[i].is_ascii_whitespace() {
151            i += 1;
152        }
153        if i >= bytes.len() {
154            break;
155        }
156
157        // Parse key: read until '=', whitespace, or end.
158        let key_start = i;
159        while i < bytes.len() && !bytes[i].is_ascii_whitespace() && bytes[i] != b'=' {
160            i += 1;
161        }
162        let key = &line[key_start..i];
163        if key.is_empty() {
164            // Skip stray garbage.
165            i += 1;
166            continue;
167        }
168
169        // Check for '='.
170        if i < bytes.len() && bytes[i] == b'=' {
171            i += 1; // consume '='
172
173            // Parse value.
174            if i < bytes.len() && bytes[i] == b'"' {
175                // Quoted value: read until closing '"', handling escapes.
176                i += 1; // skip opening quote
177                let val_start = i;
178                while i < bytes.len() && bytes[i] != b'"' {
179                    if bytes[i] == b'\\' && i + 1 < bytes.len() {
180                        i += 2;
181                    } else {
182                        i += 1;
183                    }
184                }
185                let value = unescape_logfmt_value(&line[val_start..i]);
186                if i < bytes.len() {
187                    i += 1; // skip closing quote
188                }
189                pairs.push((key.to_string(), value));
190            } else {
191                // Unquoted value: read until whitespace or end.
192                let val_start = i;
193                while i < bytes.len() && !bytes[i].is_ascii_whitespace() {
194                    i += 1;
195                }
196                pairs.push((key.to_string(), line[val_start..i].to_string()));
197            }
198        } else {
199            // Bare key (no '='): treat as boolean true.
200            pairs.push((key.to_string(), String::new()));
201        }
202    }
203
204    if pairs.is_empty() {
205        return None;
206    }
207    // Require at least one pair with '=' to distinguish from plain text.
208    if !line.contains('=') {
209        return None;
210    }
211    Some(pairs)
212}
213
214fn unescape_logfmt_value(s: &str) -> String {
215    let mut result = String::with_capacity(s.len());
216    let mut chars = s.chars();
217    while let Some(c) = chars.next() {
218        if c == '\\' {
219            if let Some(next) = chars.next() {
220                result.push(next);
221            }
222        } else {
223            result.push(c);
224        }
225    }
226    result
227}
228
229// ---------------------------------------------------------------------------
230// Field extraction with alias tables (inspired by pamburus/hl)
231// ---------------------------------------------------------------------------
232
233/// Try to extract a normalized level from common field names.
234///
235/// Handles both string values (case-insensitive) and integer values
236/// (pino/syslog style). See `normalize_level_value` for the mapping.
237fn extract_level(obj: &Map<String, Value>) -> Option<String> {
238    for key in &["level", "severity", "lvl", "PRIORITY", "@level"] {
239        if let Some(val) = obj.get(*key) {
240            if let Some(level) = normalize_level_value(val) {
241                return Some(level);
242            }
243        }
244    }
245    None
246}
247
248/// Normalize a level value to one of: error, warn, info, debug, trace.
249///
250/// String matching is case-insensitive. Integer values follow pino
251/// (10-60) and syslog/RFC 5424 (0-7) conventions.
252fn normalize_level_value(val: &Value) -> Option<String> {
253    match val {
254        Value::String(s) => normalize_level_str(s),
255        Value::Number(n) => {
256            let n = n.as_i64()?;
257            // pino: 10=trace, 20=debug, 30=info, 40=warn, 50=error, 60=fatal
258            match n {
259                50 | 60 => Some("error".into()),
260                40 => Some("warn".into()),
261                30 => Some("info".into()),
262                20 => Some("debug".into()),
263                10 => Some("trace".into()),
264                // syslog/RFC 5424: 0=emerg..7=debug
265                0..=3 => Some("error".into()),
266                4 | 5 => Some("warn".into()),
267                6 => Some("info".into()),
268                7 => Some("debug".into()),
269                _ => None,
270            }
271        }
272        _ => None,
273    }
274}
275
276pub fn normalize_level_str(s: &str) -> Option<String> {
277    let lower = s.to_ascii_lowercase();
278    match lower.as_str() {
279        "error" | "err" | "fatal" | "critical" | "panic" | "alert" | "emerg" => {
280            Some("error".into())
281        }
282        "warn" | "warning" | "wrn" => Some("warn".into()),
283        "info" | "inf" | "information" | "notice" => Some("info".into()),
284        "debug" | "dbg" => Some("debug".into()),
285        "trace" | "trc" => Some("trace".into()),
286        _ => None,
287    }
288}
289
290fn extract_first_string(obj: &Map<String, Value>, keys: &[&str]) -> Option<String> {
291    for key in keys {
292        if let Some(Value::String(s)) = obj.get(*key) {
293            return Some(s.clone());
294        }
295    }
296    None
297}
298
299/// Try to extract the human-readable message from common field names.
300fn extract_msg(obj: &Map<String, Value>) -> Option<String> {
301    extract_first_string(obj, &["msg", "message", "event", "@message"])
302}
303
304/// Try to extract the logger name from common field names.
305fn extract_logger(obj: &Map<String, Value>) -> Option<String> {
306    extract_first_string(obj, &["logger", "name", "component", "module"])
307}
308
309#[cfg(test)]
310mod tests {
311    use super::*;
312
313    #[test]
314    fn test_json_parse() {
315        let line = r#"{"level":"info","msg":"server started","port":8080}"#;
316        let parsed = parse(line, "json");
317        assert_eq!(parsed.level.as_deref(), Some("info"));
318        assert_eq!(parsed.msg.as_deref(), Some("server started"));
319        assert!(parsed.fields_json.is_some());
320    }
321
322    #[test]
323    fn test_json_level_normalization() {
324        let line = r#"{"level":"FATAL","msg":"crash"}"#;
325        let parsed = parse(line, "json");
326        assert_eq!(parsed.level.as_deref(), Some("error"));
327    }
328
329    #[test]
330    fn test_json_pino_integer_level() {
331        let line = r#"{"level":50,"msg":"error occurred"}"#;
332        let parsed = parse(line, "json");
333        assert_eq!(parsed.level.as_deref(), Some("error"));
334    }
335
336    #[test]
337    fn test_json_syslog_priority() {
338        let line = r#"{"PRIORITY":3,"msg":"system error"}"#;
339        let parsed = parse(line, "json");
340        assert_eq!(parsed.level.as_deref(), Some("error"));
341    }
342
343    #[test]
344    fn test_json_msg_aliases() {
345        // structlog uses "event"
346        let line = r#"{"event":"hello","level":"info"}"#;
347        let parsed = parse(line, "json");
348        assert_eq!(parsed.msg.as_deref(), Some("hello"));
349    }
350
351    #[test]
352    fn test_logfmt_parse() {
353        let line = r#"level=info msg="server started" port=8080"#;
354        let parsed = parse(line, "logfmt");
355        assert_eq!(parsed.level.as_deref(), Some("info"));
356        assert_eq!(parsed.msg.as_deref(), Some("server started"));
357        assert!(parsed.fields_json.is_some());
358    }
359
360    #[test]
361    fn test_logfmt_bare_key() {
362        let line = r#"level=debug ready msg="ok""#;
363        let parsed = parse(line, "logfmt");
364        assert_eq!(parsed.level.as_deref(), Some("debug"));
365        assert_eq!(parsed.msg.as_deref(), Some("ok"));
366        // "ready" is a bare key → true
367        let fields: Value = serde_json::from_str(parsed.fields_json.as_deref().unwrap()).unwrap();
368        assert_eq!(fields["ready"], Value::Bool(true));
369    }
370
371    #[test]
372    fn test_logfmt_quoted_value_with_spaces() {
373        let line = r#"level=error msg="connection refused: timeout""#;
374        let parsed = parse(line, "logfmt");
375        assert_eq!(parsed.msg.as_deref(), Some("connection refused: timeout"));
376    }
377
378    #[test]
379    fn test_text_format() {
380        let line = r#"{"level":"info"}"#;
381        let parsed = parse(line, "text");
382        assert!(parsed.level.is_none());
383        assert!(parsed.fields_json.is_none());
384        assert_eq!(parsed.message, line);
385    }
386
387    #[test]
388    fn test_json_parse_failure_falls_back() {
389        let line = "{not valid json";
390        let parsed = parse(line, "json");
391        assert!(parsed.level.is_none());
392        assert!(parsed.fields_json.is_none());
393        assert_eq!(parsed.message, line);
394    }
395
396    #[test]
397    fn test_logfmt_logger_extraction() {
398        let line = r#"level=info msg="hi" logger=myapp"#;
399        let parsed = parse(line, "logfmt");
400        assert_eq!(parsed.logger.as_deref(), Some("myapp"));
401    }
402
403    #[test]
404    fn test_json_nested_not_extracted_as_msg() {
405        // Only top-level fields are extracted; nested objects stay in fields_json.
406        let line = r#"{"level":"info","fields":{"message":"nested"}}"#;
407        let parsed = parse(line, "json");
408        assert_eq!(parsed.level.as_deref(), Some("info"));
409        assert_eq!(parsed.msg, None); // top-level has no msg/message/event
410    }
411}