From d50b2e52f82783d9cffb12a3cadc9a3f6f0e3755 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:41:00 +0200 Subject: [PATCH 1/3] feat(uniffi): expose the telemetry core `telemetry_configure` installs the process pipeline with a host-implemented `TelemetryTransport` or the `livekit-net` HTTP client, and `telemetry_configure_pulled` hands bindings with thread-bound callbacks (Dart) a bounded pull queue instead. `telemetry_scope` gives each Room its `TelemetryScope`: `set_server` at connect and on every token refresh, spans, the subscribe lifecycle, `record_peer_stats` paced by `stats_poll_interval_ms`, `track_ended`, `disconnected`, and the app API (`emit_custom`, `set_attribute`). `telemetry_disable` is the synchronous opt-out and `telemetry_is_disabled` reads it process-wide; device state, device events and log records come from the host. Dart tests cover the pull queue, and macOS links reserve header padding for Flutter's install name rewrite. --- .cargo/config.toml | 6 +- .changeset/uniffi-telemetry.md | 5 + .gitignore | 1 + Cargo.lock | 1 + livekit-uniffi/Cargo.toml | 1 + livekit-uniffi/src/lib.rs | 3 + livekit-uniffi/src/telemetry.rs | 800 ++++++++++++++++++ .../support/dart/test/telemetry_test.dart | 82 ++ 8 files changed, 897 insertions(+), 2 deletions(-) create mode 100644 .changeset/uniffi-telemetry.md create mode 100644 livekit-uniffi/src/telemetry.rs create mode 100644 livekit-uniffi/support/dart/test/telemetry_test.dart diff --git a/.cargo/config.toml b/.cargo/config.toml index 2adead06d..122b31359 100644 --- a/.cargo/config.toml +++ b/.cargo/config.toml @@ -8,16 +8,18 @@ rustflags = ["-C", "target-feature=+crt-static"] [target.aarch64-pc-windows-msvc] rustflags = ["-C", "target-feature=+crt-static"] +# `-headerpad_max_install_names` on macOS: Flutter's native-assets step rewrites the cdylib's +# install name with install_name_tool, which fails without the padding. # `-ObjC` must also be passed to rustdoc: doctest binaries are linked by rustdoc # and do not inherit `rustflags`. Without it the ObjC categories in libwebrtc's # static lib (e.g. `NSString (StdString)`) are not loaded, and static # initializers such as RTCH264ProfileLevelId.mm abort at process startup. [target.x86_64-apple-darwin] -rustflags = ["-C", "link-args=-ObjC"] +rustflags = ["-C", "link-args=-ObjC -Wl,-headerpad_max_install_names"] rustdocflags = ["-C", "link-args=-ObjC"] [target.aarch64-apple-darwin] -rustflags = ["-C", "link-args=-ObjC"] +rustflags = ["-C", "link-args=-ObjC -Wl,-headerpad_max_install_names"] rustdocflags = ["-C", "link-args=-ObjC"] [target.aarch64-apple-ios] diff --git a/.changeset/uniffi-telemetry.md b/.changeset/uniffi-telemetry.md new file mode 100644 index 000000000..278ef60c9 --- /dev/null +++ b/.changeset/uniffi-telemetry.md @@ -0,0 +1,5 @@ +--- +livekit-uniffi: minor +--- + +Expose the client telemetry core over UniFFI: `telemetry_configure` / `telemetry_configure_pulled` (a bounded pull queue for bindings with thread-bound callbacks, served with `next` or polled with `try_next`), `telemetry_scope` with `TelemetryScope.set_server` (connect and every token refresh), spans, the subscribe lifecycle, `record_peer_stats` + `stats_poll_interval_ms`, `track_ended`, `emit_custom` / `set_attribute` (the Room-scoped app API), `telemetry_disable` (a synchronous opt-out, in effect when it returns), device state and log forwarding, with a host-implemented `TelemetryTransport` or the `livekit-net` HTTP client. diff --git a/.gitignore b/.gitignore index 328df4d65..fbae334f9 100644 --- a/.gitignore +++ b/.gitignore @@ -6,3 +6,4 @@ soxr-sys/test-output.wav .env .cursor __pycache__ +/target-*/ diff --git a/Cargo.lock b/Cargo.lock index 811a82f7d..10dacf191 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3872,6 +3872,7 @@ dependencies = [ "livekit-datatrack", "livekit-net", "livekit-protocol", + "livekit-telemetry", "livekit-token", "log", "once_cell", diff --git a/livekit-uniffi/Cargo.toml b/livekit-uniffi/Cargo.toml index ed131007e..d1f261ddc 100644 --- a/livekit-uniffi/Cargo.toml +++ b/livekit-uniffi/Cargo.toml @@ -16,6 +16,7 @@ livekit-protocol = { workspace = true } livekit-common = { workspace = true, features = ["uniffi"] } livekit-token = { workspace = true } livekit-datatrack = { workspace = true, features = ["uniffi"] } +livekit-telemetry = { workspace = true, features = ["uniffi", "net"] } livekit-net = { workspace = true, features = ["uniffi"] } livekit-data-stream = { workspace = true } uniffi = { workspace = true, features = ["scaffolding-ffi-buffer-fns", "tokio"] } diff --git a/livekit-uniffi/src/lib.rs b/livekit-uniffi/src/lib.rs index c93e40cdc..8a3b5d5d2 100644 --- a/livekit-uniffi/src/lib.rs +++ b/livekit-uniffi/src/lib.rs @@ -21,6 +21,9 @@ pub mod data_stream; /// Access token generation and verification from [`livekit-api::access_token`]. pub mod access_token; +/// Client telemetry core from [`livekit-telemetry`]. +pub mod telemetry; + /// Forward log messages from Rust. pub mod log_forward; diff --git a/livekit-uniffi/src/telemetry.rs b/livekit-uniffi/src/telemetry.rs new file mode 100644 index 000000000..9274b1a59 --- /dev/null +++ b/livekit-uniffi/src/telemetry.rs @@ -0,0 +1,800 @@ +//! Client telemetry core from the [`livekit-telemetry`] crate. +//! +//! What a platform calls, and when: +//! +//! - SDK init: [`telemetry_configure`] (or [`telemetry_configure_pulled`] where foreign callbacks +//! are thread-bound) with storage and instruments. No endpoint, no headers: destinations come +//! from the rooms. +//! - Room created: [`telemetry_scope`]; connect: [`TelemetryScope::set_server`] with the server +//! URL and the participant token, and again with every refreshed token. +//! - Room life: spans ([`TelemetryScope::start`]), subscribe lifecycle, one `getStats()` report +//! per peer connection ([`TelemetryScope::record_peer_stats`]) every +//! [`TelemetryScope::stats_poll_interval_ms`], [`TelemetryScope::disconnected`]. +//! - App API: [`TelemetryScope::emit_custom`], [`TelemetryScope::set_attribute`] (room-scoped), +//! [`telemetry_disable`] (process-wide opt-out). +//! - OS signals: [`telemetry_set_device_state`], [`telemetry_device_event`], [`telemetry_log`]. +//! +//! Transport: pass a host-implemented `TelemetryTransport`, or `None` to ride the HTTP client +//! the host registered with `livekit-net` (`set_http_client`) for signaling. + +use std::{ + collections::HashMap, + sync::{ + atomic::{AtomicBool, AtomicU64, Ordering}, + Arc, Mutex, + }, +}; + +use livekit_telemetry::{ + global::{self, TelemetryInstrument}, + Attribute, AttributeValue, DeviceEvent, DeviceState, DisconnectReason, ExportError, + ExportRequest, ExportResponse, LogRecord, NetTransport, ReconnectReason, RoomIdentity, RtcStat, + SpanName, SpanOutcome, SpanStep, SpanTrack, StreamDirection, TelemetryConfig, TelemetryEvent, + TelemetryStats, TelemetryTransport, TraceContext, TrackKind, +}; +use tokio::sync::{mpsc, oneshot, watch}; + +/// Why a [`Telemetry`] pipeline could not be created. +#[derive(uniffi::Error, thiserror::Error, Debug)] +#[uniffi(flat_error)] +pub enum TelemetryError { + /// No transport was passed and no HTTP client is registered with `livekit-net`. + #[error("no telemetry transport: pass one, or register an HTTP client with livekit-net first")] + NoTransport, +} + +/// Start the process pipeline; a previous one is drained and replaced. `transport = None` uses +/// the HTTP client registered with `livekit-net`, if any. After [`telemetry_disable`] nothing +/// starts, and whatever `config.storage_dir` still holds is deleted. +#[uniffi::export] +pub fn telemetry_configure( + config: TelemetryConfig, + transport: Option>, + instruments: Vec>, +) -> Result<(), TelemetryError> { + let transport: Arc = match transport { + Some(transport) => transport, + None => Arc::new(NetTransport::from_registry().ok_or(TelemetryError::NoTransport)?), + }; + install(livekit_telemetry::Telemetry::new(config, transport), instruments); + Ok(()) +} + +/// Like [`telemetry_configure`], exporting through a queue the host drains from its own thread. +/// For bindings whose callbacks cannot be invoked from Rust threads (uniffi-dart today): such a +/// host passes no `instruments` — it starts and stops its own around configure and disable, so +/// Rust never calls back into it — and polls the queue with +/// [`try_next`](TelemetryExportQueue::try_next). +#[uniffi::export] +pub fn telemetry_configure_pulled( + config: TelemetryConfig, + instruments: Vec>, +) -> Arc { + let queue = TelemetryExportQueue::new(); + queue.process.store(true, Ordering::SeqCst); + install(livekit_telemetry::Telemetry::new(config, queue.clone()), instruments); + queue +} + +fn install( + (telemetry, exporter): (livekit_telemetry::Telemetry, livekit_telemetry::Exporter), + instruments: Vec>, +) { + let previous = global::install(telemetry, instruments); + // Refused after the opt-out (the new pipeline was purged instead, its `storage_dir` with it), + // or opted out meanwhile (every generation revoked): nothing to run, nothing to drain. + if global::is_disabled() { + return; + } + crate::runtime::runtime().spawn(exporter.run()); + if let Some(previous) = previous { + crate::runtime::runtime().spawn(async move { previous.shutdown().await }); + } +} + +/// Drain and stop the process pipeline; every call is a no-op afterwards. After +/// [`telemetry_disable`], returns once its purge is complete. +#[uniffi::export(async_runtime = "tokio")] +pub async fn telemetry_shutdown() { + purged().await; + global::shutdown().await; +} + +/// Opt-out for the rest of the process, in effect when this returns: no scope is handed out, +/// every existing scope and instrument captures nothing, and later configures are refused. +/// Everything not yet sent — queued, open, cached on disk — is then deleted in the background +/// without another upload (awaited by [`telemetry_shutdown`] and [`telemetry_flush`]). Data +/// already sent cannot be recalled. +/// +/// The installed instruments' `stop()` runs synchronously on the calling thread, under the +/// lifecycle lock, before this returns: it must not block on the caller's queue and must not +/// call back into configure, disable, flush or shutdown. +#[uniffi::export] +pub fn telemetry_disable() { + start_purge(&PURGE, global::disable()); +} + +/// Whether the process opted out with [`telemetry_disable`]: `true` for every thread and isolate +/// from the moment that call revokes collection, before it returns. A cheap read that never +/// panics; platforms with per-isolate state may check it right before each `getStats` or submit. +#[uniffi::export] +pub fn telemetry_is_disabled() -> bool { + global::is_disabled() +} + +/// The opt-out's purge in progress: done once it holds `true`. +type PurgeSlot = Mutex>>; + +/// The process-wide opt-out's purge. +static PURGE: PurgeSlot = Mutex::new(None); + +/// Run `purge` in the background and record it in `slot`, done only once every purge recorded +/// there before it is done too: a repeated opt-out never lets a waiter skip the first purge. +fn start_purge(slot: &PurgeSlot, purge: impl std::future::Future + Send + 'static) { + let (done, receiver) = watch::channel(false); + let previous = slot.lock().unwrap_or_else(|e| e.into_inner()).replace(receiver); + crate::runtime::runtime().spawn(async move { + purge.await; + if let Some(mut previous) = previous { + let _ = previous.wait_for(|done| *done).await; + } + let _ = done.send(true); + }); +} + +/// Wait for the opt-out's purge, if one was started. +async fn purged() { + wait_purge(&PURGE).await; +} + +/// Wait for the purge recorded in `slot`, if any. +async fn wait_purge(slot: &PurgeSlot) { + let receiver = slot.lock().unwrap_or_else(|e| e.into_inner()).clone(); + if let Some(mut receiver) = receiver { + let _ = receiver.wait_for(|done| *done).await; + } +} + +/// Cache everything queued and upload what the network allows. After [`telemetry_disable`], +/// returns once its purge is complete. +#[uniffi::export(async_runtime = "tokio")] +pub async fn telemetry_flush() { + purged().await; + global::flush().await; +} + +/// A scope — one room, one call — on the process pipeline; `None` while telemetry is off. +#[uniffi::export] +pub fn telemetry_scope() -> Option> { + global::scope().map(|scope| Arc::new(TelemetryScope(scope))) +} + +/// A process-level event (outside any room). +#[uniffi::export] +pub fn telemetry_emit(event: TelemetryEvent) { + global::emit(event); +} + +/// A captured log line; the core applies the per-source floor and builds the record. For the +/// platform's own and WebRTC's lines only — never a Rust entry received through the log +/// forwarder, which the core has already copied (`log_forward_bootstrap`). +#[uniffi::export] +pub fn telemetry_log(record: LogRecord) { + global::log(record); +} + +/// Audio route, interruption, denied permission: a process-level record built by the core. +#[uniffi::export] +pub fn telemetry_device_event(event: DeviceEvent) { + global::device_event(event); +} + +/// The device's current state; drives the upload cadence and yields change events. +#[uniffi::export] +pub fn telemetry_set_device_state(state: DeviceState) { + global::set_device_state(state); +} + +/// SDK metadata on every record of every scope (not an app API: app attributes are +/// room-scoped, [`TelemetryScope::set_attribute`]); `None` removes it. +#[uniffi::export] +pub fn telemetry_set_attribute(key: String, value: Option) { + global::set_attribute(&key, value); +} + +/// Pipeline health (drops by reason, uploads, backlog); `None` while telemetry is off. +#[uniffi::export] +pub fn telemetry_stats() -> Option { + global::stats() +} + +/// The stats as one line for a debug console, or `off`. +#[uniffi::export] +pub fn telemetry_diagnostics() -> String { + global::diagnostics() +} + +/// One Room's telemetry session: its trace, attributes, spans, stats and destination. +#[derive(uniffi::Object)] +pub struct TelemetryScope(livekit_telemetry::Scope); + +#[uniffi::export] +impl TelemetryScope { + /// The session's trace id as 32 hex characters. + pub fn trace_id(&self) -> String { + self.0.trace_id() + } + + /// An SDK event filed under this session. + pub fn emit(&self, event: TelemetryEvent) { + self.0.emit(event); + } + + /// The room's server URL and participant token: at connect and with every refreshed token. + /// Cheap and idempotent. The core derives the ingest URL (LiveKit Cloud only), reads the + /// grant and expiry, and uploads this room's records with this room's token. + pub fn set_server(&self, url: String, token: String) { + self.0.set_server(&url, &token); + } + + /// An app event, exported as `custom.` with the room's correlation attributes. Names + /// and keys up to 128 bytes, values up to 1024, at most 64 attributes, no `lk.*` keys; + /// anything else is rejected and counted. + pub fn emit_custom(&self, name: String, attributes: HashMap) { + let attributes = attributes.into_iter().map(|(k, v)| Attribute::new(k, v)).collect(); + self.0.emit_custom(&name, attributes); + } + + /// An app correlation attribute on every record this room captures from now on; `None` + /// removes it. Same limits as `emit_custom`, at most 64 per room. + pub fn set_attribute(&self, key: String, value: Option) { + self.0.set_attribute(&key, value.map(AttributeValue::Str)); + } + + /// One peer connection's whole `getStats()` report and its tracks (MediaStreamTrack id → + /// track sid); the core maps every RTP stream to its track. + pub fn record_peer_stats( + &self, + report: Vec, + tracks: HashMap, + timestamp_ns: Option, + ) { + self.0.record_peer_stats(report, tracks, timestamp_ns); + } + + /// When to poll `getStats()` next, in milliseconds; ask after every poll. + pub fn stats_poll_interval_ms(&self) -> u64 { + self.0.stats_poll_interval_ms() + } + + /// Start a typed span in this session's trace, stamped now; `parent` nests it. + pub fn start(&self, name: SpanName, parent: Option>) -> Arc { + Arc::new(TelemetrySpan(self.0.start(name, parent.map(|p| p.0.clone())))) + } + + /// The room and local participant, on every record of this session from now on. + pub fn set_room(&self, room: RoomIdentity) { + self.0.set_room(room); + } + + /// The session ended for good (not a reconnect). + pub fn disconnected(&self, reason: DisconnectReason) { + self.0.disconnected(reason); + } + + /// A log record filed under this session even without an ambient span. + pub fn log(&self, record: LogRecord) { + self.0.log(record); + } + + /// One track's whole `getStats()` report; the core maps it (see `record_stats`). + pub fn record_stats_report( + &self, + track_sid: String, + kind: TrackKind, + direction: StreamDirection, + report: Vec, + timestamp_ns: Option, + ) { + self.0.record_stats_report(&track_sid, kind, direction, report, timestamp_ns); + } + + /// Intent to subscribe (autoSubscribe: the remote publish; manual: the call; tracks already + /// in the room at join: at connect): opens `lk.subscribe`. + pub fn subscribe_started(&self, track: SpanTrack) { + self.0.subscribe_started(track); + } + + /// The server confirmed the subscription. Without an earlier `subscribe_started` this is the + /// intent (a fallback, measured from here): `lk.subscribe` opens and `stats_poll_interval_ms` + /// turns fast at once. + pub fn subscribed(&self, track: SpanTrack) { + self.0.subscribed(track); + } + + /// A track left the room (unpublished, unsubscribed, publisher gone): a pending subscribe + /// ends, the track's last RTC window ships, its state is forgotten. + pub fn track_ended(&self, sid: String) { + self.0.track_ended(&sid); + } + + /// The subscription failed; `error_type` is the platform's error name. + pub fn subscribe_failed(&self, sid: String, error_type: String) { + self.0.subscribe_failed(&sid, &error_type); + } +} + +/// The protocol's `DisconnectReason` number as the shared enum. +#[uniffi::export] +pub fn telemetry_disconnect_reason(proto: i32) -> DisconnectReason { + DisconnectReason::from_proto(proto) +} + +/// The protocol's `ReconnectReason` number (`RR_*`) as the shared enum. +#[uniffi::export] +pub fn telemetry_reconnect_reason(proto: i32) -> ReconnectReason { + ReconnectReason::from_proto(proto) +} + +/// One export the host has to perform on behalf of a pulled pipeline. +#[derive(uniffi::Record)] +pub struct PendingExport { + pub id: u64, + pub request: ExportRequest, +} + +struct Pending { + export: PendingExport, + done: oneshot::Sender>, +} + +/// Pull-side transport: Rust never calls into the host. The exporter queues each request; the +/// host awaits [`next`](Self::next) (a Rust future — those cross every binding) or polls +/// [`try_next`](Self::try_next), performs the HTTP call on its own thread, and reports the +/// outcome with [`complete`](Self::complete), +/// which unblocks the exporter's retry/drop/go-silent logic exactly as a direct transport would. +/// +/// Exists because uniffi-dart's foreign-trait callbacks are isolate-bound (`Pointer.fromFunction`) +/// and abort the VM when invoked from a tokio thread; Swift and Kotlin callbacks are thread-agnostic +/// and use [`TelemetryTransport`] directly. +/// One attempt at an SDK operation, owned by the core. Every call is synchronous and stamps the +/// clock inside, so the only skew is the FFI call; `describe()` is the console line on every +/// platform. Detached (no session) it still times and describes itself. +#[derive(uniffi::Object)] +pub struct TelemetrySpan(Arc); + +#[uniffi::export] +impl TelemetrySpan { + /// A checkpoint, stamped now. + pub fn step(&self, step: SpanStep) { + self.0.step(step); + } + + /// The open bag; replaces an existing key. + pub fn set_attribute(&self, key: String, value: AttributeValue) { + self.0.set_attribute(key, value); + } + + /// The track the operation is about. + pub fn set_track(&self, track: SpanTrack) { + self.0.set_track(track); + } + + /// End once; `error` becomes `error.type` and the status message. + pub fn end(&self, outcome: SpanOutcome, error: Option) { + self.0.end(outcome, error); + } + + /// End with an error; `error` becomes `error.type`. + pub fn fail(&self, error: String) { + self.0.fail(error); + } + + /// End as cancelled. + pub fn cancel(&self) { + self.0.cancel(); + } + + /// Whether the span has ended. + pub fn is_ended(&self) -> bool { + self.0.is_ended() + } + + /// `None` for a detached span. + pub fn context(&self) -> Option { + self.0.context() + } + + /// Seconds to the end, or to the last step while running. + pub fn total_secs(&self) -> f64 { + self.0.total_secs() + } + + /// `lk.connect: ws_open +1.49s, …, total 1.83s, ok` + pub fn describe(&self) -> String { + self.0.describe() + } +} + +/// The pull queue's serving side (see its type docs above [`TelemetrySpan`]). +#[derive(uniffi::Object)] +pub struct TelemetryExportQueue { + /// Bounded to one request: the exporter has one request out at a time, so nothing piles up + /// behind a host that stopped serving. `None` once finished. + tx: Mutex>>, + rx: tokio::sync::Mutex>, + inflight: Mutex>>>, + seq: AtomicU64, + finished: AtomicBool, + /// Serves the process pipeline (`telemetry_configure_pulled`): nothing is handed out once + /// the app opted out, even before the purge has cancelled the request. + process: AtomicBool, + /// Requests handed to the queue so far (tests only). + #[cfg(test)] + queued: AtomicU64, +} + +#[uniffi::export(async_runtime = "tokio")] +impl TelemetryExportQueue { + /// An empty queue; `telemetry_configure_pulled` makes the one it serves. + #[uniffi::constructor] + pub fn new() -> Arc { + let (tx, rx) = mpsc::channel(1); + Arc::new(Self { + tx: Mutex::new(Some(tx)), + rx: tokio::sync::Mutex::new(rx), + inflight: Mutex::new(HashMap::new()), + seq: AtomicU64::new(0), + finished: AtomicBool::new(false), + process: AtomicBool::new(false), + #[cfg(test)] + queued: AtomicU64::new(0), + }) + } + + /// End the serving loop: `next` resolves `None` from now on, and whatever is still queued is + /// discarded, never served. Call after `telemetry_shutdown` or `telemetry_disable`, or when a + /// later `telemetry_configure_pulled` replaced this queue. + /// (Not `close`: UniFFI's Kotlin objects already have `AutoCloseable.close`.) + pub fn finish(&self) { + self.finished.store(true, Ordering::SeqCst); + self.tx.lock().unwrap_or_else(|e| e.into_inner()).take(); + self.inflight.lock().unwrap_or_else(|e| e.into_inner()).clear(); + // Dropping what is queued tells its exporter the host dropped it. + if let Ok(mut rx) = self.rx.try_lock() { + while rx.try_recv().is_ok() {} + } + } + + /// The next request to perform; `None` once finished. A request the exporter already gave up + /// on (timed out, cancelled by shutdown or the opt-out) is never handed out. + pub async fn next(&self) -> Option { + loop { + if self.finished.load(Ordering::SeqCst) { + return None; + } + let pending = self.rx.lock().await.recv().await?; + if let Some(export) = self.serve(pending) { + return Some(export); + } + if self.finished.load(Ordering::SeqCst) { + return None; + } + } + } + + /// The request waiting right now, if any — without waiting: for hosts that must never leave + /// a pending Rust future holding one of their continuations (Dart: poll from a `Timer`). A + /// request nobody polls for still times out and is withdrawn exactly as with + /// [`next`](Self::next); `None` also once finished, or while a `next` is waiting. + pub fn try_next(&self) -> Option { + let mut rx = self.rx.try_lock().ok()?; + loop { + if self.finished.load(Ordering::SeqCst) { + return None; + } + if let Some(export) = self.serve(rx.try_recv().ok()?) { + return Some(export); + } + } + } + + /// The collector's answer to the request with `id`, whatever its status; the core classifies it. + pub fn complete(&self, id: u64, response: ExportResponse) { + self.settle(id, Ok(response)); + } + + /// The request with `id` got no answer (network error, timeout, invalid URL). + pub fn fail(&self, id: u64, error: ExportError) { + self.settle(id, Err(error)); + } +} + +impl TelemetryExportQueue { + /// Hand `pending` out, unless the queue is finished or its exporter gave up on it. + fn serve(&self, pending: Pending) -> Option { + let opted_out = self.process.load(Ordering::SeqCst) && global::is_disabled(); + if opted_out || self.finished.load(Ordering::SeqCst) || pending.done.is_closed() { + return None; + } + let mut inflight = self.inflight.lock().unwrap_or_else(|e| e.into_inner()); + inflight.retain(|_, done| !done.is_closed()); + inflight.insert(pending.export.id, pending.done); + Some(pending.export) + } + + fn settle(&self, id: u64, outcome: Result) { + let done = self.inflight.lock().unwrap_or_else(|e| e.into_inner()).remove(&id); + if let Some(done) = done { + let _ = done.send(outcome); + } + } + + /// Drop queued requests nobody waits for any more (their exporter timed out or stopped): + /// with one request out at a time, anything still queued when a new one comes is stale. + fn discard_stale(&self) { + let Ok(mut rx) = self.rx.try_lock() else { return }; + let mut live = Vec::new(); + while let Ok(pending) = rx.try_recv() { + if !pending.done.is_closed() { + live.push(pending); + } + } + drop(rx); + let tx = self.tx.lock().unwrap_or_else(|e| e.into_inner()).clone(); + for pending in live { + if let Some(tx) = &tx { + let _ = tx.try_send(pending); + } + } + } +} + +/// Forgets an in-flight entry when the exporter stops waiting for its answer (timeout, +/// shutdown, opt-out), so nothing is kept for an answer nobody will read. +struct Forget<'a> { + queue: &'a TelemetryExportQueue, + id: u64, +} + +impl Drop for Forget<'_> { + fn drop(&mut self) { + self.queue.inflight.lock().unwrap_or_else(|e| e.into_inner()).remove(&self.id); + } +} + +#[async_trait::async_trait] +impl TelemetryTransport for TelemetryExportQueue { + async fn send(&self, request: ExportRequest) -> Result { + let closed = || ExportError::Retryable { + reason: "export queue closed".into(), + retry_after_ms: None, + }; + self.discard_stale(); + let id = self.seq.fetch_add(1, Ordering::Relaxed); + let (done, wait) = oneshot::channel(); + let pending = Pending { export: PendingExport { id, request }, done }; + let _forget = Forget { queue: self, id }; + let Some(tx) = self.tx.lock().unwrap_or_else(|e| e.into_inner()).clone() else { + return Err(closed()); + }; + // Waits for room while the host is busy; cancelled with the exporter's timeout, and then + // the request was never queued. The sender is not kept while the answer is awaited. + tx.send(pending).await.map_err(|_| closed())?; + drop(tx); + #[cfg(test)] + self.queued.fetch_add(1, Ordering::SeqCst); + wait.await.unwrap_or(Err(ExportError::Retryable { + reason: "host dropped the export".into(), + retry_after_ms: None, + })) + } +} + +#[cfg(test)] +mod tests { + use std::time::Duration; + + use super::*; + + fn request() -> ExportRequest { + ExportRequest { + url: "https://p.livekit.cloud/observability/client/logs/otlp/v0".into(), + headers: HashMap::from([("Authorization".into(), "Bearer secret".into())]), + body: vec![1, 2, 3], + } + } + + fn runtime() -> tokio::runtime::Runtime { + tokio::runtime::Builder::new_current_thread().enable_time().build().expect("runtime") + } + + /// Finding r1-4: a host that stopped serving never gets a pile of stale requests — and + /// their tokens — to send later: requests the exporter gave up on are discarded, not served. + #[test] + fn a_suspended_host_is_never_served_stale_requests() { + runtime().block_on(async { + let queue = TelemetryExportQueue::new(); + for _ in 0..5 { + let sent = tokio::time::timeout(Duration::from_millis(10), queue.send(request())); + assert!(sent.await.is_err(), "the exporter's timeout gives up"); + } + assert!(queue.inflight.lock().expect("lock").is_empty(), "nothing kept for answers"); + let served = tokio::time::timeout(Duration::from_millis(50), queue.next()).await; + assert!(served.is_err(), "no stale request is handed out"); + }); + } + + /// Finding r1-4: a request cancelled before the host asked for it (shutdown, opt-out) is + /// never served; a live one still is. + #[test] + fn a_cancelled_request_is_never_served() { + runtime().block_on(async { + let queue = TelemetryExportQueue::new(); + let cancelled = { + let queue = queue.clone(); + tokio::spawn(async move { queue.send(request()).await }) + }; + tokio::task::yield_now().await; + cancelled.abort(); + let _ = cancelled.await; + let live = { + let queue = queue.clone(); + tokio::spawn(async move { queue.send(request()).await }) + }; + tokio::task::yield_now().await; + let served = queue.next().await.expect("the live request"); + queue.complete(served.id, ExportResponse::accepted()); + assert_eq!(live.await.expect("task"), Ok(ExportResponse::accepted())); + }); + } + + /// Finding r1-4: finishing the queue discards what is queued instead of serving it. + #[test] + fn finishing_discards_what_is_queued() { + runtime().block_on(async { + let queue = TelemetryExportQueue::new(); + let waiting = { + let queue = queue.clone(); + tokio::spawn(async move { queue.send(request()).await }) + }; + tokio::task::yield_now().await; + queue.finish(); + assert!(queue.next().await.is_none(), "not served after finish"); + assert!(waiting.await.expect("task").is_err(), "the exporter hears it was dropped"); + }); + } + + /// Flutter review: `try_next` never waits. It serves a live request, never one its exporter + /// gave up on while nobody polled, and nothing once finished. + #[test] + fn try_next_serves_without_waiting_and_never_a_stale_request() { + runtime().block_on(async { + let queue = TelemetryExportQueue::new(); + assert!(queue.try_next().is_none(), "empty: returns at once"); + let sent = tokio::time::timeout(Duration::from_millis(10), queue.send(request())); + assert!(sent.await.is_err(), "nobody polled: the exporter's timeout gives up"); + assert!(queue.try_next().is_none(), "the withdrawn request is never served"); + let live = { + let queue = queue.clone(); + tokio::spawn(async move { queue.send(request()).await }) + }; + tokio::task::yield_now().await; + let served = queue.try_next().expect("the live request"); + queue.complete(served.id, ExportResponse::accepted()); + assert_eq!(live.await.expect("task"), Ok(ExportResponse::accepted())); + queue.complete(served.id, ExportResponse::accepted()); // stale id: ignored + queue.fail( + u64::MAX, + ExportError::Retryable { reason: "x".into(), retry_after_ms: None }, + ); + queue.finish(); + assert!(queue.try_next().is_none()); + }); + } + + /// An unsigned participant token with the observability grant (the core reads claims, never + /// verifies them). + fn granted_token() -> String { + fn b64url(bytes: &[u8]) -> String { + const ABC: &[u8; 64] = + b"ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789-_"; + let mut out = String::new(); + for chunk in bytes.chunks(3) { + let n = chunk + .iter() + .enumerate() + .fold(0u32, |n, (i, b)| n | (*b as u32) << (16 - 8 * i)); + for i in 0..=chunk.len() { + out.push(ABC[(n >> (18 - 6 * i) & 63) as usize] as char); + } + } + out + } + let exp = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .expect("clock") + .as_secs() + + 3600; + let claims = format!(r#"{{"exp":{exp},"observability":{{"write":true}}}}"#); + format!("{}.{}.sig", b64url(br#"{"alg":"HS256"}"#), b64url(claims.as_bytes())) + } + + /// Wait (bounded) until `condition` holds. + fn eventually(condition: impl Fn() -> bool) -> bool { + (0..500).any(|_| { + condition() || { + std::thread::sleep(Duration::from_millis(10)); + false + } + }) + } + + /// Final review: a second opt-out while the first one's purge is blocked does not let a + /// waiter (`telemetry_flush`/`telemetry_shutdown`) return before that first purge finishes. + #[test] + fn a_repeated_opt_out_waits_for_the_first_purge() { + static SLOT: PurgeSlot = Mutex::new(None); + let (release, blocked) = oneshot::channel::<()>(); + start_purge(&SLOT, blocked); + start_purge(&SLOT, async {}); + runtime().block_on(async { + let early = tokio::time::timeout(Duration::from_millis(200), wait_purge(&SLOT)).await; + assert!(early.is_err(), "the first purge is still running"); + release.send(()).expect("first purge waiting"); + tokio::time::timeout(Duration::from_secs(5), wait_purge(&SLOT)) + .await + .expect("done once the first purge is"); + }); + } + + /// Finding r2-2 + Swift review: the opt-out is in effect when `telemetry_disable` returns — + /// no scope, nothing captured — and once its purge is awaited (`telemetry_shutdown`) the + /// pulled request of a pipeline replaced while its upload was pending is never handed to the + /// host; `telemetry_is_disabled` turns `true` for every thread. (The only test here that + /// touches the process-wide pipeline: the opt-out is one-way.) + #[test] + fn the_opt_out_withdraws_a_replaced_generations_pulled_request() { + let config = || TelemetryConfig { resource: Vec::new(), ..Default::default() }; + let first = telemetry_configure_pulled(config(), Vec::new()); + let room = telemetry_scope().expect("scope"); + room.set_server("wss://p.livekit.cloud".into(), granted_token()); + room.emit_custom("ping".into(), HashMap::new()); + crate::runtime::runtime().spawn(telemetry_flush()); + assert!( + eventually(|| first.queued.load(Ordering::SeqCst) == 1), + "precondition: a pulled request is pending in the first queue" + ); + let _second = telemetry_configure_pulled(config(), Vec::new()); // `first` now drains + assert!(!telemetry_is_disabled()); + telemetry_disable(); + assert!(telemetry_is_disabled(), "visible process-wide once the call returns"); + assert!(std::thread::spawn(telemetry_is_disabled).join().expect("thread")); + assert!(first.try_next().is_none(), "nothing is served once the call returns"); + assert!(telemetry_scope().is_none(), "no scope once the call returns"); + assert!(telemetry_stats().is_none()); + room.emit_custom("after".into(), HashMap::new()); // refused: its generation is revoked + // A configure after the opt-out deletes what its storage dir still holds and starts no + // exporter: the queue it hands back is held by nothing else. + let dir = std::env::temp_dir().join(format!("lk-refused-{}", std::process::id())); + let _ = std::fs::remove_dir_all(&dir); + std::fs::create_dir(&dir).expect("dir"); + let stamp = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .expect("clock") + .as_nanos(); + std::fs::write(dir.join(format!("{stamp:020}-000001-1-l.otlp")), b"cached").expect("batch"); + let refused = telemetry_configure_pulled( + TelemetryConfig { storage_dir: Some(dir.to_string_lossy().into_owned()), ..config() }, + Vec::new(), + ); + assert_eq!(Arc::strong_count(&refused), 1, "no exporter holds its transport"); + assert_eq!(std::fs::read_dir(&dir).expect("dir").count(), 0, "the cache is purged"); + let _ = std::fs::remove_dir_all(&dir); + runtime().block_on(async { + telemetry_shutdown().await; // awaits the purge + let served = tokio::time::timeout(Duration::from_millis(300), first.next()).await; + assert!(!matches!(served, Ok(Some(_))), "withdrawn, never served"); + }); + } +} diff --git a/livekit-uniffi/support/dart/test/telemetry_test.dart b/livekit-uniffi/support/dart/test/telemetry_test.dart new file mode 100644 index 000000000..e0cc9789b --- /dev/null +++ b/livekit-uniffi/support/dart/test/telemetry_test.dart @@ -0,0 +1,82 @@ +import 'dart:convert'; +import 'dart:typed_data'; + +import 'package:livekit_uniffi/livekit_telemetry.dart'; +import 'package:livekit_uniffi/livekit_uniffi.dart'; +import 'package:test/test.dart'; + +/// Rust must never call back into Dart from its own threads: uniffi-dart's foreign-trait +/// callbacks are isolate-bound (`Pointer.fromFunction`) and the VM aborts with "Cannot invoke +/// native callback outside an isolate" when the exporter invokes `TelemetryTransport.send` from +/// a tokio worker. The pull queue inverts the direction: Dart polls `tryNext()` from its own +/// timer (never leaving a Rust future holding a continuation of an isolate that may die), +/// performs the request, and reports the collector's answer back with `complete()` (or `fail()` +/// when there was none); the core classifies the status and body. +Future serve(TelemetryExportQueue queue, List sink, int count) async { + for (var i = 0; i < count; i++) { + var pending = queue.tryNext(); + while (pending == null) { + await Future.delayed(const Duration(milliseconds: 10)); + pending = queue.tryNext(); + } + sink.add(pending.request); + queue.complete(id: pending.id, response: ExportResponse(status: 200, headers: {}, body: Uint8List(0))); + } +} + +/// An unsigned participant token with the observability grant: the core reads claims, never +/// verifies them (the collector does). +String grantedToken() { + String part(Map json) => + base64Url.encode(utf8.encode(jsonEncode(json))).replaceAll('=', ''); + final exp = DateTime.now().millisecondsSinceEpoch ~/ 1000 + 3600; + return '${part({'alg': 'HS256'})}.${part({'exp': exp, 'observability': {'write': true}})}.sig'; +} + +void main() { + group('telemetry', () { + test('exports a room\'s records to its project through the pull queue', () async { + final requests = []; + final queue = telemetryConfigurePulled( + config: TelemetryConfig(resource: [], logSeverity: Severity.warn), + instruments: [], + ); + final serving = serve(queue, requests, 2); + + final room = telemetryScope()!; + final token = grantedToken(); + room.setServer(url: 'wss://my-project.livekit.cloud', token: token); + room.setAttribute(key: 'app.call_id', value: 'c1'); + room.emitCustom(name: 'checkout', attributes: {'plan': 'pro'}); + await telemetryFlush(); + expect(requests, hasLength(1)); + expect(requests.single.url, 'https://my-project.livekit.cloud/observability/client/logs/otlp/v0'); + expect(requests.single.headers['Authorization'], 'Bearer $token'); + expect(requests.single.headers['Content-Type'], 'application/x-protobuf'); + expect(requests.single.body, isNotEmpty); + expect(telemetryStats()!.uploadsSent, 1); + expect(telemetryStats()!.dropped, 0); + expect(room.statsPollIntervalMs(), 30000); + + // Shutdown leaves the session summary, a process-level batch of its own. + await telemetryShutdown(); + await serving; + expect(requests, hasLength(2)); + expect(telemetryStats(), isNull); + // The queue outlives the pipeline (Dart holds it): finishing it ends the serving loop. + queue.finish(); + expect(await queue.next(), isNull); + }); + + test('refuses to start without any transport', () { + expect( + () => telemetryConfigure( + config: TelemetryConfig(resource: [], logSeverity: Severity.warn), + transport: null, + instruments: [], + ), + throwsA(anything), + ); + }); + }); +} From 99452cba6342267dc3b73a59153a621436143a84 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:41:24 +0200 Subject: [PATCH 2/3] feat(uniffi): copy the core's own warnings and errors into telemetry Once a platform installs the log forwarder, warnings and errors from `livekit*` targets also reach the telemetry pipeline as log records, with any compact JWT and bearer or token credential values masked; the telemetry crate's own records are never copied. The platform's level still filters only what is forwarded to it, so the console output is unchanged. --- livekit-uniffi/src/log_forward.rs | 324 +++++++++++++++++++++++++++++- 1 file changed, 322 insertions(+), 2 deletions(-) diff --git a/livekit-uniffi/src/log_forward.rs b/livekit-uniffi/src/log_forward.rs index 1cb669c55..8123b5a57 100644 --- a/livekit-uniffi/src/log_forward.rs +++ b/livekit-uniffi/src/log_forward.rs @@ -12,6 +12,8 @@ // See the License for the specific language governing permissions and // limitations under the License. +use std::sync::atomic::{AtomicUsize, Ordering}; + use log::{Level, LevelFilter, Log, Record}; use once_cell::sync::OnceCell; use tokio::sync::{mpsc, Mutex}; @@ -24,11 +26,18 @@ static LOGGER: OnceCell = OnceCell::new(); /// Generally, you will invoke this once early in program execution. However, /// subsequent invocations are allowed to change the log level. /// +/// Also turns on the telemetry copy of the core's own warnings and errors (see +/// `telemetry_copy`), whatever `level` is: `level` filters only what is forwarded here. Don't +/// pass the forwarded entries on to `telemetry_log`. +/// #[uniffi::export] fn log_forward_bootstrap(level: LevelFilter) { let logger = LOGGER.get_or_init(Logger::new); _ = log::set_logger(logger); // Returns an error if already set (ignore) - log::set_max_level(level); + logger.forward.store(level as usize, Ordering::Relaxed); + // Never stricter than `Warn`: the telemetry copy needs the core's warnings even when the + // platform forwards only errors. + log::set_max_level(level.max(LevelFilter::Warn)); } /// Asynchronously receives a forwarded log entry. @@ -77,12 +86,15 @@ pub struct LogForwardEntry { struct Logger { tx: mpsc::UnboundedSender, rx: Mutex>, + /// The platform's level (a `LevelFilter` as `usize`): what is forwarded to it. The global + /// `log` max level may be looser, for the telemetry copy. + forward: AtomicUsize, } impl Logger { fn new() -> Self { let (tx, rx) = mpsc::unbounded_channel(); - Self { tx, rx: rx.into() } + Self { tx, rx: rx.into(), forward: AtomicUsize::new(LevelFilter::Trace as usize) } } } @@ -91,6 +103,15 @@ impl Log for Logger { true } fn log(&self, record: &Record) { + // Telemetry gets its copy first, built from the borrowed record; the console entry below + // is exactly what it always was. + if let Some(copy) = telemetry_copy(record) { + livekit_telemetry::global::log(copy); + } + // `Level` and `LevelFilter` share their numbering (`Error` = 1 … `Trace` = 5). + if record.level() as usize > self.forward.load(Ordering::Relaxed) { + return; + } let record = LogForwardEntry { level: record.metadata().level(), target: record.target().to_string(), @@ -103,3 +124,302 @@ impl Log for Logger { } fn flush(&self) {} } + +/// The core's own warnings and errors — `livekit*` targets at `Warn` or `Error` — as telemetry +/// log records, so no platform has to route them there itself. Works once the platform has +/// installed this forwarder (`log_forward_bootstrap`): Rust allows one logger per process, and +/// without it the core's `log` records go nowhere. Platforms must therefore not also feed these +/// forwarded entries into `telemetry_log`, or they would be counted twice. +/// +/// Never the telemetry crate's own records (`livekit_telemetry*`): a record about a failed +/// upload must not produce another upload. The pipeline applies the same guard again. +fn telemetry_copy(record: &Record) -> Option { + let severity = match record.level() { + Level::Error => livekit_telemetry::Severity::Error, + Level::Warn => livekit_telemetry::Severity::Warn, + _ => return None, + }; + let target = record.target(); + if !target.starts_with("livekit") || target.starts_with("livekit_telemetry") { + return None; + } + Some(livekit_telemetry::LogRecord { + severity, + source: livekit_telemetry::LogSource::Ffi, + body: mask_jwts(&record.args().to_string()), + logger: Some(target.to_owned()), + function: record.module_path().map(str::to_owned), + file: record.file().map(str::to_owned), + line: record.line(), + timestamp_ns: None, + span_id: None, + }) +} + +/// `body` with credentials replaced: the value after `Bearer`, or after a `token` key (`token=`, +/// `access_token=`, `"token": "…"`, quoted or not, any spacing), becomes ``, and every +/// compact JWT — three base64url segments whose header and payload decode to JSON objects, +/// whatever their formatting or the punctuation around them — becomes ``. A participant +/// token quoted in an error must not leave the device. +fn mask_jwts(body: &str) -> String { + let body = mask_credential_values(body); + let jwt_char = |c: char| c.is_ascii_alphanumeric() || matches!(c, '-' | '_' | '.'); + let mut out = String::with_capacity(body.len()); + let mut rest = body.as_str(); + while let Some(start) = rest.find(jwt_char) { + out.push_str(&rest[..start]); + let run = &rest[start..]; + let run = &run[..run.find(|c| !jwt_char(c)).unwrap_or(run.len())]; + mask_run(run, &mut out); + rest = &rest[start + run.len()..]; + } + out.push_str(rest); + out +} + +/// Push `run` (base64url characters and dots) with every `header.payload.signature` inside it — +/// the signature may be empty, dots or a word glued on with `-`/`_` may surround it — as ``. +fn mask_run(run: &str, out: &mut String) { + let segments: Vec<&str> = run.split('.').collect(); + let mut i = 0; + while i < segments.len() { + if i > 0 { + out.push('.'); + } + let token = segments.get(i + 2).and_then(|_| { + json_object(segments[i + 1]).then_some(())?; + header_start(segments[i]) + }); + let Some(at) = token else { + out.push_str(segments[i]); + i += 1; + continue; + }; + out.push_str(&segments[i][..at]); + out.push_str(""); + i += 3; + } +} + +/// Where a JWT header starts in `segment`: at its start, or after a `-`/`_` joining a word to it. +fn header_start(segment: &str) -> Option { + std::iter::once(0) + .chain(segment.match_indices(['-', '_']).map(|(i, _)| i + 1)) + .find(|&at| json_object(&segment[at..])) +} + +/// A base64url segment whose decoded bytes start, after whitespace, with `{`. +fn json_object(segment: &str) -> bool { + let sextet = |c: u8| match c { + b'A'..=b'Z' => c - b'A', + b'a'..=b'z' => c - b'a' + 26, + b'0'..=b'9' => c - b'0' + 52, + b'-' => 62, + _ => 63, // `_`: the caller only passes base64url characters + }; + let (mut acc, mut bits) = (0u32, 0); + for c in segment.bytes() { + acc = (acc << 6 | u32::from(sextet(c))) & 0xFFF; + bits += 6; + if bits >= 8 { + bits -= 8; + let byte = (acc >> bits) as u8; + if !byte.is_ascii_whitespace() { + return byte == b'{'; + } + } + } + false +} + +/// The value after a credential marker (matched case-insensitively) as ``: after +/// `bearer` and whitespace, or after `bearer`/`token`, an optional closing quote, `=` or `:`, with any +/// whitespace and an optional opening quote in between. A quoted value ends at its closing quote, +/// any other at a delimiter. +fn mask_credential_values(body: &str) -> String { + // Same byte offsets as `body`: ASCII lowercasing never changes a length. + let lower = body.to_ascii_lowercase(); + let mut out = String::with_capacity(body.len()); + let (mut at, mut from) = (0, 0); + while let Some((key, len)) = ["bearer", "token"] + .iter() + .filter_map(|k| Some((lower[from..].find(k)? + from, k.len()))) + .min() + { + from = key + len; + let Some(value) = credential_value(&body[from..], len == "bearer".len()) else { continue }; + let (start, end) = (from + value.start, from + value.end); + out.push_str(&body[at..start]); + out.push_str(""); + at = end; + from = end; + } + out.push_str(&body[at..]); + out +} + +/// The non-empty value in `rest`, which follows a `bearer` (`bearer`) or `token` key, as a byte +/// range. Only a bearer value may follow its key after whitespace alone: `token expired` is prose. +fn credential_value(rest: &str, bearer: bool) -> Option> { + let trimmed = |s: &str| s.len() - s.trim_start().len(); + let skip = if bearer && rest.starts_with(char::is_whitespace) { + trimmed(rest) + } else { + let quote = usize::from(rest.starts_with(['"', '\''])); + let after = &rest[quote..]; + let gap = trimmed(after); + if !after[gap..].starts_with(['=', ':']) { + return None; + } + quote + gap + 1 + trimmed(&after[gap + 1..]) + }; + let value = &rest[skip..]; + let (start, len) = match value.chars().next() { + Some(q @ ('"' | '\'')) => (skip + 1, value[1..].find(q).unwrap_or(value.len() - 1)), + _ => ( + skip, + value + .find(|c: char| { + c.is_whitespace() || matches!(c, '&' | '"' | '\'' | ',' | ';' | ')' | '<' | '>') + }) + .unwrap_or(value.len()), + ), + }; + (len > 0).then_some(start..start + len) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn record<'a>(level: Level, target: &'a str, args: std::fmt::Arguments<'a>) -> Record<'a> { + Record::builder() + .level(level) + .target(target) + .args(args) + .file(Some("f.rs")) + .line(Some(7)) + .build() + } + + /// The core's own warnings and errors reach telemetry; nothing else does, and never the + /// telemetry crate's own (no feedback loop). + #[test] + fn only_the_cores_own_warnings_and_errors_are_copied() { + let copied = |level, target| telemetry_copy(&record(level, target, format_args!("boom"))); + let copy = copied(Level::Warn, "livekit::rtc_engine").expect("a core warning"); + assert_eq!(copy.severity, livekit_telemetry::Severity::Warn); + assert_eq!(copy.source, livekit_telemetry::LogSource::Ffi); + assert_eq!( + (copy.body.as_str(), copy.logger.as_deref()), + ("boom", Some("livekit::rtc_engine")) + ); + assert_eq!((copy.file.as_deref(), copy.line), (Some("f.rs"), Some(7))); + assert!(copied(Level::Error, "livekit_datatrack").is_some()); + assert!(copied(Level::Info, "livekit::room").is_none(), "info stays on the console"); + assert!(copied(Level::Error, "livekit_telemetry::exporter").is_none(), "no feedback loop"); + assert!(copied(Level::Error, "hyper::client").is_none(), "not the core's"); + assert!(copied(Level::Error, "my_app").is_none()); + } + + /// A token quoted in a core error never reaches telemetry; the console still gets the line. + #[test] + fn tokens_in_copied_records_are_masked() { + let jwt = "eyJhbGciOiJIUzI1NiJ9.eyJzdWIiOiJib2IifQ.c2lnbmF0dXJl-_x"; + let copy = telemetry_copy(&record( + Level::Error, + "livekit::signal_client", + format_args!("connect failed: wss://x/rtc?access_token={jwt}&auto=1 ({jwt})"), + )) + .expect("copied"); + assert_eq!(copy.body, "connect failed: wss://x/rtc?access_token=&auto=1 ()"); + // Codex final review B3: a valid token need not start `eyJ` — a header `{ "alg":…` and a + // payload with leading whitespace encode differently — and a transport error may quote + // the bearer value itself. + let spaced = "eyAiYWxnIjoiSFMyNTYiLCJ0eXAiOiJKV1QifQ\ + .CiAgeyJzdWIiOiJib2IiLCJ2aWRlbyI6eyJyb29tIjoiciJ9fQ.c2lnbmF0dXJlLWJ5dGVz"; + let copy = telemetry_copy(&record( + Level::Warn, + "livekit_signaling", + format_args!( + "signal transport error: HTTP 401 for request with Authorization: Bearer opaque-\ + credential; token {spaced}. retrying" + ), + )) + .expect("copied"); + assert_eq!( + copy.body, + "signal transport error: HTTP 401 for request with Authorization: Bearer ; \ + token . retrying" + ); + assert_eq!(mask_jwts(r#"{"token":"abc.def"}"#), r#"{"token":""}"#); + // Final review r2 B3: a JWT inside a punctuation-delimited run, and credential values + // that are quoted or spaced. + for (logged, masked) in [ + (format!("token (...{jwt})"), "token (...)".to_owned()), + (format!("token (...{spaced})"), "token (...)".to_owned()), + (format!("id=lk_{jwt}.x"), "id=lk_.x".to_owned()), + (r#"token="opaque-secret""#.into(), r#"token="""#.into()), + (r#"{"token": "opaque-secret"}"#.into(), r#"{"token": ""}"#.into()), + ( + r#"{"token" : 'opaque secret', "k": 1}"#.into(), + r#"{"token" : '', "k": 1}"#.into(), + ), + ( + "Authorization: Bearer opaque-secret".into(), + "Authorization: Bearer ".into(), + ), + ("ACCESS_TOKEN = opaque-secret&x".into(), "ACCESS_TOKEN = &x".into()), + ] { + assert_eq!(mask_jwts(&logged), masked, "{logged}"); + } + let unsigned = "eyAiYWxnIjoiSFMyNTYiLCJ0eXAiOiJKV1QifQ.CiAgeyJzdWIiOiJib2IifQ."; + assert_eq!(mask_jwts(unsigned), "", "an unsigned token is masked too"); + for kept in [ + "eyJ", + "eyJabc.", + "eyJabc..sig", + "see eyJ-not-a-token.", + "plain", + "wss://p.livekit.cloud/rtc", + "v1.2.3 of livekit_api.client.rs", + "token expired, tokens: 3, token_type=x", + "(...eyJ)", + ] { + assert_eq!(mask_jwts(kept), kept); + } + } + + /// Swift round-2 review: a platform forwarding only errors (`OSLogger(minLevel: .error)`) + /// still gets the core's warnings into telemetry, and its console still sees only errors. + #[test] + fn the_platforms_level_filters_the_console_not_the_telemetry_copy() { + log_forward_bootstrap(LevelFilter::Error); + assert_eq!(log::max_level(), LevelFilter::Warn, "warnings still reach the logger"); + let logger = Logger::new(); + logger.forward.store(LevelFilter::Error as usize, Ordering::Relaxed); + let warning = record(Level::Warn, "livekit::room", format_args!("careful")); + assert!(telemetry_copy(&warning).is_some(), "telemetry copies the warning"); + logger.log(&warning); + logger.log(&record(Level::Error, "livekit::room", format_args!("boom"))); + let mut rx = logger.rx.try_lock().expect("rx"); + assert_eq!(rx.try_recv().expect("forwarded").level, Level::Error, "the error is forwarded"); + assert!(rx.try_recv().is_err(), "the warning is not: the console is unchanged"); + drop(rx); + log_forward_bootstrap(LevelFilter::Debug); + assert_eq!(log::max_level(), LevelFilter::Debug, "a looser level is kept as is"); + } + + /// The console entry is unchanged by the copy. + #[test] + fn the_console_entry_is_what_it_always_was() { + let logger = Logger::new(); + logger.log(&record(Level::Warn, "livekit::room", format_args!("hello {}", 1))); + let entry = logger.rx.try_lock().expect("rx").try_recv().expect("forwarded"); + assert_eq!((entry.level, entry.target.as_str()), (Level::Warn, "livekit::room")); + assert_eq!( + (entry.message.as_str(), entry.file.as_deref(), entry.line), + ("hello 1", Some("f.rs"), Some(7)) + ); + } +} From a9ec8d5296efc24d43560717e991e91f4d02661e 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:42:26 +0200 Subject: [PATCH 3/3] docs(telemetry): add the telemetry_ping example and local testing guide `telemetry_ping` records a small session (an `lk.connect` span with its checkpoints, an `lk.publish`, an `lk.subscribe` ended by first media, a custom event) and exports it to the collector `LK_TELEMETRY_ENDPOINT` names; the README shows how to run it against a local OpenTelemetry backend, and a Dart test does the same through the pull queue. --- Cargo.lock | 10 ++ Cargo.toml | 1 + examples/telemetry_ping/Cargo.toml | 11 +++ examples/telemetry_ping/src/main.rs | 91 +++++++++++++++++++ livekit-telemetry/README.md | 18 ++++ .../dart/test/live_collector_test.dart | 88 ++++++++++++++++++ 6 files changed, 219 insertions(+) create mode 100644 examples/telemetry_ping/Cargo.toml create mode 100644 examples/telemetry_ping/src/main.rs create mode 100644 livekit-uniffi/support/dart/test/live_collector_test.dart diff --git a/Cargo.lock b/Cargo.lock index 10dacf191..6ed4b1665 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -6974,6 +6974,16 @@ version = "0.13.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "adb6935a6f5c20170eeceb1a3835a49e12e19d792f6dd344ccc76a985ca5a6ca" +[[package]] +name = "telemetry_ping" +version = "0.1.0" +dependencies = [ + "env_logger 0.11.11", + "livekit-net", + "livekit-telemetry", + "tokio", +] + [[package]] name = "tempfile" version = "3.27.0" diff --git a/Cargo.toml b/Cargo.toml index 265d7a29d..3058ba0bf 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -42,6 +42,7 @@ members = [ "examples/rpc", "examples/save_to_disk", "examples/screensharing", + "examples/telemetry_ping", "examples/token_source", "examples/webhooks", ] diff --git a/examples/telemetry_ping/Cargo.toml b/examples/telemetry_ping/Cargo.toml new file mode 100644 index 000000000..dee6f7709 --- /dev/null +++ b/examples/telemetry_ping/Cargo.toml @@ -0,0 +1,11 @@ +[package] +name = "telemetry_ping" +version = "0.1.0" +edition = "2021" +publish = false + +[dependencies] +livekit-telemetry = { workspace = true, features = ["net"] } +livekit-net = { workspace = true, features = ["native", "rustls-tls-native-roots"] } +tokio = { workspace = true, features = ["rt-multi-thread", "macros"] } +env_logger = { workspace = true } diff --git a/examples/telemetry_ping/src/main.rs b/examples/telemetry_ping/src/main.rs new file mode 100644 index 000000000..79cc00a13 --- /dev/null +++ b/examples/telemetry_ping/src/main.rs @@ -0,0 +1,91 @@ +// 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. + +//! Records a small session through the telemetry core — an `lk.connect` span with its +//! checkpoints, an `lk.publish`, an `lk.subscribe` ended by first media, and a custom event — +//! and exports it. +//! +//! - A local OpenTelemetry collector: `LK_TELEMETRY_ENDPOINT=http://localhost:4318` (the core's +//! test-only override; any server URL and token are accepted). +//! - A LiveKit Cloud project: `LK_URL=wss://.livekit.cloud` and `LK_TOKEN=`. + +use std::{env, sync::Arc}; + +use livekit_telemetry::{ + Attribute, NetTransport, RoomIdentity, RtcStatsSample, SpanName, SpanOutcome, SpanStep, + SpanTrack, StreamDirection, Telemetry, TelemetryConfig, TrackKind, TrackSource, + ENDPOINT_OVERRIDE_ENV, +}; + +#[tokio::main] +async fn main() { + env_logger::init(); + let url = env::var("LK_URL").unwrap_or_else(|_| "ws://localhost:7880".to_owned()); + let token = env::var("LK_TOKEN").unwrap_or_default(); + if env::var(ENDPOINT_OVERRIDE_ENV).is_err() && env::var("LK_URL").is_err() { + eprintln!("set {ENDPOINT_OVERRIDE_ENV}=http://localhost:4318, or LK_URL and LK_TOKEN"); + return; + } + + let config = TelemetryConfig { + resource: vec![ + Attribute::new("service.name", "telemetry_ping"), + Attribute::new("os.name", env::consts::OS), + ], + // Optional on-disk cache: run once with the collector down, once with it up. + storage_dir: env::var("LK_TELEMETRY_DIR").ok(), + ..Default::default() + }; + let transport = NetTransport::from_registry().expect("livekit-net has no HTTP client"); + let (telemetry, exporter) = Telemetry::new(config, Arc::new(transport)); + tokio::spawn(exporter.run()); + + let room = telemetry.begin_scope(); + room.set_server(&url, &token); + room.set_room(RoomIdentity { name: Some("telemetry-ping".into()), ..Default::default() }); + + let connect = room.start(SpanName::Connect, None); + for step in [SpanStep::WsOpen, SpanStep::Signal, SpanStep::JoinRecv, SpanStep::PcCreated] { + connect.step(step); + } + connect.end(SpanOutcome::Ok, None); + + let microphone = SpanTrack { + sid: Some("TR_ping_mic".into()), + kind: TrackKind::Audio, + source: TrackSource::Microphone, + remote_identity: None, + }; + let publish = room.start(SpanName::Publish, None); + publish.set_track(microphone); + publish.end(SpanOutcome::Ok, None); + + let remote = SpanTrack { + sid: Some("TR_ping_remote".into()), + kind: TrackKind::Video, + source: TrackSource::Camera, + remote_identity: Some("bob".into()), + }; + room.subscribe_started(remote.clone()); + room.subscribed(remote); + let mut media = + RtcStatsSample::new("TR_ping_remote", TrackKind::Video, StreamDirection::Inbound); + media.bytes = Some(1_500); + room.record_stats(media); // first media: the lk.subscribe span ends ok + + room.emit_custom("ping", vec![Attribute::new("seq", 1i64)]); + telemetry.shutdown().await; + println!("trace {} — {}", room.trace_id(), telemetry.stats()); +} diff --git a/livekit-telemetry/README.md b/livekit-telemetry/README.md index 1809183d9..00f350d17 100644 --- a/livekit-telemetry/README.md +++ b/livekit-telemetry/README.md @@ -59,6 +59,24 @@ println!("{}", telemetry.stats()); // drops by reason, uploads, cached batches Event names, attributes, cadences and the upload policy are defined in [`SPEC.md`](SPEC.md). +## Local testing + +Set `LK_TELEMETRY_ENDPOINT` in the process environment to send every upload to an OpenTelemetry +backend of your own instead of LiveKit Cloud — the test-only override, not part of any platform +API. A base URL gets the standard OTLP paths (`/v1/logs`, `/v1/traces`). + +```sh +docker run -d --name lk-lgtm -p 3000:3000 -p 4318:4318 grafana/otel-lgtm # OTLP/HTTP :4318, UI :3000 +LK_TELEMETRY_ENDPOINT=http://localhost:4318 cargo run -p telemetry_ping # prints the trace id +``` + +`telemetry_ping` records a small session (an `lk.connect` span with its checkpoints, an +`lk.publish`, an `lk.subscribe` ended by first media, a custom event). In the UI +(http://localhost:3000 → Explore): Loki `{service_name="telemetry_ping"}` for the log records, +Tempo with the printed trace id for the spans. Add `LK_TELEMETRY_DIR=/tmp/lk-telemetry` to use +the file cache: stop the container, run once, start it again and run once more to watch the +cached batches replay. Platform end-to-end tests set the same variable before the SDK starts. + ## Design notes - **The room decides the destination.** A platform passes only the server URL and the token diff --git a/livekit-uniffi/support/dart/test/live_collector_test.dart b/livekit-uniffi/support/dart/test/live_collector_test.dart new file mode 100644 index 000000000..3ba90165b --- /dev/null +++ b/livekit-uniffi/support/dart/test/live_collector_test.dart @@ -0,0 +1,88 @@ +import 'dart:io'; +import 'dart:typed_data'; + +import 'package:http/http.dart' as http; +import 'package:livekit_uniffi/livekit_telemetry.dart'; +import 'package:livekit_uniffi/livekit_uniffi.dart'; +import 'package:test/test.dart'; + +/// Posts one custom event to a real collector through the pull queue. +/// +/// Skipped unless a collector is named: `LK_TELEMETRY_ENDPOINT` (a local OpenTelemetry collector, +/// e.g. `http://localhost:4318` — read by the Rust core itself, the test-only override), or +/// `LK_TELEMETRY_URL` + `LK_TELEMETRY_TOKEN` (a LiveKit Cloud project and a token carrying the +/// observability grant). Unlike `telemetry_test.dart`, which fakes the collector's answer, this +/// one performs the request so the Dart binding is exercised against a real ingest. +/// +/// Dart drives the transport rather than implementing `TelemetryTransport`: uniffi-dart's +/// foreign-trait callbacks are isolate-bound, so Rust calling into Dart from a tokio worker aborts +/// the VM. The queue inverts that — Dart awaits `next()` and reports the answer with `complete()`. +Future serve(TelemetryExportQueue queue, List statuses, int count) async { + final client = http.Client(); + try { + for (var i = 0; i < count; i++) { + final pending = await queue.next(); + if (pending == null) return; + final request = pending.request; + try { + final response = await client.post( + Uri.parse(request.url), + headers: request.headers, + body: request.body, + ); + statuses.add(response.statusCode); + queue.complete( + id: pending.id, + response: ExportResponse( + status: response.statusCode, + headers: {}, + body: Uint8List.fromList(response.bodyBytes), + ), + ); + } catch (error) { + statuses.add(-1); + queue.fail( + id: pending.id, + error: RetryableExportException(reason: '$error', retryAfterMs: null), + ); + } + } + } finally { + client.close(); + } +} + +void main() { + final endpoint = Platform.environment['LK_TELEMETRY_ENDPOINT']; + final url = Platform.environment['LK_TELEMETRY_URL']; + final token = Platform.environment['LK_TELEMETRY_TOKEN']; + final target = endpoint ?? url; + + test('posts a custom event to a live collector through the pull queue', () async { + final statuses = []; + final queue = telemetryConfigurePulled( + config: TelemetryConfig( + resource: [ + Attribute(key: 'service.name', value: StrAttributeValue('livekit-client-dart')), + Attribute(key: 'service.version', value: StrAttributeValue('0.0.0-local')), + Attribute(key: 'os.name', value: StrAttributeValue(Platform.operatingSystem)), + ], + logSeverity: Severity.warn, + ), + instruments: [], + ); + final serving = serve(queue, statuses, 2); + + final room = telemetryScope()!; + room.setServer(url: url ?? 'ws://localhost:7880', token: token ?? 'local'); + room.emitCustom(name: 'ping', attributes: {'seq': '1'}); + await telemetryFlush(); + await telemetryShutdown(); + await serving; + queue.finish(); + + print('dart → $target: statuses $statuses'); + expect(statuses, isNotEmpty); + expect(statuses.first, inInclusiveRange(200, 299)); + }, skip: target == null ? 'set LK_TELEMETRY_ENDPOINT, or LK_TELEMETRY_URL and LK_TELEMETRY_TOKEN' : null); +}