From 934de045871f6c406d743a66f5db575a4353468a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?B=C5=82az=CC=87ej=20Pankowski?= <86720177+pblazej@users.noreply.github.com> Date: Thu, 1 Oct 2026 12:34:34 +0200 Subject: [PATCH 1/2] feat(telemetry): add the pipeline, sessions, spans and exporter `Telemetry` is the synchronous handle (emit, log, device state, flush, shutdown, purge, stats) behind a flood guard; `Scope` is one Room on it, with its own trace id, server URL and token, attributes, RTC stats and the `lk.subscribe` lifecycle; `Span` is one typed attempt (`SpanName`, `SpanStep`). The `Exporter` actor drains the queue and finished spans, encodes and gzips one batch per project into the cache before any network is involved, then uploads oldest first: it acts on every collector answer, pauses a destination with jittered backoff or for the delay the server asked for, splits a 413, holds uploads while a Room connects or the device asks for quiet, and meters requests while a Room is in a call. --- Cargo.lock | 1 + livekit-telemetry/Cargo.toml | 2 + livekit-telemetry/README.md | 81 ++ livekit-telemetry/src/exporter.rs | 1368 ++++++++++++++++++++++++++++ livekit-telemetry/src/lib.rs | 16 +- livekit-telemetry/src/scope.rs | 403 +++++++- livekit-telemetry/src/telemetry.rs | 892 ++++++++++++++++++ livekit-telemetry/src/trace.rs | 525 +++++++++++ 8 files changed, 3285 insertions(+), 3 deletions(-) create mode 100644 livekit-telemetry/README.md create mode 100644 livekit-telemetry/src/exporter.rs create mode 100644 livekit-telemetry/src/telemetry.rs create mode 100644 livekit-telemetry/src/trace.rs diff --git a/Cargo.lock b/Cargo.lock index 53f94fe5b..811a82f7d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3818,6 +3818,7 @@ version = "0.1.0" dependencies = [ "async-trait", "base64 0.22.1", + "flate2", "livekit-net", "log", "opentelemetry-proto", diff --git a/livekit-telemetry/Cargo.toml b/livekit-telemetry/Cargo.toml index ad1fe9df6..0943e1689 100644 --- a/livekit-telemetry/Cargo.toml +++ b/livekit-telemetry/Cargo.toml @@ -21,6 +21,8 @@ base64 = "0.22" serde_json = { workspace = true } # Server URLs are parsed, never string-matched, before a token may follow them. url = "2.3" +# gzip for request bodies; the deflate backend (miniz_oxide) is already linked by data streams. +flate2 = "1" # OTLP message types only (`gen-tonic-messages` = prost structs, no tonic). Same prost as # livekit-protocol so a single prost is linked. opentelemetry-proto = { version = "0.32", default-features = false, features = ["logs", "trace", "gen-tonic-messages"] } diff --git a/livekit-telemetry/README.md b/livekit-telemetry/README.md new file mode 100644 index 000000000..1809183d9 --- /dev/null +++ b/livekit-telemetry/README.md @@ -0,0 +1,81 @@ +# LiveKit Telemetry + +**Important**: +This is an internal crate that powers client telemetry in LiveKit client SDKs (through +`livekit-uniffi`: Swift, Kotlin, Dart) and is not meant to be used directly. + +The core buffers records on the device and ships them out-of-band as standard +[OTLP/HTTP](https://opentelemetry.io/docs/specs/otlp/) logs and traces to the LiveKit Cloud +project a room belongs to. Everything hard lives here once — destination and credentials, +batching, encoding, retry, persistence, holds, loss accounting — while platforms provide only +what they are uniquely placed to: OS signals (instruments) and the byte-moving transport. + +```text + SDK / instruments ──emit()──▶ Telemetry ──▶ Store ──drain──▶ Exporter ──push──▶ BatchCache ──upload oldest-first──▶ TelemetryTransport + (source) (records) (tick · encode) (MemoryCache | FileCache) (NetTransport | host HTTP) +``` + +| Layer | OTel equivalent | Type | +|---|---|---| +| source | `Logger.emit`, `Tracer.start` | [`Telemetry`], [`Scope`] (one per room), [`Span`] | +| store | `BatchLogRecordProcessor` queue | `Store` (in-memory, drop-oldest, bounded) | +| exporter | batch processor timer + exporter | [`Exporter`] (actor, spawn `run()`) | +| cache | disk buffering | [`BatchCache`]: [`MemoryCache`], [`FileCache`] (`storage_dir`) | +| transport | exporter's HTTP client | [`TelemetryTransport`]; `NetTransport` over `livekit-net` (feature `net`) | + +## Usage + +```rust +# use std::sync::Arc; +# use livekit_telemetry::*; +# struct Discard; +# #[async_trait::async_trait] +# impl TelemetryTransport for Discard { +# async fn send(&self, _: ExportRequest) -> Result { +# Ok(ExportResponse::accepted()) +# } +# } +# #[tokio::main(flavor = "current_thread")] async fn main() { +let config = TelemetryConfig { storage_dir: None, ..Default::default() }; +let (telemetry, exporter) = Telemetry::new(config, Arc::new(Discard)); +tokio::spawn(exporter.run()); + +// One scope per room. Its server URL and token decide where its records go; hand the token +// over again on every refresh. +let room = telemetry.begin_scope(); +room.set_server("wss://my-project.livekit.cloud", ""); +room.set_attribute("app.call_id", Some("c-42".into())); +room.emit_custom("checkout", vec![Attribute::new("plan", "pro")]); + +let connect = room.start(SpanName::Connect, None); +connect.step(SpanStep::WsOpen); +connect.end(SpanOutcome::Ok, None); + +telemetry.set_device_state(DeviceState { thermal: ThermalState::Serious, ..Default::default() }); +telemetry.shutdown().await; // cache, then upload what the network allows +println!("{}", telemetry.stats()); // drops by reason, uploads, cached batches +# } +``` + +Event names, attributes, cadences and the upload policy are defined in [`SPEC.md`](SPEC.md). + +## Design notes + +- **The room decides the destination.** A platform passes only the server URL and the token + ([`Scope::set_server`]). The core derives the ingest URL (LiveKit Cloud hosts only), reads + the token's unverified claims (observability grant, expiry) and routes each batch by its + owner — project and session — so two rooms on two projects never share a token. +- **Write-ahead cache.** Every batch is encoded, gzipped and cached *before* the network is + tried, then uploaded oldest-first and removed on an answer. Tokens are attached at upload + time and never persisted. +- **The core reads every answer.** Transports return status, headers and body; the core alone + classifies OTLP partial success, 4xx, 413 (split), 429/503 `Retry-After` and `RetryInfo`, + and Cloud's disabled and credential answers. +- **Loss is counted, never silent.** Every way a record can be lost has a counter; the deltas + ride along with the next batch as `lk.telemetry.report` and are readable locally through + [`Telemetry::stats`]. +- **Device state comes from the host; the policy lives here.** Thermal, power, memory, + network and lifecycle are OS APIs the host watches and pushes through + [`Telemetry::set_device_state`]; the core stretches its cadence and holds uploads. +- **Size.** OTLP types come from `opentelemetry-proto` (`gen-tonic-messages`, no tonic); JWT + claims are read with `base64` and `serde_json`, both already linked by the UniFFI library. diff --git a/livekit-telemetry/src/exporter.rs b/livekit-telemetry/src/exporter.rs new file mode 100644 index 000000000..b7f6992b4 --- /dev/null +++ b/livekit-telemetry/src/exporter.rs @@ -0,0 +1,1368 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::{ + collections::{HashMap, VecDeque}, + future::Future, + io::Write, + pin::Pin, + sync::{atomic::AtomicU64, Arc}, + time::Duration, +}; + +use flate2::{write::GzEncoder, Compression}; +use tokio::sync::{mpsc, oneshot}; +use tokio::time::{sleep_until, timeout, Instant}; + +use crate::{ + destination::{Dead, Route, Signal, Target, PROCESS_OWNER}, + event::now_unix_nanos, + otlp, + scope::ScopeState, + stats::{Counters, Snapshot, TelemetryStatus}, + store::Queued, + telemetry::Shared, + transport::Verdict, + AppState, DeviceState, ExportError, ExportRequest, MemoryPressure, NetworkType, + TelemetryTransport, ThermalState, +}; + +/// Who a batch belongs to: its session's project (`None` until the session has a server) and +/// the session itself. +#[derive(Debug, PartialEq)] +struct Owner { + host: Option, + session: String, +} + +/// Local retry backoff: 1 s doubling per consecutive failure up to 60 s, with full jitter +/// (a uniform wait in `[0, backoff]`) so a fleet that failed together does not retry together. +/// Retrying never gives up on a batch: the cache's age and size bound what is kept. +const RETRY_BASE: Duration = Duration::from_secs(1); +const RETRY_CAP: Duration = Duration::from_secs(60); +/// Delays the server asks for (`Retry-After`, `RetryInfo`) are honored even beyond +/// [`RETRY_CAP`]; clamped only to the cache's age limit, past which there is nothing to wait for. +const MAX_SERVER_DELAY: Duration = Duration::from_secs(24 * 60 * 60); +/// A cached batch older than this is dropped (`expired`) — at start and while running. +const MAX_AGE: Duration = Duration::from_secs(24 * 60 * 60); +/// Soft holds (see [`Exporter::hold_reason`]) last at most this long before one batch goes out +/// anyway: the cap that bounds the policy when its signals lie. Hard holds have no such escape. +const MAX_HOLD: Duration = Duration::from_secs(60); +/// While one of these is open the uplink belongs to signaling and ICE/DTLS. +const SENSITIVE_SPANS: &[&str] = &["lk.connect", "lk.reconnect"]; +/// RFC 9218 request priority: lowest urgency, not incremental. A hint for HTTP/2+ hops that +/// implement it; the host transport marks the local traffic class (see `TelemetryTransport`). +const PRIORITY: &str = "u=7"; + +impl Signal { + fn tag(self) -> char { + match self { + Signal::Logs => 'l', + Signal::Traces => 't', + } + } +} + +/// A cached batch's id, which is also its ownership record: +/// `-----`. Sorting ids sorts batches oldest +/// first; the host and the session that produced it pick destination and token at upload time +/// (no host: process-level, sent to the latest project). No token is ever written to disk. +struct BatchId<'a> { + /// Creation time, unix ns. + stamp: u64, + count: u64, + signal: Signal, + owner: Option<&'a str>, + host: Option<&'a str>, +} + +impl<'a> BatchId<'a> { + fn format(signal: Signal, seq: &str, count: u64, owner: &str, host: Option<&str>) -> String { + let (stamp, tag, host) = (now_unix_nanos(), signal.tag(), host.unwrap_or_default()); + format!("{stamp:020}-{seq}-{count}-{tag}-{owner}-{host}") + } + + fn parse(id: &'a str) -> Self { + let mut parts = id.splitn(6, '-'); + let stamp = parts.next().and_then(|n| n.parse().ok()).unwrap_or(0); + let count = parts.nth(1).and_then(|n| n.parse().ok()).unwrap_or(0); + let signal = match parts.next() { + Some("t") => Signal::Traces, + _ => Signal::Logs, + }; + let owner = parts.next().filter(|o| !o.is_empty()); + Self { stamp, count, signal, owner, host: parts.next().filter(|h| !h.is_empty()) } + } + + /// The id of one half of this batch after a 413: same stamp and owner, so it keeps the + /// original's place in line, `seq` extended with `a`/`b`. + fn half(id: &str, which: char, count: u64) -> String { + let mut parts: Vec<&str> = id.splitn(6, '-').collect(); + let seq = format!("{}{which}", parts.get(1).copied().unwrap_or_default()); + let count = count.to_string(); + if parts.len() >= 3 { + parts[1] = &seq; + parts[2] = &count; + } + parts.join("-") + } +} + +/// The records a cached batch id says it holds. +pub(crate) fn events_in(id: &str) -> u64 { + BatchId::parse(id).count +} + +pub(crate) enum Command { + Flush(oneshot::Sender<()>), + Shutdown(oneshot::Sender<()>), + /// Opt-out: stop without uploading and delete everything held. + Purge(oneshot::Sender<()>), +} + +/// What one request achieved. +enum Outcome { + Answered(Verdict), + /// No answer: offline, connection or TLS failure, timeout. + NoAnswer { + timed_out: bool, + reason: String, + }, +} + +/// The request on the wire, polled by the actor's loop so commands and deadlines keep being +/// served while it is out. +struct InFlight { + id: String, + signal: Signal, + /// Which destination's pause the answer drives. + key: String, + target: Target, + body: Vec, + request: Pin + Send>>, +} + +/// One pass over the cache: the batches still to look at and the requests it may make. +#[derive(Default)] +struct Round { + queue: VecDeque, + budget: usize, + attempts: usize, + /// A batch waited for a token or a destination. + waiting: bool, + /// A batch waited for its destination's pause. + paused: bool, + /// The one-batch escape of a soft hold: the hold does not stop it. + escape: bool, + /// Its requests are charged to the interval's allowance: next to a call, and not a drain, + /// flush, background pass or escape. + metered: bool, + /// The hold policy stopped the pass. + held: bool, +} + +/// A destination's pause after a failure: uploads to it wait, others carry on. +#[derive(Debug, Default)] +struct Pause { + /// Retry after this (backoff or server delay); `None` once it has passed. + until: Option, + /// The part the server asked for: honored in full, even while draining. + server_until: Option, + /// Consecutive failures: the backoff exponent. + failing: u32, +} + +/// Background actor that turns stored events into OTLP requests. +/// +/// The role of OTel's `BatchLogRecordProcessor` + OTLP exporter in one place. Every tick it +/// [`enqueue`](Self::enqueue)s: drains the queue and the finished spans, appends an +/// `lk.telemetry.report` when something was dropped or failed since the last one, encodes and +/// gzips one batch per destination project and writes it to the [`BatchCache`](crate::BatchCache) +/// *before* any network is involved; then it [`begin_round`](Self::begin_round)s the cache oldest-first +/// through the [`TelemetryTransport`], one request at a time. What the collector answers decides +/// what happens to the batch ([`Verdict`]); a failure pauses uploads with jittered exponential +/// backoff, a throttle for as long as the collector asked, and collection never stops meanwhile. +/// +/// Telemetry must never win over media, so uploads are shaped as well as batched: at most +/// `max_batches_per_upload` per tick while a session may be live, and none at all while the room +/// is connecting or reconnecting or while the device asks for quiet ([`DeviceState::holds_uploads`]) +/// — bounded by [`MAX_HOLD`]. Yielding to media on the wire is the transport's job: every request +/// carries `Priority: u=7` (RFC 9218) and a gzipped body. +/// +/// The tick period is `flush_interval_ms × cadence factor`: device pressure and a CPU-limited +/// encoder stretch it up to 4×, and entering the background flushes once immediately (the app +/// may be suspended any moment). +/// +/// Drive it with `spawn(exporter.run())` on the consumer's runtime. It stops after +/// [`Telemetry::shutdown`](crate::Telemetry::shutdown) or when the last +/// [`Telemetry`](crate::Telemetry) handle is dropped. +pub struct Exporter { + shared: Arc, + transport: Arc, + commands: mpsc::UnboundedReceiver, + /// Per-destination pauses after failures, by destination key (the project host). + pauses: HashMap, + /// The upload pass in progress, and its request on the wire. + round: Option, + inflight: Option, + /// Another pass is due when this one ends (new batches arrived meanwhile). + rerun: bool, + /// `flush()` callers, answered when the pass ends. + flush_waiters: Vec>, + /// The next pass comes from an explicit `flush()` or entering the background: it has no + /// per-pass budget. + flushing: bool, + /// Requests left in this flush interval (`max_batches_per_upload`, refilled every tick), + /// charged while a Room is in a call: wake-ups between ticks never add requests beyond it. + allowance: usize, + /// Batches done with whose delete failed: never sent again, deletion retried each pass. + undeletable: std::collections::HashSet, + /// When this launch wrote each cached batch, on the monotonic clock (see + /// [`Exporter::expired`]). + born: HashMap, + /// Policy transitions are logged once, not per tick: the hold reason last logged and the + /// cadence factor last logged. + hold_logged: Option<&'static str>, + cadence_logged: u32, + seq: u64, + /// Counter values at the last `lk.telemetry.report`. + last_report: Snapshot, + /// When the current upload hold began, for the [`MAX_HOLD`] cap. + held_since: Option, + /// Shutting down: the call is over, so only the device's own holds still apply. + draining: bool, + /// Append an `lk.telemetry.report` to the next batch even if nothing went wrong (the + /// shutdown summary). + force_report: bool, +} + +impl Exporter { + pub(crate) fn new( + shared: Arc, + transport: Arc, + commands: mpsc::UnboundedReceiver, + ) -> Self { + Self { + shared, + transport, + commands, + pauses: HashMap::new(), + round: None, + inflight: None, + rerun: false, + flush_waiters: Vec::new(), + undeletable: std::collections::HashSet::new(), + flushing: false, + allowance: 0, + born: HashMap::new(), + hold_logged: None, + cadence_logged: 1, + seq: 0, + last_report: Snapshot::default(), + held_since: None, + draining: false, + force_report: false, + } + } + + /// Run until shut down: replay whatever the cache holds, then export on every tick and on + /// demand. A request on the wire never blocks the loop: commands, ticks and deadlines are + /// served while it is out. + pub async fn run(mut self) { + self.begin_round(); + // Deadlines rather than tickers: `next = last + period`, re-derived every loop, so a + // cadence change applies to the pending tick in both directions (pressure postpones it, + // relief brings it forward) and a missed tick never bursts. The first flush is immediate, + // the first window closes a full period after it opens. + let mut last_flush = Instant::now() - self.period(self.shared.config.flush_interval_ms); + let mut last_window = Instant::now(); + loop { + self.log_cadence(); + let next_flush = last_flush + self.period(self.shared.config.flush_interval_ms); + let next_window = last_window + self.period(self.shared.config.stats_window_ms); + let retry = self.retry_deadline(); + let subscribe = self.subscribe_deadline(); + tokio::select! { + outcome = answer(&mut self.inflight) => self.on_answer(outcome), + _ = sleep_until(next_flush) => { + // A window due at the same moment closes first, so it rides this batch. + if Instant::now() >= next_window { + self.close_windows(); + last_window = Instant::now(); + } + // A new interval: a new allowance of requests. + self.allowance = self.shared.config.max_batches_per_upload.max(1) as usize; + self.export_pending(); + last_flush = Instant::now(); + } + _ = sleep_until(next_window) => { + self.close_windows(); + last_window = Instant::now(); + } + // A pause ran out or a soft hold reached its cap: go without waiting for a tick. + _ = at(retry) => { + self.clear_expired_pauses(); + self.begin_round(); + } + _ = at(subscribe) => self.expire_subscribes(), + // Coalesced wake-ups. Only the queue crossing its threshold or the app going to + // the background encode a new batch; a lifted route, token or hold re-runs a pass + // over what is cached (within the interval's allowance); a subscribe or a cadence + // change only makes the loop re-read its deadlines. So the cadence stays one + // export per interval. + _ = self.shared.wake.notified() => { + let taken = |flag: &std::sync::atomic::AtomicBool| { + flag.swap(false, std::sync::atomic::Ordering::SeqCst) + }; + let background = taken(&self.shared.backgrounded); + let overflow = taken(&self.shared.overflowed); + let release = taken(&self.shared.released); + if background { + // The app may be suspended any moment: close the RTC windows and get + // everything into the cache and out — the whole cache, like a flush, + // whatever is left of the interval's allowance. + self.close_windows(); + self.flushing = true; + } + if background || overflow { + self.export_pending(); + } else if release { + self.begin_round(); + } + } + command = self.commands.recv() => match command { + Some(Command::Flush(done)) => { + // An explicit flush drains the cache: no per-pass budget, but every hold + // and pause still applies. + self.flushing = true; + self.export_pending(); + self.answer_when_done(done); + } + + Some(Command::Shutdown(done)) => { + self.drain().await; + let _ = done.send(()); + return; + } + Some(Command::Purge(done)) => { + // Cancel the request on the wire: its answer must not bring anything back. + self.inflight = None; + self.round = None; + self.shared.clear(); + self.answer_flushes(); + let _ = done.send(()); + return; + } + // Every `Telemetry` handle is gone. + None => { + self.drain().await; + return; + } + }, + } + } + } + + /// `base_ms × cadence factor`. + fn period(&self, base_ms: u64) -> Duration { + Duration::from_millis(base_ms.max(1)) * self.cadence_factor() + } + + fn device(&self) -> DeviceState { + self.shared.device.lock().unwrap_or_else(|e| e.into_inner()).unwrap_or_default() + } + + fn cadence_factor(&self) -> u32 { + self.shared.cadence_factor() + } + + fn cpu_limited(&self) -> bool { + self.shared.windows.lock().unwrap_or_else(|e| e.into_inner()).cpu_limited() + } + + /// Whether a Room is in a call: it has a server and has not disconnected. + fn in_call(&self) -> bool { + let scopes = self.shared.scopes.lock().unwrap_or_else(|e| e.into_inner()); + scopes.iter().any(|s| s.upgrade().is_some_and(|s| s.in_call())) + } + + /// Every live session's earliest subscribe deadline. + fn subscribe_deadline(&self) -> Option { + let scopes = self.shared.scopes.lock().unwrap_or_else(|e| e.into_inner()); + scopes.iter().filter_map(|s| s.upgrade()?.subscribe_deadline()).min() + } + + fn expire_subscribes(&self) { + let scopes: Vec> = { + let scopes = self.shared.scopes.lock().unwrap_or_else(|e| e.into_inner()); + scopes.iter().filter_map(|s| s.upgrade()).collect() + }; + for scope in scopes { + scope.expire_subscribes(); + } + } + + /// One debug line per cadence change, naming what stretched it. + fn log_cadence(&mut self) { + let factor = self.cadence_factor(); + if factor == self.cadence_logged { + return; + } + self.cadence_logged = factor; + let device = self.device(); + let mut why = Vec::new(); + if !matches!(device.thermal, ThermalState::Unknown | ThermalState::Nominal) { + why.push(format!("thermal {:?}", device.thermal).to_lowercase()); + } + if device.memory != MemoryPressure::Normal { + why.push(format!("memory {:?}", device.memory).to_lowercase()); + } + if device.low_power_mode == Some(true) { + why.push("low power mode".to_string()); + } + if device.app_state == AppState::Background { + why.push("background".to_string()); + } + if self.cpu_limited() { + why.push("encoder cpu-limited".to_string()); + } + log::debug!( + "cadence ×{factor}, flush every {}s ({})", + self.period(self.shared.config.flush_interval_ms).as_secs(), + if why.is_empty() { "pressure over".to_string() } else { why.join(", ") } + ); + } + + /// Last chance: everything queued into the cache and out, within `export_timeout_ms`. + /// Ignores the failure backoff, the batch budget and the session holds (the call is over); + /// respects server-directed delays and the device's own holds. At the deadline the request + /// on the wire is cancelled and the actor stops: its batch stays cached for the next launch. + async fn drain(&mut self) { + let deadline = + Instant::now() + Duration::from_millis(self.shared.config.export_timeout_ms.max(1)); + self.draining = true; + self.force_report = true; + // Pending subscribes end now and ride the last batch; nothing keeps a span afterwards. + self.shared.end_pending_subscribes(true); + for pause in self.pauses.values_mut() { + pause.until = pause.server_until.filter(|t| Instant::now() < *t); + } + self.close_windows(); + self.export_pending(); + let mut commands_open = true; + while self.round.is_some() || self.inflight.is_some() { + tokio::select! { + outcome = answer(&mut self.inflight) => self.on_answer(outcome), + _ = sleep_until(deadline) => { + log::debug!("shutdown deadline: cancelling the upload in flight"); + self.inflight = None; + self.round = None; + break; + } + // The opt-out reaches a draining generation too: cancel, clear, confirm. + command = self.commands.recv(), if commands_open => match command { + Some(Command::Purge(done)) => { + self.inflight = None; + self.round = None; + self.shared.clear(); + let _ = done.send(()); + break; + } + Some(Command::Flush(done)) => self.answer_when_done(done), + Some(Command::Shutdown(done)) => { + let _ = done.send(()); + } + None => commands_open = false, + }, + } + // A pass that ended with a request left to make (a split, the rerun) goes on. + if self.inflight.is_none() && self.round.is_some() { + self.advance(); + } + } + self.answer_flushes(); + let left = self.shared.cache.pending().len(); + if left > 0 { + log::debug!("{left} batches still cached at shutdown (replayed next start)"); + } + } + + /// Turn every open RTC stats window into its `lk.rtc.stats.sample` event. Windows bypass the + /// flood guard: they are the pipeline's own, bounded output. + fn close_windows(&mut self) { + let events = self.shared.windows.lock().unwrap_or_else(|e| e.into_inner()).close(); + for event in events { + self.shared.store.push(event); + } + } + + fn export_pending(&mut self) { + self.enqueue(); + self.begin_round(); + } + + /// Split records by owner — the project their session routes to, and the session — so each + /// batch has exactly one destination and one token. Records whose project receives nothing + /// are dropped here, counted as `disabled`. + fn by_owner( + &self, + items: Vec, + owner_of: impl Fn(&T) -> (&ScopeState, Option<&String>), + ) -> Vec<(Owner, Vec)> { + let destinations = self.destinations(); + let mut groups: Vec<(Owner, Vec)> = Vec::new(); + for item in items { + let (session, route) = owner_of(&item); + let process = session.trace_id == self.shared.process.trace_id; + let session = if process { PROCESS_OWNER.to_owned() } else { session.hex() }; + // The project captured with the record, never the session's current one — or, for a + // Room record captured before its first server, that Room's first project, so the + // cached batch names its project and replays after a restart. + let host = match route { + Some(route) => Some(route.clone()), + None if !process => destinations.first_project(&session), + None => None, + }; + let owner = Owner { host, session }; + if !destinations.alive(owner.host.as_deref()) { + Counters::add(&self.shared.counters.disabled, 1); + continue; + } + match groups.iter_mut().find(|(o, _)| *o == owner) { + Some((_, group)) => group.push(item), + None => groups.push((owner, vec![item])), + } + } + groups + } + + /// Encode everything queued — log records and finished spans — into the cache. No network. + fn enqueue(&mut self) { + self.enqueue_spans(); + let config = &self.shared.config; + let (max, max_bytes) = ( + config.max_batch_size.max(1) as usize, + usize::try_from(config.max_batch_bytes.max(1)).unwrap_or(usize::MAX), + ); + // One report per pass at most: a report that is itself dropped (oversized) or evicts + // another batch must not make the next report due in the same pass, or with + // `max_batch_size` 1 it would take the only place forever and no record would move. + let mut reported = false; + loop { + // Self-telemetry rides along with real data: never its own request, never its own + // cadence, and only when there is something to report — plus once at shutdown. It + // takes one of the batch's `max_batch_size` places. + let now = self.shared.counters.snapshot(); + let delta = now.since(&self.last_report); + let report_due = !reported && (delta.has_problems() || self.force_report); + let room = if report_due { max - 1 } else { max }; + let mut batch = self.shared.store.drain(room, max_bytes); + // Real data to ride along with (with `max_batch_size` 1 the report goes first, alone); + // the shutdown summary goes out even with nothing else queued. + let data = !batch.is_empty() || (room == 0 && !self.shared.store.is_empty()); + if !data && !self.force_report { + return; + } + if report_due { + let cached = self.shared.cache.pending().len() as u64; + // The host's console gets the cumulative line through the FFI log path. + log::debug!("{}", crate::stats::TelemetryStats::new(now, cached, self.status())); + let report = delta.report(cached); + batch.push(Queued::new(report, self.shared.process.clone())); + self.last_report = now; + self.force_report = false; + reported = true; + } + let global = self.shared.global.lock().unwrap_or_else(|e| e.into_inner()).clone(); + for (owner, records) in self.by_owner(batch, |q| (&q.session, q.route.as_ref())) { + let count = records.len() as u64; + let body = otlp::encode_logs(&self.shared.config.resource, &global, records); + self.push_batch(Signal::Logs, &owner, count, &body); + } + } + } + + /// Finished spans travel as their own batches on the traces signal: every one of them, so a + /// burst larger than a batch (or a shutdown) leaves nothing behind in memory. + fn enqueue_spans(&mut self) { + loop { + let (spans, dropped) = { + let mut registry = self.shared.spans.lock().unwrap_or_else(|e| e.into_inner()); + ( + registry.drain( + self.shared.config.max_batch_size.max(1) as usize, + usize::try_from(self.shared.config.max_batch_bytes.max(1)) + .unwrap_or(usize::MAX), + ), + registry.take_dropped(), + ) + }; + Counters::add(&self.shared.counters.queue_full, dropped); + if spans.is_empty() { + return; + } + let global = self.shared.global.lock().unwrap_or_else(|e| e.into_inner()).clone(); + for (owner, spans) in self.by_owner(spans, |s| (&s.session, s.route.as_ref())) { + let count = spans.len() as u64; + let body = otlp::encode_spans(&self.shared.config.resource, &global, spans); + self.push_batch(Signal::Traces, &owner, count, &body); + } + } + } + + /// Gzip the encoded batch and cache it. Compressed at rest as well as on the wire: the cache + /// holds 5–10× more, the disk write shrinks, and a replay costs no CPU. + /// + /// The record-size estimates that cut batches cannot see what encoding adds (session and + /// pipeline attributes, the self-report), so the real encoded size is checked here: a batch + /// over `max_batch_bytes` is halved until it fits, or is a single record. + fn push_batch(&mut self, signal: Signal, owner: &Owner, count: u64, body: &[u8]) { + let limit = + usize::try_from(self.shared.config.max_batch_bytes.max(1)).unwrap_or(usize::MAX); + if body.len() > limit { + let halves = match signal { + Signal::Logs => otlp::split_logs(body), + Signal::Traces => otlp::split_spans(body), + }; + if let Some([(first, first_count), (second, second_count)]) = halves { + self.push_batch(signal, owner, first_count, &first); + self.push_batch(signal, owner, second_count, &second); + } else { + // One record larger than a request may be: never cached, never sent. + log::warn!("a record of {} bytes exceeds max_batch_bytes; dropped", body.len()); + Counters::add(&self.shared.counters.oversized, count); + } + return; + } + let body = gzip(body); + self.seq += 1; + let seq = format!("{:06}", self.seq); + let id = BatchId::format(signal, &seq, count, &owner.session, owner.host.as_deref()); + self.store_batch(id, &body); + } + + /// Write one gzipped batch to the cache — unless the app opted out meanwhile. + fn store_batch(&mut self, id: String, body: &[u8]) { + if self.shared.revoked() { + return; + } + self.born.insert(id.clone(), Instant::now()); + match self.shared.cache.push(&id, body) { + Ok(evicted) => self.count_evicted(&evicted), + Err(err) => { + log::debug!("could not cache {id}: {err}"); + Counters::add(&self.shared.counters.cache_error, events_in(&id)); + self.born.remove(&id); + } + } + } + + /// Older batches pushed out by the cache's bounds: lost, but counted. Losing them mid-hold is + /// its own answer — the collector held us off longer than the cache could carry — so it is + /// counted apart from an ordinary overflow. + fn count_evicted(&mut self, evicted: &[String]) { + let now = Instant::now(); + let counters = self.shared.counters.clone(); + for id in evicted { + // `throttled` only when the evicted batch's own destination is under a server pause — + // for a host-less batch (process-level, pre-connect), the project it routes to. + let batch = BatchId::parse(id); + let host = match self.destinations().route(batch.host, batch.owner) { + Route::Send(target) => target.project.unwrap_or_default(), + _ => batch.host.unwrap_or_default().to_owned(), + }; + let held = + self.pauses.get(&host).is_some_and(|p| p.server_until.is_some_and(|t| now < t)); + let counter = if held { &counters.throttled } else { &counters.cache_full }; + Counters::add(counter, events_in(id)); + self.born.remove(id); + } + } + + /// Why uploads should wait right now, if they should. Data keeps flowing into the cache + /// meanwhile — write-ahead caching is what makes holding free. + fn hold_reason(&self) -> Option<&'static str> { + if self.device().holds_uploads() { + return Some("device asks for quiet"); + } + if self.draining { + return None; + } + if self.shared.spans.lock().unwrap_or_else(|e| e.into_inner()).any_open(SENSITIVE_SPANS) { + return Some("connecting"); + } + None + } + + fn status(&self) -> TelemetryStatus { + *self.shared.status.lock().unwrap_or_else(|e| e.into_inner()) + } + + /// Record the policy state; true when it changed (log once, not per tick). + fn set_status(&self, status: TelemetryStatus) -> bool { + let mut current = self.shared.status.lock().unwrap_or_else(|e| e.into_inner()); + let changed = *current != status; + *current = status; + changed + } + + /// Log a hold transition once. + fn log_hold(&mut self, reason: Option<&'static str>, backlog: usize) { + match (self.hold_logged, reason) { + (None, Some(now)) => log::debug!("held ({now}), backlog {backlog}"), + (Some(was), Some(now)) if was != now => { + log::debug!("held ({now}, was {was}), backlog {backlog}") + } + (Some(was), None) => match self.held_since { + Some(since) => log::debug!( + "resumed after {}s ({was}), backlog {backlog}", + since.elapsed().as_secs() + ), + None => log::debug!("resumed ({was}), backlog {backlog}"), + }, + _ => {} + } + self.hold_logged = reason; + } + + /// How many batches this call may send: none while held — a hard hold (offline) for as long + /// as it lasts, a soft one up to [`MAX_HOLD`], then one — otherwise what is left of the + /// interval's allowance when `metered`, everything when not. + fn budget(&mut self, backlog: usize, metered: bool) -> (usize, bool) { + if self.offline() { + self.log_hold(Some("offline"), backlog); + return (0, false); + } + let reason = self.hold_reason(); + self.log_hold(reason, backlog); + let Some(reason) = reason else { + self.held_since = None; + return (if metered { self.allowance } else { usize::MAX }, false); + }; + let since = *self.held_since.get_or_insert_with(Instant::now); + if since.elapsed() < MAX_HOLD { + log::trace!("holding {backlog} batches: {reason}"); + return (0, false); + } + log::warn!("held {}s ({reason}): sending one batch anyway", MAX_HOLD.as_secs()); + // Held long enough: one batch goes out, then the hold starts over. + self.held_since = Some(Instant::now()); + Counters::add(&self.shared.counters.hold_cap_hits, 1); + (1, true) + } + + /// A hard hold: nothing is attempted while the device is offline. + fn offline(&self) -> bool { + self.device().network == NetworkType::Unavailable + } + + /// Whether the pass may make its next request: the hold policy is evaluated before every + /// request, not once per pass. A hard hold stops the pass; a soft hold that began meanwhile + /// stops it too — unless this pass is the hold's one-batch escape (or a shutdown drain). + fn may_send(&mut self) -> bool { + // The opt-out stops a pass mid-way (a second guard: `advance` checks it as well). + if self.offline() || self.shared.revoked() { + return false; + } + if self.draining || self.round.as_ref().is_some_and(|r| r.escape) { + return true; + } + if self.hold_reason().is_some() { + self.held_since.get_or_insert_with(Instant::now); + return false; + } + true + } + + /// Whether a cached batch is past [`MAX_AGE`]. Wall-clock age decides for batches from a + /// previous launch; for this launch's own the monotonic clock must agree, so a clock set + /// forward cannot expire a fresh backlog (and one set back only delays expiry). + fn expired(&self, id: &str) -> bool { + let wall = Duration::from_nanos(now_unix_nanos().saturating_sub(BatchId::parse(id).stamp)); + wall > MAX_AGE && self.born.get(id).is_none_or(|born| born.elapsed() > MAX_AGE) + } + + /// When the loop must wake to upload without a tick: the earliest pause to run out, or a + /// soft hold reaching its cap while batches wait. + fn retry_deadline(&self) -> Option { + let pause = self.pauses.values().filter_map(|p| p.until).min(); + // A hard hold suppresses the soft hold's cap: waking for it could send nothing. It comes + // back when the hard hold lifts (a device change wakes the loop). + let hold = self + .held_since + .filter(|_| self.round.is_none() && !self.offline()) + .filter(|_| !self.shared.cache.pending().is_empty()) + .map(|since| since + MAX_HOLD); + pause.into_iter().chain(hold).min() + } + + fn clear_expired_pauses(&mut self) { + let now = Instant::now(); + for pause in self.pauses.values_mut() { + pause.until = pause.until.filter(|t| now < *t); + pause.server_until = pause.server_until.filter(|t| now < *t); + } + } + + /// Answer a `flush()` once the pass it started is over. + fn answer_when_done(&mut self, done: oneshot::Sender<()>) { + if self.round.is_some() || self.inflight.is_some() { + // Callers that stopped waiting are not kept. + self.flush_waiters.retain(|waiter| !waiter.is_closed()); + self.flush_waiters.push(done); + } else { + let _ = done.send(()); + } + } + + fn answer_flushes(&mut self) { + for done in self.flush_waiters.drain(..) { + let _ = done.send(()); + } + } + + /// Start a pass over the cache: oldest first, within this tick's budget, one request at a + /// time. A pass already running picks the new batches up in a second pass. + fn begin_round(&mut self) { + #[cfg(test)] + self.shared.passes.fetch_add(1, std::sync::atomic::Ordering::Relaxed); + if self.shared.revoked() { + return; + } + if self.round.is_some() || self.inflight.is_some() { + self.rerun = true; + return; + } + // An explicit flush applies to this pass only, whatever it finds. + let flushing = std::mem::take(&mut self.flushing); + let counters = self.shared.counters.clone(); + self.retry_deletes(); + self.bind_pre_connect(); + for id in self.shared.cache.pending() { + if self.expired(&id) { + self.forget(&id, Some(&counters.expired)); + } + } + self.retire_credentials(); + if self.nobody_listens() { + self.purge_dead(); + self.set_status(TelemetryStatus::Off); + return self.answer_flushes(); + } + let pending = self.shared.cache.pending(); + if pending.is_empty() { + self.held_since = None; + self.set_status(TelemetryStatus::Ok); + return self.answer_flushes(); + } + // Only next to a call is a pass charged to the interval's allowance. + let metered = !(self.draining || flushing) && self.in_call(); + let (budget, escape) = self.budget(pending.len(), metered); + if budget == 0 { + if self.held_since.is_some() || self.offline() { + self.set_status(TelemetryStatus::Held); + } + return self.answer_flushes(); + } + let metered = metered && !escape; + self.round = + Some(Round { queue: pending.into(), budget, escape, metered, ..Round::default() }); + self.advance(); + } + + /// Take the pass to its next request, or to its end. + fn advance(&mut self) { + if self.inflight.is_some() || self.shared.revoked() { + return; + } + let counters = self.shared.counters.clone(); + loop { + let next = match self.round.as_mut() { + None => return, + Some(round) if round.attempts >= round.budget => None, + Some(round) => round.queue.pop_front(), + }; + let Some(id) = next else { return self.finish_round() }; + if self.undeletable.contains(&id) { + continue; + } + if !self.may_send() { + self.mark(|round| round.held = true); + return self.finish_round(); + } + let (route, signal) = { + let batch = BatchId::parse(&id); + (self.destinations().route(batch.host, batch.owner), batch.signal) + }; + let target = match route { + Route::Send(target) => target, + // A hard hold: no token that may be used (missing, expired, grant-less, refused). + Route::Wait => { + self.mark(|round| round.waiting = true); + continue; + } + Route::Drop => { + self.forget(&id, Some(&counters.disabled)); + continue; + } + }; + // The project the request goes to — and nothing else — is what its answer is about. + let key = target.project.clone().unwrap_or_default(); + if self.pauses.get(&key).is_some_and(|p| p.until.is_some_and(|t| Instant::now() < t)) { + self.mark(|round| round.paused = true); + continue; + } + let body = match self.shared.cache.read(&id) { + Ok(body) if intact(&body) => body, + Ok(_) => { + log::debug!("cached batch {id} is corrupt; dropped"); + self.forget(&id, Some(&counters.corrupt)); + continue; + } + Err(err) if err.kind() == std::io::ErrorKind::NotFound => { + self.born.remove(&id); + continue; + } + // There but unreadable right now (file protection on a locked device, a + // permission): not corrupt — it is left alone and tried again later. + Err(err) => { + log::debug!("cached batch {id} is inaccessible ({err}); trying later"); + self.mark(|round| round.waiting = true); + continue; + } + }; + self.mark(|round| round.attempts += 1); + if self.round.as_ref().is_some_and(|r| r.metered) { + self.allowance = self.allowance.saturating_sub(1); + } + let request = self.request(&body, signal, &target); + self.inflight = Some(InFlight { id, signal, key, target, body, request }); + return; + } + } + + fn mark(&mut self, f: impl FnOnce(&mut Round)) { + if let Some(round) = self.round.as_mut() { + f(round); + } + } + + /// The pass is over: report where uploads stand, then start the next pass if one is due. + fn finish_round(&mut self) { + let Some(round) = self.round.take() else { return }; + let now = Instant::now(); + let paused = |server: bool| { + self.pauses.values().any(|p| { + let until = if server { p.server_until } else { p.until }; + until.is_some_and(|t| now < t) + }) + }; + let status = if paused(true) { + TelemetryStatus::Throttled + } else if paused(false) { + TelemetryStatus::Paused + } else if round.held { + TelemetryStatus::Held + } else if round.waiting && round.attempts == 0 { + TelemetryStatus::Waiting + } else if self.held_since.is_some() { + TelemetryStatus::Held + } else { + TelemetryStatus::Ok + }; + if self.set_status(status) && status == TelemetryStatus::Waiting { + log::debug!("waiting for a destination or a usable token"); + } + if std::mem::take(&mut self.rerun) { + self.begin_round(); + } + if self.round.is_none() && self.inflight.is_none() { + self.answer_flushes(); + } + } + + /// What the collector answered to the request on the wire, acted on; then the pass goes on. + fn on_answer(&mut self, outcome: Outcome) { + let Some(InFlight { id, signal, key, target, body, .. }) = self.inflight.take() else { + return; + }; + // The app opted out while the request was out: whatever came back, nothing is written, + // retried or rescheduled. + if self.shared.revoked() { + self.round = None; + return; + } + let counters = self.shared.counters.clone(); + match outcome { + Outcome::Answered(Verdict::Accepted { rejected, reason }) => { + self.forget(&id, None); + Counters::add(&counters.uploads_sent, 1); + Counters::add(&counters.upload_bytes, body.len() as u64); + // The batch's other records were accepted; the refused ones are counted, never + // retried (OTLP partial success). + if rejected > 0 { + log::warn!("collector refused {rejected} records: {reason}"); + Counters::add(&counters.rejected, rejected); + } + self.recovered(&key); + } + Outcome::Answered(Verdict::TooLarge) => { + // Split in two; the halves go at once, down to single records — each one a + // request of its own against this pass's budget. + match self.split(&id, &body, signal) { + Ok(halves) if halves.is_empty() => { + log::warn!("a single record is larger than the collector accepts; dropped"); + self.forget(&id, Some(&counters.oversized)); + } + Ok(halves) => { + if let Some(round) = self.round.as_mut() { + for half in halves.into_iter().rev() { + round.queue.push_front(half); + } + } + } + Err(err) => { + // The batch stays whole where it was; try the split again later. + self.back_off(&key, None, &format!("cannot store the split batch: {err}")); + } + } + } + Outcome::Answered(Verdict::Rejected(reason)) => { + log::warn!("batch rejected by the collector, data lost: {reason}"); + self.forget(&id, Some(&counters.rejected)); + self.recovered(&key); + } + Outcome::Answered( + verdict @ (Verdict::Unauthorized(_) | Verdict::NotFound | Verdict::Disabled), + ) if self.has_override() => { + log::warn!("the collector at {} refused a batch: {verdict:?}", target.logs); + self.forget(&id, Some(&counters.rejected)); + } + Outcome::Answered(Verdict::Unauthorized(reason)) => { + // A credential problem, not a data problem: this token is not used again and the + // batch waits for the next one. + log::debug!("upload unauthorized: {reason}"); + Counters::add(&counters.auth_denied, 1); + if let Some(token) = &target.token { + self.destinations().refuse(token); + } + if let Some(round) = self.round.as_mut() { + round.waiting = true; + } + } + Outcome::Answered(Verdict::NotFound) => { + self.kill(target.project.as_deref(), Dead::NotFound) + } + Outcome::Answered(Verdict::Disabled) => { + self.kill(target.project.as_deref(), Dead::Disabled) + } + Outcome::Answered(Verdict::Throttled { delay_ms, reason }) => { + Counters::add(&counters.upload_failures, 1); + self.back_off(&key, Some(Duration::from_millis(delay_ms)), &reason); + } + Outcome::Answered(Verdict::Retry { delay_ms, reason }) => { + Counters::add(&counters.upload_failures, 1); + self.back_off(&key, delay_ms.map(Duration::from_millis), &reason); + } + Outcome::NoAnswer { timed_out, reason } => { + let counter = + if timed_out { &counters.upload_timeouts } else { &counters.upload_failures }; + Counters::add(counter, 1); + self.back_off(&key, None, &reason); + } + } + self.advance(); + } + + /// 413: re-encode the batch as two halves and replace it in the cache transactionally — + /// both halves committed where the batch was (on disk stays on disk) before it goes, nothing + /// evicted in between. `Ok` holds the halves' ids, oldest first, or none for a single record + /// (it cannot shrink); `Err` means the cache could not take the halves and the batch is + /// still there, whole. + fn split(&mut self, id: &str, body: &[u8], signal: Signal) -> std::io::Result> { + let Some(encoded) = gunzip(body) else { return Ok(Vec::new()) }; + let halves = match signal { + Signal::Logs => otlp::split_logs(&encoded), + Signal::Traces => otlp::split_spans(&encoded), + }; + let Some([(first, first_count), (second, second_count)]) = halves else { + return Ok(Vec::new()); + }; + let new = [ + (BatchId::half(id, 'a', first_count), gzip(&first)), + (BatchId::half(id, 'b', second_count), gzip(&second)), + ]; + let evicted = self.shared.cache.replace(id, &new)?; + self.count_evicted(&evicted); + self.born.remove(id); + let now = Instant::now(); + for (half, _) in &new { + self.born.insert(half.clone(), now); + } + log::debug!( + "413: split a batch of {} into {first_count} + {second_count}", + first_count + second_count + ); + Ok(new.into_iter().map(|(half, _)| half).collect()) + } + + /// A Room batch cached before the Room had a server names no project; once the Room's first + /// project is known it is rewritten under an id that names it (a journaled replace, as + /// crash-safe as a push), so it replays after a restart, when this launch's Room → project + /// map is gone. Recovery rolls an interrupted rewrite forward once its journal is durable + /// ([`FileCache`](crate::cache::FileCache)); a crash before that leaves the batch unbound. + fn bind_pre_connect(&mut self) { + for id in self.shared.cache.pending() { + let batch = BatchId::parse(&id); + let (None, Some(owner)) = (batch.host, batch.owner) else { continue }; + // Accepted already, its delete pending: a new id would be sent again. + if owner == PROCESS_OWNER || self.undeletable.contains(&id) { + continue; + } + let Some(project) = self.destinations().first_project(owner) else { continue }; + let Ok(body) = self.shared.cache.read(&id) else { continue }; + // A host-less id ends with the empty host: appending one is the bound id. + let bound = format!("{id}{project}"); + if let Ok(evicted) = self.shared.cache.replace(&id, &[(bound.clone(), body)]) { + if let Some(born) = self.born.remove(&id) { + self.born.insert(bound, born); + } + self.count_evicted(&evicted); + } + } + } + + /// Keep credentials only for Rooms still alive or with batches still cached. + fn retire_credentials(&self) { + let live: HashMap> = { + let scopes = self.shared.scopes.lock().unwrap_or_else(|e| e.into_inner()); + scopes.iter().filter_map(|s| s.upgrade()).map(|s| (s.hex(), s.route())).collect() + }; + let backlog = self + .shared + .cache + .pending() + .iter() + .filter_map(|id| { + let batch = BatchId::parse(id); + Some((batch.host.map(str::to_owned), batch.owner?.to_owned())) + }) + .collect(); + self.destinations().retain_owners(&live, &backlog); + } + + fn destinations(&self) -> std::sync::MutexGuard<'_, crate::destination::Destinations> { + self.shared.destinations.lock().unwrap_or_else(|e| e.into_inner()) + } + + fn has_override(&self) -> bool { + self.destinations().has_override() + } + + /// Every project the process talked to receives nothing: telemetry is off. + fn nobody_listens(&self) -> bool { + matches!(self.destinations().route(None, None), Route::Drop) + } + + /// The project on `host` receives nothing more: its cached batches go now. + fn kill(&mut self, host: Option<&str>, why: Dead) { + if let Some(host) = host { + self.destinations().kill(host, why); + } + self.purge_dead(); + } + + /// Delete every cached batch whose project receives nothing, counted as `disabled`. + fn purge_dead(&mut self) { + let counters = self.shared.counters.clone(); + for id in self.shared.cache.pending() { + let batch = BatchId::parse(&id); + if self.destinations().route(batch.host, batch.owner) == Route::Drop { + self.forget(&id, Some(&counters.disabled)); + } + } + } + + /// Remove a cached batch; when it was lost, count its records against that reason. A + /// delete that fails leaves the batch pending deletion: never uploaded again, deleted as soon + /// as the storage lets us (a restart before then may send it once more). + fn forget(&mut self, id: &str, lost: Option<&AtomicU64>) { + if let Err(err) = self.shared.cache.remove(id) { + log::warn!("cannot delete cached batch {id} ({err}); it will not be sent again"); + self.undeletable.insert(id.to_owned()); + } + self.born.remove(id); + if let Some(counter) = lost { + Counters::add(counter, events_in(id)); + } + } + + /// Try the deletes that failed before. + fn retry_deletes(&mut self) { + let cache = self.shared.cache.clone(); + self.undeletable.retain(|id| cache.remove(id).is_err()); + } + + fn recovered(&mut self, key: &str) { + if let Some(pause) = self.pauses.remove(key) { + log::info!( + "uploads recovered after {} failures, backlog {}", + pause.failing, + self.shared.cache.pending().len() + ); + } + } + + /// Pause uploads to one destination after a failure — the others carry on: for the delay the + /// server named, honored in full (neither shutdown nor a hold's escape hatch cuts it short), + /// else for a jittered exponential backoff. + fn back_off(&mut self, key: &str, server_delay: Option, reason: &str) { + let pause = self.pauses.entry(key.to_owned()).or_default(); + pause.failing += 1; + let wait = match server_delay { + Some(delay) => { + let delay = delay.min(MAX_SERVER_DELAY); + pause.server_until = Some(Instant::now() + delay); + delay + } + None => backoff(pause.failing), + }; + pause.until = Some(Instant::now() + wait); + let failing = pause.failing; + let backlog = self.shared.cache.pending().len(); + let wait = wait.as_secs_f64(); + if failing == 1 { + log::warn!("upload failed: {reason}; retrying in {wait:.1}s, backlog {backlog}"); + } else { + log::debug!( + "upload failed ({failing} in a row): {reason}; retrying in {wait:.1}s, backlog {backlog}" + ); + } + } + + /// One request, bounded by `export_timeout_ms`, as a future the loop polls. Never retried + /// here: the answer decides. + fn request( + &self, + body: &[u8], + signal: Signal, + target: &Target, + ) -> Pin + Send>> { + let mut headers = HashMap::from([ + ("Content-Type".to_owned(), otlp::CONTENT_TYPE.to_owned()), + ("Content-Encoding".to_owned(), "gzip".to_owned()), + ("Priority".to_owned(), PRIORITY.to_owned()), + ]); + if let Some(token) = &target.token { + headers.insert("Authorization".to_owned(), format!("Bearer {token}")); + } + let request = + ExportRequest { url: target.url(signal).to_owned(), headers, body: body.to_vec() }; + let bound = Duration::from_millis(self.shared.config.export_timeout_ms.max(1)); + let transport = self.transport.clone(); + Box::pin(async move { + match timeout(bound, transport.send(request)).await { + Ok(Ok(response)) => Outcome::Answered(Verdict::of(&response)), + Ok(Err(ExportError::Retryable { reason, retry_after_ms: Some(ms) })) => { + Outcome::Answered(Verdict::Throttled { delay_ms: ms, reason }) + } + Ok(Err(ExportError::Retryable { reason, retry_after_ms: None })) => { + Outcome::NoAnswer { timed_out: false, reason } + } + Ok(Err(ExportError::Rejected { reason })) => { + Outcome::Answered(Verdict::Rejected(reason)) + } + Ok(Err(ExportError::Disabled)) => Outcome::Answered(Verdict::Disabled), + Err(_) => Outcome::NoAnswer { timed_out: true, reason: "timed out".to_owned() }, + } + }) + } +} + +/// The answer to the request on the wire, or never while none is out. +async fn answer(inflight: &mut Option) -> Outcome { + match inflight { + Some(inflight) => inflight.request.as_mut().await, + None => std::future::pending().await, + } +} + +/// Sleep until `deadline`, or forever without one. +async fn at(deadline: Option) { + match deadline { + Some(deadline) => sleep_until(deadline).await, + None => std::future::pending().await, + } +} + +/// The `failures`-th consecutive failure's wait: [`RETRY_BASE`] doubling up to [`RETRY_CAP`], +/// fully jittered: uniform in `[0, backoff]`. +fn backoff(failures: u32) -> Duration { + let full = RETRY_BASE.saturating_mul(1 << failures.saturating_sub(1).min(16)).min(RETRY_CAP); + full.mul_f64(rand::random::()) +} + +/// A cached batch is intact when it is one complete gzip member (RFC 1952) whose CRC checks +/// out: a truncated or bit-flipped file is caught here, never sent. +fn intact(body: &[u8]) -> bool { + body.starts_with(&[0x1f, 0x8b]) && gunzip(body).is_some() +} + +fn gunzip(body: &[u8]) -> Option> { + use std::io::Read; + let mut out = Vec::with_capacity(body.len() * 4); + flate2::read::GzDecoder::new(body).read_to_end(&mut out).ok()?; + Some(out) +} + +/// Level 1: protobuf with repeated attribute keys shrinks 5–10× already; higher levels buy little +/// for more CPU. +fn gzip(body: &[u8]) -> Vec { + let mut encoder = GzEncoder::new(Vec::with_capacity(body.len() / 4), Compression::fast()); + // Writing into a Vec cannot fail. + let _ = encoder.write_all(body); + encoder.finish().unwrap_or_default() +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn batch_ids_carry_their_owner() { + let id = + BatchId::format(Signal::Traces, "000007", 12, "abc", Some("my-proj.livekit.cloud")); + let parsed = BatchId::parse(&id); + assert_eq!((parsed.count, parsed.signal), (12, Signal::Traces)); + assert_eq!(parsed.owner, Some("abc")); + assert_eq!(parsed.host, Some("my-proj.livekit.cloud"), "hosts may contain dashes"); + let process = BatchId::format(Signal::Logs, "000008", 3, "def", None); + assert_eq!(BatchId::parse(&process).host, None); + assert_eq!(events_in(&process), 3); + assert!(id < process, "ids sort oldest first"); + let (a, b) = (BatchId::half(&id, 'a', 6), BatchId::half(&id, 'b', 6)); + assert!(id < a && a < b && b < process, "halves keep their place in line"); + assert_eq!((events_in(&a), BatchId::parse(&b).host), (6, parsed.host)); + } + + #[test] + fn backoff_doubles_with_full_jitter_up_to_a_minute() { + for (failures, full) in [(1, 1), (2, 2), (3, 4), (6, 32), (7, 60), (40, 60)] { + let full = Duration::from_secs(full); + let waits: Vec = (0..64).map(|_| backoff(failures)).collect(); + assert!(waits.iter().all(|w| *w <= full), "{failures}: {waits:?} vs {full:?}"); + assert!(waits.iter().any(|w| *w < full / 2), "{failures}: jitter spans the range"); + } + } +} diff --git a/livekit-telemetry/src/lib.rs b/livekit-telemetry/src/lib.rs index 6971ebdd4..c4045e13e 100644 --- a/livekit-telemetry/src/lib.rs +++ b/livekit-telemetry/src/lib.rs @@ -12,8 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. -// Internals the pipeline (`Telemetry`, `Exporter`) consumes once it is in place. -#![allow(dead_code)] +#![doc = include_str!("../README.md")] /// Event data model: what SDKs push in. mod event; @@ -32,8 +31,12 @@ mod rtc; /// Spans: one attempt at an operation, with explicit handles across the FFI. mod scope; +#[cfg_attr(test, allow(dead_code))] // `Spans::open_count` serves the device contract tests mod span; +/// Batch exporter actor: timer, OTLP encoding, retry policy. +mod exporter; + /// OTLP/HTTP protobuf encoding of a batch. mod otlp; @@ -49,13 +52,22 @@ mod transport; /// Where batches go: server URL + token → ingest URL, grant, expiry, per-project routing. mod destination; +/// Entry point and configuration. +#[allow(dead_code)] // `weak_commands` serves `global`; test hooks serve the pipeline tests +mod telemetry; +mod trace; + pub use cache::{BatchCache, FileCache, MemoryCache}; pub use destination::ENDPOINT_OVERRIDE_ENV; pub use device::*; pub use event::*; +pub use exporter::Exporter; pub use rtc::{RtcStat, RtcStatsSample, StreamDirection, TrackKind}; +pub use scope::{DisconnectReason, RoomIdentity, Scope}; pub use span::SpanOutcome; pub use stats::{TelemetryStats, TelemetryStatus}; +pub use telemetry::*; +pub use trace::*; pub use transport::*; #[cfg(feature = "uniffi")] diff --git a/livekit-telemetry/src/scope.rs b/livekit-telemetry/src/scope.rs index bf282f015..ab497c979 100644 --- a/livekit-telemetry/src/scope.rs +++ b/livekit-telemetry/src/scope.rs @@ -13,11 +13,18 @@ // limitations under the License. use std::{ + collections::HashMap, fmt, sync::{Arc, Mutex}, + time::Duration, }; -use crate::{Attribute, AttributeValue}; +use tokio::time::Instant; + +use crate::{ + rtc::RtcStat, Attribute, AttributeValue, LogRecord, RtcStatsSample, Span, SpanName, + SpanOutcome, SpanStep, SpanTrack, StreamDirection, Telemetry, TelemetryEvent, TrackKind, +}; /// One session's identity: the trace id every one of its records carries, and the attributes /// attached to them at export time (`lk.room.sid`, `lk.participant.identity`, …). @@ -29,6 +36,10 @@ pub(crate) struct ScopeState { /// The app's correlation attributes, copied into each record when it is captured, so a /// later change never rewrites what is already queued. custom: Mutex>, + /// Open `lk.subscribe` spans by track sid, from intent to first media. + subscribes: Mutex, Instant)>>, + /// Published tracks awaiting their first outbound reading, and since when. + publishing: Mutex>, /// The project host this session's batches go to (`Scope::set_server`); `None` until then. route: Mutex>, /// The last `(url, token)` handed over, so handing the same pair again costs a comparison; @@ -47,6 +58,8 @@ impl ScopeState { trace_id, attributes: Mutex::new(Vec::new()), custom: Mutex::new(Vec::new()), + subscribes: Mutex::new(HashMap::new()), + publishing: Mutex::new(HashMap::new()), route: Mutex::new(None), server: Mutex::new(None), }) @@ -75,6 +88,58 @@ impl ScopeState { *self.route.lock().unwrap_or_else(|e| e.into_inner()) = Some(host); } + /// A track was published: poll fast until its first outbound reading (at most + /// [`Scope::SUBSCRIBE_TIMEOUT`]). + pub fn await_first_outbound(&self, sid: &str) { + self.publishing + .lock() + .unwrap_or_else(|e| e.into_inner()) + .insert(sid.to_owned(), Instant::now()); + } + + /// Whether a published track still waits for its first outbound reading (entries older than + /// the subscribe timeout are forgotten, so a track that never sends cannot keep polling fast). + fn awaiting_outbound(&self) -> bool { + let mut publishing = self.publishing.lock().unwrap_or_else(|e| e.into_inner()); + publishing.retain(|_, since| since.elapsed() < Scope::SUBSCRIBE_TIMEOUT); + !publishing.is_empty() + } + + /// When the oldest open subscribe runs out of time. + pub fn subscribe_deadline(&self) -> Option { + let open = self.subscribes.lock().unwrap_or_else(|e| e.into_inner()); + open.values().map(|(_, since)| *since + Scope::SUBSCRIBE_TIMEOUT).min() + } + + /// The pipeline stops: end every pending subscribe (`timed_out` past its deadline, + /// `cancelled` before) so it ships with the last batch, or, on opt-out, just let go of them. + pub fn end_subscribes(&self, export: bool) { + let open: Vec<_> = + self.subscribes.lock().unwrap_or_else(|e| e.into_inner()).drain().collect(); + if export { + for (_, (span, since)) in open { + end_unfinished(&span, since); + } + } + } + + /// End every subscribe past its deadline as `timed_out` — driven by the exporter's clock, so + /// a track that never produces a single RTP reading still times out. + pub fn expire_subscribes(&self) { + let expired: Vec> = { + let mut open = self.subscribes.lock().unwrap_or_else(|e| e.into_inner()); + let sids: Vec = open + .iter() + .filter(|(_, (_, since))| since.elapsed() >= Scope::SUBSCRIBE_TIMEOUT) + .map(|(sid, _)| sid.clone()) + .collect(); + sids.iter().filter_map(|sid| open.remove(sid)).map(|(span, _)| span).collect() + }; + for span in expired { + span.fail("timed_out".to_owned()); + } + } + /// The trace id as 32 hex characters. pub fn hex(&self) -> String { format!("{:032x}", u128::from_be_bytes(self.trace_id)) @@ -163,3 +228,339 @@ impl fmt::Debug for ScopeState { write!(f, "Scope({})", self.hex()) } } + +/// Who this session is: attached to every record once the room is joined. `None` clears. +#[cfg_attr(feature = "uniffi", derive(uniffi::Record))] +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct RoomIdentity { + #[cfg_attr(feature = "uniffi", uniffi(default))] + pub sid: Option, + #[cfg_attr(feature = "uniffi", uniffi(default))] + pub name: Option, + #[cfg_attr(feature = "uniffi", uniffi(default))] + pub participant_sid: Option, + #[cfg_attr(feature = "uniffi", uniffi(default))] + pub participant_identity: Option, +} + +/// Why a session ended: the protocol's `DisconnectReason`, plus the client giving up. +#[cfg_attr(feature = "uniffi", derive(uniffi::Enum))] +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum DisconnectReason { + Unknown, + ClientInitiated, + DuplicateIdentity, + ServerShutdown, + ParticipantRemoved, + RoomDeleted, + StateMismatch, + JoinFailure, + Migration, + SignalClose, + RoomClosed, + UserUnavailable, + UserRejected, + SipTrunkFailure, + ConnectionTimeout, + MediaFailure, + AgentError, + /// The reconnect policy ran out of attempts. + ReconnectFailed, +} + +impl DisconnectReason { + /// The protocol's `DisconnectReason` number, so no SDK keeps its own switch. + pub fn from_proto(value: i32) -> Self { + match value { + 1 => Self::ClientInitiated, + 2 => Self::DuplicateIdentity, + 3 => Self::ServerShutdown, + 4 => Self::ParticipantRemoved, + 5 => Self::RoomDeleted, + 6 => Self::StateMismatch, + 7 => Self::JoinFailure, + 8 => Self::Migration, + 9 => Self::SignalClose, + 10 => Self::RoomClosed, + 11 => Self::UserUnavailable, + 12 => Self::UserRejected, + 13 => Self::SipTrunkFailure, + 14 => Self::ConnectionTimeout, + 15 => Self::MediaFailure, + 16 => Self::AgentError, + _ => Self::Unknown, + } + } +} + +/// A session on the shared pipeline: its own trace id and attributes, the same queue, cache, +/// cadence and exporter as every other session in the process. +/// +/// The pipeline starts once, at SDK init, so nothing that happens before the first room is +/// lost; a `Scope` is what a room — one call — gets from it, and what its spans, RTC windows +/// and events are filed under. Everything emitted outside a session (device state, pre-room +/// errors, self-telemetry) belongs to the pipeline's own process session. Cheap to clone. +#[derive(Clone)] +pub struct Scope { + pub(crate) telemetry: Telemetry, + pub(crate) state: Arc, +} + +impl Scope { + /// The session's trace id as 32 hex characters — print it (`lkt_…`) so support can find + /// the call. + pub fn trace_id(&self) -> String { + self.state.hex() + } + + /// The LiveKit server this room talks to and the participant token it holds. Call it at + /// connect and again with every refreshed token: the core derives the ingest URL (LiveKit + /// Cloud only), reads the token's grant and expiry, and routes this session's records to its + /// own project with its own token. Cheap and idempotent — the same pair again is a no-op. + pub fn set_server(&self, url: &str, token: &str) { + self.telemetry.set_server(&self.state, url, token); + } + + /// Queue an event or log record under this session. + pub fn emit(&self, event: TelemetryEvent) { + self.telemetry.emit_in(event, &self.state); + } + + /// An app-defined event, exported as `custom.` under this session with the session's + /// correlation attributes (its own attributes win over them). Rejected and counted as + /// `invalid` — never truncated — when the name is empty or over 128 bytes, an attribute is + /// over the limits or in the SDK's namespace (`lk.*`, `session.id`), or there are more than + /// 64 attributes. + pub fn emit_custom(&self, name: &str, attributes: Vec) { + let valid = !name.is_empty() + && name.len() <= crate::event::MAX_NAME_BYTES + && attributes.len() <= crate::event::MAX_CUSTOM_ATTRIBUTES + && attributes.iter().all(|a| crate::event::valid_custom(&a.key, Some(&a.value))); + if !valid { + self.telemetry.count_invalid(); + return; + } + self.emit(TelemetryEvent::custom(name, attributes)); + } + + /// A log record filed under this session even without an ambient span — for platforms with + /// no task-local context (Dart outside a zone). Same floor and filters as `Telemetry::log`. + pub fn log(&self, record: LogRecord) { + if let Some(event) = self.telemetry.log_event(record) { + self.emit(event); + } + } + + /// The session ended for good (not a reconnect): the `lk.room.disconnected` record. + pub fn disconnected(&self, reason: DisconnectReason) { + let open: Vec<_> = + self.state.subscribes.lock().unwrap_or_else(|e| e.into_inner()).drain().collect(); + for (_, (span, since)) in open { + end_unfinished(&span, since); + } + self.telemetry.retire_stats(&self.state, None); + // Out of the call: uploads no longer yield to it. + self.state.server.lock().unwrap_or_else(|e| e.into_inner()).take(); + let severity = if reason == DisconnectReason::ClientInitiated { + crate::Severity::Info + } else { + crate::Severity::Warn + }; + let reason = crate::device::snake(reason); + self.emit( + TelemetryEvent::new("lk.room.disconnected") + .with_severity(severity) + .with_body(format!("disconnected: {reason}")) + .with_attribute("lk.disconnect.reason", reason), + ); + } + + /// Set (or, with `None`, remove) an app correlation attribute — `app.call_id`, + /// `enduser.id` — on every record this session captures from now on: logs, spans, events and + /// RTC windows. Records already captured keep the value they had. Rejected and counted as + /// `invalid` when the key is empty or over 128 bytes, a string value is over 1024 bytes, the + /// key is the SDK's (`lk.*`, `session.id`), or the session already has 64. + pub fn set_attribute(&self, key: &str, value: Option) { + // Under the windows lock: the open windows close under the old value and the change + // lands before any reading can open a new one (readings record under the same lock). + let check = value.clone(); + let accepted = self.telemetry.with_windows_split( + &self.state, + || self.state.accepts_custom(key, check.as_ref()), + || { + self.state.set_custom(key, value); + }, + ); + if !accepted { + self.telemetry.count_invalid(); + } + } + + /// Push one `getStats()` reading; its window ships under this session. The first inbound + /// reading with bytes is a subscribed track's first media. + pub fn record_stats(&self, sample: RtcStatsSample) { + if sample.direction == StreamDirection::Inbound && sample.bytes.unwrap_or(0) > 0 { + self.first_media(&sample.track_sid); + } + if sample.direction == StreamDirection::Outbound { + self.state + .publishing + .lock() + .unwrap_or_else(|e| e.into_inner()) + .remove(&sample.track_sid); + } + self.telemetry.record_stats_in(sample, &self.state); + } + + /// A whole `getStats()` report for one track, as the platform got it. The core picks the RTP + /// streams, resolves codec and RTT, converts units and records one sample per stream (outbound + /// ones tagged with their layer), so no SDK maps stats fields itself. + pub fn record_stats_report( + &self, + track_sid: &str, + kind: TrackKind, + direction: StreamDirection, + report: Vec, + timestamp_ns: Option, + ) { + for sample in + crate::rtc::samples_from_report(track_sid, kind, direction, &report, timestamp_ns) + { + self.record_stats(sample); + } + } + + /// A whole peer connection's `getStats()` report, as the platform got it, and the tracks it + /// carries (MediaStreamTrack id → track sid, local and remote). The core finds each track's + /// RTP streams and records them as [`record_stats_report`](Self::record_stats_report) would: + /// one call per peer connection per poll, whatever the number of participants and tracks. + pub fn record_peer_stats( + &self, + report: Vec, + tracks: HashMap, + timestamp_ns: Option, + ) { + for sample in crate::rtc::samples_from_peer_report(&report, &tracks, timestamp_ns) { + self.record_stats(sample); + } + } + + /// How long the platform should wait before its next `getStats()` poll of this room: 1 s + /// while a subscribe waits for its first media (the core sees first media in the readings) or + /// a newly published track for its first outbound reading (known from its `lk.publish` + /// span's `set_track` with a sid; at most 30 s), else half the current stats window — two readings per window, at the pressure-stretched + /// cadence. Ask again after every poll. + pub fn stats_poll_interval_ms(&self) -> u64 { + let subscribing = + !self.state.subscribes.lock().unwrap_or_else(|e| e.into_inner()).is_empty(); + if subscribing || self.state.awaiting_outbound() { + return 1000; + } + self.telemetry.stats_window_ms() / 2 + } + + /// No media within this long ends `lk.subscribe` with `error.type = timed_out`. + pub const SUBSCRIBE_TIMEOUT: Duration = Duration::from_secs(30); + + /// The intent to subscribe exists (autoSubscribe: at the remote publish; manual: at the + /// subscribe call; a track already in the room at join: at connect): `lk.subscribe` opens, + /// once per track. Its end is an RTC fact the core + /// sees itself — the first inbound reading with bytes — so no SDK keeps this state. + pub fn subscribe_started(&self, track: SpanTrack) { + let Some(sid) = track.sid.clone() else { return }; + if self.telemetry.shared.revoked() { + return; + } + { + let mut open = self.state.subscribes.lock().unwrap_or_else(|e| e.into_inner()); + if open.contains_key(&sid) { + return; + } + let span = self.start(SpanName::Subscribe, None); + span.set_track(track); + open.insert(sid, (span, Instant::now())); + } + // The exporter owns the clock that enforces the deadline. + self.telemetry.wake(); + } + + /// The server confirmed the subscription: the `subscribed` step. Without an earlier intent + /// this is the intent (a fallback, measured from here): the span opens and + /// [`stats_poll_interval_ms`](Self::stats_poll_interval_ms) turns fast at once. + pub fn subscribed(&self, track: SpanTrack) { + self.subscribe_started(track.clone()); + let Some(sid) = &track.sid else { return }; + if let Some((span, _)) = + self.state.subscribes.lock().unwrap_or_else(|e| e.into_inner()).get(sid) + { + span.step(SpanStep::Subscribed); + } + } + + /// A track left this session — unpublished, unsubscribed, its publisher gone. A pending + /// `lk.subscribe` ends `cancelled` (or `timed_out` past its deadline), the track's last + /// partial RTC window ships now, and the core forgets everything it kept for the track. + pub fn track_ended(&self, sid: &str) { + self.state.publishing.lock().unwrap_or_else(|e| e.into_inner()).remove(sid); + if let Some((span, since)) = self.take_subscribe(sid) { + end_unfinished(&span, since); + } + self.telemetry.retire_stats(&self.state, Some(sid)); + } + + /// The subscription failed (`error.type` = the platform's error name). + pub fn subscribe_failed(&self, sid: &str, error_type: &str) { + if let Some((span, _)) = self.take_subscribe(sid) { + span.fail(error_type.to_owned()); + } + } + + /// First media ends the subscribe `ok` — unless its deadline already passed and the sweep + /// has not run yet: then it is `timed_out`, whoever looks first. + fn first_media(&self, sid: &str) { + let Some((span, since)) = self.take_subscribe(sid) else { return }; + if since.elapsed() >= Self::SUBSCRIBE_TIMEOUT { + span.fail("timed_out".to_owned()); + return; + } + span.step(SpanStep::FirstMedia); + span.end(SpanOutcome::Ok, None); + } + + fn take_subscribe(&self, sid: &str) -> Option<(Arc, Instant)> { + self.state.subscribes.lock().unwrap_or_else(|e| e.into_inner()).remove(sid) + } + + /// Start a typed span in this session's trace, stamped now. `parent` nests it. + pub fn start(&self, name: SpanName, parent: Option>) -> Arc { + let parent = parent.and_then(|p| p.context()).map(|c| c.span_id); + Span::bound(name, parent, self.telemetry.clone(), &self.state) + } + + /// The room and local participant, as `lk.room.*` / `lk.participant.*` on every record. + pub fn set_room(&self, room: RoomIdentity) { + for (key, value) in [ + ("lk.room.sid", room.sid), + ("lk.room.name", room.name), + ("lk.participant.sid", room.participant_sid), + ("lk.participant.identity", room.participant_identity), + ] { + // Identifiers, not free text: an over-long one is not attached (counted as invalid). + if value.as_ref().is_some_and(|v| v.len() > crate::event::MAX_VALUE_BYTES) { + self.telemetry.count_invalid(); + continue; + } + self.state.set_attribute(key, value.map(AttributeValue::Str)); + } + } +} + +/// A subscribe that ended without media: `timed_out` once past its deadline — even when the +/// cleanup comes later — `cancelled` before. +fn end_unfinished(span: &Span, since: Instant) { + if since.elapsed() >= Scope::SUBSCRIBE_TIMEOUT { + span.fail("timed_out".to_owned()); + } else { + span.cancel(); + } +} diff --git a/livekit-telemetry/src/telemetry.rs b/livekit-telemetry/src/telemetry.rs new file mode 100644 index 000000000..a72359419 --- /dev/null +++ b/livekit-telemetry/src/telemetry.rs @@ -0,0 +1,892 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::{ + sync::{ + atomic::{AtomicBool, Ordering}, + Arc, Mutex, Weak, + }, + time::Duration, +}; + +use tokio::sync::{mpsc, oneshot}; +use tokio::time::{timeout, Instant}; + +use crate::span::SpanKind; +use crate::{ + cache::SpillCache, + destination::{Destinations, ENDPOINT_OVERRIDE_ENV}, + event::now_unix_nanos, + exporter::Command, + rtc::StatsWindows, + scope::{Scope, ScopeState}, + span::Spans, + stats::{Counters, TelemetryStatus}, + store::{Queued, Store}, + Attribute, AttributeValue, BatchCache, DeviceState, Exporter, FileCache, LogRecord, LogSource, + MemoryCache, RtcStatsSample, Severity, SpanOutcome, TelemetryEvent, TelemetryStats, + TelemetryTransport, +}; +use crate::{DeviceEvent, Span, SpanName}; + +/// Pipeline configuration: storage and tuning. Where batches go is not configurable — the core +/// derives it from the server URL and token each room hands over ([`Scope::set_server`]). +/// +/// Defaults are conservative, sized so a fleet cannot overload the collector: one export and +/// one RTC stats window per minute, OTel's `BatchLogRecordProcessor` queue (2048) and batch (512). +#[cfg_attr(feature = "uniffi", derive(uniffi::Record))] +#[derive(Debug, Clone)] +pub struct TelemetryConfig { + /// Resource attributes describing the emitter (`service.name`, `os.name`, + /// `device.model.identifier`, `session.id`, …). `telemetry.sdk.*` are filled in by the core. + #[cfg_attr(feature = "uniffi", uniffi(default = []))] + pub resource: Vec, + /// Who is reporting, typed; the core owns the semconv keys. Extra attributes go in `resource`. + #[cfg_attr(feature = "uniffi", uniffi(default))] + pub sdk: Option, + /// Directory for the on-disk batch cache (created if missing; its parent must exist). + /// `None` keeps batches in memory only: they survive failed uploads, not the process. + #[cfg_attr(feature = "uniffi", uniffi(default))] + pub storage_dir: Option, + /// Cap on cached batches, in memory or on disk; the oldest are evicted first. + #[cfg_attr(feature = "uniffi", uniffi(default = 4194304))] + pub max_cache_bytes: u64, + /// Base export cadence; stretched up to 4× by [`DeviceState::cadence_factor`]. + #[cfg_attr(feature = "uniffi", uniffi(default = 60000))] + pub flush_interval_ms: u64, + /// Events buffered before the oldest are dropped. + #[cfg_attr(feature = "uniffi", uniffi(default = 2048))] + pub max_queue_size: u32, + /// Events per export request. + #[cfg_attr(feature = "uniffi", uniffi(default = 512))] + pub max_batch_size: u32, + /// Bound on a single transport attempt, and on `shutdown`. + #[cfg_attr(feature = "uniffi", uniffi(default = 10000))] + pub export_timeout_ms: u64, + /// RTC stats window: readings pushed with [`Scope::record_stats`] are summarised into one + /// `lk.rtc.stats.sample` per track and direction every window (stretched like the cadence). + #[cfg_attr(feature = "uniffi", uniffi(default = 60000))] + pub stats_window_ms: u64, + /// Flood guard for discrete events: beyond this many `emit`s per 10 minutes the rest are + /// dropped and counted as `rate_limited`. RTC windows and self-telemetry are exempt; 0 = off. + #[cfg_attr(feature = "uniffi", uniffi(default = 300))] + pub max_events_per_10min: u32, + /// Cached batches uploaded per tick while a session may be live — bounds how fast a backlog + /// (offline period, previous launch) replays next to a call: 4 × ~20 KB gzipped per minute + /// is ~10 kbps. `shutdown` drains without the budget. + #[cfg_attr(feature = "uniffi", uniffi(default = 4))] + pub max_batches_per_upload: u32, + /// Export as soon as the queue holds about this many bytes, without waiting for the tick + /// (design doc: "flush on the tick or at 256 KB"). + #[cfg_attr(feature = "uniffi", uniffi(default = 262144))] + pub flush_threshold_bytes: u64, + /// Cap on one request's payload before compression (design doc: "single POST ≤ 1 MB"). + #[cfg_attr(feature = "uniffi", uniffi(default = 1048576))] + pub max_batch_bytes: u64, + /// Lowest severity a plain log record (an event with no name) needs to leave the device. + /// Events are not subject to it. Design doc: warn. + /// `None` is `Warn`. (Optional because UniFFI 0.31 cannot default an enum literal.) + #[cfg_attr(feature = "uniffi", uniffi(default))] + pub log_severity: Option, + /// Platform instruments not to run; all run by default. + #[cfg_attr(feature = "uniffi", uniffi(default = []))] + pub disabled_instruments: Vec, +} + +/// The platform instruments a config can switch off. +#[cfg_attr(feature = "uniffi", derive(uniffi::Enum))] +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum Instrument { + /// Thermal, power, memory, network, battery, audio session. + Device, + /// Warn/error lines from the SDK, the core and WebRTC. + Logs, + /// `getStats` windows per track and subscribe spans. + Rtc, + /// Connect, reconnect and publish spans. + Room, +} + +impl Default for TelemetryConfig { + fn default() -> Self { + Self { + resource: Vec::new(), + sdk: None, + storage_dir: None, + max_cache_bytes: 4 * 1024 * 1024, + flush_interval_ms: 60_000, + max_queue_size: 2048, + max_batch_size: 512, + export_timeout_ms: 10_000, + stats_window_ms: 60_000, + max_events_per_10min: 300, + max_batches_per_upload: 4, + flush_threshold_bytes: 256 * 1024, + max_batch_bytes: 1024 * 1024, + log_severity: None, + disabled_instruments: Vec::new(), + } + } +} + +/// Entry point: the synchronous, never-blocking side SDKs push into. +/// +/// Fail-open by design: [`emit`](Self::emit) cannot fail or block — when the queue is full the +/// oldest event is dropped and counted in [`stats`](Self::stats). Cheap to clone; every clone +/// feeds the same pipeline. +/// +/// ``` +/// # use std::sync::Arc; +/// # use livekit_telemetry::*; +/// # struct Discard; +/// # #[async_trait::async_trait] +/// # impl TelemetryTransport for Discard { +/// # async fn send(&self, _: ExportRequest) -> Result { +/// # Ok(ExportResponse::accepted()) +/// # } +/// # } +/// # #[tokio::main(flavor = "current_thread")] async fn main() { +/// let (telemetry, exporter) = Telemetry::new(TelemetryConfig::default(), Arc::new(Discard)); +/// tokio::spawn(exporter.run()); +/// +/// // One scope per room: the server URL and token route its records to the room's project. +/// let room = telemetry.begin_scope(); +/// room.set_server("wss://my-project.livekit.cloud", ""); +/// room.emit_custom("checkout", vec![Attribute::new("plan", "pro")]); +/// telemetry.emit(TelemetryEvent::new("lk.ping")); +/// telemetry.shutdown().await; +/// # } +/// ``` +#[derive(Clone)] +pub struct Telemetry { + pub(crate) shared: Arc, + guard: Arc>, + commands: mpsc::UnboundedSender, +} + +/// How much longer than `export_timeout_ms` [`Telemetry::shutdown`] waits for the exporter to +/// confirm it stopped. +const SHUTDOWN_GRACE: Duration = Duration::from_secs(1); + +/// Everything the synchronous side ([`Telemetry`]) and the exporter actor share. +pub(crate) struct Shared { + pub store: Store, + pub config: TelemetryConfig, + pub cache: Arc, + pub counters: Arc, + /// Read synchronously by the exporter, so a state pushed right before `emit` already governs + /// the flush that follows. + pub device: Mutex>, + pub windows: Mutex, + pub spans: Mutex, + /// The pipeline's own session: whatever is emitted outside a room session. + pub process: Arc, + /// Attributes attached to every record of every session. + pub global: Mutex>, + /// Where each session's batches go, and with which token. + pub destinations: Mutex, + pub status: Mutex, + /// Every room session, for the subscribe deadlines the exporter enforces. + pub scopes: Mutex>>, + /// The queue crossed `flush_threshold_bytes` since the exporter last looked: export now. + pub overflowed: AtomicBool, + /// A destination, a token or a hold changed: re-evaluate what waits (nothing is encoded). + pub released: AtomicBool, + /// Test-only: runs after a capture's early opt-out check, before it takes the lock it commits + /// under (to hold a producer exactly where a racing purge could slip in). + #[cfg(test)] + pub pause: Mutex>>, + /// Upload passes started, for tests that must see the exporter stay idle. + #[cfg(test)] + pub passes: std::sync::atomic::AtomicU64, + /// Coalesced wake-ups for the exporter (see [`Exporter`]): notifications never queue up. + pub wake: tokio::sync::Notify, + /// The app went to the background since the exporter last looked. + pub backgrounded: AtomicBool, + /// The app opted out ([`Telemetry::purge`]). One-way: checked before anything is captured, + /// written or sent, so a request still in flight cannot bring anything back. + pub revoked: Arc, +} + +impl Shared { + /// The test hook between a capture's early check and its commit; nothing in production. + fn before_commit(&self) { + #[cfg(test)] + { + let pause = self.pause.lock().unwrap_or_else(|e| e.into_inner()).clone(); + if let Some(pause) = pause { + pause(); + } + } + } + + /// The queue crossed its threshold: wake the exporter to export now. + pub fn overflow(&self) { + self.overflowed.store(true, Ordering::SeqCst); + self.wake.notify_one(); + } + + /// Something that held batches back changed: wake the exporter to re-evaluate. + pub fn release(&self) { + self.released.store(true, Ordering::SeqCst); + self.wake.notify_one(); + } + + /// Record a checkpoint inside an open span (`ws_open`, `join_recv`, …), stamped now. + pub fn add_span_event(&self, span: u64, name: &str, attributes: Vec) { + if self.revoked() { + return; + } + self.spans.lock().unwrap_or_else(|e| e.into_inner()).add_event(span, name, attributes); + } + + /// End a span with its outcome; `error_type` becomes `error.type` and the status message. + /// The span is exported with the next batch. Ending twice, or an unknown handle, is a no-op. + pub fn end_span( + &self, + span: u64, + outcome: SpanOutcome, + error_type: Option, + attributes: Vec, + ) { + if self.revoked() { + return; + } + self.spans + .lock() + .unwrap_or_else(|e| e.into_inner()) + .end(span, outcome, error_type, attributes); + } + + /// End every Room's pending subscribes (shutdown) or just forget them (opt-out): nothing + /// is left referencing a span once the pipeline stops. + pub fn end_pending_subscribes(&self, export: bool) { + let scopes: Vec> = { + let scopes = self.scopes.lock().unwrap_or_else(|e| e.into_inner()); + scopes.iter().filter_map(|s| s.upgrade()).collect() + }; + for scope in scopes { + scope.end_subscribes(export); + } + } + + /// Device pressure, doubled again while WebRTC reports the encoder CPU-limited; capped at 4×. + pub fn cadence_factor(&self) -> u32 { + let device = self.device.lock().unwrap_or_else(|e| e.into_inner()).unwrap_or_default(); + let cpu = self.windows.lock().unwrap_or_else(|e| e.into_inner()).cpu_limited(); + (device.cadence_factor() * if cpu { 2 } else { 1 }).min(4) + } + + pub fn revoked(&self) -> bool { + self.revoked.load(Ordering::SeqCst) + } + + /// Drop every record and batch, in memory and on disk, counted as `purged`. `false` when + /// the storage refused to delete something: that data is still on disk. + pub fn clear(&self) -> bool { + self.end_pending_subscribes(false); + let queued = self.store.clear(); + let spans = self.spans.lock().unwrap_or_else(|e| e.into_inner()).clear(); + let windows = self.windows.lock().unwrap_or_else(|e| e.into_inner()).clear(); + // Count what really went, once: what was there before and is gone after (an + // unreadable directory lists nothing, so nothing is counted for it). + let before = self.cache.pending(); + let deleted = self.cache.clear(); + let after: std::collections::HashSet = self.cache.pending().into_iter().collect(); + let cached: u64 = before + .iter() + .filter(|id| !after.contains(*id)) + .map(|id| crate::exporter::events_in(id)) + .sum(); + if let Err(err) = &deleted { + log::warn!("opt-out: cached telemetry could not all be deleted ({err})"); + } + Counters::add(&self.counters.purged, queued + spans + windows + cached); + deleted.is_ok() + } +} + +/// Fixed-window cap on discrete events (the design doc's ~300 per 10 min). +struct FloodGuard { + max: u32, + window_start: Instant, + count: u32, + /// The first drop of a window logs; the rest are counted. + warned: bool, +} + +impl FloodGuard { + const WINDOW: Duration = Duration::from_secs(10 * 60); + + fn new(max: u32) -> Self { + Self { max, window_start: Instant::now(), count: 0, warned: false } + } + + fn admit(&mut self) -> bool { + if self.max == 0 { + return true; + } + let now = Instant::now(); + if now.duration_since(self.window_start) >= Self::WINDOW { + self.window_start = now; + self.count = 0; + self.warned = false; + } + if self.count >= self.max { + if !self.warned { + self.warned = true; + log::warn!( + "flood: {} records in 10 min; dropping until the window moves", + self.max + ); + } + return false; + } + self.count += 1; + true + } +} + +impl Telemetry { + /// Build the pipeline with the cache the config asks for: a [`FileCache`] in `storage_dir`, + /// or a [`MemoryCache`] when unset — or when the directory is unusable (logged, never an + /// error). Spawn the returned [`Exporter`] with `exporter.run()` on your runtime. + pub fn new( + config: TelemetryConfig, + transport: Arc, + ) -> (Self, Exporter) { + let mut fell_back = false; + let cache: Arc = match config.storage_dir.as_deref() { + Some(dir) => match FileCache::open(dir, config.max_cache_bytes) { + Ok(cache) => Arc::new(cache), + Err(err) => { + log::warn!("cannot use storage dir {dir}: {err}; caching in memory"); + fell_back = true; + Arc::new(MemoryCache::new(config.max_cache_bytes)) + } + }, + None => Arc::new(MemoryCache::new(config.max_cache_bytes)), + }; + let (telemetry, exporter) = Self::with_cache(config, transport, cache); + if fell_back { + // Reported like any other refused write: batches survive failures, not the process. + Counters::add(&telemetry.shared.counters.cache_write_errors, 1); + } + (telemetry, exporter) + } + + /// Build the pipeline around a caller-provided [`BatchCache`] (`storage_dir` is ignored). + pub fn with_cache( + mut config: TelemetryConfig, + transport: Arc, + cache: Arc, + ) -> (Self, Exporter) { + add_sdk_resource(&mut config.resource, config.sdk.as_ref()); + let counters = Arc::new(Counters::default()); + let revoked = Arc::new(AtomicBool::new(false)); + let cache: Arc = Arc::new(SpillCache::new( + cache, + config.max_cache_bytes, + counters.clone(), + revoked.clone(), + )); + // Batches the cache evicted while opening (past the age or size limit) are lost too. + let evicted: u64 = + cache.take_evicted().iter().map(|id| crate::exporter::events_in(id)).sum(); + Counters::add(&counters.cache_full, evicted); + let endpoint_override = std::env::var(ENDPOINT_OVERRIDE_ENV).ok(); + if let Some(endpoint) = endpoint_override.as_deref().filter(|e| !e.is_empty()) { + log::info!("{ENDPOINT_OVERRIDE_ENV} set: every upload goes to {endpoint}"); + } + // Unbounded, but only request/response commands travel here (flush, shutdown, purge), + // each with a caller awaiting its answer: the queue is bounded by the callers waiting. + // Notifications go through the coalescing `Shared::wake` instead. + let (commands, receiver) = mpsc::unbounded_channel(); + let shared = Arc::new(Shared { + store: Store::new( + config.max_queue_size.max(1) as usize, + usize::try_from(config.flush_threshold_bytes.max(1)).unwrap_or(usize::MAX), + counters.clone(), + ) + .with_consent(revoked.clone()), + cache, + counters, + device: Mutex::new(None), + windows: Mutex::new(StatsWindows::with_consent(revoked.clone())), + spans: Mutex::new( + Spans::new(config.max_queue_size.max(1) as usize).with_consent(revoked.clone()), + ), + // One pipeline per process; sessions (rooms) carry their own trace ids. This is the + // pipeline's own session, for everything emitted outside a room. + process: ScopeState::new(), + global: Mutex::new(Vec::new()), + destinations: Mutex::new(Destinations::new(endpoint_override.as_deref())), + status: Mutex::new(TelemetryStatus::Ok), + scopes: Mutex::new(Vec::new()), + revoked, + #[cfg(test)] + passes: Default::default(), + #[cfg(test)] + pause: Mutex::new(None), + wake: tokio::sync::Notify::new(), + overflowed: AtomicBool::new(false), + released: AtomicBool::new(false), + backgrounded: AtomicBool::new(false), + config, + }); + let guard = Arc::new(Mutex::new(FloodGuard::new(shared.config.max_events_per_10min))); + let exporter = Exporter::new(shared.clone(), transport, receiver); + (Self { shared, guard, commands }, exporter) + } + + /// Route a room's records: its server URL names the project (LiveKit Cloud only), its token + /// authorizes the upload. Called at connect and again with every refreshed token; cheap and + /// idempotent. Everything cached for the project starts uploading once the token allows. + pub(crate) fn set_server(&self, session: &ScopeState, url: &str, token: &str) { + if session.same_server(url, token) { + // A reconnect with the same pair still retries an ingest that answered 404 — one + // map lookup, no parsing. + let route = session.route(); + let revived = route.is_some_and(|host| { + self.shared.destinations.lock().unwrap_or_else(|e| e.into_inner()).revive(&host) + }); + if revived { + self.shared.release(); + } + return; + } + let host = self.shared.destinations.lock().unwrap_or_else(|e| e.into_inner()).set( + url, + token, + &session.hex(), + ); + let Some(host) = host else { + log::warn!("server url has no host; this session's telemetry stays cached"); + return; + }; + // A project change closes the session's RTC windows first: a window belongs to one project. + if session.route().as_deref() != Some(host.as_str()) { + self.with_windows_split(session, || true, || session.set_route(host.clone())); + } + // A new route or a new token may release what waits in the cache. + self.shared.release(); + } + + /// The RTC stats window right now: `stats_window_ms` stretched by device pressure and a + /// CPU-limited encoder (at most 4×). + pub(crate) fn stats_window_ms(&self) -> u64 { + self.shared.config.stats_window_ms.max(1) * u64::from(self.shared.cadence_factor()) + } + + /// An app handed over something over the limits (a custom event or attribute). + pub(crate) fn count_invalid(&self) { + Counters::add(&self.shared.counters.invalid, 1); + } + + /// Nudge the exporter to re-read its deadlines (a subscribe started). + pub(crate) fn wake(&self) { + self.shared.wake.notify_one(); + } + + /// Point every upload at `endpoint`, as [`ENDPOINT_OVERRIDE_ENV`] does. + #[cfg(test)] + pub(crate) fn override_endpoint(&self, endpoint: &str) { + *self.shared.destinations.lock().unwrap_or_else(|e| e.into_inner()) = + Destinations::new(Some(endpoint)); + } + + /// Start a session — one room, one call — with its own trace id and attributes on this + /// pipeline. Sessions do not need ending: a room's last record is simply its last. + pub fn begin_scope(&self) -> Scope { + let state = ScopeState::new(); + let mut scopes = self.shared.scopes.lock().unwrap_or_else(|e| e.into_inner()); + scopes.retain(|scope| scope.strong_count() > 0); + scopes.push(Arc::downgrade(&state)); + Scope { telemetry: self.clone(), state } + } + + /// A span in the pipeline's own trace: app-defined work outside any room, or the SDK before a + /// room exists. Stamped now; `parent` nests it. + pub fn start(&self, name: SpanName, parent: Option>) -> Arc { + let parent = parent.and_then(|p| p.context()).map(|c| c.span_id); + Span::bound(name, parent, self.clone(), &self.shared.process) + } + + /// Queue an event or log record for export. Stamps it with the current time unless it + /// carries one. + /// + /// A record with an empty `name` is a plain log line: only `Warn` and `Error` ones leave the + /// device (design doc: debug/info logs never do). Discrete events are subject to the flood + /// guard (`max_events_per_10min`); what it drops is counted as `rate_limited`. + /// Something happened to the device mid-call (audio route, interruption, a denied permission): + /// a process-level record with a display body, built here so every platform files it alike. + pub fn device_event(&self, event: DeviceEvent) { + self.emit_in(event.into_event(), &self.shared.process); + } + + /// A captured log line. WebRTC only counts at error; the SDK and the core at the configured + /// floor; the core's own telemetry module never (a rejected batch that produced a record that + /// produced a batch would never end). + pub fn log(&self, record: LogRecord) { + if let Some(event) = self.log_event(record) { + self.emit(event); + } + } + + /// The record as an event, or nothing when it is below the floor (WebRTC: error only) or is + /// telemetry's own. + pub(crate) fn log_event(&self, record: LogRecord) -> Option { + let floor = match record.source { + LogSource::WebRtc => self.log_severity().max(Severity::Error), + _ => self.log_severity(), + }; + if record.severity < floor { + return None; + } + if record.source == LogSource::Ffi + && record.logger.as_deref().is_some_and(|l| l.starts_with("livekit_telemetry")) + { + return None; + } + Some(record.into()) + } + + /// Queue an event or log record, filed under the session of the span it names (`span_id`), + /// else under the pipeline's own process session. + pub fn emit(&self, event: TelemetryEvent) { + // A record emitted inside a room's span belongs to that room's session; anything else + // is the process's own. + let session = event + .span_id + .and_then(|id| self.shared.spans.lock().unwrap_or_else(|e| e.into_inner()).scope_of(id)) + .unwrap_or_else(|| self.shared.process.clone()); + self.emit_in(event, &session); + } + + fn log_severity(&self) -> Severity { + self.shared.config.log_severity.unwrap_or(Severity::Warn) + } + + pub(crate) fn emit_in(&self, mut event: TelemetryEvent, session: &Arc) { + if event.name.is_empty() && event.severity < self.log_severity() { + return; + } + if self.shared.revoked() { + return; + } + if !self.collects(session) { + Counters::add(&self.shared.counters.disabled, 1); + return; + } + if !self.guard.lock().unwrap_or_else(|e| e.into_inner()).admit() { + Counters::add(&self.shared.counters.rate_limited, 1); + return; + } + if event.timestamp_ns.is_none() { + event.timestamp_ns = Some(now_unix_nanos()); + } + if self.shared.store.push(Queued::new(event, session.clone())) { + self.shared.overflow(); + } + } + + /// Queue a consumer-defined event, exported as `custom.` (see + /// [`TelemetryEvent::custom`]): the stringly-typed escape hatch next to the `lk.*` + /// catalogue. Same flood guard, same pipeline. + pub fn emit_custom(&self, name: &str, attributes: Vec) { + Scope { telemetry: self.clone(), state: self.shared.process.clone() } + .emit_custom(name, attributes); + } + + /// Set a pipeline-wide attribute (a consumer's `enduser.id`, an `acme.tenant`), attached to + /// every record of every session from now on unless the record — or its session — already + /// carries the key. `None` removes it. Scope-level identity goes through + /// [`Scope::set_attribute`]. + pub fn set_attribute(&self, key: &str, value: Option) { + let mut global = self.shared.global.lock().unwrap_or_else(|e| e.into_inner()); + global.retain(|a| a.key != key); + if let Some(value) = value { + global.push(Attribute::new(key, value)); + } + } + + /// The session's trace id as 32 hex characters — what every span and log record of this + /// pipeline carries. Print it (`lkt_…`) so support can find the session. + pub fn trace_id(&self) -> String { + self.shared.process.hex() + } + + /// Open a span: one attempt at an operation (`lk.connect`, `lk.publish`, …). Returns the + /// handle to record checkpoints and to end it with; `parent` nests it under another open span. + #[cfg(test)] + pub(crate) fn begin_span(&self, name: &str, kind: SpanKind, parent: Option) -> u64 { + self.begin_span_in(name, kind, parent, &self.shared.process) + } + + pub(crate) fn begin_span_in( + &self, + name: &str, + kind: SpanKind, + parent: Option, + session: &Arc, + ) -> u64 { + // After the opt-out nothing opens: 0 is never a span id, so the span is detached. + if self.shared.revoked() { + return 0; + } + self.shared.before_commit(); + self.shared.spans.lock().unwrap_or_else(|e| e.into_inner()).begin_in( + name, + kind, + parent, + session.clone(), + ) + } + + #[cfg(test)] + pub(crate) fn add_span_event(&self, span: u64, name: &str, attributes: Vec) { + self.shared.add_span_event(span, name, attributes); + } + + #[cfg(test)] + pub(crate) fn end_span( + &self, + span: u64, + outcome: SpanOutcome, + error_type: Option, + attributes: Vec, + ) { + self.shared.end_span(span, outcome, error_type, attributes); + } + + /// Push one `getStats()` reading. Readings are windowed on device into `lk.rtc.stats.sample` + /// events (see `stats_window_ms`); they never count against the flood guard. + pub fn record_stats(&self, sample: RtcStatsSample) { + self.record_stats_in(sample, &self.shared.process); + } + + pub(crate) fn record_stats_in(&self, sample: RtcStatsSample, session: &Arc) { + if self.shared.revoked() || !self.collects(session) { + return; + } + self.shared.before_commit(); + self.shared.windows.lock().unwrap_or_else(|e| e.into_inner()).record_in(sample, session); + } + + /// Whether `session` still collects: not once its project is known to receive nothing + /// (self-hosted server, no observability grant, disabled by the collector) — its records are + /// dropped at the door, never written to the cache. + fn collects(&self, session: &ScopeState) -> bool { + let route = session.route(); + self.shared.destinations.lock().unwrap_or_else(|e| e.into_inner()).alive(route.as_deref()) + } + + /// Change something a session's RTC windows captured when they opened (its correlation + /// attributes, its project), atomically with respect to readings, which record under the same + /// lock: under the windows lock, check the change is `accept`ed, close the session's open + /// windows, then `apply` it. Returns whether it was accepted. + pub(crate) fn with_windows_split( + &self, + session: &ScopeState, + accept: impl FnOnce() -> bool, + apply: impl FnOnce(), + ) -> bool { + let closed = { + let mut windows = self.shared.windows.lock().unwrap_or_else(|e| e.into_inner()); + if !accept() { + return false; + } + let closed = windows.split(session); + apply(); + closed + }; + for window in closed { + if self.shared.store.push(window) { + self.shared.overflow(); + } + } + true + } + + /// Close and forget the RTC windows of one track (or, with `None`, of the whole session): + /// the last partial window is queued now instead of at the next tick. + pub(crate) fn retire_stats(&self, session: &Arc, track_sid: Option<&str>) { + let closed = self + .shared + .windows + .lock() + .unwrap_or_else(|e| e.into_inner()) + .retire(session, track_sid); + for window in closed { + // Windows bypass the flood guard: they are the pipeline's own, bounded output. + if self.shared.store.push(window) { + self.shared.overflow(); + } + } + } + + /// Tell the pipeline what the device looks like. Emits the `lk.device.*.changed` events for + /// whatever differs from the last state (everything, the first time) and re-tunes the + /// pipeline: pressure stretches the cadence up to 4× ([`DeviceState::cadence_factor`]), a + /// constrained network or a nearly empty battery holds uploads + /// ([`DeviceState::holds_uploads`]), and entering the background flushes once right away. + pub fn set_device_state(&self, state: DeviceState) { + let mut previous = self.shared.device.lock().unwrap_or_else(|e| e.into_inner()); + for event in state.change_events(previous.as_ref()) { + self.emit(event); + } + let before = previous.replace(state).unwrap_or_default(); + drop(previous); + // Only three things are worth a wake-up: entering the background (export now, the app + // may be suspended), a hold changing (re-evaluate; nothing new is encoded for it) and the + // cadence changing (the pending tick is re-derived: relief brings it forward). + let offline = |s: &DeviceState| s.network == crate::NetworkType::Unavailable; + if state.app_state == crate::AppState::Background && before.app_state != state.app_state { + self.shared.backgrounded.store(true, Ordering::SeqCst); + self.shared.wake.notify_one(); + } else if before.holds_uploads() != state.holds_uploads() + || offline(&before) != offline(&state) + { + self.shared.release(); + } else if before.cadence_factor() != state.cadence_factor() { + self.shared.wake.notify_one(); + } + } + + /// Cache everything queued and upload the whole cache — no per-pass budget — as far as the + /// holds and pauses in force allow. Returns when that pass is over. + pub async fn flush(&self) { + self.command(Command::Flush).await; + } + + /// Flush, then stop the exporter, and return once it has stopped. The exporter gives the + /// network `export_timeout_ms`, then cancels the request in flight and exits; what did not + /// go out stays cached (with a [`FileCache`], for the next launch). Events emitted afterwards + /// are never exported. If the exporter is not running at all, this gives up a second after + /// that bound. + pub async fn shutdown(&self) { + let bound = Duration::from_millis(self.shared.config.export_timeout_ms.max(1)); + let _ = timeout(bound + SHUTDOWN_GRACE, self.command(Command::Shutdown)).await; + } + + /// Pipeline health: drops by reason, uploads, cached batches. The same numbers ride to the + /// backend as `lk.telemetry.report` events whenever something went wrong. + pub fn stats(&self) -> TelemetryStats { + TelemetryStats::new( + self.shared.counters.snapshot(), + self.shared.cache.pending().len() as u64, + *self.shared.status.lock().unwrap_or_else(|e| e.into_inner()), + ) + } + + /// Opt-out: stop the exporter without another upload and delete everything it holds — the + /// queue, open and finished spans, RTC windows and every cached batch, on disk included. + /// Events emitted afterwards go nowhere. Returns `false` when the storage refused to delete + /// something (it is logged; the files stay until a later purge or configure succeeds). + pub async fn purge(&self) -> bool { + self.shared.revoked.store(true, Ordering::SeqCst); + self.shared.clear(); + let bound = Duration::from_millis(self.shared.config.export_timeout_ms.max(1)); + let _ = timeout(bound, self.command(Command::Purge)).await; + // Whatever an upload in flight put back while the exporter wound down. + self.shared.clear() + } + + /// A handle that can reach the exporter without keeping it alive (the opt-out uses it to + /// cancel and await every generation). + pub(crate) fn weak_commands(&self) -> mpsc::WeakUnboundedSender { + self.commands.downgrade() + } + + async fn command(&self, make: impl FnOnce(oneshot::Sender<()>) -> Command) { + let (done, wait) = oneshot::channel(); + if self.commands.send(make(done)).is_ok() { + let _ = wait.await; + } + } +} + +/// Which LiveKit client SDK is reporting: `service.name` becomes `livekit-client-`. +#[cfg_attr(feature = "uniffi", derive(uniffi::Enum))] +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum Sdk { + Swift, + Android, + Flutter, + ReactNative, + Unity, + Rust, +} + +impl Sdk { + fn as_str(self) -> &'static str { + match self { + Self::Swift => "swift", + Self::Android => "android", + Self::Flutter => "flutter", + Self::ReactNative => "react-native", + Self::Unity => "unity", + Self::Rust => "rust", + } + } +} + +/// The reporting SDK and the device it runs on. Lowered to semconv: `service.name`, +/// `service.version`, `os.name`, `os.version`, `device.model.identifier`. +#[cfg_attr(feature = "uniffi", derive(uniffi::Record))] +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct TelemetryResource { + pub sdk: Sdk, + pub sdk_version: String, + pub os_name: String, + pub os_version: String, + #[cfg_attr(feature = "uniffi", uniffi(default))] + pub device_model: Option, +} + +impl TelemetryResource { + pub(crate) fn attributes(&self) -> Vec { + let mut out = vec![ + Attribute::new("service.name", format!("livekit-client-{}", self.sdk.as_str())), + Attribute::new("service.version", self.sdk_version.as_str()), + Attribute::new("os.name", self.os_name.as_str()), + Attribute::new("os.version", self.os_version.as_str()), + ]; + if let Some(model) = &self.device_model { + out.push(Attribute::new("device.model.identifier", model.as_str())); + } + out + } +} + +/// Lower the typed resource, then fill in the `telemetry.sdk.*` attributes and a fallback +/// `service.name`. Attributes already present (the open bag) win. +fn add_sdk_resource(resource: &mut Vec, sdk: Option<&TelemetryResource>) { + for attribute in sdk.map(TelemetryResource::attributes).unwrap_or_default() { + if !resource.iter().any(|a| a.key == attribute.key) { + resource.push(attribute); + } + } + let defaults = [ + ("service.name", "livekit-client"), + ("telemetry.sdk.name", env!("CARGO_PKG_NAME")), + ("telemetry.sdk.language", "rust"), + ("telemetry.sdk.version", env!("CARGO_PKG_VERSION")), + ]; + for (key, value) in defaults { + if !resource.iter().any(|a| a.key == key) { + resource.push(Attribute::new(key, value)); + } + } +} diff --git a/livekit-telemetry/src/trace.rs b/livekit-telemetry/src/trace.rs new file mode 100644 index 000000000..72aaf97a5 --- /dev/null +++ b/livekit-telemetry/src/trace.rs @@ -0,0 +1,525 @@ +//! The SDK's span vocabulary and the span itself, owned by the core so every platform names, +//! times and describes an operation the same way. Every call is synchronous and stamps the clock +//! inside, so the only skew is the FFI call itself; context propagation (the "current" span) +//! stays with the platform runtime, which is the one thing a core cannot do. + +use crate::device::snake; +use std::{ + sync::{Arc, Mutex, Weak}, + time::Duration, +}; + +use tokio::time::Instant; + +use crate::span::SpanKind; +use crate::{ + scope::ScopeState, telemetry::Shared, Attribute, AttributeValue, SpanOutcome, Telemetry, + TrackKind, +}; + +/// What triggered a reconnect cycle: the protocol's `ReconnectReason` values, plus the ones only +/// a client knows. +#[cfg_attr(feature = "uniffi", derive(uniffi::Enum))] +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum ReconnectReason { + /// The signaling socket closed (`RR_SIGNAL_DISCONNECTED`). + SignalDisconnected, + PublisherFailed, + SubscriberFailed, + /// A peer connection failed and the platform did not say which. + TransportFailed, + SwitchCandidate, + /// The device moved to another network (Wi-Fi ↔ cellular): the client noticed first. + NetworkChanged, + /// A test or a debug menu asked for it. + Debug, + Unknown, +} + +impl ReconnectReason { + /// The protocol's `ReconnectReason` number (`RR_*`). + pub fn from_proto(value: i32) -> Self { + match value { + 1 => Self::SignalDisconnected, + 2 => Self::PublisherFailed, + 3 => Self::SubscriberFailed, + 4 => Self::SwitchCandidate, + _ => Self::Unknown, + } + } +} + +/// What an SDK operation is. The kind follows from the name: connects talk to the server +/// (`client`), the rest is internal work. +#[cfg_attr(feature = "uniffi", derive(uniffi::Enum))] +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum SpanName { + Connect, + Reconnect { + reason: ReconnectReason, + }, + Publish, + Subscribe, + /// An app-defined span; its name is the consumer's. + Custom { + name: String, + }, +} + +impl SpanName { + /// The span name on the wire (`lk.connect`, …). + pub fn label(&self) -> &str { + match self { + Self::Connect => "lk.connect", + Self::Reconnect { .. } => "lk.reconnect", + Self::Publish => "lk.publish", + Self::Subscribe => "lk.subscribe", + Self::Custom { name } => name, + } + } + + fn kind(&self) -> SpanKind { + match self { + Self::Connect | Self::Reconnect { .. } => SpanKind::Client, + _ => SpanKind::Internal, + } + } + + fn attributes(&self) -> Vec { + match self { + Self::Reconnect { reason } => { + vec![Attribute::new("lk.reconnect.reason", snake(reason))] + } + // The user-initiated connect is attempt 1; a platform that retries sets it again. + Self::Connect => vec![Attribute::new("lk.connect.attempt", 1i64)], + _ => Vec::new(), + } + } +} + +/// A checkpoint inside a span. One vocabulary for all spans; the core does not police which step +/// belongs to which span, the dashboard does. +#[cfg_attr(feature = "uniffi", derive(uniffi::Enum))] +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum SpanStep { + WsOpen, + Signal, + JoinRecv, + PcCreated, + OfferSent, + AnswerSent, + Engine, + PcConnected, + RoomConnected, + Subscribed, + FirstMedia, + /// One reconnect attempt; also sets `lk.reconnect.attempts` and `lk.reconnect.mode`. + Attempt { + number: u32, + full: bool, + }, + Custom { + name: String, + }, +} + +impl SpanStep { + fn label(&self) -> String { + match self { + Self::WsOpen => "ws_open".into(), + Self::Signal => "signal".into(), + Self::JoinRecv => "join_recv".into(), + Self::PcCreated => "pc_created".into(), + Self::OfferSent => "offer_sent".into(), + Self::AnswerSent => "answer_sent".into(), + Self::Engine => "engine".into(), + Self::PcConnected => "pc_connected".into(), + Self::RoomConnected => "room_connected".into(), + Self::Subscribed => "subscribed".into(), + Self::FirstMedia => "first_media".into(), + Self::Attempt { number, full } => { + format!("attempt {number} {}", if *full { "full" } else { "quick" }) + } + Self::Custom { name } => name.clone(), + } + } + + fn attributes(&self) -> Vec { + match self { + Self::Attempt { number, full } => vec![ + Attribute::new("lk.reconnect.attempts", *number as i64), + Attribute::new("lk.reconnect.mode", if *full { "full" } else { "quick" }), + ], + _ => Vec::new(), + } + } +} + +/// The protocol's `TrackSource`. +#[cfg_attr(feature = "uniffi", derive(uniffi::Enum))] +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum TrackSource { + Camera, + Microphone, + ScreenShare, + ScreenShareAudio, + Unknown, +} + +/// The track a publish or subscribe span is about. +#[cfg_attr(feature = "uniffi", derive(uniffi::Record))] +#[derive(Debug, Clone, PartialEq)] +pub struct SpanTrack { + /// Unknown until the server assigns it (publish): set the track again once it is. + #[cfg_attr(feature = "uniffi", uniffi(default))] + pub sid: Option, + pub kind: TrackKind, + pub source: TrackSource, + /// The publisher, for subscribe spans. + #[cfg_attr(feature = "uniffi", uniffi(default))] + pub remote_identity: Option, +} + +impl SpanTrack { + fn attributes(&self) -> Vec { + let mut out = vec![ + Attribute::new("lk.track.kind", format!("{:?}", self.kind).to_lowercase()), + Attribute::new("lk.track.source", snake(self.source)), + ]; + if let Some(sid) = &self.sid { + out.push(Attribute::new("lk.track.sid", sid.as_str())); + } + if let Some(identity) = &self.remote_identity { + out.push(Attribute::new("lk.participant.remote_identity", identity.as_str())); + } + out + } +} + +/// A span's identity in its session's trace, for log correlation. +#[cfg_attr(feature = "uniffi", derive(uniffi::Record))] +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct TraceContext { + pub trace_id: String, + pub span_id: u64, +} + +/// A span's tie to its pipeline. Weak: a span a Room still holds (a pending subscribe) never +/// keeps the pipeline — and its exporter — alive. +struct Bound { + shared: Weak, + /// The session the span belongs to (weak, like `shared`). + session: Weak, + trace_id: String, + id: u64, +} + +#[derive(Default)] +struct State { + /// (label, offset from start) + steps: Vec<(String, Duration)>, + attributes: Vec, + ended: Option<(SpanOutcome, Duration)>, +} + +/// One attempt at an SDK operation: a timed interval with checkpoints, attributes and an outcome. +/// Bound to a session it is exported as an OTLP span when it ends; detached it still times and +/// describes itself, so the console line looks the same with telemetry off. +pub struct Span { + name: SpanName, + started: Instant, + bound: Option, + state: Mutex, +} + +/// What a span keeps before it is exported: the OTel default limits (128 events, 128 +/// attributes), and no label or value longer than an app may hand over elsewhere. +const MAX_STEPS: usize = 128; +const MAX_ATTRIBUTES: usize = 128; + +/// Set an attribute; a new key beyond [`MAX_ATTRIBUTES`] is ignored (the OTel limit). Returns +/// `false` when the key or a string value is over its length limit: not kept, the caller counts it. +fn upsert(attributes: &mut Vec, attribute: Attribute) -> bool { + let too_long = attribute.key.len() > crate::event::MAX_KEY_BYTES + || matches!(&attribute.value, AttributeValue::Str(v) if v.len() > crate::event::MAX_VALUE_BYTES); + if too_long { + return false; + } + let known = attributes.iter().any(|a| a.key == attribute.key); + if !known && attributes.len() >= MAX_ATTRIBUTES { + return true; + } + attributes.retain(|a| a.key != attribute.key); + attributes.push(attribute); + true +} + +impl Span { + /// Timings and a description only; nothing is exported. + pub fn detached(name: SpanName) -> Arc { + Arc::new(Self::new(name, None)) + } + + pub(crate) fn bound( + name: SpanName, + parent: Option, + telemetry: Telemetry, + session: &Arc, + ) -> Arc { + // A caller-named span over the name limit is not recorded (counted as invalid). + if name.label().len() > crate::event::MAX_NAME_BYTES { + telemetry.count_invalid(); + return Arc::new(Self::new(name, None)); + } + let id = telemetry.begin_span_in(name.label(), name.kind(), parent, session); + if id == 0 { + return Arc::new(Self::new(name, None)); // opted out: detached + } + let shared = Arc::downgrade(&telemetry.shared); + let bound = Bound { shared, session: Arc::downgrade(session), trace_id: session.hex(), id }; + Arc::new(Self::new(name, Some(bound))) + } + + fn new(name: SpanName, bound: Option) -> Self { + // A name over the limit is not kept, detached spans included: only a fixed placeholder. + let name = match name { + SpanName::Custom { name } if name.len() > crate::event::MAX_NAME_BYTES => { + SpanName::Custom { name: "invalid".to_owned() } + } + name => name, + }; + let state = State { attributes: name.attributes(), ..State::default() }; + Self { name, started: Instant::now(), bound, state: Mutex::new(state) } + } + + fn lock(&self) -> std::sync::MutexGuard<'_, State> { + self.state.lock().unwrap_or_else(|e| e.into_inner()) + } + + /// What this span is. + pub fn name(&self) -> SpanName { + self.name.clone() + } + + /// The span name on the wire. + pub fn label(&self) -> String { + self.name.label().to_owned() + } + + /// A checkpoint, stamped now. Ignored once the span has ended. + pub fn step(&self, step: SpanStep) { + let at = self.started.elapsed(); + let label = step.label(); + if label.len() > crate::event::MAX_NAME_BYTES { + return self.count_invalid(1); // a caller-named step over the name limit: not kept + } + { + let mut state = self.lock(); + if state.ended.is_some() || state.steps.len() >= MAX_STEPS { + return; + } + state.steps.push((label.clone(), at)); + for attribute in step.attributes() { + upsert(&mut state.attributes, attribute); // the core's own: always within limits + } + } + if let Some(bound) = &self.bound { + if let Some(shared) = bound.shared.upgrade() { + shared.add_span_event(bound.id, &label, Vec::new()); + } + } + } + + /// The open bag, for app-defined spans and one-off details. Replaces an existing key. + pub fn set_attribute(&self, key: String, value: AttributeValue) { + if !upsert(&mut self.lock().attributes, Attribute::new(key, value)) { + self.count_invalid(1); + } + } + + /// Count caller strings this span did not keep (over their length limits) as `invalid`. + fn count_invalid(&self, n: u64) { + if n == 0 { + return; + } + if let Some(shared) = self.bound.as_ref().and_then(|b| b.shared.upgrade()) { + crate::stats::Counters::add(&shared.counters.invalid, n); + } + } + + /// The track the operation is about (its sid may arrive later: set it again). + pub fn set_track(&self, track: SpanTrack) { + let invalid = { + let mut state = self.lock(); + track + .attributes() + .into_iter() + .filter(|a| !upsert(&mut state.attributes, a.clone())) + .count() + }; + self.count_invalid(invalid as u64); + // A published track's sid is known: its session polls fast until the first outbound + // reading arrives (see `Scope::stats_poll_interval_ms`). + if let (SpanName::Publish, Some(sid)) = (&self.name, &track.sid) { + if let Some(session) = self.bound.as_ref().and_then(|b| b.session.upgrade()) { + session.await_first_outbound(sid); + } + } + } + + /// End once; later calls are no-ops. `error` becomes `error.type` and the status message. + pub fn end(&self, outcome: SpanOutcome, error: Option) { + let at = self.started.elapsed(); + let attributes = { + let mut state = self.lock(); + if state.ended.is_some() { + return; + } + state.ended = Some((outcome, at)); + state.attributes.clone() + }; + if let Some(bound) = &self.bound { + if let Some(shared) = bound.shared.upgrade() { + // `error.type` is a type name, not a message: an over-long one is replaced. + let error = error.map(|e| { + if e.len() > crate::event::MAX_NAME_BYTES { + crate::stats::Counters::add(&shared.counters.invalid, 1); + "invalid".to_owned() + } else { + e + } + }); + shared.end_span(bound.id, outcome, error, attributes); + } + // The platform's console line, the same on every SDK (debug: FFI log path). + log::debug!("{}", self.describe()); + } + } + + /// End with an error; `error` becomes `error.type`. + pub fn fail(&self, error: String) { + self.end(SpanOutcome::Error, Some(error)); + } + + /// End as cancelled (status unset, `lk.outcome = cancelled`). + pub fn cancel(&self) { + self.end(SpanOutcome::Cancelled, None); + } + + /// Whether the span has ended. + pub fn is_ended(&self) -> bool { + self.lock().ended.is_some() + } + + /// How it ended, once it has. + pub fn outcome(&self) -> Option { + self.lock().ended.map(|(outcome, _)| outcome) + } + + /// `None` for a detached span. + pub fn context(&self) -> Option { + self.bound.as_ref().map(|b| TraceContext { trace_id: b.trace_id.clone(), span_id: b.id }) + } + + /// Seconds from start to the end, or to the last step while still running. + pub fn total_secs(&self) -> f64 { + let state = self.lock(); + state + .ended + .map(|(_, at)| at) + .or_else(|| state.steps.last().map(|(_, at)| *at)) + .unwrap_or_default() + .as_secs_f64() + } + + /// `lk.connect: ws_open +1.49s, signal +0.03s, total 1.83s, ok` — the same line on every + /// platform, for the console when a span ends. + pub fn describe(&self) -> String { + let state = self.lock(); + let mut parts = Vec::with_capacity(state.steps.len() + 2); + let mut previous = Duration::ZERO; + for (label, at) in &state.steps { + parts.push(format!("{label} +{:.2}s", at.saturating_sub(previous).as_secs_f64())); + previous = *at; + } + let total = state.ended.map(|(_, at)| at).unwrap_or(previous); + parts.push(format!("total {:.2}s", total.as_secs_f64())); + if let Some((outcome, _)) = state.ended { + parts.push(outcome.as_str().to_owned()); + } + format!("{}: {}", self.name.label(), parts.join(", ")) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test(start_paused = true)] + async fn a_detached_span_times_and_describes_itself() { + let span = + Span::detached(SpanName::Reconnect { reason: ReconnectReason::SignalDisconnected }); + tokio::time::advance(Duration::from_millis(1490)).await; + span.step(SpanStep::Attempt { number: 1, full: false }); + tokio::time::advance(Duration::from_millis(30)).await; + span.step(SpanStep::WsOpen); + assert_eq!( + span.describe(), + "lk.reconnect: attempt 1 quick +1.49s, ws_open +0.03s, total 1.52s" + ); + tokio::time::advance(Duration::from_millis(10)).await; + span.end(SpanOutcome::Ok, None); + span.fail("late".into()); + assert_eq!(span.outcome(), Some(SpanOutcome::Ok), "ends once"); + assert_eq!( + span.describe(), + "lk.reconnect: attempt 1 quick +1.49s, ws_open +0.03s, total 1.53s, ok" + ); + assert!(span.context().is_none()); + let attributes = span.lock().attributes.clone(); + let get = |key: &str| attributes.iter().find(|a| a.key == key).map(|a| a.value.clone()); + assert_eq!( + get("lk.reconnect.reason"), + Some(AttributeValue::Str("signal_disconnected".into())) + ); + assert_eq!(get("lk.reconnect.attempts"), Some(AttributeValue::Int(1))); + assert_eq!(get("lk.reconnect.mode"), Some(AttributeValue::Str("quick".into()))); + } + + #[test] + fn a_track_sets_its_attributes_and_can_learn_its_sid_later() { + let span = Span::detached(SpanName::Publish); + span.set_track(SpanTrack { + sid: None, + kind: TrackKind::Video, + source: TrackSource::Camera, + remote_identity: None, + }); + span.set_track(SpanTrack { + sid: Some("TR_1".into()), + kind: TrackKind::Video, + source: TrackSource::Camera, + remote_identity: None, + }); + let attributes = span.lock().attributes.clone(); + let get = |key: &str| attributes.iter().find(|a| a.key == key).map(|a| a.value.clone()); + assert_eq!(get("lk.track.kind"), Some(AttributeValue::Str("video".into()))); + assert_eq!(get("lk.track.sid"), Some(AttributeValue::Str("TR_1".into()))); + assert_eq!(attributes.iter().filter(|a| a.key == "lk.track.source").count(), 1); + } + + /// Finding r1-12: what a span retains before export is bounded, however long it lives. + #[test] + fn retained_span_state_is_bounded() { + let span = Span::detached(SpanName::Custom { name: "work".into() }); + for n in 0..1000 { + span.step(SpanStep::Custom { name: format!("step {n}") }); + span.set_attribute(format!("k{n}"), AttributeValue::Str("v".into())); + } + span.step(SpanStep::Custom { name: "x".repeat(1000) }); + span.set_attribute("k0".into(), AttributeValue::Str("x".repeat(2000))); + let state = span.lock(); + assert_eq!((state.steps.len(), state.attributes.len()), (128, 128)); + assert!(state.attributes.iter().all(|a| a.value.size_hint() <= 1024)); + } +} From c16912c32668c2b5118527ed91dbf701f0effae8 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?B=C5=82az=CC=87ej=20Pankowski?= <86720177+pblazej@users.noreply.github.com> Date: Thu, 1 Oct 2026 12:34:34 +0200 Subject: [PATCH 2/2] docs(telemetry): specify the upload policy, app API and typed surface --- livekit-telemetry/SPEC.md | 242 ++++++++++++++++++++++++++++++++++++++ 1 file changed, 242 insertions(+) diff --git a/livekit-telemetry/SPEC.md b/livekit-telemetry/SPEC.md index c2428964e..0da641e92 100644 --- a/livekit-telemetry/SPEC.md +++ b/livekit-telemetry/SPEC.md @@ -18,6 +18,15 @@ Set once per pipeline (`TelemetryConfig.resource`): ## Pipeline, scopes and destination +One pipeline per process — started at SDK init, so audio pre-initialization, permission failures +and connect attempts that never reach a server are captured — and one **scope** per room (one +call). A scope is a trace id plus the attributes attached to its records (`lk.room.sid`, +`lk.participant.identity`, …); spans, RTC windows and events are filed under the scope that +produced them, and `session.id` (OTel semconv) is written on every record as an attribute. A log +record emitted inside a room's span is filed under that room's scope; anything emitted outside +a scope — device state, pre-room errors, self-telemetry — belongs to the pipeline's own process +scope. Scopes are not ended: a room's last record is simply its last. + ### Destination and credentials The platform passes exactly two things per room, through `Scope::set_server(url, token)`: the @@ -228,6 +237,104 @@ judgement and `qualityLimitationReason` is WebRTC's. The core also paces the pla `getStats()` polling (`Scope::stats_poll_interval_ms`): every second while a subscribe waits for its first media, else twice per (stretched) window — 30 s by default. +## Upload policy — telemetry never wins over media + +Uploads are shaped, not just batched: + +- **One request in flight**, oldest batch first, polled by the exporter's loop: commands, ticks, + subscribe deadlines and wake-ups are served while it is out. Each answer is classified by the + core (see *Collector answers*); a transport returns status, headers and body and fails only + without a response. A pause stops uploads **to that destination** — other projects carry on — + never collection: records keep landing in the cache and ship when it lifts. +- **Retry:** local backoff starts at 1 s and doubles per consecutive failure up to 60 s, with + full jitter (a uniform wait in `[0, backoff]`). A delay the server names — `Retry-After` + (delay-seconds or HTTP-date, RFC 9110 §10.2.3) or `RetryInfo.retry_delay` — is honored in + full even when longer, validated (garbage ignored, negative → now, overflow ignored) and + clamped only to the 24 h age limit; neither shutdown nor a hold's escape cuts it short. + Running out of patience never deletes: only the cache's age and size bound what is kept. +- **Budget:** while a Room is in a call (it has a server and has not disconnected), at most + `max_batches_per_upload` (default 4) requests per flush interval — an allowance refilled at + each tick, 413 halves included, each request at most `max_batch_bytes` of protobuf before gzip + — however many wake-ups happen; a backlog (offline period, previous launch) waits its turn + oldest first. With no Room in a call nothing is metered: each pass sends the whole cache. A new + batch is encoded only at the tick, when the queue crosses `flush_threshold_bytes`, or when the + app enters the background; a lifted route, token or hold only re-runs a pass over what is + cached (within the allowance); a subscribe or a cadence change only re-reads deadlines. + Entering the background and an explicit `flush()` drain the whole cache without the allowance, + within every hold and pause. `shutdown` drains without the budget within + `export_timeout_ms`, then cancels the request on the wire and stops the exporter; what did not + go out stays cached. +- **Holds** — nothing is sent, everything keeps flowing into the write-ahead cache. The policy is + evaluated before every request, not once per pass, and a hold changing (or the app entering the + background, or the cadence changing) wakes the exporter: + - *hard* (no escape, for as long as they last): the device is offline; no usable token + (missing, expired, grant-less, refused); the project disabled data recording; the opt-out. + - *soft* (at most 60 s, on its own clock rather than the next tick, then one batch goes out and + the hold starts over — the cap that bounds the policy when its signals lie): an `lk.connect` or `lk.reconnect` span is open (signaling + and ICE/DTLS own the uplink); Low Data Mode / Data Saver; battery ≤ 10 % unplugged. The cap is + not scheduled while a hard hold is on; it resumes when the hard hold lifts. + `qualityLimitationDurations.bandwidth` is deliberately *not* a hold: WebRTC reports it for + minutes during a normal ramp-up and for as long as an encoder stalls. +- **Bytes:** bodies are gzipped (level 1, `Content-Encoding: gzip`) when cached, so a batch is + 5–10× smaller on disk and on the wire and a replay costs no CPU; a cached batch is CRC-checked + before it is sent. A request never carries more than `max_batch_size` (512) records, nor more + than `max_batch_bytes` (1 MiB) of encoded protobuf — checked on the encoded body, session + attributes and the self-report included (the report takes one of the `max_batch_size` + places); a batch over it is halved until it fits, and a single record over it is dropped and + counted as `oversized`. Caller strings a span or scope retains are bounded: names, keys and + error types ≤ 128 bytes, values and Room identities ≤ 1 KiB (anything longer is not kept and + is counted as `invalid`). When the queue reaches + `flush_threshold_bytes` (256 KiB) it is exported at once instead of at the next tick. +- **Backlog:** nested bounds, every eviction counted — the queue (`max_queue_size`, 2048 + records), the cache's size (`max_cache_bytes`, 4 MiB compressed) and file count (512 + batches), and age (24 h, enforced at start and while running; the monotonic clock must agree + for this launch's own batches, so a clock set forward cannot expire them). Oldest goes first. +- **Durability:** a batch is committed when the cache's `push` returns: `FileCache` writes a + `.tmp`, `fsync`s it, renames it into place and `fsync`s the directory (on Unix; elsewhere a + directory cannot be synced and that step is skipped). If any step fails the file is removed + and the batch is kept in memory instead (counted as a write error) — never claimed committed. + A 413 split is journaled (halves written, synced and renamed to `@.pend`, + directory synced, parent deleted = commit point, halves published); an interrupted split is + finished or rolled back when the cache opens, before anything is evicted, so every record is + there exactly once. A split that fails after its commit point has succeeded: its halves wait in + the journal and the next listing of the cache publishes them (so they upload, and an opt-out + purges and counts them). Splits and recovery share a process-wide lock, so a reconfigure's + cache on the same directory never sees the draining pipeline's split half-done; stray files + are swept only when a cache opens. A crash loses what was not committed yet — records since the last tick + (≤ one flush interval, ≤ 2048 queued), open spans, open RTC windows. Duplicates: a crash + between an answer and its delete sends that one batch again at the next launch; a delete that + fails without a crash leaves the batch pending deletion (never sent again in this launch, but + sent once more after a restart if it is still there). A disk that refuses writes gets batches + kept in memory (bounded by `max_cache_bytes`, oldest evicted and counted), which survive + failures but not the process. A batch that exists but cannot be read right now (file + protection, permissions) is kept, not counted corrupt. +- **Priority:** every request carries `Priority: u=7` (RFC 9218, lowest urgency) for HTTP/2+ + hops that implement it. +- **Redirects:** a 3xx the transport returns drops the batch. Transports must not forward + `Authorization` across origins; the `livekit-net` native client strips it when a redirect + changes the host or the port (tested). A scheme-only change on the same explicit port is not + covered by those tests. +- **Threads:** pipeline work runs on the SDK's runtime; none of it is on a media or UI thread and + `emit`/`record_stats` never block on it. + +### Collector answers + +| Answer | Batch | Pipeline | +|---|---|---| +| 2xx | removed | — | +| 2xx with OTLP `partial_success` | removed; rejected records counted, never retried | — | +| 400, 3xx, other 4xx | dropped, counted `rejected` | — | +| 413 | replaced by two halves in one cache transaction (both committed where the batch was, on disk stays on disk, nothing evicted in between) and retried at once, down to one record; if the cache cannot take the halves the batch stays whole; a lone oversized record is dropped, counted `oversized` | — | +| 401/403 "data recording is disabled by owner" | purged with the project's whole cache | project silent for the process | +| other 401/403 | kept | that token is never sent again; the project's batches wait for the next token | +| 404 | purged with the project's cache | project silent until its next token (a new connect or refresh) | +| 429 | kept | paused for `Retry-After`, else `RetryInfo`, else 60 s | +| 503 with `Retry-After`/`RetryInfo` | kept | paused for that long | +| 502, 503, 504; any error with `RetryInfo` (Cloud's retryable 500) | kept | paused for the named delay, else backoff | +| other 5xx | dropped, counted `rejected` | — | +| no answer: timeout, connection, DNS, TLS | kept | backoff; never a capability decision | +| invalid request (transport-side) | dropped, counted `rejected` | — | + ## Log records A `TelemetryEvent` with an empty `name` is a plain log record (OTLP log without `event_name`): @@ -250,6 +357,105 @@ forwarder the core's records go nowhere. Platforms must not feed forwarded Rust them, so they would be counted twice. `telemetry_log` is for the platform's own and WebRTC's lines. +## App API: custom events, correlation attributes, opt-out + +Three things an app can do; everything else is the SDK's. + +- **Custom events** are room-scoped: `Scope::emit_custom(name, attributes)` (Swift: + `room.emitTelemetryEvent(_:attributes:)`). The core prefixes the name with `custom.` + (`checkout` → `custom.checkout`), so a custom event can never collide with, or spoof, an + `lk.*` event and the backend can filter or quota the namespace as a whole. Severity is `info`; + custom events count against the flood guard like any discrete event. +- **Correlation attributes** are room-scoped too: `Scope::set_attribute(key, value | nil)` — + `app.call_id`, `enduser.id` — inherited by the room's subsequent logs, spans, events and RTC + windows. Values are snapshotted when a record is captured: a later change never rewrites what + is queued, and a change closes the Room's open RTC windows first, so a window never mixes + readings taken under two values. A custom event's own attributes override the room's. There is no public + process-wide attribute (a mutable global would leak across simultaneous rooms); the + process-level set is internal, for SDK metadata. +- **Limits:** names and keys ≤ 128 UTF-8 bytes, string values ≤ 1024 bytes, ≤ 64 attributes per + room and per event. SDK-owned keys are reserved: `lk.*` (room, participant, track, outcome …) + and `session.id`. Anything else is rejected — never truncated, a truncated id collides — and + counted as `lk.telemetry.dropped.invalid`. +- **Opt-out** is process-wide and synchronous: when `telemetry_disable()` returns, no scope is + handed out, every existing scope and instrument captures nothing and later configures are + refused; the purge — cancel scheduled work, delete everything not yet sent (queue, open spans + and windows, every cached batch on disk) — then runs on the core's runtime + (`telemetry_shutdown` / `telemetry_flush` await it; `global::disable()` returns it as a + future). + Install, instrument start/stop and the opt-out run under one lifecycle lock, so an instrument + never starts after the opt-out stopped it. Every pipeline generation still alive (one replaced + but still draining included) is revoked — the flag is checked under the lock of every structure + that commits data (queue, spans, RTC windows, the cache's write gate), which the purge also + takes, so nothing lands after it — then its exporter's request on the wire is cancelled and + the exporter awaited, and later configures in the process are refused. The pull queue never + hands out a request its exporter gave up on, whether the host awaits `next()` or polls + `try_next()` (Dart, from a timer: no Rust future ever holds a continuation of an isolate that + may die; such hosts pass no instruments either, so Rust never calls into them). Data already sent cannot be recalled: at most the + one batch per generation already on the wire, or already handed to a pull-queue host. If the + storage refuses a delete the purge is reported incomplete (logged; `global::disable` and + `Telemetry::purge` return `false`; the exported `telemetry_disable` returns nothing). Platforms + mark the call TODO pending the token discussion. + +## Flood guard + +Discrete events (`emit`) are capped at `max_events_per_10min` (default 300, design doc); what +exceeds it is dropped and reported as `lk.telemetry.dropped.rate_limited`. `lk.rtc.stats.sample` +windows and `lk.telemetry.report` are exempt. + +```yaml +event: lk.rtc.stats.sample +area: rtc +severity: info +cadence: one per track and direction per stats window (default 60 s, stretched by the cadence + factor); closed early on background, when the track leaves (`track_ended`), at + disconnect and at shutdown. Produced by the core from raw readings — one getStats() + report per peer connection (`record_peer_stats`), polled when the core says + (`stats_poll_interval_ms`). Keyed by session and track: two rooms receiving the same + track keep two windows. +attributes: + lk.track.sid: string + lk.track.kind: enum(audio | video) + lk.track.direction: enum(inbound | outbound) + lk.rtc.codec: string # mimeType, when known + lk.rtc.window_ms: int # actual window length + lk.rtc.samples: int # readings in the window + # cumulative counters — the last reading's value, monotonic (W3C webrtc-stats model) + lk.rtc.bytes: int + lk.rtc.packets: int + lk.rtc.packets_lost: int # inbound + lk.rtc.freeze_count: int # inbound video + lk.rtc.freezes_duration_ms: int # inbound video + lk.rtc.concealed_samples: int # inbound audio + lk.rtc.concealment_events: int # inbound audio + lk.rtc.jitter_buffer_delay_ms: int # inbound + lk.rtc.jitter_buffer_emitted_count: int # inbound + lk.rtc.quality_limitation.bandwidth_ms: int # outbound video + lk.rtc.quality_limitation.cpu_ms: int # outbound video + lk.rtc.quality_limitation.other_ms: int # outbound video + lk.rtc.pause_count: int # inbound video + lk.rtc.pauses_duration_ms: int # inbound video + lk.rtc.silent_concealed_samples: int # inbound audio + lk.rtc.interruption_count: int # inbound audio + lk.rtc.interruptions_duration_ms: int # inbound audio + # gauges — min / max / avg over the window + lk.rtc.jitter_ms.{min,max,avg}: double + lk.rtc.rtt_ms.{min,max,avg}: double # remote-inbound RTT for outbound, candidate-pair for inbound + lk.rtc.fps.{min,max,avg}: double # video + lk.rtc.audio_level.{min,max,avg}: double # audio +platforms: all +``` + +```yaml +event: lk.room.disconnected +area: session +severity: info (client_initiated) | warn (anything else) +cadence: once, when the Room leaves connected for good — never on a reconnect +attributes: + lk.disconnect.reason: enum(client_initiated | duplicate_identity | server_shutdown | participant_removed | room_deleted | state_mismatch | join_failure | migration | signal_close | room_closed | user_unavailable | user_rejected | sip_trunk_failure | connection_timeout | media_failure | agent_error | reconnect_failed | unknown) # the protocol's DisconnectReason, plus the client giving up +platforms: all +``` + ## Spans A span is **one attempt** at an operation. The scope (one Room connection lifetime, across @@ -325,3 +531,39 @@ attributes: lk.participant.remote_identity: string checkpoints: subscribed, first_media ``` + +## Typed surface + +Everything the SDKs have in common enters the core typed; the core owns the keys, the bodies and +the policy. `Attribute { key, value }` survives only as the open bag: `emit_custom`, +`set_attribute`, and `Span::set_attribute` for app-defined spans. + +| Platform calls | The core produces | +|---|---| +| `Scope::set_server(url, token)` — at connect and on every token refresh | the room's destination: `https:///observability/client/{logs,traces}/otlp/v0` for Cloud hosts, the token's grant and expiry read, its batches routed to its project with its own token | +| `Scope::emit_custom(name, attributes)`, `Scope::set_attribute(key, value)` | `custom.` events and correlation attributes, validated and snapshotted (see *App API*) | +| `telemetry_disable()` | the opt-out, in effect when it returns: capture stops; everything unsent is then deleted | +| `TelemetryConfig.sdk: TelemetryResource { sdk: Sdk, sdk_version, os_name, os_version, device_model }` | `service.name = livekit-client-`, `service.version`, `os.*`, `device.model.identifier`, plus `telemetry.sdk.*` | +| `log(LogRecord { severity, source: LogSource, body, logger, function, file, line, timestamp_ns, span_id })` | a record with `code.function.name`, `code.file.path`, `code.line.number`, `lk.log.source`, `lk.log.logger`; the per-source floor (WebRTC at `error`, own module never) | +| `Scope::set_room(RoomIdentity { sid, name, participant_sid, participant_identity })` | `lk.room.*`, `lk.participant.*` on every record of the scope | +| `Scope::start(SpanName, parent) -> Span`; `Span::detached(name)` | an OTLP span (`lk.connect` / `lk.reconnect` are `client`, the rest `internal`); `Reconnect { reason }` sets `lk.reconnect.reason` | +| `Span::step(SpanStep)` | a span event named `ws_open` … `room_connected`, `subscribed`, `first_media`, `attempt N quick|full` (which also sets `lk.reconnect.attempts` / `.mode`) | +| `Span::set_track(SpanTrack { sid, kind, source, remote_identity })` | `lk.track.sid`, `lk.track.kind`, `lk.track.source`, `lk.participant.remote_identity` | +| `Span::end(outcome, error)` / `fail(error)` / `cancel()` | status, `error.type`, `lk.outcome`; ending twice is a no-op | +| `Span::describe()` | `lk.connect: ws_open +1.49s, signal +0.03s, total 1.83s, ok` — the console line, identical on every platform | +| `Span::context()` | `TraceContext { trace_id, span_id }` for log correlation; `None` when detached | +| `device_event(DeviceEvent::{AudioRouteChanged, AudioInterruption, CaptureFailed})` | `lk.device.audio_route.changed`, `lk.device.audio.interruption`, `lk.device.capture.failed` with display bodies; every value is a shared enum (`AudioOutput`, `CaptureDevice`, `CaptureFailure`) | +| `Scope::subscribe_started(SpanTrack)`, `subscribed(SpanTrack)`, `subscribe_failed(sid, error_type)` | the `lk.subscribe` span, ended by the core at the first inbound reading with bytes, or `timed_out` after 30 s; tracks already in the room at join: `subscribe_started` at connect (a lone `subscribed` still opens the span, as a fallback) | +| `Scope::track_ended(sid)` | a pending subscribe ends (`cancelled`, or `timed_out` past its deadline); the track's last RTC window ships; its per-track state is retired | +| `Scope::record_peer_stats(Vec, tracks: {MediaStreamTrack id → sid}, ts)` | one peer connection's raw `getStats()` report → the core finds each track's RTP streams (`trackIdentifier`, or the `media-source` an `outbound-rtp` names), resolves codec / RTT, converts seconds to ms and records one sample per stream | +| `Scope::stats_poll_interval_ms()` | when to poll `getStats()` next: 1 s while a subscribe awaits first media, else half the stretched window | +| `Scope::record_stats_report(sid, kind, direction, Vec, ts)` | the same mapping for one track's report (platforms that poll per track) | +| `DisconnectReason::from_proto(i32)`, `ReconnectReason::from_proto(i32)` | the protocol numbers → the shared enums | +| `Scope::log(LogRecord)` | a record filed under the session without an ambient span (Dart has no task-local outside a zone); same floor and filters as `Telemetry::log` | +| `Scope::disconnected(DisconnectReason)` | `lk.room.disconnected` with `lk.disconnect.reason` — info when the client hung up, warn otherwise; pending subscribes end, the session's RTC windows ship and its state is retired | +| `RtcStatsSample.layer` (rid, ssrc or stats id) | simulcast layers folded into one monotonic series per track before windowing | + +Timing rule: span calls are synchronous and stamp the clock inside the core, so the only skew is +the FFI call. Anything that may cross an executor hop before reaching the core (a log record) carries +its own `timestamp_ns` from capture. Context propagation — the "current" span — stays with the +platform runtime (task-local, coroutine context, zone); that is the one piece a core cannot own.