1use serde_json::{Map, Value};
8
9#[derive(Debug, Clone, Default)]
15pub struct ParsedLog {
16 pub message: String,
18 pub level: Option<String>,
20 pub msg: Option<String>,
22 pub logger: Option<String>,
24 pub fields_json: Option<String>,
28}
29
30impl ParsedLog {
31 fn plain(message: impl Into<String>) -> Self {
33 Self {
34 message: message.into(),
35 ..Default::default()
36 }
37 }
38}
39
40const MAX_PARSE_LINE_LEN: usize = 65536;
44
45pub 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
60fn 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 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
85fn parse_logfmt(line: &str) -> Option<ParsedLog> {
90 let pairs = parse_logfmt_pairs(line)?;
91
92 let mut obj = Map::new();
94 for (key, maybe_value) in &pairs {
95 let json_val = match maybe_value {
98 None => Value::Bool(true),
99 Some(value) if value.is_empty() => Value::String(value.clone()),
100 Some(value) => {
101 if let Ok(n) = value.parse::<i64>() {
102 Value::Number(n.into())
103 } else if let Ok(n) = value.parse::<f64>() {
104 serde_json::Number::from_f64(n)
105 .map(Value::Number)
106 .unwrap_or_else(|| Value::String(value.clone()))
107 } else if value.eq_ignore_ascii_case("true") {
108 Value::Bool(true)
109 } else if value.eq_ignore_ascii_case("false") {
110 Value::Bool(false)
111 } else if value.eq_ignore_ascii_case("null") {
112 Value::Null
113 } else {
114 Value::String(value.clone())
115 }
116 }
117 };
118 obj.insert(key.clone(), json_val);
119 }
120
121 let level = extract_level(&obj);
122 let msg = extract_msg(&obj);
123 let logger = extract_logger(&obj);
124 let value = Value::Object(obj);
125 let fields_json = serde_json::to_string(&value).ok()?;
126
127 Some(ParsedLog {
128 message: line.to_string(),
129 level,
130 msg,
131 logger,
132 fields_json: Some(fields_json),
133 })
134}
135
136fn parse_logfmt_pairs(line: &str) -> Option<Vec<(String, Option<String>)>> {
150 let bytes = line.as_bytes();
151 let mut pairs = Vec::new();
152 let mut bare_key_count = 0;
153 let mut i = 0;
154
155 while i < bytes.len() {
156 while i < bytes.len() && bytes[i].is_ascii_whitespace() {
158 i += 1;
159 }
160 if i >= bytes.len() {
161 break;
162 }
163
164 let key_start = i;
166 while i < bytes.len() && !bytes[i].is_ascii_whitespace() && bytes[i] != b'=' {
167 i += 1;
168 }
169 let key = &line[key_start..i];
170 if key.is_empty() {
171 i += 1;
173 continue;
174 }
175 if !key
180 .bytes()
181 .all(|b| b.is_ascii_alphanumeric() || b == b'_' || b == b'.' || b == b'-' || b == b'@')
182 {
183 while i < bytes.len() && bytes[i] == b'=' {
186 i += 1;
187 if i < bytes.len() && bytes[i] == b'"' {
188 i += 1;
189 while i < bytes.len() && bytes[i] != b'"' {
190 if bytes[i] == b'\\' && i + 1 < bytes.len() {
191 i += 2;
192 } else {
193 i += 1;
194 }
195 }
196 if i < bytes.len() {
197 i += 1;
198 }
199 } else {
200 while i < bytes.len() && !bytes[i].is_ascii_whitespace() {
201 i += 1;
202 }
203 }
204 }
205 continue;
206 }
207
208 if i < bytes.len() && bytes[i] == b'=' {
210 i += 1; if i < bytes.len() && bytes[i] == b'"' {
214 i += 1; let val_start = i;
217 while i < bytes.len() && bytes[i] != b'"' {
218 if bytes[i] == b'\\' && i + 1 < bytes.len() {
219 i += 2;
220 } else {
221 i += 1;
222 }
223 }
224 let value = unescape_logfmt_value(&line[val_start..i]);
225 if i < bytes.len() {
226 i += 1; }
228 pairs.push((key.to_string(), Some(value)));
229 } else {
230 let val_start = i;
232 while i < bytes.len() && !bytes[i].is_ascii_whitespace() {
233 i += 1;
234 }
235 pairs.push((key.to_string(), Some(line[val_start..i].to_string())));
236 }
237 } else {
238 pairs.push((key.to_string(), None));
240 bare_key_count += 1;
241 }
242 }
243
244 if pairs.is_empty() {
245 return None;
246 }
247 if !line.contains('=') {
249 return None;
250 }
251 if bare_key_count * 2 > pairs.len() {
255 return None;
256 }
257 Some(pairs)
258}
259
260fn unescape_logfmt_value(s: &str) -> String {
261 let mut result = String::with_capacity(s.len());
262 let mut chars = s.chars();
263 while let Some(c) = chars.next() {
264 if c == '\\' {
265 if let Some(next) = chars.next() {
266 result.push(next);
267 }
268 } else {
269 result.push(c);
270 }
271 }
272 result
273}
274
275fn extract_level(obj: &Map<String, Value>) -> Option<String> {
284 for key in &["level", "severity", "lvl", "PRIORITY", "@level"] {
285 if let Some(val) = obj.get(*key)
286 && let Some(level) = normalize_level_value(val)
287 {
288 return Some(level);
289 }
290 }
291 None
292}
293
294fn normalize_level_value(val: &Value) -> Option<String> {
299 match val {
300 Value::String(s) => normalize_level_str(s),
301 Value::Number(n) => {
302 let n = n.as_i64()?;
303 match n {
305 50 | 60 => Some("error".into()),
306 40 => Some("warn".into()),
307 30 => Some("info".into()),
308 20 => Some("debug".into()),
309 10 => Some("trace".into()),
310 0..=3 => Some("error".into()),
312 4 | 5 => Some("warn".into()),
313 6 => Some("info".into()),
314 7 => Some("debug".into()),
315 _ => None,
316 }
317 }
318 _ => None,
319 }
320}
321
322pub fn normalize_level_str(s: &str) -> Option<String> {
323 let lower = s.to_ascii_lowercase();
324 match lower.as_str() {
325 "error" | "err" | "fatal" | "critical" | "panic" | "alert" | "emerg" => {
326 Some("error".into())
327 }
328 "warn" | "warning" | "wrn" => Some("warn".into()),
329 "info" | "inf" | "information" | "notice" => Some("info".into()),
330 "debug" | "dbg" => Some("debug".into()),
331 "trace" | "trc" => Some("trace".into()),
332 _ => None,
333 }
334}
335
336fn extract_first_string(obj: &Map<String, Value>, keys: &[&str]) -> Option<String> {
337 for key in keys {
338 if let Some(Value::String(s)) = obj.get(*key) {
339 return Some(s.clone());
340 }
341 }
342 None
343}
344
345fn extract_msg(obj: &Map<String, Value>) -> Option<String> {
347 extract_first_string(obj, &["msg", "message", "event", "@message"])
348}
349
350fn extract_logger(obj: &Map<String, Value>) -> Option<String> {
352 extract_first_string(obj, &["logger", "name", "component", "module"])
353}
354
355#[cfg(test)]
356mod tests {
357 use super::*;
358
359 #[test]
360 fn test_json_parse() {
361 let line = r#"{"level":"info","msg":"server started","port":8080}"#;
362 let parsed = parse(line, "json");
363 assert_eq!(parsed.level.as_deref(), Some("info"));
364 assert_eq!(parsed.msg.as_deref(), Some("server started"));
365 assert!(parsed.fields_json.is_some());
366 }
367
368 #[test]
369 fn test_json_level_normalization() {
370 let line = r#"{"level":"FATAL","msg":"crash"}"#;
371 let parsed = parse(line, "json");
372 assert_eq!(parsed.level.as_deref(), Some("error"));
373 }
374
375 #[test]
376 fn test_json_pino_integer_level() {
377 let line = r#"{"level":50,"msg":"error occurred"}"#;
378 let parsed = parse(line, "json");
379 assert_eq!(parsed.level.as_deref(), Some("error"));
380 }
381
382 #[test]
383 fn test_json_syslog_priority() {
384 let line = r#"{"PRIORITY":3,"msg":"system error"}"#;
385 let parsed = parse(line, "json");
386 assert_eq!(parsed.level.as_deref(), Some("error"));
387 }
388
389 #[test]
390 fn test_json_msg_aliases() {
391 let line = r#"{"event":"hello","level":"info"}"#;
393 let parsed = parse(line, "json");
394 assert_eq!(parsed.msg.as_deref(), Some("hello"));
395 }
396
397 #[test]
398 fn test_logfmt_parse() {
399 let line = r#"level=info msg="server started" port=8080"#;
400 let parsed = parse(line, "logfmt");
401 assert_eq!(parsed.level.as_deref(), Some("info"));
402 assert_eq!(parsed.msg.as_deref(), Some("server started"));
403 assert!(parsed.fields_json.is_some());
404 }
405
406 #[test]
407 fn test_logfmt_bare_key() {
408 let line = r#"level=debug ready msg="ok""#;
409 let parsed = parse(line, "logfmt");
410 assert_eq!(parsed.level.as_deref(), Some("debug"));
411 assert_eq!(parsed.msg.as_deref(), Some("ok"));
412 let fields: Value = serde_json::from_str(parsed.fields_json.as_deref().unwrap()).unwrap();
414 assert_eq!(fields["ready"], Value::Bool(true));
415 }
416
417 #[test]
418 fn test_logfmt_quoted_value_with_spaces() {
419 let line = r#"level=error msg="connection refused: timeout""#;
420 let parsed = parse(line, "logfmt");
421 assert_eq!(parsed.msg.as_deref(), Some("connection refused: timeout"));
422 }
423
424 #[test]
425 fn test_text_format() {
426 let line = r#"{"level":"info"}"#;
427 let parsed = parse(line, "text");
428 assert!(parsed.level.is_none());
429 assert!(parsed.fields_json.is_none());
430 assert_eq!(parsed.message, line);
431 }
432
433 #[test]
434 fn test_json_parse_failure_falls_back() {
435 let line = "{not valid json";
436 let parsed = parse(line, "json");
437 assert!(parsed.level.is_none());
438 assert!(parsed.fields_json.is_none());
439 assert_eq!(parsed.message, line);
440 }
441
442 #[test]
443 fn test_logfmt_logger_extraction() {
444 let line = r#"level=info msg="hi" logger=myapp"#;
445 let parsed = parse(line, "logfmt");
446 assert_eq!(parsed.logger.as_deref(), Some("myapp"));
447 }
448
449 #[test]
450 fn test_logfmt_multiple_bare_keys_still_parsed() {
451 let line = r#"level=debug ready enabled msg="ok""#;
453 let parsed = parse(line, "logfmt");
454 assert_eq!(parsed.level.as_deref(), Some("debug"));
455 assert_eq!(parsed.msg.as_deref(), Some("ok"));
456 let fields: Value = serde_json::from_str(parsed.fields_json.as_deref().unwrap()).unwrap();
457 assert_eq!(fields["ready"], Value::Bool(true));
458 assert_eq!(fields["enabled"], Value::Bool(true));
459 }
460
461 #[test]
462 fn test_logfmt_rejects_go_standard_log() {
463 let line = "2026/07/23 22:02:39 INFO acquired instance lock path=/foo/bar";
466 let parsed = parse(line, "logfmt");
467 assert!(parsed.level.is_none());
469 assert!(parsed.fields_json.is_none());
470 assert_eq!(parsed.message, line);
471 }
472
473 #[test]
474 fn test_logfmt_explicit_empty_value() {
475 let line = r#"level=debug msg="" extra="""#;
478 let parsed = parse(line, "logfmt");
479 assert_eq!(parsed.level.as_deref(), Some("debug"));
480 assert_eq!(parsed.msg.as_deref(), Some(""));
481 let fields: Value = serde_json::from_str(parsed.fields_json.as_deref().unwrap()).unwrap();
482 assert_eq!(fields["msg"], Value::String("".into()));
483 assert_eq!(fields["extra"], Value::String("".into()));
484 }
485
486 #[test]
487 fn test_logfmt_hyphenated_key() {
488 let line = r#"level=info msg="ok" request-id=abc123"#;
490 let parsed = parse(line, "logfmt");
491 assert_eq!(parsed.level.as_deref(), Some("info"));
492 let fields: Value = serde_json::from_str(parsed.fields_json.as_deref().unwrap()).unwrap();
493 assert_eq!(fields["request-id"], Value::String("abc123".into()));
494 }
495
496 #[test]
497 fn test_logfmt_at_prefixed_key() {
498 let line = r#"@level=info @message="server started" port=8080"#;
500 let parsed = parse(line, "logfmt");
501 assert_eq!(parsed.level.as_deref(), Some("info"));
502 assert_eq!(parsed.msg.as_deref(), Some("server started"));
503 let fields: Value = serde_json::from_str(parsed.fields_json.as_deref().unwrap()).unwrap();
504 assert_eq!(fields["port"], Value::Number(8080.into()));
505 }
506
507 #[test]
508 fn test_json_nested_not_extracted_as_msg() {
509 let line = r#"{"level":"info","fields":{"message":"nested"}}"#;
511 let parsed = parse(line, "json");
512 assert_eq!(parsed.level.as_deref(), Some("info"));
513 assert_eq!(parsed.msg, None); }
515}