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

@@ -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<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`
@@ -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<Ctx>, 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<String>, removed: usize, kept: usize, total: usize },
}
#[tracing::instrument(name = "feed", skip_all, fields(feed = id))]
async fn scan_one(
ctx: &Arc<Ctx>,
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<Ctx>,
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
/// wherever it sits in the queue.
#[tracing::instrument(skip(ctx))]
async fn download_one(ctx: &Arc<Ctx>, id: i64) -> Result<()> {
let cfg = ctx.cfg();
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
/// 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<Ctx>,
feed_id: &str,