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:
@@ -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.
|
||||
|
||||
@@ -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.
|
||||
|
||||
13
Cargo.lock
generated
13
Cargo.lock
generated
@@ -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]]
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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 |
|
||||
|
||||
@@ -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<j>\\\\{.*\\\\})$" | line_format "{{.j}}" | json | __error__=""'
|
||||
HTTP = (SEL + ' |= "ipx::http: " | regexp "ipx::http: (?P<method>[A-Z]+) (?P<path>\\\\S+) -> (?P<status>\\\\d+) in (?P<ms>\\\\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<v>\\\\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})
|
||||
|
||||
|
||||
30
src/ipc.rs
30
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;
|
||||
}
|
||||
|
||||
28
src/main.rs
28
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=<that token>",
|
||||
config_path.display()
|
||||
);
|
||||
} else {
|
||||
println!(
|
||||
tracing::info!(
|
||||
"web ui at http://{bind}/ (the sign-in token is [web] token in {})",
|
||||
config_path.display()
|
||||
);
|
||||
|
||||
14
src/web.rs
14
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::<MatchedPath>().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
|
||||
|
||||
Reference in New Issue
Block a user