diff --git a/Cargo.lock b/Cargo.lock index 4c73ce507..53f94fe5b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3816,11 +3816,19 @@ dependencies = [ name = "livekit-telemetry" version = "0.1.0" dependencies = [ + "async-trait", + "base64 0.22.1", + "livekit-net", "log", "opentelemetry-proto", "prost 0.14.4", + "prost-types 0.14.4", "rand 0.9.5", + "serde_json", + "thiserror 2.0.19", + "tokio", "uniffi", + "url", ] [[package]] diff --git a/livekit-telemetry/Cargo.toml b/livekit-telemetry/Cargo.toml index 7364e25b3..ad1fe9df6 100644 --- a/livekit-telemetry/Cargo.toml +++ b/livekit-telemetry/Cargo.toml @@ -8,17 +8,34 @@ edition.workspace = true repository.workspace = true [dependencies] +tokio = { workspace = true, default-features = false, features = ["macros", "sync", "time"] } log = { workspace = true } +thiserror = { workspace = true } +async-trait = "0.1" prost = { workspace = true } +# `google.protobuf.Any`/`Duration` for the `google.rpc.Status` error body (RetryInfo). +prost-types = { workspace = true } rand = { workspace = true } +# Unverified JWT claims (observability grant, expiry); both already linked by livekit-api. +base64 = "0.22" +serde_json = { workspace = true } +# Server URLs are parsed, never string-matched, before a token may follow them. +url = "2.3" # 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"] } +livekit-net = { workspace = true, optional = true } uniffi = { workspace = true, features = ["scaffolding-ffi-buffer-fns"], optional = true } [features] +# Default HTTP transport over the pluggable `livekit-net` client (native backend, or one +# the host registered with `livekit_net::set_http_client`). +net = ["dep:livekit-net"] uniffi = ["dep:uniffi"] +[dev-dependencies] +tokio = { workspace = true, default-features = false, features = ["rt", "rt-multi-thread", "macros", "sync", "time", "test-util", "net", "io-util"] } + # How CI checks this crate's features, read by # `.github/workflows/feature-combinations-curated.yml` via `cargo metadata`. [package.metadata.feature-combinations] diff --git a/livekit-telemetry/SPEC.md b/livekit-telemetry/SPEC.md index 07a350f87..fabef938d 100644 --- a/livekit-telemetry/SPEC.md +++ b/livekit-telemetry/SPEC.md @@ -16,6 +16,55 @@ Set once per pipeline (`TelemetryConfig.resource`): | `device.model.identifier` | platform SDK | `iPhone16,1` | | `telemetry.sdk.name/language/version` | core | `livekit-telemetry`, `rust`, `0.1.0` | +## Pipeline, scopes and destination + +### Destination and credentials + +The platform passes exactly two things per room, through `Scope::set_server(url, token)`: the +LiveKit server URL the room connects to and the participant token it connects with — at connect, +and again with **every** refreshed token (the SFU sends one right after join, then every few +minutes). The call is cheap and idempotent. There is no client-side endpoint, header or sink. + +- **Ingest URL:** derived from the server URL's host, `https:///observability/client/{logs,traces}/otlp/v0`, + for LiveKit Cloud hosts only (`*.livekit.cloud`, including `*.staging.livekit.cloud`). Any other + host (self-hosted, OSS) has no ingest: one local warning, and that room's records are dropped + at the door instead of cached. +- **Token:** `Authorization: Bearer `. The core reads the token's *unverified* claims: the + observability grant (`observability.write` or `observability.clientWrite`) and `exp`. It never + sends a token known to be expired or refused; batches wait for the next token instead + (a hard hold). A project whose first token has no grant never opted in: nothing is collected + for it. A refreshed token without the grant (today's SFU drops it) does not replace a granted + one that is still valid. Tokens live in memory only — never in a batch, never on disk. +- **Server URL validation:** parsed with WHATWG URL rules, never string-matched. A token is only + sent to `https://.livekit.cloud/…` built from the parsed host alone: TLS scheme + (`wss`/`https`), a domain under the Cloud suffix with a label of its own, the default port, no + userinfo. +- **Ownership:** every record captures its owner when it is captured — its session and the + project that session is routed to at that moment — and keeps it: a Room that reconnects to + another project takes nothing queued or cached along. Credentials are keyed by (project, + session): a live Room uploads with its **own** latest token for that project, never another + Room's. A Room's records captured before it had a server go to its own first project, never to + another Room's. Process-level records (device state, pre-room errors, self-telemetry) go to + the project most recently handed a token, with that project's latest token; so do batches from + a previous launch, which wait — up to the 24 h age limit — for a token of the same project. + The answer to a request is attributed to the project it was sent to — its 404, disable or + pause never lands on another project. A session's credentials stay while the session is alive + (its Room, or records, windows or spans still referencing it) — for every project it was routed + to, so records captured for an earlier project can still be sent; that is the residual: a live + Room keeps one credential per project it has used — and while cached batches need them; + project-level copies only while a live session is routed there or the backlog has batches for + it. A token the collector refused is recorded by identity (a hash) and never sent again from + any slot, until it expires (a token without `exp` stays refused for the process; like the + per-host project table, which keeps one small entry per host ever seen, that grows only with + what a process meets — bounded in practice, not by a cap). A past, + negative or non-numeric `exp` counts as expired. +- **Waiting** for a first destination or a usable token is uncapped, bounded only by the cache. + +Local end-to-end tests point everything at an OpenTelemetry collector of their own with the +`LK_TELEMETRY_ENDPOINT` environment variable, read by the core at start (a base URL gets +`/v1/logs` and `/v1/traces`; a URL ending in `logs` is used as is). It is not part of any platform +API. With it, every batch goes there without Cloud rules or tokens. + ## Events An event with no `body` is exported with its name as the body as well as in `event_name`: log diff --git a/livekit-telemetry/src/destination.rs b/livekit-telemetry/src/destination.rs new file mode 100644 index 000000000..b8045ab60 --- /dev/null +++ b/livekit-telemetry/src/destination.rs @@ -0,0 +1,768 @@ +// 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. + +//! Where batches go and with which credential — decided here, once, for every platform. +//! +//! A platform hands over exactly two things: the LiveKit server URL a room connects to and the +//! participant token it connects with (again on every token refresh). The core derives the +//! ingest URL, checks that the host is LiveKit Cloud, reads the token's unverified claims (the +//! observability grant, the expiry), and routes every batch to the project its session belongs +//! to, so two rooms talking to two projects never share a token or a destination. + +use std::{ + collections::HashMap, + time::{Duration, SystemTime, UNIX_EPOCH}, +}; + +use base64::{ + alphabet, + engine::{DecodePaddingMode, GeneralPurpose, GeneralPurposeConfig}, + Engine, +}; +use tokio::time::Instant; + +/// The environment variable that points every upload at a collector of your own — a local +/// OpenTelemetry collector for end-to-end tests. Not part of any platform API: set it in the +/// test process's environment before the pipeline starts. Either a base URL +/// (`http://localhost:4318`, OTLP paths `/v1/logs` and `/v1/traces` are appended) or a full logs +/// URL ending in `logs` (its traces URL is derived by replacing that segment). +pub const ENDPOINT_OVERRIDE_ENV: &str = "LK_TELEMETRY_ENDPOINT"; + +/// The longest a token is believed to stay valid, whatever its `exp` says. +const MAX_TOKEN_LIFETIME: Duration = Duration::from_secs(366 * 24 * 60 * 60); + +/// LiveKit Cloud project hosts: `.livekit.cloud` in production, +/// `.staging.livekit.cloud` in staging. Self-hosted servers have no ingest. +const CLOUD_SUFFIX: &str = ".livekit.cloud"; + +/// Which OTLP signal a batch carries; picks the ingest path. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum Signal { + Logs, + Traces, +} + +/// A request target: URL per signal and the bearer token, when one is due. +#[derive(Debug, Clone, PartialEq)] +pub(crate) struct Target { + /// The project the batch is sent to (`None` for the test override): what its answer is + /// attributed to — a 404, a disable, a pause. + pub project: Option, + pub logs: String, + pub traces: String, + pub token: Option, +} + +impl Target { + pub fn url(&self, signal: Signal) -> &str { + match signal { + Signal::Logs => &self.logs, + Signal::Traces => &self.traces, + } + } +} + +/// What the unverified claims of a participant token say — enough to not send with a token the +/// collector is known to refuse. The signature is the collector's business. +#[derive(Debug, Clone, PartialEq)] +pub(crate) struct Token { + raw: String, + /// `observability.write` or `observability.clientWrite`. + grant: bool, + /// From `exp`, on the monotonic clock (so tests can move it). + expires: Option, + /// The credential's identity (a hash of the raw token): what a refusal is recorded against. + id: u64, +} + +impl Token { + pub fn parse(raw: &str) -> Self { + let claims = raw + .split('.') + .nth(1) + .and_then(|payload| { + const URL_SAFE: GeneralPurpose = GeneralPurpose::new( + &alphabet::URL_SAFE, + GeneralPurposeConfig::new() + .with_decode_padding_mode(DecodePaddingMode::Indifferent), + ); + URL_SAFE.decode(payload).ok() + }) + .and_then(|json| serde_json::from_slice::(&json).ok()); + let grant = claims.as_ref().is_some_and(|c| { + let observability = &c["observability"]; + observability["write"].as_bool() == Some(true) + || observability["clientWrite"].as_bool() == Some(true) + }); + // Untrusted input, read conservatively: a past or negative `exp` is expired, and so is + // one that is not a number at all; a fractional one counts in whole seconds; one beyond + // any real token lifetime is clamped (an `Instant` that far out would overflow). Only a + // token without `exp` has no known expiry. + let expires = claims.as_ref().and_then(|c| c.get("exp")).map(|exp| { + let now = + SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs_f64(); + let left = match exp.as_f64() { + Some(exp) if exp.is_finite() && exp > now => { + Duration::from_secs_f64((exp - now).min(MAX_TOKEN_LIFETIME.as_secs_f64())) + } + _ => Duration::ZERO, + }; + Instant::now().checked_add(left).unwrap_or_else(Instant::now) + }); + let id = { + use std::hash::{Hash, Hasher}; + let mut hasher = std::collections::hash_map::DefaultHasher::new(); + raw.hash(&mut hasher); + hasher.finish() + }; + Self { raw: raw.to_owned(), grant, expires, id } + } + + fn expired(&self) -> bool { + self.expires.is_some_and(|at| Instant::now() >= at) + } + + /// May be sent: carries the grant, not expired, not refused. + fn usable(&self, refused: &Refusals) -> bool { + self.grant && !self.expired() && !refused.contains(self) + } +} + +/// Why a project receives nothing (anymore). Each is logged once. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum Dead { + /// Not a LiveKit Cloud host (self-hosted, a typo): there is no ingest to talk to. + NotCloud, + /// The project's own token carries no observability grant: the customer did not opt in. + NoGrant, + /// The collector answered 404: no client ingest on this host. Revived by the next token + /// (a new connect or refresh), so a misrouted answer cannot silence a project for good. + NotFound, + /// The collector said the project's data recording is disabled by its owner. + Disabled, +} + +/// One LiveKit Cloud project, keyed by its host. +#[derive(Debug, Default)] +struct Project { + /// The latest token any room handed over for the project: what process-level batches and + /// batches from a previous launch are sent with. + latest: Option, + /// A token with the grant was seen: the customer opted in. Refreshed tokens may lack the + /// grant (today's server drops it on refresh) without withdrawing that. + consented: bool, + dead: Option, +} + +/// Offer `new` for a token slot: a refused token is never taken, into any slot; a token with the +/// grant wins; one without it does not replace a usable granted token (today's refreshes drop +/// the grant). +fn offer(slot: &mut Option, new: &Token, refused: &Refusals) { + if refused.contains(new) { + return; + } + let keep = slot.as_ref().is_some_and(|t| t.raw == new.raw || (!new.grant && t.usable(refused))); + if !keep { + *slot = Some(new.clone()); + } +} + +/// Credentials the collector refused, by identity, shared by every slot: a refused token is never +/// sent again, whichever Room or project hands it over. An entry lives until its token expires; +/// tokens without `exp` stay refused for the life of the process. +// ponytail: grows with refusals of tokens without `exp` (LiveKit tokens always carry one); a cap +// if a server ever mints tokens without expiry. +#[derive(Debug, Default)] +struct Refusals(HashMap>); + +impl Refusals { + fn contains(&self, token: &Token) -> bool { + self.0.contains_key(&token.id) + } + + /// Returns whether the refusal is new. + fn add(&mut self, token: &Token) -> bool { + let now = Instant::now(); + self.0.retain(|_, expires| expires.is_none_or(|at| at > now)); + self.0.insert(token.id, token.expires).is_none() + } +} + +/// The owner of process-level records (device state, pre-room errors, self-telemetry): no Room, +/// so they go to the latest project with its latest token. +pub(crate) const PROCESS_OWNER: &str = "process"; + +/// Where a batch is headed. +#[derive(Debug, Clone, PartialEq)] +pub(crate) enum Route { + /// Ready: send now. + Send(Target), + /// Keep it cached: no destination yet, or no usable token (missing, expired, grant-less, + /// refused) until the platform hands over a fresh one. A hard hold: no escape hatch sends + /// it anyway. + Wait, + /// Drop it: its project receives nothing. + Drop, +} + +/// Every project and session this process talks to, plus the test override. Tokens live here, +/// in memory only: never persisted, never in a batch. +#[derive(Debug, Default)] +pub(crate) struct Destinations { + projects: HashMap, + /// Each Room session's own latest token per project, by (host, session id): its batches for + /// that project are sent with it and never with another Room's, nor with the token it later + /// got for another project. + sessions: HashMap<(String, String), Token>, + /// Each Room session's first project: where the records it captured before it had a server + /// go. A Room's records never fall back to another Room's project. + rooms: HashMap, + refused: Refusals, + /// The host process-level batches (device state, pre-room logs, self-telemetry) go to: the + /// project most recently handed a token, while it is alive. + latest: Option, + /// [`ENDPOINT_OVERRIDE_ENV`]: every batch goes there, no Cloud rules, no token. + endpoint_override: Option, +} + +impl Destinations { + pub fn new(endpoint_override: Option<&str>) -> Self { + Self { + endpoint_override: endpoint_override.filter(|e| !e.is_empty()).map(override_target), + ..Self::default() + } + } + + /// The first project a Room session was routed to: where what it captured before it had a + /// server belongs. + pub fn first_project(&self, owner: &str) -> Option { + self.rooms.get(owner).cloned() + } + + pub fn has_override(&self) -> bool { + self.endpoint_override.is_some() + } + + /// Session `owner`'s server URL and token (at connect, and again on every refresh). Returns + /// the project host its batches carry, or `None` when the URL has no host. + pub fn set(&mut self, url: &str, token: &str, owner: &str) -> Option { + let Server { key: host, cloud } = Server::parse(url)?; + if self.endpoint_override.is_some() { + return Some(host); + } + let project = self.projects.entry(host.clone()).or_default(); + if project.dead.is_none() && !cloud { + project.dead = Some(Dead::NotCloud); + log::warn!("{host} is not LiveKit Cloud: client telemetry stays on this device"); + } + if matches!(project.dead, Some(Dead::NotCloud | Dead::Disabled)) { + return Some(host); + } + let token = Token::parse(token); + if project.dead == Some(Dead::NotFound) { + log::debug!("new token for {host}: trying its ingest again"); + project.dead = None; + } + if token.grant { + project.consented = true; + if project.dead == Some(Dead::NoGrant) { + project.dead = None; + } + } else if !project.consented && project.dead.is_none() { + project.dead = Some(Dead::NoGrant); + log::info!("the token for {host} has no observability grant: telemetry stays local"); + } + offer(&mut project.latest, &token, &self.refused); + let key = (host.clone(), owner.to_owned()); + let mut own = self.sessions.remove(&key); + offer(&mut own, &token, &self.refused); + self.sessions.extend(own.map(|own| (key, own))); + self.rooms.entry(owner.to_owned()).or_insert_with(|| host.clone()); + self.latest = Some(host.clone()); + Some(host) + } + + /// Where a batch goes: `host` is its project as captured with its records, `owner` the + /// session that produced them. A Room known to this process sends with its own latest token + /// for that project and waits when it has none; process-level batches and a previous + /// launch's use the project's latest. A Room's records captured before it had a server go + /// to its first project, never to another Room's. + pub fn route(&self, host: Option<&str>, owner: Option<&str>) -> Route { + if let Some(target) = &self.endpoint_override { + return Route::Send(target.clone()); + } + let room = owner.filter(|o| *o != PROCESS_OWNER); + let host = match (host, room) { + (Some(host), _) => host, + (None, Some(room)) => match self.rooms.get(room) { + Some(first) => first.as_str(), + None => return Route::Wait, + }, + (None, None) => match self.fallback() { + Some(host) => host, + // Nothing to talk to yet: wait. Only dead projects so far: nobody will ever + // take process-level data, so it is not kept either. + None if self.all_dead() => return Route::Drop, + None => return Route::Wait, + }, + }; + let Some(project) = self.projects.get(host) else { return Route::Wait }; + if project.dead.is_some() { + return Route::Drop; + } + let token = match room { + Some(room) if self.rooms.contains_key(room) => { + self.sessions.get(&(host.to_owned(), room.to_owned())) + } + _ => project.latest.as_ref(), + }; + match token { + Some(token) if token.usable(&self.refused) => Route::Send(Target { + project: Some(host.to_owned()), + logs: ingest_url(host, Signal::Logs), + traces: ingest_url(host, Signal::Traces), + token: Some(token.raw.clone()), + }), + _ => Route::Wait, + } + } + + /// The host of the project process-level batches go to: the latest, else any alive one. + fn fallback(&self) -> Option<&str> { + let alive = |host: &&String| self.projects.get(*host).is_some_and(|p| p.dead.is_none()); + self.latest + .as_ref() + .filter(alive) + .or_else(|| self.projects.keys().find(alive)) + .map(String::as_str) + } + + fn all_dead(&self) -> bool { + !self.projects.is_empty() && self.projects.values().all(|p| p.dead.is_some()) + } + + /// Whether a session routed to `host` should still collect: `false` once its project is + /// known to receive nothing, so its records are dropped at the door instead of cached. + pub fn alive(&self, host: Option<&str>) -> bool { + self.endpoint_override.is_some() || !matches!(self.route(host, None), Route::Drop) + } + + /// The collector refused `token` (401/403): it is not sent again, from any slot; batches that + /// would use it wait for the next one. + pub fn refuse(&mut self, token: &str) { + if self.refused.add(&Token::parse(token)) { + log::warn!("the collector refused a token; waiting for a new one"); + } + } + + /// Keep credentials only where something still needs them. `live` maps every session still + /// alive — held by its Room, or by records, windows or spans not yet cached — to its current + /// project; its credentials stay (for each project it was routed to: records it captured for + /// an earlier project may still be queued). `backlog` is (project or none, owner) of every + /// cached batch: their credentials stay too. Project-level copies stay while a live session + /// is routed there or the backlog has batches for the project (process-level batches: the + /// latest project). + pub fn retain_owners( + &mut self, + live: &HashMap>, + backlog: &std::collections::HashSet<(Option, String)>, + ) { + // A Room's batch captured before it had a server is cached without a project: it goes + // to that Room's first project, so it backs that project's credential. + let backlog: std::collections::HashSet<(Option, String)> = backlog + .iter() + .map(|(host, owner)| match (host, self.rooms.get(owner)) { + (None, Some(first)) if owner != PROCESS_OWNER => { + (Some(first.clone()), owner.clone()) + } + _ => (host.clone(), owner.clone()), + }) + .collect(); + let backed = + |host: &str, owner: &str| backlog.contains(&(Some(host.to_owned()), owner.to_owned())); + self.sessions.retain(|(host, owner), _| backed(host, owner) || live.contains_key(owner)); + self.rooms + .retain(|owner, _| live.contains_key(owner) || backlog.iter().any(|(_, o)| o == owner)); + let unrouted = backlog.iter().any(|(host, _)| host.is_none()); + let latest = self.latest.clone(); + for (host, project) in &mut self.projects { + let needed = live.values().any(|route| route.as_deref() == Some(host)) + || backlog.iter().any(|(h, _)| h.as_deref() == Some(host)) + || (unrouted && latest.as_deref() == Some(host)); + if !needed { + project.latest = None; + } + } + } + + /// How many credentials are held, Room and project copies together. + #[cfg(test)] + pub fn credentials(&self) -> usize { + self.sessions.len() + self.projects.values().filter(|p| p.latest.is_some()).count() + } + + /// Try `host`'s ingest again if it answered 404 (a reconnect). Returns whether it did. + pub fn revive(&mut self, host: &str) -> bool { + match self.projects.get_mut(host) { + Some(project) if project.dead == Some(Dead::NotFound) => { + project.dead = None; + true + } + _ => false, + } + } + + /// The project on `host` receives nothing (see [`Dead`] for how long). + pub fn kill(&mut self, host: &str, why: Dead) { + let project = self.projects.entry(host.to_owned()).or_default(); + if project.dead.is_none() { + match why { + Dead::NotFound => log::warn!("{host} has no client telemetry ingest; going silent"), + Dead::Disabled => log::warn!("{host} disabled data recording; going silent"), + Dead::NotCloud | Dead::NoGrant => {} + } + project.dead = Some(why); + } + } +} + +/// A server URL, parsed (WHATWG URL rules, the same a browser or `URLSession` applies): its +/// routing key — the canonical host, with the port when one is given — and whether it is a +/// LiveKit Cloud project a token may be sent to. +struct Server { + key: String, + cloud: bool, +} + +impl Server { + fn parse(url: &str) -> Option { + let url = url::Url::parse(url).ok()?; + let host = url.host_str().filter(|h| !h.is_empty())?.to_ascii_lowercase(); + let key = match url.port() { + Some(port) => format!("{host}:{port}"), + None => host.clone(), + }; + // Credentials only ever go to `https://.livekit.cloud/…`, built from the + // validated host alone: a TLS scheme, a domain under the Cloud suffix with a label of its + // own, the default port, no userinfo. Anything else — a look-alike that only a string + // check would accept (`evil.example\.livekit.cloud` is host `evil.example`), a + // plaintext scheme, an odd port — is not Cloud. + let cloud = matches!(url.scheme(), "wss" | "https") + && matches!(url.host(), Some(url::Host::Domain(_))) + && host.strip_suffix(CLOUD_SUFFIX).is_some_and(|label| !label.is_empty()) + && url.port().is_none() + && url.username().is_empty() + && url.password().is_none(); + Some(Self { key, cloud }) + } +} + +/// `https:///observability/client/{logs,traces}/otlp/v0`: the client table set. +fn ingest_url(host: &str, signal: Signal) -> String { + let signal = match signal { + Signal::Logs => "logs", + Signal::Traces => "traces", + }; + format!("https://{host}/observability/client/{signal}/otlp/v0") +} + +fn override_target(endpoint: &str) -> Target { + let endpoint = endpoint.trim_end_matches('/'); + let (logs, traces) = match endpoint.rsplit_once("logs") { + Some((before, after)) => (endpoint.to_owned(), format!("{before}traces{after}")), + None => (format!("{endpoint}/v1/logs"), format!("{endpoint}/v1/traces")), + }; + Target { project: None, logs, traces, token: None } +} + +#[cfg(test)] +pub(crate) mod tests { + use super::*; + + /// An unsigned JWT with these claims — enough for the core, which never verifies. + pub(crate) fn jwt(claims: serde_json::Value) -> String { + let engine = base64::engine::general_purpose::URL_SAFE_NO_PAD; + format!( + "{}.{}.sig", + engine.encode(br#"{"alg":"HS256"}"#), + engine.encode(claims.to_string()) + ) + } + + pub(crate) fn granted(ttl_secs: u64) -> String { + let exp = SystemTime::now().duration_since(UNIX_EPOCH).expect("clock").as_secs() + ttl_secs; + jwt(serde_json::json!({ "exp": exp, "observability": { "write": true } })) + } + + pub(crate) fn grantless(ttl_secs: u64) -> String { + let exp = SystemTime::now().duration_since(UNIX_EPOCH).expect("clock").as_secs() + ttl_secs; + jwt(serde_json::json!({ "exp": exp, "video": { "roomJoin": true } })) + } + + fn send(route: Route) -> Target { + match route { + Route::Send(target) => target, + other => panic!("expected a target, got {other:?}"), + } + } + + #[test] + fn only_livekit_cloud_hosts_get_an_ingest_url() { + let mut d = Destinations::default(); + let token = granted(600); + let host = d.set("wss://Proj.LiveKit.Cloud/rtc?access_token=x#f", &token, "s"); + assert_eq!(host.as_deref(), Some("proj.livekit.cloud")); + let target = send(d.route(host.as_deref(), Some("s"))); + assert_eq!(target.logs, "https://proj.livekit.cloud/observability/client/logs/otlp/v0"); + assert_eq!(target.traces, "https://proj.livekit.cloud/observability/client/traces/otlp/v0"); + assert_eq!(target.token.as_deref(), Some(token.as_str())); + + let staging = d.set("https://p.staging.livekit.cloud", &token, "s"); + assert!(matches!(d.route(staging.as_deref(), None), Route::Send(_)), "staging is Cloud"); + + for url in + ["ws://localhost:7880", "wss://livekit.example.com", "wss://livekit.cloud.evil.io"] + { + let host = d.set(url, &token, "s"); + assert_eq!(d.route(host.as_deref(), None), Route::Drop, "{url}: no ingest"); + } + assert_eq!(d.set("nonsense", &token, "s"), None); + let empty = d.set("wss:///rtc", &token, "s"); // the parser reads host `rtc` + assert_eq!(d.route(empty.as_deref(), None), Route::Drop); + } + + /// Finding r1-1: whatever a string check would accept, the token only goes to a host the URL + /// parser agrees is `.livekit.cloud`, over TLS, on the default port, without + /// userinfo — and the endpoint is built from that host alone. + #[test] + fn look_alike_server_urls_never_get_the_token() { + let token = granted(600); + for url in [ + "wss://evil.example\\.livekit.cloud", + "wss://evil.example\\x.livekit.cloud/rtc", + "wss://user:pass@p.livekit.cloud", + "wss://evil.example@p.livekit.cloud", + "wss://p.livekit.cloud:8443", + "ws://p.livekit.cloud", + "http://p.livekit.cloud", + "wss://livekit.cloud", + "wss://.livekit.cloud", + "wss://p.livekit.cloud.evil.io", + "wss://p.livekit.cloud%2eevil.io", + "wss://evil.io#.livekit.cloud", + "wss://evil.io?.livekit.cloud", + "wss://[::1]", + "wss://1.2.3.4", + "wss://p.livekit.cloud\u{0}.evil", + ] { + let mut d = Destinations::default(); + let host = d.set(url, &token, "s"); + let route = d.route(host.as_deref(), Some("s")); + assert!(!matches!(route, Route::Send(_)), "{url} → {host:?} must not get the token"); + } + let mut d = Destinations::default(); + let host = d.set("WSS://P.LiveKit.Cloud:443/rtc?x=1#y", &token, "s"); + let target = send(d.route(host.as_deref(), Some("s"))); + assert_eq!(target.logs, "https://p.livekit.cloud/observability/client/logs/otlp/v0"); + } + + #[test] + fn a_token_without_the_grant_is_not_consent() { + let mut d = Destinations::default(); + let host = d.set("wss://p.livekit.cloud", &grantless(600), "s"); + assert_eq!(d.route(host.as_deref(), Some("s")), Route::Drop, "never opted in"); + assert_eq!(d.route(None, None), Route::Drop, "and no project takes process data"); + let host = d.set("wss://p.livekit.cloud", &granted(600), "s"); + assert!(matches!(d.route(host.as_deref(), Some("s")), Route::Send(_)), "a grant opts in"); + } + + #[tokio::test(start_paused = true)] + async fn a_refresh_that_drops_the_grant_keeps_the_granted_token_until_it_expires() { + let mut d = Destinations::default(); + let join = granted(60); + let host = d.set("wss://p.livekit.cloud", &join, "s"); + d.set("wss://p.livekit.cloud", &grantless(600), "s"); + assert_eq!(send(d.route(host.as_deref(), Some("s"))).token, Some(join), "still granted"); + tokio::time::advance(Duration::from_secs(61)).await; + assert_eq!(d.route(host.as_deref(), Some("s")), Route::Wait, "expired: hold, never drop"); + let refreshed = granted(600); + d.set("wss://p.livekit.cloud", &refreshed, "s"); + assert_eq!(send(d.route(host.as_deref(), Some("s"))).token, Some(refreshed)); + } + + #[test] + fn a_refused_token_is_never_sent_again() { + let mut d = Destinations::default(); + let first = granted(600); + let host = d.set("wss://p.livekit.cloud", &first, "s").expect("host"); + d.refuse(&first); + assert_eq!(d.route(Some(&host), Some("s")), Route::Wait); + d.set("wss://p.livekit.cloud", &first, "s"); + assert_eq!(d.route(Some(&host), Some("s")), Route::Wait, "handing it over again: still no"); + d.set("wss://p.livekit.cloud", &granted(900), "s"); + assert!(matches!(d.route(Some(&host), Some("s")), Route::Send(_))); + } + + #[test] + fn rooms_never_borrow_each_others_tokens() { + let mut d = Destinations::default(); + let (a, b, c) = (granted(600), granted(700), granted(800)); + let host_a = d.set("wss://a.livekit.cloud", &a, "room-a"); + let host_b = d.set("wss://b.livekit.cloud", &b, "room-b"); + d.set("wss://a.livekit.cloud", &c, "room-c"); // a second room on project a + let target = send(d.route(host_a.as_deref(), Some("room-a"))); + assert!(target.logs.starts_with("https://a.") && target.token == Some(a.clone())); + let target = send(d.route(host_a.as_deref(), Some("room-c"))); + assert_eq!(target.token, Some(c.clone()), "same project, its own token"); + let target = send(d.route(host_b.as_deref(), Some("room-b"))); + assert!(target.logs.starts_with("https://b.") && target.token == Some(b)); + d.refuse(&a); + assert_eq!(d.route(host_a.as_deref(), Some("room-a")), Route::Wait, "no borrowing c"); + let previous_launch = send(d.route(host_a.as_deref(), Some("gone"))); + assert_eq!(previous_launch.token, Some(c), "a previous launch's batch: project's latest"); + d.kill("b.livekit.cloud", Dead::Disabled); + assert!(send(d.route(None, None)).logs.starts_with("https://a."), "process: a live one"); + } + + /// Finding r1-2: a Room's batches keep the project they were captured for — after it + /// reconnects to another project they still go to the first with the first project's token — + /// and records it captured before it had a server never go to another Room's project. + #[test] + fn ownership_survives_a_room_changing_projects() { + let mut d = Destinations::default(); + let (a, b) = (granted(600), granted(700)); + let host_a = d.set("wss://a.livekit.cloud", &a, "room").expect("a"); + d.set("wss://b.livekit.cloud", &b, "room"); + let target = send(d.route(Some(&host_a), Some("room"))); + assert!(target.logs.starts_with("https://a.") && target.token == Some(a)); + + d.set("wss://c.livekit.cloud", &granted(800), "other"); + assert_eq!(d.route(None, Some("unconnected")), Route::Wait, "no borrowing c"); + let early = send(d.route(None, Some("room"))); + assert!(early.logs.starts_with("https://a."), "a Room's first project"); + assert!(matches!(d.route(None, Some(PROCESS_OWNER)), Route::Send(_))); + } + + /// Findings r1-13 / r2-4: a refusal belongs to the credential — no slot, Room, project or + /// process route can make it usable again — and credentials follow the live Rooms and their + /// backlog, project-level copies included. + #[test] + fn refusals_stick_and_credentials_follow_live_rooms() { + let mut d = Destinations::default(); + let tokens: Vec = (0..40).map(|n| granted(600 + n)).collect(); + for (n, token) in tokens.iter().enumerate() { + d.set("wss://p.livekit.cloud", token, &format!("room-{n}")); + d.refuse(token); + } + for (n, token) in tokens.iter().enumerate() { + let room = format!("room-{n}"); + d.set("wss://p.livekit.cloud", token, &room); // handed over again + assert_eq!(d.route(Some("p.livekit.cloud"), Some(&room)), Route::Wait, "{n}"); + d.set("wss://p.livekit.cloud", token, &format!("new-{n}")); // by a new Room + assert_eq!(d.route(Some("p.livekit.cloud"), Some(&format!("new-{n}"))), Route::Wait); + assert_eq!(d.route(None, Some(PROCESS_OWNER)), Route::Wait, "{n}: process route"); + assert_eq!(d.route(Some("p.livekit.cloud"), Some("gone")), Route::Wait, "{n}: restart"); + } + // Refuse A, let B become the latest, hand A over again: A stays refused everywhere. + let (a, b) = (granted(900), granted(901)); + d.set("wss://q.livekit.cloud", &a, "ra"); + d.refuse(&a); + d.set("wss://q.livekit.cloud", &b, "rb"); + d.set("wss://q.livekit.cloud", &a, "ra"); + let process = send(d.route(Some("q.livekit.cloud"), Some("gone"))); + assert_eq!(process.token, Some(b), "the project copy is b, never the refused a"); + + // Churn: 50 projects, their Rooms gone, nothing cached: nothing is retained. + for n in 0..50 { + d.set(&format!("wss://c{n}.livekit.cloud"), &granted(1000 + n), &format!("churn-{n}")); + } + let live: HashMap> = + [("rb".to_owned(), Some("q.livekit.cloud".to_owned()))].into(); + d.retain_owners(&live, &Default::default()); + assert_eq!(d.credentials(), 2, "rb's token and q's project copy, nothing else"); + let backlog = [(Some("q.livekit.cloud".to_owned()), "ra".to_owned())].into(); + d.set("wss://q.livekit.cloud", &granted(902), "ra"); + d.retain_owners(&Default::default(), &backlog); + assert_eq!(d.credentials(), 2, "a gone Room's backlog keeps its credential and q's copy"); + } + + #[test] + fn not_found_recovers_with_the_next_token_disabled_does_not() { + let mut d = Destinations::default(); + let host = d.set("wss://p.livekit.cloud", &granted(600), "s").expect("host"); + d.kill(&host, Dead::NotFound); + assert_eq!(d.route(Some(&host), Some("s")), Route::Drop); + d.set("wss://p.livekit.cloud", &granted(700), "s"); + assert!(matches!(d.route(Some(&host), Some("s")), Route::Send(_)), "a new connect retries"); + d.kill(&host, Dead::Disabled); + d.set("wss://p.livekit.cloud", &granted(800), "s"); + assert_eq!(d.route(Some(&host), Some("s")), Route::Drop, "the owner's switch stands"); + } + + #[test] + fn the_override_takes_everything_without_a_token() { + let d = Destinations::new(Some("http://localhost:4318/")); + let target = send(d.route(None, None)); + assert_eq!(target.logs, "http://localhost:4318/v1/logs"); + assert_eq!(target.traces, "http://localhost:4318/v1/traces"); + assert_eq!(target.token, None); + let full = Destinations::new(Some("https://h/observability/client/logs/otlp/v0")); + assert_eq!( + send(full.route(Some("ignored"), None)).traces, + "https://h/observability/client/traces/otlp/v0" + ); + let mut local = Destinations::new(Some("http://localhost:4318")); + let host = local.set("ws://localhost:7880", "not-a-jwt", "s"); + assert!(matches!(local.route(host.as_deref(), Some("s")), Route::Send(_))); + } + + /// Finding r2-11: a past, negative, fractional-past or malformed `exp` is expired (a hard + /// hold); only a missing one means no known expiry. + #[tokio::test(start_paused = true)] + async fn untrusted_expiry_claims_fail_closed() { + let now = SystemTime::now().duration_since(UNIX_EPOCH).expect("clock").as_secs_f64(); + let token = |exp: serde_json::Value| { + jwt(serde_json::json!({ "exp": exp, "observability": { "write": true } })) + }; + for exp in [ + serde_json::json!(-1), + serde_json::json!(0), + serde_json::json!(now - 0.5), + serde_json::json!(-1e308), + serde_json::json!("tomorrow"), + serde_json::json!(true), + serde_json::json!(null), + serde_json::json!({}), + ] { + let mut d = Destinations::default(); + let host = d.set("wss://p.livekit.cloud", &token(exp.clone()), "s"); + assert_eq!(d.route(host.as_deref(), Some("s")), Route::Wait, "exp {exp}: expired"); + } + let mut d = Destinations::default(); + let host = d.set("wss://p.livekit.cloud", &token(serde_json::json!(now + 60.9)), "s"); + assert!(matches!(d.route(host.as_deref(), Some("s")), Route::Send(_)), "fractional future"); + let host = d.set("wss://p.livekit.cloud", &token(serde_json::json!(1e300)), "s"); + assert!(matches!(d.route(host.as_deref(), Some("s")), Route::Send(_)), "extreme: clamped"); + } + + #[test] + fn tokens_are_read_not_verified() { + assert!(Token::parse(&granted(60)).grant); + let client_write = jwt(serde_json::json!({ "observability": { "clientWrite": true } })); + let parsed = Token::parse(&client_write); + assert!(parsed.grant && parsed.expires.is_none(), "no exp: no known expiry"); + assert!(!Token::parse("garbage").grant); + assert!(!Token::parse("a.!!!.c").grant); + } +} diff --git a/livekit-telemetry/src/lib.rs b/livekit-telemetry/src/lib.rs index ac69c5014..67fa4b169 100644 --- a/livekit-telemetry/src/lib.rs +++ b/livekit-telemetry/src/lib.rs @@ -34,9 +34,17 @@ mod otlp; /// OTLP protobuf types (re-exported from `opentelemetry-proto`). mod proto; +/// Transport seam: how encoded batches leave the device. +mod transport; + +/// Where batches go: server URL + token → ingest URL, grant, expiry, per-project routing. +mod destination; + +pub use destination::ENDPOINT_OVERRIDE_ENV; pub use event::*; pub use span::SpanOutcome; pub use stats::{TelemetryStats, TelemetryStatus}; +pub use transport::*; #[cfg(feature = "uniffi")] uniffi::setup_scaffolding!(); diff --git a/livekit-telemetry/src/transport.rs b/livekit-telemetry/src/transport.rs new file mode 100644 index 000000000..060d5eb36 --- /dev/null +++ b/livekit-telemetry/src/transport.rs @@ -0,0 +1,568 @@ +// 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; + +use prost::Message; + +/// One fully composed OTLP/HTTP export request. +/// +/// The core fills in the URL, the headers (content type, auth, …) and the protobuf body; +/// a transport only moves the bytes. Non-HTTP transports (e.g. a data channel) may ignore +/// `url`/`headers` and forward `body`, which is a standard `ExportLogsServiceRequest`. +#[cfg_attr(feature = "uniffi", derive(uniffi::Record))] +#[derive(Debug, Clone, PartialEq)] +pub struct ExportRequest { + pub url: String, + pub headers: HashMap, + pub body: Vec, +} + +/// What the collector answered. A transport moves bytes both ways and never interprets them: +/// status classification, `Retry-After`, the `google.rpc.Status` body and the Cloud "disabled" +/// contract are read once, here, for every platform ([`ExportError::from_response`]). +#[cfg_attr(feature = "uniffi", derive(uniffi::Record))] +#[derive(Debug, Clone, PartialEq, Default)] +pub struct ExportResponse { + pub status: u16, + pub headers: HashMap, + pub body: Vec, +} + +impl ExportResponse { + /// `2xx`, nothing else: what an accepting collector answers. + pub fn accepted() -> Self { + Self { status: 200, ..Self::default() } + } +} + +/// LiveKit Cloud's answer when the project owner switched data recording off (401 or 403, plain +/// text, no machine-readable code yet). Matched as a phrase, not a word. +const DISABLED_ANSWER: &str = "data recording is disabled by owner"; + +/// How long a 429 that names no delay holds uploads: LiveKit Cloud's quota answer names none. +const THROTTLE_DEFAULT_MS: u64 = 60_000; + +/// Why a transport could not get an answer. A 4xx/5xx is an answer, not an error: return it as an +/// [`ExportResponse`] and the core classifies it. +#[cfg_attr(feature = "uniffi", derive(uniffi::Error))] +#[derive(Debug, Clone, PartialEq, thiserror::Error)] +pub enum ExportError { + /// No response: connection refused or reset, DNS or TLS failure, offline, timed out. The + /// batch stays cached and is retried with backoff. `retry_after_ms` is honored when set. + #[error("retryable export error: {reason}")] + Retryable { reason: String, retry_after_ms: Option }, + /// The request can never succeed (an invalid URL): the batch is dropped. + #[error("export rejected: {reason}")] + Rejected { reason: String }, + /// Telemetry is off for this project: the project goes silent. + #[error("telemetry disabled by the collector")] + Disabled, +} + +/// A foreign transport threw something its binding did not declare (a Swift `Error`, a Kotlin +/// `RuntimeException`): no answer, retried with backoff — never a panic in the exporter. +#[cfg(feature = "uniffi")] +impl From for ExportError { + fn from(error: uniffi::UnexpectedUniFFICallbackError) -> Self { + Self::Retryable { + reason: format!("transport threw: {}", error.reason), + retry_after_ms: None, + } + } +} + +/// What a collector's answer means for the batch: OTLP/HTTP failure semantics, plus the answers +/// LiveKit Cloud gives on top of them. +#[derive(Debug, Clone, PartialEq)] +pub(crate) enum Verdict { + /// 2xx. `rejected` counts the records an OTLP partial success refused (never retried). + Accepted { rejected: u64, reason: String }, + /// 400, 3xx (a credential-bearing request is never redirected), and every other 4xx/5xx the + /// OTLP spec does not call retryable: drop the batch. + Rejected(String), + /// 413: split the batch and retry the halves, down to single records. + TooLarge, + /// 401/403 that is not the "disabled" answer: the credential is the problem, not the data — + /// hold the project's batches until the platform hands over a new token. + Unauthorized(String), + /// 404: the host has no client ingest; the project goes silent. + NotFound, + /// 401/403 "data recording is disabled by owner": the project goes silent and its cache is + /// purged. + Disabled, + /// 429, or 503 naming a delay: pause every upload for `delay_ms`, keep collecting. + Throttled { delay_ms: u64, reason: String }, + /// 502/503/504, or any error carrying `RetryInfo`: the server failed transiently. Retry the + /// batch after `delay_ms` or a backoff; a batch that keeps failing is dropped eventually. + Retry { delay_ms: Option, reason: String }, +} + +impl Verdict { + /// Classify a collector response. + pub(crate) fn of(response: &ExportResponse) -> Self { + let status = response.status; + if (200..300).contains(&status) { + return match PartialSuccessResponse::decode(response.body.as_slice()) { + Ok(PartialSuccessResponse { partial_success: Some(p) }) if p.rejected > 0 => { + Self::Accepted { rejected: p.rejected as u64, reason: p.error_message } + } + _ => Self::Accepted { rejected: 0, reason: String::new() }, + }; + } + // OTLP/HTTP errors are a protobuf `google.rpc.Status`; Cloud's auth answers are text. + let rpc = RpcStatus::decode(response.body.as_slice()).ok(); + let text = match rpc.as_ref().map(|s| s.message.trim()).filter(|m| !m.is_empty()) { + Some(message) => message.to_owned(), + None => String::from_utf8_lossy(&response.body).trim().chars().take(200).collect(), + }; + // `reason`, not `message`: a UniFFI field named `message` collides with + // `Throwable.message` in Kotlin. + let reason = if text.is_empty() { + format!("HTTP {status}") + } else { + format!("HTTP {status}: {text}") + }; + if matches!(status, 401 | 403) { + return if text.to_ascii_lowercase().contains(DISABLED_ANSWER) { + Self::Disabled + } else { + Self::Unauthorized(reason) + }; + } + let header = response + .headers + .iter() + .find(|(name, _)| name.eq_ignore_ascii_case("retry-after")) + .and_then(|(_, value)| retry_after_ms(value)); + let retry_info = rpc.as_ref().and_then(retry_info_ms); + let delay = header.or(retry_info); + // The status decides whether a batch is retried; a delay hint only says when. A hint on a + // final status (400, 422, 501, a redirect) changes nothing. + match status { + 404 => Self::NotFound, + 413 => Self::TooLarge, + // A rate limit is an instruction to stop, so 429 always yields a wait: LiveKit Cloud + // answers a quota check with a bare `ResourceExhausted` — no `Retry-After`, no + // `RetryInfo` — and without one the exporter would spend its retries hammering the + // endpoint that just asked for quiet. + 429 => Self::Throttled { delay_ms: delay.unwrap_or(THROTTLE_DEFAULT_MS), reason }, + 503 if delay.is_some() => { + Self::Throttled { delay_ms: delay.unwrap_or_default(), reason } + } + 502..=504 => Self::Retry { delay_ms: delay, reason }, + // LiveKit Cloud's one retryable 500 carries `RetryInfo` (custom: OTLP retries only + // 429/502/503/504). + 500 if retry_info.is_some() => Self::Retry { delay_ms: retry_info, reason }, + _ => Self::Rejected(reason), + } + } +} + +/// `ExportLogsServiceResponse` and `ExportTraceServiceResponse` share one shape: field 1 is the +/// partial success, whose field 1 counts the rejected records and field 2 says why. +#[derive(Clone, PartialEq, prost::Message)] +struct PartialSuccessResponse { + #[prost(message, optional, tag = "1")] + partial_success: Option, +} + +#[derive(Clone, PartialEq, prost::Message)] +struct PartialSuccess { + #[prost(int64, tag = "1")] + rejected: i64, + #[prost(string, tag = "2")] + error_message: String, +} + +/// `google.rpc.Status`, the OTLP/HTTP error body. Two messages, three fields: not worth a +/// generated crate. +#[derive(Clone, PartialEq, prost::Message)] +struct RpcStatus { + #[prost(int32, tag = "1")] + code: i32, + #[prost(string, tag = "2")] + message: String, + #[prost(message, repeated, tag = "3")] + details: Vec, +} + +/// `google.rpc.RetryInfo`. +#[derive(Clone, PartialEq, prost::Message)] +struct RetryInfo { + #[prost(message, optional, tag = "1")] + retry_delay: Option, +} + +/// RFC 9110 §10.2.3 `Retry-After`: delay-seconds or an IMF-fixdate. Untrusted input: garbage is +/// ignored, a date in the past means now; the caller clamps the far end. +fn retry_after_ms(value: &str) -> Option { + let value = value.trim(); + // delay-seconds: digits only; more than a u64 holds saturates (the exporter clamps it). + if !value.is_empty() && value.bytes().all(|b| b.is_ascii_digit()) { + return Some(value.parse::().unwrap_or(u64::MAX).saturating_mul(1000)); + } + let at = imf_fixdate_secs(value)?; + let now = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).ok()?.as_secs(); + Some(at.saturating_sub(now).saturating_mul(1000)) +} + +/// `Sun, 06 Nov 1994 08:49:37 GMT` → unix seconds. The obsolete RFC 850 / asctime forms are not +/// accepted (RFC 9110 lets a recipient ignore them). +fn imf_fixdate_secs(value: &str) -> Option { + let mut parts = value.split_ascii_whitespace(); + let (_weekday, day, month, year, time, zone) = + (parts.next()?, parts.next()?, parts.next()?, parts.next()?, parts.next()?, parts.next()?); + if zone != "GMT" || parts.next().is_some() { + return None; + } + const MONTHS: [&str; 12] = + ["Jan", "Feb", "Mar", "Apr", "May", "Jun", "Jul", "Aug", "Sep", "Oct", "Nov", "Dec"]; + let month = MONTHS.iter().position(|m| *m == month)? as i64 + 1; + let (day, year): (i64, i64) = (day.parse().ok()?, year.parse().ok()?); + let mut hms = time.split(':').map(|n| n.parse::().ok()); + let (h, m, s) = (hms.next()??, hms.next()??, hms.next()??); + // IMF-fixdate has a four-digit year; bounding it keeps the arithmetic below in range. + if !(1970..=9999).contains(&year) + || hms.next().is_some() + || !(1..=31).contains(&day) + || !(0..24).contains(&h) + || !(0..60).contains(&m) + || !(0..61).contains(&s) + { + return None; + } + // Days from the civil date (proleptic Gregorian), after Howard Hinnant's algorithm. + let y = if month <= 2 { year - 1 } else { year }; + let era = y.div_euclid(400); + let yoe = y - era * 400; + let doy = (153 * (month + if month > 2 { -3 } else { 9 }) + 2) / 5 + day - 1; + let doe = yoe * 365 + yoe / 4 - yoe / 100 + doy; + let days = era * 146_097 + doe - 719_468; + u64::try_from(days * 86_400 + h * 3600 + m * 60 + s).ok() +} + +fn retry_info_ms(status: &RpcStatus) -> Option { + status + .details + .iter() + .filter(|detail| detail.type_url.ends_with("google.rpc.RetryInfo")) + .find_map(|detail| RetryInfo::decode(detail.value.as_slice()).ok()?.retry_delay) + .map(|delay| { + (delay.seconds.max(0) as u64) + .saturating_mul(1000) + .saturating_add(delay.nanos.clamp(0, 999_999_999) as u64 / 1_000_000) + }) +} + +/// Moves an encoded batch off the device and hands back whatever came back. +/// +/// Implemented in Rust ([`NetTransport`], feature `net`) or by the host platform through UniFFI. +/// The contract: +/// +/// - POST `url` with `headers` and `body` as given, and return the response whatever its status; +/// fail only when there was no response (connection, DNS, TLS, timeout) — +/// [`ExportError::Retryable`] — or the request cannot be made at all +/// ([`ExportError::Rejected`], e.g. an invalid URL). +/// - Never retry: the core owns the retry policy. +/// - Never forward the `Authorization` header across origins: do not follow a redirect to +/// another scheme, host or port with it (returning the 3xx is fine — the core drops the batch). +/// `NetTransport` over `livekit-net`'s native client strips it on cross-origin redirects. +/// - Keep the request off the call's critical path: lowest priority where the platform allows it. +#[cfg_attr(feature = "uniffi", uniffi::export(with_foreign))] +#[async_trait::async_trait] +pub trait TelemetryTransport: Send + Sync { + async fn send(&self, request: ExportRequest) -> Result; +} + +#[cfg(feature = "net")] +mod net { + use std::sync::Arc; + + use livekit_net::{Header, HttpClient, HttpClientExt, TransportError}; + + use super::{ExportError, ExportRequest, ExportResponse, TelemetryTransport}; + + /// Default transport: HTTP POST through a `livekit-net` client — the native backend, or + /// whatever the host registered with `livekit_net::set_http_client`. + pub struct NetTransport(Arc); + + impl NetTransport { + /// Post through `client`. + pub fn new(client: Arc) -> Self { + Self(client) + } + + /// Resolve the process-wide `livekit-net` client; `None` when none is available. + pub fn from_registry() -> Option { + livekit_net::http_client().map(Self) + } + } + + #[async_trait::async_trait] + impl TelemetryTransport for NetTransport { + async fn send(&self, request: ExportRequest) -> Result { + let headers = + request.headers.into_iter().map(|(name, value)| Header { name, value }).collect(); + match self.0.post(request.url, headers, request.body).await { + Ok(response) => Ok(ExportResponse { + status: response.status, + headers: response.headers.into_iter().map(|h| (h.name, h.value)).collect(), + body: response.body, + }), + // The client swallowed the body with the status; the core still classifies it. + Err(TransportError::Http { status }) => { + Ok(ExportResponse { status, ..Default::default() }) + } + Err(other) => { + Err(ExportError::Retryable { reason: other.to_string(), retry_after_ms: None }) + } + } + } + } +} + +#[cfg(feature = "net")] +pub use net::NetTransport; + +#[cfg(test)] +mod tests { + use prost::Message; + + use super::*; + + fn status_body(message: &str, retry: Option<(i64, i32)>) -> Vec { + let details = retry + .map(|(seconds, nanos)| prost_types::Any { + type_url: "type.googleapis.com/google.rpc.RetryInfo".into(), + value: RetryInfo { retry_delay: Some(prost_types::Duration { seconds, nanos }) } + .encode_to_vec(), + }) + .into_iter() + .collect(); + RpcStatus { code: 8, message: message.into(), details }.encode_to_vec() + } + + fn response(status: u16, headers: &[(&str, &str)], body: Vec) -> ExportResponse { + ExportResponse { + status, + headers: headers.iter().map(|(k, v)| (k.to_string(), v.to_string())).collect(), + body, + } + } + + fn verdict(status: u16, headers: &[(&str, &str)], body: Vec) -> Verdict { + Verdict::of(&response(status, headers, body)) + } + + #[test] + fn success_and_partial_success() { + assert_eq!( + Verdict::of(&ExportResponse::accepted()), + Verdict::Accepted { rejected: 0, reason: String::new() } + ); + assert!(matches!(verdict(204, &[], vec![]), Verdict::Accepted { rejected: 0, .. })); + let partial = PartialSuccessResponse { + partial_success: Some(PartialSuccess { rejected: 3, error_message: "too old".into() }), + } + .encode_to_vec(); + assert_eq!( + verdict(200, &[], partial), + Verdict::Accepted { rejected: 3, reason: "too old".into() } + ); + assert!(matches!(verdict(200, &[], b"not protobuf".to_vec()), Verdict::Accepted { .. })); + } + + #[test] + fn client_errors() { + assert_eq!( + verdict(400, &[], status_body("bad batch", None)), + Verdict::Rejected("HTTP 400: bad batch".into()) + ); + assert_eq!(verdict(413, &[], vec![]), Verdict::TooLarge); + assert!(matches!(verdict(307, &[], vec![]), Verdict::Rejected(_)), "never follow"); + assert!(matches!(verdict(422, &[], vec![]), Verdict::Rejected(_))); + assert_eq!(verdict(404, &[], vec![]), Verdict::NotFound); + for body in ["invalid token", "operation requires observability write grant"] { + assert_eq!( + verdict(401, &[], body.as_bytes().to_vec()), + Verdict::Unauthorized(format!("HTTP 401: {body}")) + ); + } + assert!(matches!(verdict(403, &[], vec![]), Verdict::Unauthorized(_))); + } + + #[test] + fn disabled_by_owner() { + let text = b"project data recording is disabled by owner".to_vec(); + assert_eq!(verdict(401, &[], text), Verdict::Disabled); + let rpc = status_body("Project data recording is disabled by owner", None); + assert_eq!(verdict(403, &[], rpc), Verdict::Disabled); + // Only the documented phrase: another answer that merely mentions "disabled" is a + // credential problem, not a project switch. + for body in ["account disabled", "token disabled for this room"] { + assert!(matches!( + verdict(401, &[], body.as_bytes().to_vec()), + Verdict::Unauthorized(_) + )); + } + } + + /// Finding r1-9: the status decides; a delay hint on a final status is ignored. Only 429, + /// 502–504 and Cloud's 500-with-RetryInfo are retried. + #[test] + fn delay_hints_never_make_a_final_status_retryable() { + type Hint<'a> = (&'a [(&'a str, &'a str)], Vec); + let hints: [Hint; 3] = [ + (&[("Retry-After", "5")], vec![]), + (&[], status_body("later", Some((5, 0)))), + (&[("Retry-After", "5")], status_body("later", Some((5, 0)))), + ]; + for (headers, body) in &hints { + for status in [300, 302, 307, 400, 405, 409, 410, 422, 501, 505, 599] { + assert!( + matches!(verdict(status, headers, body.clone()), Verdict::Rejected(_)), + "{status} with {headers:?}" + ); + } + assert_eq!(verdict(413, headers, body.clone()), Verdict::TooLarge); + assert_eq!(verdict(404, headers, body.clone()), Verdict::NotFound); + assert!(matches!(verdict(429, headers, body.clone()), Verdict::Throttled { .. })); + assert!(matches!(verdict(503, headers, body.clone()), Verdict::Throttled { .. })); + for status in [502, 504] { + assert!(matches!( + verdict(status, headers, body.clone()), + Verdict::Retry { delay_ms: Some(5_000), .. } + )); + } + } + // 500: only with RetryInfo in the body, never on a Retry-After header alone. + assert!(matches!(verdict(500, &[("Retry-After", "5")], vec![]), Verdict::Rejected(_))); + assert!(matches!( + verdict(500, &[], status_body("later", Some((5, 0)))), + Verdict::Retry { delay_ms: Some(5_000), .. } + )); + } + + /// Finding r1-11: untrusted time values never panic or wrap; oversized ones saturate. + #[test] + fn untrusted_time_values_are_bounded() { + assert_eq!(retry_after_ms("99999999999999999999999"), Some(u64::MAX)); + assert_eq!(retry_after_ms("18446744073709551615"), Some(u64::MAX)); + assert_eq!(retry_after_ms("-5"), None); + assert_eq!(retry_after_ms(""), None); + for date in [ + "Sun, 06 Nov 99999999999999 08:49:37 GMT", + "Sun, 06 Nov -9223372036854775808 08:49:37 GMT", + "Sun, 06 Nov 1969 08:49:37 GMT", + "Sun, 06 Nov 10000 08:49:37 GMT", + "Sun, 32 Nov 2026 08:49:37 GMT", + "Sun, 06 Nov 2026 08:49:37:99 GMT", + "Sun, 06 Nov 2026 99:49:37 GMT", + ] { + assert_eq!(imf_fixdate_secs(date), None, "{date}"); + } + assert_eq!(imf_fixdate_secs("Fri, 31 Dec 9999 23:59:59 GMT"), Some(253_402_300_799)); + let huge = verdict(429, &[("Retry-After", "99999999999999999999")], vec![]); + assert!(matches!(huge, Verdict::Throttled { delay_ms: u64::MAX, .. }), "clamped later"); + let far = verdict(429, &[], status_body("q", Some((i64::MAX, 999_999_999)))); + assert!(matches!(far, Verdict::Throttled { delay_ms: u64::MAX, .. })); + } + + #[test] + fn throttling_takes_the_header_then_the_body_then_the_default() { + let both = verdict(429, &[("Retry-After", "7")], status_body("quota", Some((30, 0)))); + assert!(matches!(both, Verdict::Throttled { delay_ms: 7_000, .. })); + let body = verdict(429, &[], status_body("quota", Some((30, 500_000_000)))); + assert!(matches!(body, Verdict::Throttled { delay_ms: 30_500, .. })); + // What LiveKit Cloud answers over quota: ResourceExhausted, no header, no RetryInfo. + assert_eq!( + verdict(429, &[], status_body("QuotaStatusExceeded", None)), + Verdict::Throttled { delay_ms: 60_000, reason: "HTTP 429: QuotaStatusExceeded".into() } + ); + let unavailable = verdict(503, &[("retry-after", " 12 ")], vec![]); + assert!(matches!(unavailable, Verdict::Throttled { delay_ms: 12_000, .. })); + let past = verdict(429, &[("Retry-After", "Wed, 21 Oct 2015 07:28:00 GMT")], vec![]); + assert!(matches!(past, Verdict::Throttled { delay_ms: 0, .. }), "a past date: now"); + let garbage = verdict(429, &[("Retry-After", "soon-ish")], vec![]); + assert!(matches!(garbage, Verdict::Throttled { delay_ms: 60_000, .. }), "ignored"); + let negative = verdict(429, &[], status_body("q", Some((-5, -1)))); + assert!(matches!(negative, Verdict::Throttled { delay_ms: 0, .. }), "clamped"); + } + + #[test] + fn retry_after_http_dates() { + assert_eq!(imf_fixdate_secs("Sun, 06 Nov 1994 08:49:37 GMT"), Some(784_111_777)); + assert_eq!(imf_fixdate_secs("Thu, 01 Jan 1970 00:00:00 GMT"), Some(0)); + assert_eq!(imf_fixdate_secs("Tue, 29 Feb 2028 12:00:00 GMT"), Some(1_835_438_400)); + assert_eq!(imf_fixdate_secs("Sunday, 06-Nov-94 08:49:37 GMT"), None, "obsolete form"); + assert_eq!(imf_fixdate_secs("Sun, 06 Nov 1994 08:49:37 PST"), None); + let now = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .expect("clock") + .as_secs(); + let date = |secs: u64| { + let days = secs / 86_400; + let (h, m, s) = ((secs % 86_400) / 3600, (secs % 3600) / 60, secs % 60); + // Civil from days (inverse of the parser's algorithm) for the fixture only. + let z = days as i64 + 719_468; + let era = z.div_euclid(146_097); + let doe = z - era * 146_097; + let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365; + let doy = doe - (365 * yoe + yoe / 4 - yoe / 100); + let mp = (5 * doy + 2) / 153; + let d = doy - (153 * mp + 2) / 5 + 1; + let mo = if mp < 10 { mp + 3 } else { mp - 9 }; + let y = yoe + era * 400 + if mo <= 2 { 1 } else { 0 }; + const M: [&str; 12] = [ + "Jan", "Feb", "Mar", "Apr", "May", "Jun", "Jul", "Aug", "Sep", "Oct", "Nov", "Dec", + ]; + format!("Mon, {d:02} {} {y} {h:02}:{m:02}:{s:02} GMT", M[mo as usize - 1]) + }; + let ms = retry_after_ms(&date(now + 120)).expect("date"); + assert!((118_000..=120_000).contains(&ms), "{ms}"); + } + + #[test] + fn server_errors() { + for status in [502, 503, 504] { + assert!(matches!(verdict(status, &[], vec![]), Verdict::Retry { delay_ms: None, .. })); + } + assert!(matches!( + verdict(502, &[("Retry-After", "3")], vec![]), + Verdict::Retry { delay_ms: Some(3_000), .. } + )); + // Cloud's retryable 500 carries RetryInfo; a bare 500 is final, per OTLP/HTTP. + let internal = verdict(500, &[], status_body("try later", Some((5, 0)))); + assert!(matches!(internal, Verdict::Retry { delay_ms: Some(5_000), .. })); + assert!(matches!(verdict(500, &[], vec![]), Verdict::Rejected(_))); + assert!(matches!(verdict(501, &[], vec![]), Verdict::Rejected(_))); + } + + /// Finding r1-10: an exception the foreign transport did not declare becomes a retryable + /// transport error. Without the conversion UniFFI panics on the exporter's thread. + #[cfg(feature = "uniffi")] + #[test] + fn an_undeclared_foreign_exception_is_a_retryable_failure() { + let lifted = as uniffi::LiftReturn< + crate::UniFfiTag, + >>::handle_callback_unexpected_error( + uniffi::UnexpectedUniFFICallbackError::new("java.lang.IllegalStateException: boom"), + ); + assert!(matches!( + lifted, + Err(ExportError::Retryable { retry_after_ms: None, ref reason }) if reason.contains("boom") + )); + } +}