use serde_json::{Map, Value};
#[derive(Debug, Clone, Default)]
pub struct ParsedLog {
pub message: String,
pub level: Option<String>,
pub msg: Option<String>,
pub logger: Option<String>,
pub fields_json: Option<String>,
}
impl ParsedLog {
fn plain(message: impl Into<String>) -> Self {
Self {
message: message.into(),
..Default::default()
}
}
}
const MAX_PARSE_LINE_LEN: usize = 65536;
pub fn parse(line: &str, format: &str) -> ParsedLog {
if line.len() > MAX_PARSE_LINE_LEN {
return ParsedLog::plain(line);
}
match format {
"json" => parse_json(line).unwrap_or_else(|| ParsedLog::plain(line)),
"logfmt" => parse_logfmt(line).unwrap_or_else(|| ParsedLog::plain(line)),
_ => ParsedLog::plain(line),
}
}
fn parse_json(line: &str) -> Option<ParsedLog> {
let value: Value = serde_json::from_str(line.trim()).ok()?;
let obj = value.as_object()?;
let level = extract_level(obj);
let msg = extract_msg(obj);
let logger = extract_logger(obj);
let fields_json = serde_json::to_string(&value).ok()?;
Some(ParsedLog {
message: line.to_string(),
level,
msg,
logger,
fields_json: Some(fields_json),
})
}
fn parse_logfmt(line: &str) -> Option<ParsedLog> {
let pairs = parse_logfmt_pairs(line)?;
let mut obj = Map::new();
for (key, value) in &pairs {
let json_val = if value.is_empty() {
Value::Bool(true)
} else if let Ok(n) = value.parse::<i64>() {
Value::Number(n.into())
} else if let Ok(n) = value.parse::<f64>() {
serde_json::Number::from_f64(n)
.map(Value::Number)
.unwrap_or_else(|| Value::String(value.clone()))
} else if value.eq_ignore_ascii_case("true") {
Value::Bool(true)
} else if value.eq_ignore_ascii_case("false") {
Value::Bool(false)
} else if value.eq_ignore_ascii_case("null") {
Value::Null
} else {
Value::String(value.clone())
};
obj.insert(key.clone(), json_val);
}
let level = extract_level(&obj);
let msg = extract_msg(&obj);
let logger = extract_logger(&obj);
let value = Value::Object(obj);
let fields_json = serde_json::to_string(&value).ok()?;
Some(ParsedLog {
message: line.to_string(),
level,
msg,
logger,
fields_json: Some(fields_json),
})
}
fn parse_logfmt_pairs(line: &str) -> Option<Vec<(String, String)>> {
let bytes = line.as_bytes();
let mut pairs = Vec::new();
let mut i = 0;
while i < bytes.len() {
while i < bytes.len() && bytes[i].is_ascii_whitespace() {
i += 1;
}
if i >= bytes.len() {
break;
}
let key_start = i;
while i < bytes.len() && !bytes[i].is_ascii_whitespace() && bytes[i] != b'=' {
i += 1;
}
let key = &line[key_start..i];
if key.is_empty() {
i += 1;
continue;
}
if i < bytes.len() && bytes[i] == b'=' {
i += 1;
if i < bytes.len() && bytes[i] == b'"' {
i += 1; let val_start = i;
while i < bytes.len() && bytes[i] != b'"' {
if bytes[i] == b'\\' && i + 1 < bytes.len() {
i += 2;
} else {
i += 1;
}
}
let value = unescape_logfmt_value(&line[val_start..i]);
if i < bytes.len() {
i += 1; }
pairs.push((key.to_string(), value));
} else {
let val_start = i;
while i < bytes.len() && !bytes[i].is_ascii_whitespace() {
i += 1;
}
pairs.push((key.to_string(), line[val_start..i].to_string()));
}
} else {
pairs.push((key.to_string(), String::new()));
}
}
if pairs.is_empty() {
return None;
}
if !line.contains('=') {
return None;
}
Some(pairs)
}
fn unescape_logfmt_value(s: &str) -> String {
let mut result = String::with_capacity(s.len());
let mut chars = s.chars();
while let Some(c) = chars.next() {
if c == '\\' {
if let Some(next) = chars.next() {
result.push(next);
}
} else {
result.push(c);
}
}
result
}
fn extract_level(obj: &Map<String, Value>) -> Option<String> {
for key in &["level", "severity", "lvl", "PRIORITY", "@level"] {
if let Some(val) = obj.get(*key)
&& let Some(level) = normalize_level_value(val)
{
return Some(level);
}
}
None
}
fn normalize_level_value(val: &Value) -> Option<String> {
match val {
Value::String(s) => normalize_level_str(s),
Value::Number(n) => {
let n = n.as_i64()?;
match n {
50 | 60 => Some("error".into()),
40 => Some("warn".into()),
30 => Some("info".into()),
20 => Some("debug".into()),
10 => Some("trace".into()),
0..=3 => Some("error".into()),
4 | 5 => Some("warn".into()),
6 => Some("info".into()),
7 => Some("debug".into()),
_ => None,
}
}
_ => None,
}
}
pub fn normalize_level_str(s: &str) -> Option<String> {
let lower = s.to_ascii_lowercase();
match lower.as_str() {
"error" | "err" | "fatal" | "critical" | "panic" | "alert" | "emerg" => {
Some("error".into())
}
"warn" | "warning" | "wrn" => Some("warn".into()),
"info" | "inf" | "information" | "notice" => Some("info".into()),
"debug" | "dbg" => Some("debug".into()),
"trace" | "trc" => Some("trace".into()),
_ => None,
}
}
fn extract_first_string(obj: &Map<String, Value>, keys: &[&str]) -> Option<String> {
for key in keys {
if let Some(Value::String(s)) = obj.get(*key) {
return Some(s.clone());
}
}
None
}
fn extract_msg(obj: &Map<String, Value>) -> Option<String> {
extract_first_string(obj, &["msg", "message", "event", "@message"])
}
fn extract_logger(obj: &Map<String, Value>) -> Option<String> {
extract_first_string(obj, &["logger", "name", "component", "module"])
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_json_parse() {
let line = r#"{"level":"info","msg":"server started","port":8080}"#;
let parsed = parse(line, "json");
assert_eq!(parsed.level.as_deref(), Some("info"));
assert_eq!(parsed.msg.as_deref(), Some("server started"));
assert!(parsed.fields_json.is_some());
}
#[test]
fn test_json_level_normalization() {
let line = r#"{"level":"FATAL","msg":"crash"}"#;
let parsed = parse(line, "json");
assert_eq!(parsed.level.as_deref(), Some("error"));
}
#[test]
fn test_json_pino_integer_level() {
let line = r#"{"level":50,"msg":"error occurred"}"#;
let parsed = parse(line, "json");
assert_eq!(parsed.level.as_deref(), Some("error"));
}
#[test]
fn test_json_syslog_priority() {
let line = r#"{"PRIORITY":3,"msg":"system error"}"#;
let parsed = parse(line, "json");
assert_eq!(parsed.level.as_deref(), Some("error"));
}
#[test]
fn test_json_msg_aliases() {
let line = r#"{"event":"hello","level":"info"}"#;
let parsed = parse(line, "json");
assert_eq!(parsed.msg.as_deref(), Some("hello"));
}
#[test]
fn test_logfmt_parse() {
let line = r#"level=info msg="server started" port=8080"#;
let parsed = parse(line, "logfmt");
assert_eq!(parsed.level.as_deref(), Some("info"));
assert_eq!(parsed.msg.as_deref(), Some("server started"));
assert!(parsed.fields_json.is_some());
}
#[test]
fn test_logfmt_bare_key() {
let line = r#"level=debug ready msg="ok""#;
let parsed = parse(line, "logfmt");
assert_eq!(parsed.level.as_deref(), Some("debug"));
assert_eq!(parsed.msg.as_deref(), Some("ok"));
let fields: Value = serde_json::from_str(parsed.fields_json.as_deref().unwrap()).unwrap();
assert_eq!(fields["ready"], Value::Bool(true));
}
#[test]
fn test_logfmt_quoted_value_with_spaces() {
let line = r#"level=error msg="connection refused: timeout""#;
let parsed = parse(line, "logfmt");
assert_eq!(parsed.msg.as_deref(), Some("connection refused: timeout"));
}
#[test]
fn test_text_format() {
let line = r#"{"level":"info"}"#;
let parsed = parse(line, "text");
assert!(parsed.level.is_none());
assert!(parsed.fields_json.is_none());
assert_eq!(parsed.message, line);
}
#[test]
fn test_json_parse_failure_falls_back() {
let line = "{not valid json";
let parsed = parse(line, "json");
assert!(parsed.level.is_none());
assert!(parsed.fields_json.is_none());
assert_eq!(parsed.message, line);
}
#[test]
fn test_logfmt_logger_extraction() {
let line = r#"level=info msg="hi" logger=myapp"#;
let parsed = parse(line, "logfmt");
assert_eq!(parsed.logger.as_deref(), Some("myapp"));
}
#[test]
fn test_json_nested_not_extracted_as_msg() {
let line = r#"{"level":"info","fields":{"message":"nested"}}"#;
let parsed = parse(line, "json");
assert_eq!(parsed.level.as_deref(), Some("info"));
assert_eq!(parsed.msg, None); }
}