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, value) in &pairs {
95 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
132fn 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 while i < bytes.len() && bytes[i].is_ascii_whitespace() {
151 i += 1;
152 }
153 if i >= bytes.len() {
154 break;
155 }
156
157 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 i += 1;
166 continue;
167 }
168
169 if i < bytes.len() && bytes[i] == b'=' {
171 i += 1; if i < bytes.len() && bytes[i] == b'"' {
175 i += 1; 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; }
189 pairs.push((key.to_string(), value));
190 } else {
191 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 pairs.push((key.to_string(), String::new()));
201 }
202 }
203
204 if pairs.is_empty() {
205 return None;
206 }
207 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
229fn 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 && let Some(level) = normalize_level_value(val)
241 {
242 return Some(level);
243 }
244 }
245 None
246}
247
248fn 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 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 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
299fn extract_msg(obj: &Map<String, Value>) -> Option<String> {
301 extract_first_string(obj, &["msg", "message", "event", "@message"])
302}
303
304fn 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 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 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 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); }
411}