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

@@ -24,12 +24,16 @@ $Q trace <trace id> # one trace as a tree, with its log lin
Every line has `timestamp`, `level`, `message` and `target`; lines inside a span have `span` (the
innermost: `{"name":"feed","feed":"x"}`). Loki's `| json` flattens it to `span_name`, `span_feed`.
Lines inside a traced span, requests and scans, also carry `trace_id` and `span_id`: give the
`trace_id` to `$Q trace` to see the whole request or scan. (From 2026-09-29 16:30 UTC; before that,
lines had no trace id, the access log's time was `ms`, and each event was logged twice, words under
`ipx::scan` and fields under `ipx::io`.)
| target | fields | what |
|---|---|---|
| `ipx::http` | `method`, `path`, `route`, `status`, `ms` | one per web request; `route` is the pattern, empty for an unrouted path |
| `ipx::io` | `ev` and the event's own: `feed`, `new`, `downloaded`, `failed`, `bytes`, `msg`, `url`, `feeds`, `pending`, `reason` | the daemon's events as they go on the wire |
| `ipx::scan` | message only | the same events in words; warnings are feed and download failures |
| `ipx::http` | `method`, `path`, `route`, `status`, `duration_ms` | one per web request; `route` is the pattern, empty for an unrouted path |
| `ipx::scan` | `ev` and the event's own: `feed`, `new`, `downloaded`, `failed`, `bytes`, `msg`, `url`, `feeds`, `reason`; on a failure `error.type` and, from an HTTP error, `http.response.status_code` | the daemon's events, one line each, in words; warnings are feed and download failures |
| `ipx::io` | `ev`, `feeds`, `pending`, `downloaded` on the `status` reply | commands arriving (`-> {...}`) and the healthcheck's answer |
| `ipx` | message, sometimes fields | start-up, shutdown, account and config messages |
Events (`ev`): `feed_start`, `feed_done` (new, downloaded, failed, torrents), `feed_skip` (not due,
@@ -37,7 +41,10 @@ routine), `feed_error` (msg), `download_done` (bytes), `download_error` (msg, ur
`torrent_deferred`, `reaped`, `scan_done` (feeds checked), `reap_done`, `status` (feeds, pending,
downloaded: the healthcheck's, every 30s), `error` (msg).
Filter on the text before `| json` where you can (`|= "\"target\":\"ipx::io\""`): it is much
`error.type` is the HTTP status (`404`, `503`) or one of `dns`, `redirect_loop`, `timeout`, `tls`,
`not_a_feed`, `site_message`, `connect`, `parse`, `other`; Loki's `| json` names it `error_type`.
Filter on the text before `| json` where you can (`|= "\"ev\":\"feed_error\""`): it is much
cheaper than parsing every line. Lines before 2026-09-29 14:00 UTC are text, not JSON, and
`| json | __error__=""` drops them.
@@ -53,8 +60,8 @@ a count says something happened, the lines and traces say why.
Feed and download failures name the feed in the message; group them in the next check instead.
2. **Failing feeds and downloads.**
```
$Q metric 'sum by (feed, msg) (count_over_time({container="iPX"} |= "\"target\":\"ipx::io\"" | json | __error__="" | ev="feed_error" [$range]))' 7d
$Q metric 'sum by (feed, msg) (count_over_time({container="iPX"} |= "\"target\":\"ipx::io\"" | json | __error__="" | ev="download_error" [$range]))' 7d
$Q metric 'sum by (feed, error_type) (count_over_time({container="iPX"} |= "\"ev\":\"feed_error\"" | json | __error__="" [$range]))' 7d
$Q metric 'sum by (feed, error_type) (count_over_time({container="iPX"} |= "\"ev\":\"download_error\"" | json | __error__="" [$range]))' 7d
```
Tell the publisher's problems from ipx's. A 404, 410, DNS failure or 503 from the feed's own
server is the publisher (worth saying, since the feed may have moved; one issue for a feed
@@ -63,7 +70,7 @@ a count says something happened, the lines and traces say why.
3. **Server errors and slow requests.**
```
$Q metric 'sum by (method, route, status) (count_over_time({container="iPX"} |= "\"target\":\"ipx::http\"" | json | __error__="" | status >= 500 [$range]))'
$Q metric 'topk(10, quantile_over_time(0.95, {container="iPX"} |= "\"target\":\"ipx::http\"" | json | __error__="" | route != "" | route != "/api/events" | unwrap ms [$range]) by (method, route))'
$Q metric 'topk(10, quantile_over_time(0.95, {container="iPX"} |= "\"target\":\"ipx::http\"" | json | __error__="" | route != "" | route != "/api/events" | unwrap duration_ms [$range]) by (method, route))'
```
Any 5xx is worth a look. 401s are people signing in, not a problem unless one address is
hammering. For a slow route, find its traces (check 5) and see which span holds the time.

View File

@@ -23,6 +23,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
### Changed
- The JSON log carries each line's `trace_id` and `span_id`, logs each scan event once instead of
twice, names a failure's kind in `error.type` (and its HTTP status in
`http.response.status_code`), and calls a request's time `duration_ms` instead of `ms`.
- `ipx list` shows each feed's id on a line of its own, labelled, under its title.
- Each browser keeps its own theme, so a phone and a desktop can differ. A browser that has not
chosen one yet starts from the theme your account had.

View File

@@ -171,8 +171,8 @@ Non-trivial logic leaves one runnable check behind. Pure functions (`merge_polic
* Only one daemon per socket. Removing the socket file defeats the guard and you get two daemons
fighting over the database, with the stale one still holding the port.
* **The Grafana dashboard reads the log's fields** (`grafana/dashboard.py`). Production logs JSON
(`IPX_LOG_FORMAT=json`); the access log's `method`, `path`, `route`, `status`, `ms` and the
events' `ev`, `feed`, `new`, `bytes`, `msg` (`log_wire` in ipc.rs) are what the panels query.
(`IPX_LOG_FORMAT=json`); the access log's `method`, `path`, `route`, `status`, `duration_ms` and
the events' `ev`, `feed`, `new`, `bytes`, `msg` (`log_event` in ipc.rs) are what the panels query.
Rename one and its panels go blank without an error; regenerate the dashboard to match.
* `/api/settings` answering `200` does **not** mean the daemon is well — the web server is a
different task. `ipx status` checks the control socket and the database; to see the worker

View File

@@ -4,7 +4,7 @@
Grafana reads that file on its own within a minute; edits made in Grafana are refused. The panels
read the fields of ipx's JSON log (IPX_LOG_FORMAT=json): renaming a field in web.rs access_log or
ipc.rs log_wire has to be matched here. `dashboard.py queries` prints each
ipc.rs log_event has to be matched here. `dashboard.py queries` prints each
query, to try against Loki.
"""
import json, sys
@@ -13,8 +13,10 @@ LOKI = {"type": "loki", "uid": "${loki}"}
TEMPO = {"type": "tempo", "uid": "${tempo}"}
SEL = '{container="iPX"}'
# ipx logs one JSON object a line (IPX_LOG_FORMAT=json). Each scan and download event carries its
# fields (ev, feed, new, bytes, msg, ...); each request its method, path, route, status and ms.
EV = SEL + ' |= "\\"target\\":\\"ipx::io\\"" | json | __error__="" | ev != ""'
# fields (ev, feed, new, bytes, msg, error.type, ...); each request its method, path, route, status
# and duration_ms. Events are logged under ipx::scan, the healthcheck's status under ipx::io, so
# the filter is on the ev field, not the target.
EV = SEL + ' |= "\\"ev\\":\\"" | json | __error__="" | ev != ""'
HTTP = SEL + ' |= "\\"target\\":\\"ipx::http\\"" | json | __error__="" | path != "/api/events"'
BAD = SEL + ' | json | __error__="" | level =~ "WARN|ERROR"'
TEXT = ' | line_format "{{.level}} {{.target}}: {{.message}}"'
@@ -122,12 +124,12 @@ y += 8
row("Web")
ts("Requests by status", [loki(f'sum by (status) (count_over_time({HTTP} [$__interval]))', legend="{{status}}")],
0, bars=True, stack=True, desc="The event stream the page keeps open is left out.")
ts("Response time", [loki(f'quantile_over_time(0.5, {HTTP} | unwrap ms [$__interval]) by ()', legend="median"),
loki(f'quantile_over_time(0.95, {HTTP} | unwrap ms [$__interval]) by ()', ref="B", legend="95th percentile"),
loki(f'max_over_time({HTTP} | unwrap ms [$__interval]) by ()', ref="C", legend="slowest")],
ts("Response time", [loki(f'quantile_over_time(0.5, {HTTP} | unwrap duration_ms [$__interval]) by ()', legend="median"),
loki(f'quantile_over_time(0.95, {HTTP} | unwrap duration_ms [$__interval]) by ()', ref="B", legend="95th percentile"),
loki(f'max_over_time({HTTP} | unwrap duration_ms [$__interval]) by ()', ref="C", legend="slowest")],
12, unit="ms")
y += 8
table("Slowest routes", [loki(f'topk(15, avg_over_time({HTTP} | route != "" | unwrap ms [$__range]) by (method, route))', instant=True)],
table("Slowest routes", [loki(f'topk(15, avg_over_time({HTTP} | route != "" | unwrap duration_ms [$__range]) by (method, route))', instant=True)],
0, rename={"method": "Method", "route": "Route", "Value": "Average ms"}, sort=[{"displayName": "Average ms", "desc": True}])
table("Busiest routes", [loki(f'topk(15, sum by (method, route) (count_over_time({HTTP} | route != "" [$__range])))', instant=True)],
12, rename={"method": "Method", "route": "Route", "Value": "Requests"}, sort=[{"displayName": "Requests", "desc": True}])

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