use tracing::{Level as TracingLevel, debug, error, info, trace, warn};
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum ParsedLevel {
Trace,
Debug,
Info,
Warn,
Error,
}
impl ParsedLevel {
fn from_token(token: &str) -> Option<Self> {
match token {
"TRACE" => Some(Self::Trace),
"DEBUG" => Some(Self::Debug),
"INFO" => Some(Self::Info),
"WARN" => Some(Self::Warn),
"ERROR" => Some(Self::Error),
_ => None,
}
}
}
#[derive(Debug, PartialEq, Eq)]
pub struct TelemetryLine<'a> {
pub timestamp: &'a str,
pub level: ParsedLevel,
pub component: &'a str,
pub file_line: &'a str,
pub message: &'a str,
pub fields: Vec<(&'a str, &'a str)>,
}
pub fn parse_line(line: &str) -> Option<TelemetryLine<'_>> {
let line = line.trim_end_matches(['\r', '\n']);
let (timestamp, rest) = split_optional_timestamp(line);
let rest = rest.trim_start();
let rest = rest.strip_prefix('[')?;
let close = rest.find(']')?;
let level = ParsedLevel::from_token(&rest[..close])?;
let mut rest = &rest[close + 1..];
let mut component = "";
if let Some(after_open) = rest.strip_prefix('[')
&& let Some(close) = after_open.find(']')
{
component = &after_open[..close];
rest = &after_open[close + 1..];
}
let rest = rest.trim_start();
let (file_line, rest_after_loc) = split_optional_file_line(rest);
let rest = rest_after_loc.trim_start();
let (message, fields_block) = split_message_and_fields(rest);
let fields = fields_block.map(parse_fields).unwrap_or_default();
Some(TelemetryLine {
timestamp,
level,
component,
file_line,
message: message.trim_end(),
fields,
})
}
fn split_optional_timestamp(line: &str) -> (&str, &str) {
if line.len() >= 24
&& line.as_bytes().get(10) == Some(&b'T')
&& line.as_bytes().get(23) == Some(&b'Z')
&& line
.as_bytes()
.get(24)
.map(|c| c.is_ascii_whitespace())
.unwrap_or(false)
{
(&line[..24], &line[24..])
} else {
("", line)
}
}
fn split_optional_file_line(rest: &str) -> (&str, &str) {
let Some(end) = rest.find(char::is_whitespace) else {
return ("", rest);
};
let token = &rest[..end];
if let Some(colon) = token.rfind(':')
&& token[colon + 1..].chars().all(|c| c.is_ascii_digit())
&& !token[colon + 1..].is_empty()
{
return (token, &rest[end..]);
}
("", rest)
}
fn split_message_and_fields(rest: &str) -> (&str, Option<&str>) {
if !rest.ends_with(']') {
return (rest, None);
}
let Some(open) = rest.rfind(" [") else {
return (rest, None);
};
let block = &rest[open + 2..rest.len() - 1];
if block.contains('=') {
(&rest[..open], Some(block))
} else {
(rest, None)
}
}
fn parse_fields(block: &str) -> Vec<(&str, &str)> {
block
.split(", ")
.filter_map(|kv| kv.split_once('='))
.collect()
}
pub fn emit_as_tracing(line: &TelemetryLine<'_>) {
let fields_json = fields_to_json(&line.fields);
let component = line.component;
let file_line = line.file_line;
let message = line.message;
match line.level {
ParsedLevel::Trace => trace!(
provider = %component,
source = %file_line,
fields = %fields_json,
"{message}"
),
ParsedLevel::Debug => debug!(
provider = %component,
source = %file_line,
fields = %fields_json,
"{message}"
),
ParsedLevel::Info => info!(
provider = %component,
source = %file_line,
fields = %fields_json,
"{message}"
),
ParsedLevel::Warn => warn!(
provider = %component,
source = %file_line,
fields = %fields_json,
"{message}"
),
ParsedLevel::Error => error!(
provider = %component,
source = %file_line,
fields = %fields_json,
"{message}"
),
}
let _ = TracingLevel::INFO;
}
fn fields_to_json(fields: &[(&str, &str)]) -> String {
let mut out = String::with_capacity(fields.len() * 16);
out.push('{');
for (i, (k, v)) in fields.iter().enumerate() {
if i > 0 {
out.push(',');
}
out.push('"');
out.push_str(&escape_json(k));
out.push_str("\":\"");
out.push_str(&escape_json(v));
out.push('"');
}
out.push('}');
out
}
fn escape_json(s: &str) -> String {
s.chars()
.flat_map(|c| match c {
'"' => vec!['\\', '"'],
'\\' => vec!['\\', '\\'],
'\n' => vec!['\\', 'n'],
'\r' => vec!['\\', 'r'],
'\t' => vec!['\\', 't'],
c if c.is_control() => format!("\\u{:04x}", c as u32).chars().collect(),
c => vec![c],
})
.collect()
}
use std::io::Write as _;
use std::pin::Pin;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use std::task::{Context, Poll};
use tokio::io::AsyncWrite;
use wasmtime_wasi::cli::{IsTerminal, StdoutStream};
#[derive(Clone, Copy, Debug)]
pub enum StreamKind {
Stdout,
Stderr,
}
#[derive(Clone)]
pub struct TelemetryStream {
kind: StreamKind,
forward_unmatched: Arc<AtomicBool>,
}
impl TelemetryStream {
pub fn stdout() -> Self {
Self {
kind: StreamKind::Stdout,
forward_unmatched: Arc::new(AtomicBool::new(true)),
}
}
pub fn stderr() -> Self {
Self {
kind: StreamKind::Stderr,
forward_unmatched: Arc::new(AtomicBool::new(true)),
}
}
pub fn drop_unmatched(self) -> Self {
self.forward_unmatched.store(false, Ordering::Relaxed);
self
}
}
impl IsTerminal for TelemetryStream {
fn is_terminal(&self) -> bool {
false
}
}
impl StdoutStream for TelemetryStream {
fn async_stream(&self) -> Box<dyn AsyncWrite + Send + Sync> {
Box::new(TelemetryLineSink {
kind: self.kind,
forward_unmatched: self.forward_unmatched.clone(),
buffer: Mutex::new(Vec::new()),
})
}
}
struct TelemetryLineSink {
kind: StreamKind,
forward_unmatched: Arc<AtomicBool>,
buffer: Mutex<Vec<u8>>,
}
impl TelemetryLineSink {
fn dispatch_lines(&self, complete: &[u8]) {
for raw in complete.split(|b| *b == b'\n') {
if raw.is_empty() {
continue;
}
let raw = if raw.last() == Some(&b'\r') {
&raw[..raw.len() - 1]
} else {
raw
};
let line = match std::str::from_utf8(raw) {
Ok(s) => s,
Err(_) => {
if self.forward_unmatched.load(Ordering::Relaxed) {
forward_bytes(self.kind, raw);
forward_bytes(self.kind, b"\n");
}
continue;
}
};
if let Some(parsed) = parse_line(line) {
emit_as_tracing(&parsed);
} else if self.forward_unmatched.load(Ordering::Relaxed) {
forward_bytes(self.kind, raw);
forward_bytes(self.kind, b"\n");
}
}
}
}
fn forward_bytes(kind: StreamKind, bytes: &[u8]) {
match kind {
StreamKind::Stdout => {
let _ = std::io::stdout().write_all(bytes);
}
StreamKind::Stderr => {
let _ = std::io::stderr().write_all(bytes);
}
}
}
impl AsyncWrite for TelemetryLineSink {
fn poll_write(
self: Pin<&mut Self>,
_cx: &mut Context<'_>,
buf: &[u8],
) -> Poll<Result<usize, std::io::Error>> {
let drained = {
let mut guard = self.buffer.lock().unwrap_or_else(|e| e.into_inner());
guard.extend_from_slice(buf);
let split_at = guard.iter().rposition(|b| *b == b'\n').map(|i| i + 1);
split_at.map(|i| guard.drain(..i).collect::<Vec<_>>())
};
if let Some(complete) = drained {
self.dispatch_lines(&complete);
}
Poll::Ready(Ok(buf.len()))
}
fn poll_flush(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<Result<(), std::io::Error>> {
Poll::Ready(Ok(()))
}
fn poll_shutdown(
self: Pin<&mut Self>,
_cx: &mut Context<'_>,
) -> Poll<Result<(), std::io::Error>> {
let leftover = {
let mut guard = self.buffer.lock().unwrap_or_else(|e| e.into_inner());
std::mem::take(&mut *guard)
};
if !leftover.is_empty() {
self.dispatch_lines(&leftover);
}
Poll::Ready(Ok(()))
}
}
#[cfg(test)]
mod tests {
use super::*;
use futures::task::noop_waker;
use std::pin::Pin;
use std::task::Context;
#[test]
fn parses_full_line_with_all_segments() {
let line = "2026-05-25T10:30:42.123Z [DEBUG][messaging.webchat-gui] components/foo/src/lib.rs:160 span-start: send_payload [id=1, event_kind=send_payload]";
let p = parse_line(line).expect("should parse");
assert_eq!(p.timestamp, "2026-05-25T10:30:42.123Z");
assert_eq!(p.level, ParsedLevel::Debug);
assert_eq!(p.component, "messaging.webchat-gui");
assert_eq!(p.file_line, "components/foo/src/lib.rs:160");
assert_eq!(p.message, "span-start: send_payload");
assert_eq!(p.fields, vec![("id", "1"), ("event_kind", "send_payload")]);
}
#[test]
fn parses_line_without_component() {
let line = "[INFO] doing work [k=v]";
let p = parse_line(line).expect("should parse");
assert_eq!(p.timestamp, "");
assert_eq!(p.level, ParsedLevel::Info);
assert_eq!(p.component, "");
assert_eq!(p.file_line, "");
assert_eq!(p.message, "doing work");
assert_eq!(p.fields, vec![("k", "v")]);
}
#[test]
fn parses_line_without_fields_block() {
let line = "[WARN][svc] src/lib.rs:1 some warning";
let p = parse_line(line).expect("should parse");
assert_eq!(p.level, ParsedLevel::Warn);
assert_eq!(p.component, "svc");
assert_eq!(p.file_line, "src/lib.rs:1");
assert_eq!(p.message, "some warning");
assert!(p.fields.is_empty());
}
#[test]
fn parses_legacy_line_format() {
let line = "[DEBUG] span-start: send_payload [event_kind=send_payload, provider=messaging.webchat-gui]";
let p = parse_line(line).expect("should parse");
assert_eq!(p.timestamp, "");
assert_eq!(p.level, ParsedLevel::Debug);
assert_eq!(p.component, "");
assert_eq!(p.file_line, "");
assert_eq!(p.message, "span-start: send_payload");
assert_eq!(
p.fields,
vec![
("event_kind", "send_payload"),
("provider", "messaging.webchat-gui"),
]
);
}
#[test]
fn rejects_non_telemetry_lines() {
assert!(parse_line("[ws notifier:memory] publish tenant=demo").is_none());
assert!(parse_line("secrets: backend=dev-store").is_none());
assert!(parse_line("").is_none());
assert!(parse_line("plain text").is_none());
}
#[test]
fn handles_trailing_newline() {
let line = "[INFO] hello\n";
let p = parse_line(line).expect("should parse");
assert_eq!(p.message, "hello");
}
#[test]
fn message_with_bracket_not_followed_by_equals_stays_in_message() {
let line = "[INFO] config loaded [final state]";
let p = parse_line(line).expect("should parse");
assert_eq!(p.message, "config loaded [final state]");
assert!(p.fields.is_empty());
}
#[test]
fn file_line_with_relative_path() {
let line = "[ERROR][svc] crates/foo/src/lib.rs:42 oops [code=500]";
let p = parse_line(line).expect("should parse");
assert_eq!(p.file_line, "crates/foo/src/lib.rs:42");
assert_eq!(p.message, "oops");
assert_eq!(p.fields, vec![("code", "500")]);
}
#[test]
fn rejects_unknown_level_and_unclosed_tokens() {
assert!(parse_line("[NOTICE] hello").is_none());
assert!(parse_line("[INFO hello").is_none());
assert!(parse_line("[INFO][component hello").is_some());
}
#[test]
fn timestamp_requires_separator_after_z() {
let line = "2026-05-25T10:30:42.123Z[INFO] no separator";
let (timestamp, rest) = split_optional_timestamp(line);
assert_eq!(timestamp, "");
assert_eq!(rest, line);
}
#[test]
fn malformed_file_line_stays_in_message() {
let line = "[INFO] src/lib.rs:not-a-line still message";
let p = parse_line(line).expect("should parse");
assert_eq!(p.file_line, "");
assert_eq!(p.message, "src/lib.rs:not-a-line still message");
}
#[test]
fn fields_ignore_entries_without_equals() {
let line = "[INFO] event [good=value, missing, also=ok]";
let p = parse_line(line).expect("should parse");
assert_eq!(p.fields, vec![("good", "value"), ("also", "ok")]);
}
#[test]
fn fields_json_escapes_special_characters() {
let json = fields_to_json(&[
("quote", "a\"b"),
("slash", "a\\b"),
("line", "a\nb"),
("tab", "a\tb"),
("ctrl", "\u{1f}"),
]);
assert_eq!(
json,
"{\"quote\":\"a\\\"b\",\"slash\":\"a\\\\b\",\"line\":\"a\\nb\",\"tab\":\"a\\tb\",\"ctrl\":\"\\u001f\"}"
);
}
#[test]
fn telemetry_stream_buffers_partial_lines_until_shutdown() {
let stream = TelemetryStream::stdout().drop_unmatched();
let mut sink = TelemetryLineSink {
kind: stream.kind,
forward_unmatched: stream.forward_unmatched,
buffer: Default::default(),
};
let waker = noop_waker();
let mut cx = Context::from_waker(&waker);
assert!(
Pin::new(&mut sink)
.poll_write(&mut cx, b"[INFO][svc] partial")
.is_ready()
);
assert!(Pin::new(&mut sink).poll_shutdown(&mut cx).is_ready());
}
#[test]
fn telemetry_stream_handles_invalid_utf8_without_forwarding() {
let stream = TelemetryStream::stderr().drop_unmatched();
let mut sink = TelemetryLineSink {
kind: stream.kind,
forward_unmatched: stream.forward_unmatched,
buffer: Default::default(),
};
let waker = noop_waker();
let mut cx = Context::from_waker(&waker);
assert!(
Pin::new(&mut sink)
.poll_write(&mut cx, &[0xff, b'\n'])
.is_ready()
);
assert!(Pin::new(&mut sink).poll_flush(&mut cx).is_ready());
}
}