Trace ids, failure kinds and one line per event in the JSON log (#91)

From Dash0's structured logging guide, what applies here:

- Each JSON line inside a traced span ends with its trace_id and span_id, so a line in Loki leads
  to its trace in Tempo; the access log is written inside its request's span so it has one too.
  The JSON formatter takes no extra fields, so WithTrace appends them to the object it writes.
- A feed or download failure carries error.type (the HTTP status, or dns, redirect_loop,
  timeout, ...) and http.response.status_code, from failure_kind beside explain_failure, so
  failures group by kind without a regex over msg.
- Each event was logged twice: words under ipx::scan and fields under ipx::io. It is now one
  line under ipx::scan with both; the wire copy is at debug, for the admin page's Daemon I/O tab,
  and out of production's log. The healthcheck's status reply stays under ipx::io.
- The access log's ms is duration_ms. The dashboard and the prod-check skill follow.
- error fields are Display with the anyhow chain everywhere, not a mix of Debug and Display.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
2026-09-29 16:25:53 +00:00
parent 36b16c7f5d
commit 448e557272
8 changed files with 166 additions and 67 deletions

View File

@@ -99,6 +99,37 @@ pub struct Failure {
pub new_url: Option<String>,
}
/// A failure's kind, for the log's `error.type` (#91): the HTTP status where there is one, as
/// OpenTelemetry names an HTTP error, and otherwise a word for what went wrong. Matches the
/// same wording as `explain_failure`; a message it does not know is "other", never a wrong kind.
pub fn failure_kind(msg: &str) -> (String, Option<u16>) {
let low = msg.to_ascii_lowercase();
let code = low.split("http ").skip(1).find_map(|r| r.get(..3)?.parse::<u16>().ok());
if let Some(c) = code.filter(|c| (100..600).contains(c)) {
return (c.to_string(), Some(c));
}
let kind = if low.contains("dns error") || low.contains("failed to lookup address") || low.contains("no address associated") {
"dns"
} else if low.contains("too many redirects") {
"redirect_loop"
} else if low.contains("timed out") || low.contains("timeout") {
"timeout"
} else if low.contains("certificate") || low.contains("tls") {
"tls"
} else if low.contains("got a web page") {
"not_a_feed"
} else if low.contains("the site sent ") {
"site_message"
} else if low.contains("connect") {
"connect"
} else if low.contains("pars") {
"parse"
} else {
"other"
};
(kind.into(), None)
}
/// Reads a `last_error` the same way `set_feed_error` received it (`format!("{e:#}")` on the
/// anyhow chain from `fetch` or `parse`) and says what it means, for the errors worth telling
/// someone about. Everything else -- a timeout, a 5xx, a 429, a feed that is simply garbled --
@@ -1016,6 +1047,20 @@ mod tests {
);
}
#[test]
fn failures_are_named_by_kind_for_the_log() {
let k = |m: &str| failure_kind(m);
assert_eq!(k("HTTP 404 Not Found"), ("404".into(), Some(404)));
assert_eq!(k("HTTP 503 Service Unavailable"), ("503".into(), Some(503)));
assert_eq!(
k("connecting: error following redirect for url (https://www.toddstashwick.com/): too many redirects").0,
"redirect_loop"
);
assert_eq!(k("connecting: dns error: failed to lookup address information").0, "dns");
assert_eq!(k("operation timed out").0, "timeout");
assert_eq!(k("something new").0, "other");
}
#[test]
fn explain_failure_translates_the_errors_the_ui_should_flag() {
assert_eq!(

View File

@@ -131,34 +131,8 @@ impl Emitter {
// Level by how much it matters. With 80-odd feeds in an OPML subscription, one
// line per feed per tick for "not due yet" would push everything worth reading
// out of the buffer within a few minutes.
let routine = match &e {
Event::Progress { .. } | Event::FeedSkip { .. } | Event::FeedStart { .. } => true,
Event::FeedDone { new, downloaded, failed, torrents, .. } => {
*new == 0 && *downloaded == 0 && *failed == 0 && *torrents == 0
}
_ => false,
};
let bad = matches!(
&e,
Event::FeedError { .. } | Event::DownloadError { .. } | Event::Error { .. }
);
if let Some(line) = e.human() {
let line = line.trim();
if bad {
tracing::warn!(target: "ipx::scan", "{line}");
} else if routine {
tracing::debug!(target: "ipx::scan", "{line}");
} else {
tracing::info!(target: "ipx::scan", "{line}");
}
}
log_event(&e, self.tx.is_some());
if let Some(tx) = &self.tx {
// The outbound half of the protocol, as it goes on the wire. Progress is the
// high-volume one, so it sits at debug.
if let Ok(json) = serde_json::to_string(&e) {
log_wire(&json, matches!(e, Event::Progress { .. }));
}
// An error here only means nobody is listening yet.
let _ = tx.send(e.clone());
}
@@ -168,26 +142,61 @@ impl Emitter {
}
}
/// An event as it goes on the wire, `<- {json}`, with its fields as the log line's own as well,
/// so Loki reads `ev`, `feed`, `new` and the rest from the JSON log without parsing the message
/// (#91). A field an event lacks is left out.
fn log_wire(json: &str, debug: bool) {
let v: serde_json::Value = serde_json::from_str(json).unwrap_or_default();
/// An event, logged once (#91): its words as the message, and `ev`, `feed`, `new` and the rest as
/// fields, so Loki reads them without parsing the message. A failure also gets `error.type` and,
/// from an HTTP error, `http.response.status_code`, so failures group by kind without a regex.
/// Level by how much it matters: with 80-odd feeds in an OPML subscription, a line per feed per
/// tick for "not due yet" would push everything worth reading out of the log view in minutes,
/// and Progress fires on every whole percent. `wire` also logs the event as it goes on the
/// socket, at debug: the admin page's Daemon I/O tab shows it, production's log leaves it out.
fn log_event(e: &Event, wire: bool) {
let json = serde_json::to_string(e).unwrap_or_default();
let v: serde_json::Value = serde_json::from_str(&json).unwrap_or_default();
let s = |k: &str| v.get(k).and_then(|x| x.as_str());
let n = |k: &str| v.get(k).and_then(|x| x.as_u64());
macro_rules! wire {
($level:ident) => {
let (kind, code) = match e {
Event::FeedError { msg, .. } | Event::DownloadError { msg, .. } | Event::Error { msg } => {
let (k, c) = crate::feed::failure_kind(msg);
(Some(k), c)
}
_ => (None, None),
};
let routine = match e {
Event::Progress { .. } | Event::FeedSkip { .. } | Event::FeedStart { .. } => true,
Event::FeedDone { new, downloaded, failed, torrents, .. } => {
*new == 0 && *downloaded == 0 && *failed == 0 && *torrents == 0
}
_ => false,
};
macro_rules! line {
($level:ident, $target:literal, $text:expr) => {
tracing::$level!(
target: "ipx::io",
target: $target,
ev = s("ev"), feed = s("feed"), msg = s("msg"), url = s("url"), reason = s("reason"),
new = n("new"), downloaded = n("downloaded"), failed = n("failed"),
torrents = n("torrents"), bytes = n("bytes"), feeds = n("feeds"),
pending = n("pending"), enclosure = n("enclosure"), files = n("files"),
"<- {json}"
"error.type" = kind, "http.response.status_code" = code,
"{}", $text
)
};
}
if debug { wire!(debug) } else { wire!(info) }
// The healthcheck's answer, every 30s: a reply on the socket rather than work done, so it
// stays with the rest of the conversation, and the Scans tab stays about scans.
if matches!(e, Event::Status { .. }) {
return line!(info, "ipx::io", format!("<- {json}"));
}
if wire {
tracing::debug!(target: "ipx::io", "<- {json}");
}
let text = e.human().map(|l| l.trim().to_owned()).unwrap_or_else(|| json.clone());
if kind.is_some() {
line!(warn, "ipx::scan", text)
} else if routine {
line!(debug, "ipx::scan", text)
} else {
line!(info, "ipx::scan", text)
}
}
/// True when something is already listening -- i.e. a daemon owns this socket.
@@ -275,9 +284,7 @@ async fn handle(
Ok(Command::Status) => {
tracing::info!(target: "ipx::io", "-> {line}");
let ev = status().await;
if let Ok(json) = serde_json::to_string(&ev) {
log_wire(&json, false);
}
log_event(&ev, true);
let _ = reply.send(ev).await;
}
Ok(cmd) => {

View File

@@ -196,10 +196,10 @@ async fn main() -> Result<()> {
}))
.with(json.then(|| {
tracing_subscriber::fmt::layer()
.json()
.flatten_event(true)
.with_current_span(true)
.with_span_list(false)
.fmt_fields(tracing_subscriber::fmt::format::JsonFields::new())
.event_format(WithTrace(
tracing_subscriber::fmt::format().json().flatten_event(true).with_current_span(true).with_span_list(false),
))
.with_writer(std::io::stderr)
.with_filter(stderr_filter())
}))
@@ -275,6 +275,39 @@ async fn main() -> Result<()> {
result
}
/// The JSON log line with the trace and span it belongs to (#91), so a line in Loki leads to its
/// trace in Tempo. The JSON formatter cannot take a field of its own, so the ids go on the end of
/// the object it writes. A line outside any traced span, or with no trace exporter, is unchanged.
struct WithTrace<F>(F);
impl<S, N, F> tracing_subscriber::fmt::FormatEvent<S, N> for WithTrace<F>
where
F: tracing_subscriber::fmt::FormatEvent<S, N>,
S: tracing::Subscriber + for<'a> tracing_subscriber::registry::LookupSpan<'a>,
N: for<'w> tracing_subscriber::fmt::FormatFields<'w> + 'static,
{
fn format_event(
&self,
ctx: &tracing_subscriber::fmt::FmtContext<'_, S, N>,
mut w: tracing_subscriber::fmt::format::Writer<'_>,
ev: &tracing::Event<'_>,
) -> std::fmt::Result {
use opentelemetry::trace::TraceContextExt;
use tracing_opentelemetry::OpenTelemetrySpanExt;
let current = tracing::Span::current();
let sc = if current.is_none() { None } else { Some(current.context().span().span_context().clone()) };
let Some(sc) = sc.filter(|c| c.is_valid()) else {
return self.0.format_event(ctx, w, ev);
};
let mut line = String::new();
self.0.format_event(ctx, tracing_subscriber::fmt::format::Writer::new(&mut line), ev)?;
match line.trim_end().strip_suffix('}') {
Some(body) => writeln!(w, r#"{body},"trace_id":"{}","span_id":"{}"}}"#, sc.trace_id(), sc.span_id()),
None => w.write_str(&line),
}
}
}
/// Traces over OTLP, to Tempo for one, when OTEL_EXPORTER_OTLP_ENDPOINT names a collector
/// (`http://host:4318`: the exporter speaks OTLP over HTTP and adds `/v1/traces`). The exporter
/// reads that and the other `OTEL_` variables itself.
@@ -449,13 +482,13 @@ async fn daemon(
match ctx.db.requeue_interrupted().await {
Ok(n) if n > 0 => tracing::info!(count = n, "requeued downloads interrupted by a restart"),
Ok(_) => {}
Err(e) => tracing::warn!(error = ?e, "could not requeue interrupted downloads"),
Err(e) => tracing::warn!(error = %format!("{e:#}"), "could not requeue interrupted downloads"),
}
match retire_stranded(&ctx).await {
Ok(0) => {}
Ok(n) => tracing::info!(feeds = n, "retired feeds whose OPML is no longer in config"),
Err(e) => tracing::warn!(error = ?e, "could not retire feeds whose OPML is no longer in config"),
Err(e) => tracing::warn!(error = %format!("{e:#}"), "could not retire feeds whose OPML is no longer in config"),
}
let (tx_cmd, mut rx_cmd) = mpsc::channel::<Cmd>(64);
@@ -603,7 +636,7 @@ async fn start_web(
};
Ok(Some(tokio::spawn(async move {
if let Err(e) = web::serve(state, &bind).await {
tracing::error!(error = ?e, "web ui stopped");
tracing::error!(error = %format!("{e:#}"), "web ui stopped");
}
})))
}
@@ -1649,7 +1682,7 @@ fn spawn_torrent(ctx: &Arc<Ctx>, feed_id: String, enclosure: i64, url: String, d
match outcome {
Ok((path, bytes)) => {
if let Err(e) = db.mark_downloaded(&url, &path, bytes).await {
tracing::warn!(error = ?e, "could not record the finished torrent");
tracing::warn!(error = %format!("{e:#}"), "could not record the finished torrent");
}
ctx.out.emit(Event::DownloadDone {
feed: feed_id,

View File

@@ -2007,15 +2007,17 @@ async fn access_log(req: Request, next: Next) -> Response {
let resp = tracing::Instrument::instrument(next.run(req), span.clone()).await;
span.record("http.response.status_code", resp.status().as_u16());
if !quiet {
let ms = started.elapsed().as_millis() as u64;
let duration_ms = started.elapsed().as_millis() as u64;
let status = resp.status().as_u16();
// The same as fields, for the JSON log (#91). Unrouted, a request has no route.
let route = resp.extensions().get::<MatchedPath>().map(|r| r.as_str().to_owned());
let method = method.as_str();
// Inside the request's span, so the line carries its trace id and leads to its trace.
let _in = span.enter();
if resp.status().is_success() || resp.status().is_redirection() {
tracing::info!(target: "ipx::http", method, path, route, status, ms, "{method} {path} -> {status} in {ms}ms");
tracing::info!(target: "ipx::http", method, path, route, status, duration_ms, "{method} {path} -> {status} in {duration_ms}ms");
} else {
tracing::warn!(target: "ipx::http", method, path, route, status, ms, "{method} {path} -> {status} in {ms}ms");
tracing::warn!(target: "ipx::http", method, path, route, status, duration_ms, "{method} {path} -> {status} in {duration_ms}ms");
}
}
resp