From 1724346da78f11787307080dbbea1521dc52ad96 Mon Sep 17 00:00:00 2001 From: rays Date: Tue, 29 Sep 2026 13:22:45 +0000 Subject: [PATCH] Send OpenTelemetry traces over OTLP (#87) ipx had no spans, only log lines, so there was no way to see where a slow scan, download or request spent its time. With OTEL_EXPORTER_OTLP_ENDPOINT set, the daemon now exports traces over OTLP/HTTP (Tempo on Tower): a scan, each feed in it, the feed fetch and site icon lookup, downloads, torrents, reaps, and web requests. Log lines inside a span ride along as its events. Only the daemon exports: the healthcheck runs ipx status every 30s and would bury everything else. The web event stream and the log view's polling get no span, for the same reason. The exporter shares ipx's reqwest 0.13, so no second HTTP stack comes in. The stderr log now prefixes lines inside a span with it, as tracing-subscriber's fmt layer does (scan{only=None force=false}: ...). Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 2 + Cargo.lock | 112 ++++++++++++++++++++++++++++++++++++++++++ Cargo.toml | 4 ++ docker-compose.yml | 2 + docs/configuration.md | 1 + src/feed.rs | 2 + src/main.rs | 42 +++++++++++++++- src/web.rs | 17 ++++++- 8 files changed, 180 insertions(+), 2 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 2ddae31..4064b5c 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 +- With `OTEL_EXPORTER_OTLP_ENDPOINT` set, the daemon sends traces of its scans, downloads and web + requests to a collector such as Tempo. - A feed that has no artwork of its own shows its website's icon instead. - The refresh button turns while its feed is being checked, and the check-every-feed buttons while any of yours is. diff --git a/Cargo.lock b/Cargo.lock index 7635929..c07a50d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1832,6 +1832,9 @@ dependencies = [ "futures-util", "jsonwebtoken", "librqbit", + "opentelemetry", + "opentelemetry-otlp", + "opentelemetry_sdk", "opml", "percent-encoding", "quick-xml 0.42.0", @@ -1845,6 +1848,7 @@ dependencies = [ "tower", "tower-http 0.7.1", "tracing", + "tracing-opentelemetry", "tracing-subscriber", "url", ] @@ -2710,6 +2714,76 @@ version = "0.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7c87def4c32ab89d880effc9e097653c8da5d6ef28e6b539d313baaacfbafcbe" +[[package]] +name = "opentelemetry" +version = "0.33.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6cdb0b1b267eb9db3331b434ed9ddab10d50e280a9adf9d13e5233e2002b61b5" +dependencies = [ + "futures-core", + "futures-sink", + "js-sys", + "pin-project-lite", + "thiserror 2.0.20", +] + +[[package]] +name = "opentelemetry-http" +version = "0.33.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ee2c3625b8aa04209f7e01e513bc8044da9687fa38a9eacf59628ca5f3b87300" +dependencies = [ + "async-trait", + "bytes", + "http", + "opentelemetry", + "reqwest", +] + +[[package]] +name = "opentelemetry-otlp" +version = "0.33.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "699a67345e21962a231955b9059a157e428b219fa5df8efe11f08346fb23a344" +dependencies = [ + "http", + "httpdate", + "opentelemetry", + "opentelemetry-http", + "opentelemetry-proto", + "opentelemetry_sdk", + "prost", + "reqwest", + "thiserror 2.0.20", +] + +[[package]] +name = "opentelemetry-proto" +version = "0.33.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "25da1ac11a0aeccf38d7f77ee0348715adaf8340f65ad46c94a02c6b20e2f65d" +dependencies = [ + "opentelemetry", + "opentelemetry_sdk", + "prost", +] + +[[package]] +name = "opentelemetry_sdk" +version = "0.33.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cb39533d9d1c912123efd7d41d7e0c29d16917b60ce15b4c8d87cb1af7f67520" +dependencies = [ + "futures-channel", + "futures-executor", + "futures-util", + "opentelemetry", + "percent-encoding", + "portable-atomic", + "rand 0.9.5", + "thiserror 2.0.20", +] + [[package]] name = "opml" version = "1.1.6" @@ -2925,6 +2999,29 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "prost" +version = "0.14.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "528ac67416ff8646872a3c02cad9cc4ee5dc9f9540c9b10771855c95cb2e5ae1" +dependencies = [ + "bytes", + "prost-derive", +] + +[[package]] +name = "prost-derive" +version = "0.14.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b570b25f7617e43d59005d0990ccb79e950a423952cea19671b7a876da390adf" +dependencies = [ + "anyhow", + "itertools 0.14.0", + "proc-macro2", + "quote", + "syn 2.0.119", +] + [[package]] name = "quanta" version = "0.12.6" @@ -3225,6 +3322,7 @@ dependencies = [ "base64 0.23.1", "bytes", "encoding_rs", + "futures-channel", "futures-core", "futures-util", "h2", @@ -4563,6 +4661,20 @@ dependencies = [ "tracing-core", ] +[[package]] +name = "tracing-opentelemetry" +version = "0.34.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0a904802a1b902f43638b677ff2a650847e3b4404101b6c586d648e8c1e3e8fe" +dependencies = [ + "js-sys", + "opentelemetry", + "tracing", + "tracing-core", + "tracing-subscriber", + "web-time", +] + [[package]] name = "tracing-subscriber" version = "0.3.23" diff --git a/Cargo.toml b/Cargo.toml index 0b338ae..3f27673 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -14,6 +14,9 @@ clap = { version = "4.6.6", features = ["derive"] } futures-util = { version = "0.3.34", default-features = false, features = ["std"] } jsonwebtoken = { version = "11.1.0", default-features = false, features = ["aws_lc_rs"] } librqbit = { version = "9.0.1", default-features = false, features = ["rust-tls", "http-api-client"] } +opentelemetry = { version = "0.33", default-features = false, features = ["trace"] } +opentelemetry-otlp = { version = "0.33", default-features = false, features = ["trace", "http-proto", "reqwest-blocking-client"] } +opentelemetry_sdk = { version = "0.33", default-features = false, features = ["trace"] } opml = "1.1.6" percent-encoding = "2.3.2" quick-xml = { version = "0.42.0", features = ["escape-html"] } @@ -27,5 +30,6 @@ toml = "1.1.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"] } url = "2.5.8" diff --git a/docker-compose.yml b/docker-compose.yml index d46e88a..3457fad 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -8,6 +8,8 @@ services: PGID: "100" TZ: "America/Toronto" IPX_LOG: "ipx=info" + # Traces to Tempo, in the monitoring project; ipx sends none without it. + OTEL_EXPORTER_OTLP_ENDPOINT: "http://192.168.1.130:4318" # IPX_DATABASE_URL=postgres://... to use Postgres; without it, /data/state.db (SQLite). env_file: - ipx.env # relative: Arcane resolves it inside its own container diff --git a/docs/configuration.md b/docs/configuration.md index 9908cc0..69b92dd 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -143,4 +143,5 @@ and they are re-derived on every scan. Editing one in the UI promotes it to a ca | `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_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/src/feed.rs b/src/feed.rs index 3680e1d..5fa2ba0 100644 --- a/src/feed.rs +++ b/src/feed.rs @@ -55,6 +55,7 @@ pub enum Fetched { /// Conditional GET. reqwest handles gzip and redirects; the original's hand-rolled /// CONNECT/socket.ssl proxy path is gone -- `system-proxy` reads http_proxy/https_proxy. +#[tracing::instrument(skip_all, fields(url = %cfg.url))] pub async fn fetch( client: &reqwest::Client, cfg: &FeedCfg, @@ -371,6 +372,7 @@ pub async fn feed_behind_page(client: &reqwest::Client, url: &str) -> String { /// Artwork for a feed that has none: the icon its site's page names, or else the site's /// `/favicon.ico`. None if neither is there. +#[tracing::instrument(skip_all, fields(site = site))] pub async fn site_icon(client: &reqwest::Client, site: &str) -> Option { let timeout = std::time::Duration::from_secs(20); let resp = client.get(site).timeout(timeout).send().await.ok()?; diff --git a/src/main.rs b/src/main.rs index bf58c8a..236e9be 100644 --- a/src/main.rs +++ b/src/main.rs @@ -169,6 +169,12 @@ impl Ctx { #[tokio::main] async fn main() -> Result<()> { let cli = Cli::parse(); + // The daemon only: `ipx status` runs every half minute as the healthcheck, and a trace + // for each would bury the ones worth looking at. + let otel = match cli.command { + Command::Daemon { .. } => otel_provider()?, + _ => None, + }; // Everything goes to stderr as before, and is mirrored into a ring the UI can read. { use tracing_subscriber::layer::SubscriberExt; @@ -188,6 +194,12 @@ async fn main() -> Result<()> { .with_filter(stderr_filter), ) .with(logbuf::RingLayer.with_filter(ui_filter)) + .with(otel.as_ref().map(|p| { + use opentelemetry::trace::TracerProvider; + tracing_opentelemetry::layer() + .with_tracer(p.tracer("ipx")) + .with_filter(tracing_subscriber::EnvFilter::new("ipx=info")) + })) .init(); } @@ -241,7 +253,7 @@ async fn main() -> Result<()> { detach_torrents: is_daemon, }); - match cli.command { + let result = match cli.command { Command::List => list(&ctx).await, Command::Daemon { web } => daemon(ctx, config_path, web, events).await, Command::Add { url, folder, keywords } => { @@ -253,7 +265,29 @@ async fn main() -> Result<()> { Command::Export { file } => export(&ctx, &file).await, Command::CopyDb { from } => copy_db(&ctx, &from).await, _ => run(&ctx, wire_cmd.expect("only List and Daemon have no wire form")).await, + }; + // The batch exporter holds the last few seconds of spans; without this they are lost. + if let Some(p) = otel { + let _ = p.shutdown(); } + result +} + +/// 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. +fn otel_provider() -> Result> { + if std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT").unwrap_or_default().is_empty() { + return Ok(None); + } + let exporter = opentelemetry_otlp::SpanExporter::builder().with_http().build()?; + let resource = opentelemetry_sdk::Resource::builder().with_service_name("ipx").build(); + Ok(Some( + opentelemetry_sdk::trace::SdkTracerProvider::builder() + .with_batch_exporter(exporter) + .with_resource(resource) + .build(), + )) } /// Accounts. Passwords come in on stdin so they never reach a shell history or a `ps` @@ -872,6 +906,7 @@ async fn list(ctx: &Ctx) -> Result<()> { /// `standalone` false means this is the sweep that runs before a scan: it reports what it /// deleted, but must not emit the terminal ReapDone, or a client waiting on its `fetch` /// would stop reading before the scan had even started. +#[tracing::instrument(name = "reap", skip_all, fields(dry_run = dry_run))] async fn reap(ctx: &Ctx, dry_run: bool, standalone: bool) -> Result<()> { let r = retention::run(&ctx.cfg(), &ctx.db, dry_run).await?; for c in r.aged_out.iter().chain(r.over_quota.iter()) { @@ -890,6 +925,7 @@ async fn reap(ctx: &Ctx, dry_run: bool, standalone: bool) -> Result<()> { } /// `scope`, when not empty, narrows the scan to those feeds and the feeds inside any of them. +#[tracing::instrument(name = "scan", skip_all, fields(only = ?only, force = force))] async fn fetch(ctx: &Arc, only: Option<&str>, force: bool, scope: &[String]) -> Result<()> { let cfg = ctx.cfg(); let subs = subscriptions(ctx).await?; @@ -1127,6 +1163,7 @@ enum Outcome { Opml { added: Vec, removed: usize, kept: usize, total: usize }, } +#[tracing::instrument(name = "feed", skip_all, fields(feed = id))] async fn scan_one( ctx: &Arc, id: &str, @@ -1545,6 +1582,7 @@ fn merge_policy(subs: &[db::Sub], feed_cfg: &config::Feed, global: usize) -> Pol policy } +#[tracing::instrument(name = "download", skip_all, fields(feed = feed_id, enclosure = enclosure, url = url))] async fn fetch_one( ctx: &Arc, feed_id: &str, @@ -1627,6 +1665,7 @@ fn spawn_torrent(ctx: &Arc, feed_id: String, enclosure: i64, url: String, d /// Downloads one specific enclosure immediately, whatever the per-scan cap says and /// wherever it sits in the queue. +#[tracing::instrument(skip(ctx))] async fn download_one(ctx: &Arc, id: i64) -> Result<()> { let cfg = ctx.cfg(); let enc = ctx @@ -1696,6 +1735,7 @@ async fn download_one(ctx: &Arc, id: i64) -> Result<()> { /// Torrent progress is reported the same way an HTTP download's is, throttled to whole /// percents so a UI is not flooded. +#[tracing::instrument(name = "torrent", skip_all, fields(feed = feed_id, enclosure = enclosure))] async fn torrent_one( ctx: &Arc, feed_id: &str, diff --git a/src/web.rs b/src/web.rs index 1110c06..4d0bd1b 100644 --- a/src/web.rs +++ b/src/web.rs @@ -1938,8 +1938,23 @@ async fn access_log(req: Request, next: Next) -> Response { let path = req.uri().path().to_owned(); let method = req.method().clone(); let quiet = path.starts_with("/api/logs"); + // No trace for the event stream either: an open page asks for it again and again, and Tempo + // would fill with nothing else. + let span = if quiet || path == "/api/events" { + tracing::Span::none() + } else { + tracing::info_span!( + target: "ipx::http", + "http", + otel.name = %format!("{method} {path}"), + http.request.method = %method, + url.path = %path, + http.response.status_code = tracing::field::Empty, + ) + }; let started = std::time::Instant::now(); - let resp = next.run(req).await; + 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 status = resp.status().as_u16();