Log as JSON when IPX_LOG_FORMAT=json (#91)

The log was text, so the Grafana dashboard picked lines apart with
regular expressions, and a change of wording would have blanked its
panels. With IPX_LOG_FORMAT=json each line is one JSON object: the
access log carries method, path, route, status and ms as fields (the
route passed from the routing layer in the response's extensions), and
each wire event its ev, feed, new, downloaded, failed, bytes, msg and
the rest (log_wire), beside the old message. The two startup lines that
were println! are logged, so no line breaks the JSON. Text stays the
default, for a terminal. The dashboard reads the fields with Loki's json
parser, and groups requests by route rather than path.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
2026-09-29 13:55:50 +00:00
parent 8ce0a4cb27
commit 4f8b3d6a1d
10 changed files with 95 additions and 35 deletions

View File

@@ -9,6 +9,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
### Added ### Added
- `IPX_LOG_FORMAT=json` logs one JSON object a line, with a request's method, route, status and time
and a scan's or download's details as fields, for Loki and the like.
- A Grafana dashboard for iPX, from its log in Loki and its traces in Tempo (`grafana/dashboard.py`). - A Grafana dashboard for iPX, from its log in Loki and its traces in Tempo (`grafana/dashboard.py`).
- With `OTEL_EXPORTER_OTLP_ENDPOINT` set, the daemon sends traces of its scans, downloads and web - With `OTEL_EXPORTER_OTLP_ENDPOINT` set, the daemon sends traces of its scans, downloads and web
requests to a collector such as Tempo. requests to a collector such as Tempo.

View File

@@ -170,9 +170,10 @@ Non-trivial logic leaves one runnable check behind. Pure functions (`merge_polic
watch the shutdown channel itself; the daemon ignored SIGTERM for exactly this reason. watch the shutdown channel itself; the daemon ignored SIGTERM for exactly this reason.
* Only one daemon per socket. Removing the socket file defeats the guard and you get two daemons * 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. fighting over the database, with the stale one still holding the port.
* **The Grafana dashboard parses the log** (`grafana/dashboard.py`): the access log's * **The Grafana dashboard reads the log's fields** (`grafana/dashboard.py`). Production logs JSON
`GET /path -> 200 in 3ms` and the events' `ipx::io: <- {json}`. Change either and the panels go (`IPX_LOG_FORMAT=json`); the access log's `method`, `path`, `route`, `status`, `ms` and the
blank without an error; regenerate the dashboard with the new pattern. events' `ev`, `feed`, `new`, `bytes`, `msg` (`log_wire` 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 * `/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 different task. `ipx status` checks the control socket and the database; to see the worker
getting through its jobs, watch for `scan complete` in the log. getting through its jobs, watch for `scan complete` in the log.

13
Cargo.lock generated
View File

@@ -4675,6 +4675,16 @@ dependencies = [
"web-time", "web-time",
] ]
[[package]]
name = "tracing-serde"
version = "0.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "704b1aeb7be0d0a84fc9828cae51dab5970fee5088f83d1dd7ee6f6246fc6ff1"
dependencies = [
"serde",
"tracing-core",
]
[[package]] [[package]]
name = "tracing-subscriber" name = "tracing-subscriber"
version = "0.3.23" version = "0.3.23"
@@ -4685,12 +4695,15 @@ dependencies = [
"nu-ansi-term", "nu-ansi-term",
"once_cell", "once_cell",
"regex-automata", "regex-automata",
"serde",
"serde_json",
"sharded-slab", "sharded-slab",
"smallvec", "smallvec",
"thread_local", "thread_local",
"tracing", "tracing",
"tracing-core", "tracing-core",
"tracing-log", "tracing-log",
"tracing-serde",
] ]
[[package]] [[package]]

View File

@@ -31,5 +31,5 @@ tower = { version = "0.5.3", features = ["util"] }
tower-http = { version = "0.7.1", features = ["fs"] } tower-http = { version = "0.7.1", features = ["fs"] }
tracing = "0.1.44" tracing = "0.1.44"
tracing-opentelemetry = { version = "0.34", default-features = false } tracing-opentelemetry = { version = "0.34", default-features = false }
tracing-subscriber = { version = "0.3.23", features = ["env-filter"] } tracing-subscriber = { version = "0.3.23", features = ["env-filter", "json"] }
url = "2.5.8" url = "2.5.8"

View File

@@ -8,6 +8,8 @@ services:
PGID: "100" PGID: "100"
TZ: "America/Toronto" TZ: "America/Toronto"
IPX_LOG: "ipx=info" IPX_LOG: "ipx=info"
# One JSON object a line, which Loki (through Alloy) and the Grafana dashboard read.
IPX_LOG_FORMAT: "json"
# Traces to Tempo, in the monitoring project; ipx sends none without it. # Traces to Tempo, in the monitoring project; ipx sends none without it.
OTEL_EXPORTER_OTLP_ENDPOINT: "http://192.168.1.130:4318" OTEL_EXPORTER_OTLP_ENDPOINT: "http://192.168.1.130:4318"
# What the dashboard filters on, so a daemon run by hand for testing stays out of it. # What the dashboard filters on, so a daemon run by hand for testing stays out of it.

View File

@@ -142,6 +142,7 @@ and they are re-derived on every scan. Editing one in the UI promotes it to a ca
| `IPX_DATABASE_URL` | A `postgres://user:password@host:port/database` URL: use that database instead of `state.db` | | `IPX_DATABASE_URL` | A `postgres://user:password@host:port/database` URL: use that database instead of `state.db` |
| `IPX_TEST_DATABASE_URL` | For `cargo test`: run the database tests on this Postgres database too, each in a schema of its own | | `IPX_TEST_DATABASE_URL` | For `cargo test`: run the database tests on this Postgres database too, each in a schema of its own |
| `IPX_LOG` | What reaches stderr (`ipx=debug`, `ipx::scan=debug`, …) | | `IPX_LOG` | What reaches stderr (`ipx=debug`, `ipx::scan=debug`, …) |
| `IPX_LOG_FORMAT` | `json` for one JSON object a line, with each request's and event's fields as its own (for Loki and the like); text otherwise |
| `IPX_UI_LOG` | What the in-process log buffer captures for the UI's Log view | | `IPX_UI_LOG` | What the in-process log buffer captures for the UI's Log view |
| `OTEL_EXPORTER_OTLP_ENDPOINT` | An OTLP/HTTP collector, such as Tempo at `http://host:4318`: the daemon sends it traces of scans, downloads and web requests. The other `OTEL_EXPORTER_OTLP_*` variables apply too | | `OTEL_EXPORTER_OTLP_ENDPOINT` | An OTLP/HTTP collector, such as Tempo at `http://host:4318`: the daemon sends it traces of scans, downloads and web requests. The other `OTEL_EXPORTER_OTLP_*` variables apply too |
| `http_proxy` / `https_proxy` | Honoured for feed and enclosure fetches | | `http_proxy` / `https_proxy` | Honoured for feed and enclosure fetches |

View File

@@ -3,8 +3,8 @@
python3 grafana/dashboard.py > /mnt/fast/arcane/projects/monitoring/grafana-provisioning/dashboards/ipx.json python3 grafana/dashboard.py > /mnt/fast/arcane/projects/monitoring/grafana-provisioning/dashboards/ipx.json
Grafana reads that file on its own within a minute; edits made in Grafana are refused. The panels Grafana reads that file on its own within a minute; edits made in Grafana are refused. The panels
parse ipx's log lines, so a change to what the access log or the event log prints (web.rs read the fields of ipx's JSON log (IPX_LOG_FORMAT=json): renaming a field in web.rs access_log or
access_log, ipc.rs Emitter::emit) has to be matched here. `dashboard.py queries` prints each ipc.rs log_wire has to be matched here. `dashboard.py queries` prints each
query, to try against Loki. query, to try against Loki.
""" """
import json, sys import json, sys
@@ -12,15 +12,17 @@ import json, sys
LOKI = {"type": "loki", "uid": "${loki}"} LOKI = {"type": "loki", "uid": "${loki}"}
TEMPO = {"type": "tempo", "uid": "${tempo}"} TEMPO = {"type": "tempo", "uid": "${tempo}"}
SEL = '{container="iPX"}' SEL = '{container="iPX"}'
# Every scan and download event is logged as its wire JSON: "ipx::io: <- {...}". # ipx logs one JSON object a line (IPX_LOG_FORMAT=json). Each scan and download event carries its
EV = SEL + ' |= "ipx::io: <- {" | regexp "<- (?P<j>\\\\{.*\\\\})$" | line_format "{{.j}}" | json | __error__=""' # fields (ev, feed, new, bytes, msg, ...); each request its method, path, route, status and ms.
HTTP = (SEL + ' |= "ipx::http: " | regexp "ipx::http: (?P<method>[A-Z]+) (?P<path>\\\\S+) -> (?P<status>\\\\d+) in (?P<ms>\\\\d+)ms"' EV = SEL + ' |= "\\"target\\":\\"ipx::io\\"" | json | __error__="" | ev != ""'
' | path != "/api/events"') HTTP = SEL + ' |= "\\"target\\":\\"ipx::http\\"" | json | __error__="" | path != "/api/events"'
BAD = SEL + ' | json | __error__="" | level =~ "WARN|ERROR"'
TEXT = ' | line_format "{{.level}} {{.target}}: {{.message}}"'
TRACES = '{resource.service.name="ipx" && resource.deployment.environment.name="production"' TRACES = '{resource.service.name="ipx" && resource.deployment.environment.name="production"'
def status(field): def status(field):
return f'max(max_over_time({SEL} |= "\\"ev\\":\\"status\\"" | regexp "\\"{field}\\":(?P<v>\\\\d+)" | unwrap v [10m]))' return f'max(max_over_time({EV} | ev = "status" | unwrap {field} [10m]))'
QUERIES = {} QUERIES = {}
@@ -93,7 +95,7 @@ stat("Feed failures", f'sum(count_over_time({EV} | ev="feed_error" [$__range]))
desc="Failed feed checks in the time range.") desc="Failed feed checks in the time range.")
stat("Download failures", f'sum(count_over_time({EV} | ev="download_error" [$__range])) or vector(0)', 16, stat("Download failures", f'sum(count_over_time({EV} | ev="download_error" [$__range])) or vector(0)', 16,
thresholds=[{"color": "green", "value": None}, {"color": "orange", "value": 1}]) thresholds=[{"color": "green", "value": None}, {"color": "orange", "value": 1}])
stat("Warnings and errors", f'sum(count_over_time({SEL} |~ "^\\\\S+\\\\s+(WARN|ERROR) " [$__range])) or vector(0)', 20, stat("Warnings and errors", f'sum(count_over_time({BAD} [$__range])) or vector(0)', 20,
thresholds=[{"color": "green", "value": None}, {"color": "orange", "value": 1}, {"color": "red", "value": 50}]) thresholds=[{"color": "green", "value": None}, {"color": "orange", "value": 1}, {"color": "red", "value": 50}])
y += 4 y += 4
@@ -125,10 +127,10 @@ ts("Response time", [loki(f'quantile_over_time(0.5, {HTTP} | unwrap ms [$__inter
loki(f'max_over_time({HTTP} | unwrap ms [$__interval]) by ()', ref="C", legend="slowest")], loki(f'max_over_time({HTTP} | unwrap ms [$__interval]) by ()', ref="C", legend="slowest")],
12, unit="ms") 12, unit="ms")
y += 8 y += 8
table("Slowest requests", [loki(f'topk(15, avg_over_time({HTTP} | unwrap ms [$__range]) by (method, path))', instant=True)], table("Slowest routes", [loki(f'topk(15, avg_over_time({HTTP} | route != "" | unwrap ms [$__range]) by (method, route))', instant=True)],
0, rename={"method": "Method", "path": "Path", "Value": "Average ms"}, sort=[{"displayName": "Average ms", "desc": True}]) 0, rename={"method": "Method", "route": "Route", "Value": "Average ms"}, sort=[{"displayName": "Average ms", "desc": True}])
table("Busiest paths", [loki(f'topk(15, sum by (method, path) (count_over_time({HTTP} [$__range])))', instant=True)], table("Busiest routes", [loki(f'topk(15, sum by (method, route) (count_over_time({HTTP} | route != "" [$__range])))', instant=True)],
12, rename={"method": "Method", "path": "Path", "Value": "Requests"}, sort=[{"displayName": "Requests", "desc": True}]) 12, rename={"method": "Method", "route": "Route", "Value": "Requests"}, sort=[{"displayName": "Requests", "desc": True}])
y += 8 y += 8
# ---- Traces # ---- Traces
@@ -143,11 +145,12 @@ y += 10
# ---- Log # ---- Log
row("Log") row("Log")
panel("logs", "Warnings and errors", 24, 10, 0, [loki(f'{SEL} |~ "^\\\\S+\\\\s+(WARN|ERROR) "')], panel("logs", "Warnings and errors", 24, 10, 0, [loki(BAD + TEXT)],
options={"showTime": True, "wrapLogMessage": True, "sortOrder": "Descending", "enableLogDetails": True}) options={"showTime": True, "wrapLogMessage": True, "sortOrder": "Descending", "enableLogDetails": True})
y += 10 y += 10
panel("logs", "Log", 24, 12, 0, panel("logs", "Log", 24, 12, 0,
[loki(f'{SEL} != "GET /api/events" != "\\"ev\\":\\"feed_skip\\"" != "{{\\"cmd\\":\\"status\\"}}"')], [loki(SEL + ' | json | __error__="" | path != "/api/events" | ev != "feed_skip" | ev != "status"'
' | message != "-> {\\"cmd\\":\\"status\\"}"' + TEXT)],
description="Without the event stream's requests, not-due skips and healthcheck status calls.", description="Without the event stream's requests, not-due skips and healthcheck status calls.",
options={"showTime": True, "wrapLogMessage": True, "sortOrder": "Descending", "enableLogDetails": True}) options={"showTime": True, "wrapLogMessage": True, "sortOrder": "Descending", "enableLogDetails": True})

