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 <noreply@anthropic.com>
This commit is contained in:
2026-09-29 13:22:45 +00:00
parent 8c5eddd783
commit 1724346da7
8 changed files with 180 additions and 2 deletions

View File

@@ -9,6 +9,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
### Added ### 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. - 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 - The refresh button turns while its feed is being checked, and the check-every-feed buttons
while any of yours is. while any of yours is.

112
Cargo.lock generated
View File

@@ -1832,6 +1832,9 @@ dependencies = [
"futures-util", "futures-util",
"jsonwebtoken", "jsonwebtoken",
"librqbit", "librqbit",
"opentelemetry",
"opentelemetry-otlp",
"opentelemetry_sdk",
"opml", "opml",
"percent-encoding", "percent-encoding",
"quick-xml 0.42.0", "quick-xml 0.42.0",
@@ -1845,6 +1848,7 @@ dependencies = [
"tower", "tower",
"tower-http 0.7.1", "tower-http 0.7.1",
"tracing", "tracing",
"tracing-opentelemetry",
"tracing-subscriber", "tracing-subscriber",
"url", "url",
] ]
@@ -2710,6 +2714,76 @@ version = "0.2.1"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7c87def4c32ab89d880effc9e097653c8da5d6ef28e6b539d313baaacfbafcbe" 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]] [[package]]
name = "opml" name = "opml"
version = "1.1.6" version = "1.1.6"
@@ -2925,6 +2999,29 @@ dependencies = [
"unicode-ident", "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]] [[package]]
name = "quanta" name = "quanta"
version = "0.12.6" version = "0.12.6"
@@ -3225,6 +3322,7 @@ dependencies = [
"base64 0.23.1", "base64 0.23.1",
"bytes", "bytes",
"encoding_rs", "encoding_rs",
"futures-channel",
"futures-core", "futures-core",
"futures-util", "futures-util",
"h2", "h2",
@@ -4563,6 +4661,20 @@ dependencies = [
"tracing-core", "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]] [[package]]
name = "tracing-subscriber" name = "tracing-subscriber"
version = "0.3.23" version = "0.3.23"

View File

@@ -14,6 +14,9 @@ clap = { version = "4.6.6", features = ["derive"] }
futures-util = { version = "0.3.34", default-features = false, features = ["std"] } futures-util = { version = "0.3.34", default-features = false, features = ["std"] }
jsonwebtoken = { version = "11.1.0", default-features = false, features = ["aws_lc_rs"] } 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"] } 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" opml = "1.1.6"
percent-encoding = "2.3.2" percent-encoding = "2.3.2"
quick-xml = { version = "0.42.0", features = ["escape-html"] } 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 = { 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-subscriber = { version = "0.3.23", features = ["env-filter"] } tracing-subscriber = { version = "0.3.23", features = ["env-filter"] }
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"
# 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). # IPX_DATABASE_URL=postgres://... to use Postgres; without it, /data/state.db (SQLite).
env_file: env_file:
- ipx.env # relative: Arcane resolves it inside its own container - ipx.env # relative: Arcane resolves it inside its own container

View File

@@ -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_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_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 |
| `http_proxy` / `https_proxy` | Honoured for feed and enclosure fetches | | `http_proxy` / `https_proxy` | Honoured for feed and enclosure fetches |

View File

@@ -55,6 +55,7 @@ pub enum Fetched {
/// Conditional GET. reqwest handles gzip and redirects; the original's hand-rolled /// 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. /// 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( pub async fn fetch(
client: &reqwest::Client, client: &reqwest::Client,
cfg: &FeedCfg, 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 /// 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. /// `/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<String> { pub async fn site_icon(client: &reqwest::Client, site: &str) -> Option<String> {
let timeout = std::time::Duration::from_secs(20); let timeout = std::time::Duration::from_secs(20);
let resp = client.get(site).timeout(timeout).send().await.ok()?; let resp = client.get(site).timeout(timeout).send().await.ok()?;

View File

@@ -169,6 +169,12 @@ impl Ctx {
#[tokio::main] #[tokio::main]
async fn main() -> Result<()> { async fn main() -> Result<()> {
let cli = Cli::parse(); 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. // Everything goes to stderr as before, and is mirrored into a ring the UI can read.
{ {
use tracing_subscriber::layer::SubscriberExt; use tracing_subscriber::layer::SubscriberExt;
@@ -188,6 +194,12 @@ async fn main() -> Result<()> {
.with_filter(stderr_filter), .with_filter(stderr_filter),
) )
.with(logbuf::RingLayer.with_filter(ui_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(); .init();
} }
@@ -241,7 +253,7 @@ async fn main() -> Result<()> {
detach_torrents: is_daemon, detach_torrents: is_daemon,
}); });
match cli.command { let result = match cli.command {
Command::List => list(&ctx).await, Command::List => list(&ctx).await,
Command::Daemon { web } => daemon(ctx, config_path, web, events).await, Command::Daemon { web } => daemon(ctx, config_path, web, events).await,
Command::Add { url, folder, keywords } => { Command::Add { url, folder, keywords } => {
@@ -253,7 +265,29 @@ async fn main() -> Result<()> {
Command::Export { file } => export(&ctx, &file).await, Command::Export { file } => export(&ctx, &file).await,
Command::CopyDb { from } => copy_db(&ctx, &from).await, Command::CopyDb { from } => copy_db(&ctx, &from).await,
_ => run(&ctx, wire_cmd.expect("only List and Daemon have no wire form")).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<Option<opentelemetry_sdk::trace::SdkTracerProvider>> {
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` /// 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 /// `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` /// deleted, but must not emit the terminal ReapDone, or a client waiting on its `fetch`
/// would stop reading before the scan had even started. /// 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<()> { async fn reap(ctx: &Ctx, dry_run: bool, standalone: bool) -> Result<()> {
let r = retention::run(&ctx.cfg(), &ctx.db, dry_run).await?; let r = retention::run(&ctx.cfg(), &ctx.db, dry_run).await?;
for c in r.aged_out.iter().chain(r.over_quota.iter()) { 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. /// `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<Ctx>, only: Option<&str>, force: bool, scope: &[String]) -> Result<()> { async fn fetch(ctx: &Arc<Ctx>, only: Option<&str>, force: bool, scope: &[String]) -> Result<()> {
let cfg = ctx.cfg(); let cfg = ctx.cfg();
let subs = subscriptions(ctx).await?; let subs = subscriptions(ctx).await?;
@@ -1127,6 +1163,7 @@ enum Outcome {
Opml { added: Vec<String>, removed: usize, kept: usize, total: usize }, Opml { added: Vec<String>, removed: usize, kept: usize, total: usize },
} }
#[tracing::instrument(name = "feed", skip_all, fields(feed = id))]
async fn scan_one( async fn scan_one(
ctx: &Arc<Ctx>, ctx: &Arc<Ctx>,
id: &str, id: &str,
@@ -1545,6 +1582,7 @@ fn merge_policy(subs: &[db::Sub], feed_cfg: &config::Feed, global: usize) -> Pol
policy policy
} }
#[tracing::instrument(name = "download", skip_all, fields(feed = feed_id, enclosure = enclosure, url = url))]
async fn fetch_one( async fn fetch_one(
ctx: &Arc<Ctx>, ctx: &Arc<Ctx>,
feed_id: &str, feed_id: &str,
@@ -1627,6 +1665,7 @@ fn spawn_torrent(ctx: &Arc<Ctx>, feed_id: String, enclosure: i64, url: String, d
/// Downloads one specific enclosure immediately, whatever the per-scan cap says and /// Downloads one specific enclosure immediately, whatever the per-scan cap says and
/// wherever it sits in the queue. /// wherever it sits in the queue.
#[tracing::instrument(skip(ctx))]
async fn download_one(ctx: &Arc<Ctx>, id: i64) -> Result<()> { async fn download_one(ctx: &Arc<Ctx>, id: i64) -> Result<()> {
let cfg = ctx.cfg(); let cfg = ctx.cfg();
let enc = ctx let enc = ctx
@@ -1696,6 +1735,7 @@ async fn download_one(ctx: &Arc<Ctx>, id: i64) -> Result<()> {
/// Torrent progress is reported the same way an HTTP download's is, throttled to whole /// Torrent progress is reported the same way an HTTP download's is, throttled to whole
/// percents so a UI is not flooded. /// percents so a UI is not flooded.
#[tracing::instrument(name = "torrent", skip_all, fields(feed = feed_id, enclosure = enclosure))]
async fn torrent_one( async fn torrent_one(
ctx: &Arc<Ctx>, ctx: &Arc<Ctx>,
feed_id: &str, feed_id: &str,

View File

@@ -1938,8 +1938,23 @@ async fn access_log(req: Request, next: Next) -> Response {
let path = req.uri().path().to_owned(); let path = req.uri().path().to_owned();
let method = req.method().clone(); let method = req.method().clone();
let quiet = path.starts_with("/api/logs"); 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 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 { if !quiet {
let ms = started.elapsed().as_millis(); let ms = started.elapsed().as_millis();
let status = resp.status().as_u16(); let status = resp.status().as_u16();