diff --git a/CHANGELOG.md b/CHANGELOG.md index 117d02d..ae80029 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### 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`). - With `OTEL_EXPORTER_OTLP_ENDPOINT` set, the daemon sends traces of its scans, downloads and web requests to a collector such as Tempo. diff --git a/CLAUDE.md b/CLAUDE.md index 131ad3b..731ede4 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -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. * 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 parses the log** (`grafana/dashboard.py`): the access log's - `GET /path -> 200 in 3ms` and the events' `ipx::io: <- {json}`. Change either and the panels go - blank without an error; regenerate the dashboard with the new pattern. +* **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. + 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 getting through its jobs, watch for `scan complete` in the log. diff --git a/Cargo.lock b/Cargo.lock index c07a50d..e7223da 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4675,6 +4675,16 @@ dependencies = [ "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]] name = "tracing-subscriber" version = "0.3.23" @@ -4685,12 +4695,15 @@ dependencies = [ "nu-ansi-term", "once_cell", "regex-automata", + "serde", + "serde_json", "sharded-slab", "smallvec", "thread_local", "tracing", "tracing-core", "tracing-log", + "tracing-serde", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index 3f27673..bbded01 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -31,5 +31,5 @@ tower = { version = "0.5.3", features = ["util"] } tower-http = { version = "0.7.1", features = ["fs"] } tracing = "0.1.44" 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" diff --git a/docker-compose.yml b/docker-compose.yml index dbbccbe..a123953 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -8,6 +8,8 @@ services: PGID: "100" TZ: "America/Toronto" 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. 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. diff --git a/docs/configuration.md b/docs/configuration.md index 69b92dd..49d8666 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -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_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_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 | | `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 | diff --git a/grafana/dashboard.py b/grafana/dashboard.py index db4ecd4..db46d24 100644 --- a/grafana/dashboard.py +++ b/grafana/dashboard.py @@ -3,8 +3,8 @@ 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 -parse ipx's log lines, so a change to what the access log or the event log prints (web.rs -access_log, ipc.rs Emitter::emit) has to be matched here. `dashboard.py queries` prints each +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 query, to try against Loki. """ import json, sys @@ -12,15 +12,17 @@ import json, sys LOKI = {"type": "loki", "uid": "${loki}"} TEMPO = {"type": "tempo", "uid": "${tempo}"} SEL = '{container="iPX"}' -# Every scan and download event is logged as its wire JSON: "ipx::io: <- {...}". -EV = SEL + ' |= "ipx::io: <- {" | regexp "<- (?P\\\\{.*\\\\})$" | line_format "{{.j}}" | json | __error__=""' -HTTP = (SEL + ' |= "ipx::http: " | regexp "ipx::http: (?P[A-Z]+) (?P\\\\S+) -> (?P\\\\d+) in (?P\\\\d+)ms"' - ' | path != "/api/events"') +# 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 != ""' +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"' def status(field): - return f'max(max_over_time({SEL} |= "\\"ev\\":\\"status\\"" | regexp "\\"{field}\\":(?P\\\\d+)" | unwrap v [10m]))' + return f'max(max_over_time({EV} | ev = "status" | unwrap {field} [10m]))' 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.") 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}]) -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}]) 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")], 12, unit="ms") y += 8 -table("Slowest requests", [loki(f'topk(15, avg_over_time({HTTP} | unwrap ms [$__range]) by (method, path))', instant=True)], - 0, rename={"method": "Method", "path": "Path", "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)], - 12, rename={"method": "Method", "path": "Path", "Value": "Requests"}, sort=[{"displayName": "Requests", "desc": 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", "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}]) y += 8 # ---- Traces @@ -143,11 +145,12 @@ y += 10 # ---- 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}) y += 10 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.", options={"showTime": True, "wrapLogMessage": True, "sortOrder": "Descending", "enableLogDetails": True}) diff --git a/src/ipc.rs b/src/ipc.rs index 5ce36ec..14af6e2 100644 --- a/src/ipc.rs +++ b/src/ipc.rs @@ -157,11 +157,7 @@ impl Emitter { // 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) { - if matches!(e, Event::Progress { .. }) { - tracing::debug!(target: "ipx::io", "<- {json}"); - } else { - tracing::info!(target: "ipx::io", "<- {json}"); - } + log_wire(&json, matches!(e, Event::Progress { .. })); } // An error here only means nobody is listening yet. 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. pub async fn daemon_is_live(path: &Path) -> bool { UnixStream::connect(path).await.is_ok() @@ -258,7 +276,7 @@ async fn handle( tracing::info!(target: "ipx::io", "-> {line}"); let ev = status().await; if let Ok(json) = serde_json::to_string(&ev) { - tracing::info!(target: "ipx::io", "<- {json}"); + log_wire(&json, false); } let _ = reply.send(ev).await; } diff --git a/src/main.rs b/src/main.rs index 6511acb..a539754 100644 --- a/src/main.rs +++ b/src/main.rs @@ -183,19 +183,32 @@ async fn main() -> Result<()> { // 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 // skips that would be noise on a terminal. IPX_UI_LOG overrides it. - let stderr_filter = tracing_subscriber::EnvFilter::try_from_env("IPX_LOG") - .unwrap_or_else(|_| "ipx=info".into()); + let stderr_filter = || { + 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") .unwrap_or_else(|_| "ipx=debug".into()); tracing_subscriber::registry() - .with( + .with((!json).then(|| { tracing_subscriber::fmt::layer() .with_writer(std::io::stderr) // Colour for a terminal only: in docker logs and Loki the escapes are noise // every query has to strip (#88). .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(otel.as_ref().map(|p| { use opentelemetry::trace::TracerProvider; @@ -575,12 +588,13 @@ async fn start_web( ctx.set_cfg(fresh.clone()); // 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. - 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=", config_path.display() ); } else { - println!( + tracing::info!( "web ui at http://{bind}/ (the sign-in token is [web] token in {})", config_path.display() ); diff --git a/src/web.rs b/src/web.rs index b2729da..2b1ac92 100644 --- a/src/web.rs +++ b/src/web.rs @@ -1941,7 +1941,10 @@ async fn name_span(route: MatchedPath, req: Request, next: Next) -> Response { let span = tracing::Span::current(); span.context().span().update_name(format!("{} {}", req.method(), 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. @@ -1971,12 +1974,15 @@ 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(); + let 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::().map(|r| r.as_str().to_owned()); + let method = method.as_str(); 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 { - 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