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();