View File

@@ -157,11 +157,7 @@ impl Emitter {
// The outbound half of the protocol, as it goes on the wire. Progress is the // The outbound half of the protocol, as it goes on the wire. Progress is the
// high-volume one, so it sits at debug. // high-volume one, so it sits at debug.
if let Ok(json) = serde_json::to_string(&e) { if let Ok(json) = serde_json::to_string(&e) {
if matches!(e, Event::Progress { .. }) { log_wire(&json, matches!(e, Event::Progress { .. }));
tracing::debug!(target: "ipx::io", "<- {json}");
} else {
tracing::info!(target: "ipx::io", "<- {json}");
}
} }
// An error here only means nobody is listening yet. // An error here only means nobody is listening yet.
let _ = tx.send(e.clone()); let _ = tx.send(e.clone());
@@ -172,6 +168,28 @@ 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();
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) => {
tracing::$level!(
target: "ipx::io",
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}"
)
};
}
if debug { wire!(debug) } else { wire!(info) }
}
/// True when something is already listening -- i.e. a daemon owns this socket. /// True when something is already listening -- i.e. a daemon owns this socket.
pub async fn daemon_is_live(path: &Path) -> bool { pub async fn daemon_is_live(path: &Path) -> bool {
UnixStream::connect(path).await.is_ok() UnixStream::connect(path).await.is_ok()
@@ -258,7 +276,7 @@ async fn handle(
tracing::info!(target: "ipx::io", "-> {line}"); tracing::info!(target: "ipx::io", "-> {line}");
let ev = status().await; let ev = status().await;
if let Ok(json) = serde_json::to_string(&ev) { if let Ok(json) = serde_json::to_string(&ev) {
tracing::info!(target: "ipx::io", "<- {json}"); log_wire(&json, false);
} }
let _ = reply.send(ev).await; let _ = reply.send(ev).await;
} }

View File

@@ -183,19 +183,32 @@ async fn main() -> Result<()> {
// Two filters, deliberately different. stderr follows IPX_LOG; the in-app buffer // Two filters, deliberately different. stderr follows IPX_LOG; the in-app buffer
// keeps debug as well, so the log view can show protocol traffic and routine // keeps debug as well, so the log view can show protocol traffic and routine
// skips that would be noise on a terminal. IPX_UI_LOG overrides it. // skips that would be noise on a terminal. IPX_UI_LOG overrides it.
let stderr_filter = tracing_subscriber::EnvFilter::try_from_env("IPX_LOG") let stderr_filter = || {
.unwrap_or_else(|_| "ipx=info".into()); tracing_subscriber::EnvFilter::try_from_env("IPX_LOG").unwrap_or_else(|_| "ipx=info".into())
};
// One JSON object a line for Loki (#91), with each event's fields as its own; text
// otherwise, for someone reading a terminal.
let json = std::env::var("IPX_LOG_FORMAT").is_ok_and(|f| f.eq_ignore_ascii_case("json"));
let ui_filter = tracing_subscriber::EnvFilter::try_from_env("IPX_UI_LOG") let ui_filter = tracing_subscriber::EnvFilter::try_from_env("IPX_UI_LOG")
.unwrap_or_else(|_| "ipx=debug".into()); .unwrap_or_else(|_| "ipx=debug".into());
tracing_subscriber::registry() tracing_subscriber::registry()
.with( .with((!json).then(|| {
tracing_subscriber::fmt::layer() tracing_subscriber::fmt::layer()
.with_writer(std::io::stderr) .with_writer(std::io::stderr)
// Colour for a terminal only: in docker logs and Loki the escapes are noise // Colour for a terminal only: in docker logs and Loki the escapes are noise
// every query has to strip (#88). // every query has to strip (#88).
.with_ansi(std::io::IsTerminal::is_terminal(&std::io::stderr())) .with_ansi(std::io::IsTerminal::is_terminal(&std::io::stderr()))
.with_filter(stderr_filter), .with_filter(stderr_filter())
) }))
.with(json.then(|| {
tracing_subscriber::fmt::layer()
.json()
.flatten_event(true)
.with_current_span(true)
.with_span_list(false)
.with_writer(std::io::stderr)
.with_filter(stderr_filter())
}))
.with(logbuf::RingLayer.with_filter(ui_filter)) .with(logbuf::RingLayer.with_filter(ui_filter))
.with(otel.as_ref().map(|p| { .with(otel.as_ref().map(|p| {
use opentelemetry::trace::TracerProvider; use opentelemetry::trace::TracerProvider;
@@ -575,12 +588,13 @@ async fn start_web(
ctx.set_cfg(fresh.clone()); ctx.set_cfg(fresh.clone());
// The token signs in as the admin, and whatever reads this process's output (docker logs, // The token signs in as the admin, and whatever reads this process's output (docker logs,
// for one) is wider than who reads config.toml. So say where it is, never what it is. // for one) is wider than who reads config.toml. So say where it is, never what it is.
println!( // Logged, not printed, so a JSON log stays one object a line (#91).
tracing::info!(
"web ui token generated and saved to {} as [web] token. Open http://{bind}/?token=<that token>", "web ui token generated and saved to {} as [web] token. Open http://{bind}/?token=<that token>",
config_path.display() config_path.display()
); );
} else { } else {
println!( tracing::info!(
"web ui at http://{bind}/ (the sign-in token is [web] token in {})", "web ui at http://{bind}/ (the sign-in token is [web] token in {})",
config_path.display() config_path.display()
); );

View File

@@ -1941,7 +1941,10 @@ async fn name_span(route: MatchedPath, req: Request, next: Next) -> Response {
let span = tracing::Span::current(); let span = tracing::Span::current();
span.context().span().update_name(format!("{} {}", req.method(), route.as_str())); span.context().span().update_name(format!("{} {}", req.method(), route.as_str()));
span.record("http.route", route.as_str()); span.record("http.route", route.as_str());
next.run(req).await let mut resp = next.run(req).await;
// For access_log, which runs outside routing and cannot see it otherwise.
resp.extensions_mut().insert(route);
resp
} }
/// One line per HTTP request, so the web side shows up in the same log as the daemon. /// One line per HTTP request, so the web side shows up in the same log as the daemon.
@@ -1971,12 +1974,15 @@ async fn access_log(req: Request, next: Next) -> Response {
let resp = tracing::Instrument::instrument(next.run(req), span.clone()).await; let resp = tracing::Instrument::instrument(next.run(req), span.clone()).await;
span.record("http.response.status_code", resp.status().as_u16()); span.record("http.response.status_code", resp.status().as_u16());
if !quiet { if !quiet {
let ms = started.elapsed().as_millis(); let ms = started.elapsed().as_millis() as u64;
let status = resp.status().as_u16(); 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();
if resp.status().is_success() || resp.status().is_redirection() { if resp.status().is_success() || resp.status().is_redirection() {
tracing::info!(target: "ipx::http", "{method} {path} -> {status} in {ms}ms"); tracing::info!(target: "ipx::http", method, path, route, status, ms, "{method} {path} -> {status} in {ms}ms");
} else { } else {
tracing::warn!(target: "ipx::http", "{method} {path} -> {status} in {ms}ms"); tracing::warn!(target: "ipx::http", method, path, route, status, ms, "{method} {path} -> {status} in {ms}ms");
} }
} }
resp resp