diff --git a/.changeset/add-livekit-telemetry-crate.md b/.changeset/add-livekit-telemetry-crate.md new file mode 100644 index 000000000..40a48016e --- /dev/null +++ b/.changeset/add-livekit-telemetry-crate.md @@ -0,0 +1,5 @@ +--- +livekit-telemetry: minor +--- + +Add `livekit-telemetry`, the internal client telemetry core shared by the LiveKit client SDKs. diff --git a/Cargo.lock b/Cargo.lock index afcc5a91c..4c73ce507 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3812,6 +3812,17 @@ dependencies = [ "url", ] +[[package]] +name = "livekit-telemetry" +version = "0.1.0" +dependencies = [ + "log", + "opentelemetry-proto", + "prost 0.14.4", + "rand 0.9.5", + "uniffi", +] + [[package]] name = "livekit-token" version = "0.2.1" @@ -5037,6 +5048,46 @@ dependencies = [ "vcpkg", ] +[[package]] +name = "opentelemetry" +version = "0.32.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b0142c63252a9e054e68a4c61a5778f7b14f576274d593f8ce883d191a099682" +dependencies = [ + "futures-core", + "futures-sink", + "js-sys", + "pin-project-lite", + "thiserror 2.0.19", +] + +[[package]] +name = "opentelemetry-proto" +version = "0.32.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "56d658ba1faf63f7b9c492cfbe6e0ec365440a16132d3270c1065f7b33f1b638" +dependencies = [ + "opentelemetry", + "opentelemetry_sdk", + "prost 0.14.4", +] + +[[package]] +name = "opentelemetry_sdk" +version = "0.32.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9b59f80e1ac4d5ff7a2db8fb6c80badb7f0f3f858211fba08dd9aaec750894f9" +dependencies = [ + "futures-channel", + "futures-executor", + "futures-util", + "opentelemetry", + "percent-encoding", + "portable-atomic", + "rand 0.9.5", + "thiserror 2.0.19", +] + [[package]] name = "orbclient" version = "0.3.55" diff --git a/Cargo.toml b/Cargo.toml index 2bb66801a..265d7a29d 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -11,6 +11,7 @@ members = [ "livekit-ffi", "livekit-uniffi", "livekit-datatrack", + "livekit-telemetry", "livekit-token", "livekit-token-source", "livekit-ffi-node-bindings", @@ -60,6 +61,7 @@ livekit-api = { version = "0.8.1", path = "livekit-api" } livekit-capture = { version = "0.1.3", path = "livekit-capture" } livekit-ffi = { version = "0.12.81", path = "livekit-ffi" } livekit-datatrack = { version = "0.2.1", path = "livekit-datatrack" } +livekit-telemetry = { version = "0.1.0", path = "livekit-telemetry" } livekit-signaling = { version = "0.1.4", path = "livekit-signaling" } livekit-token = { version = "0.2.1", path = "livekit-token" } livekit-token-source = { version = "0.1.3", path = "livekit-token-source" } diff --git a/knope.toml b/knope.toml index bf337b338..039ed4fed 100644 --- a/knope.toml +++ b/knope.toml @@ -178,6 +178,14 @@ versioned_files = [ ] changelog = "livekit-datatrack/CHANGELOG.md" +[packages.livekit-telemetry] +versioned_files = [ + "livekit-telemetry/Cargo.toml", + "Cargo.lock", + { path = "Cargo.toml", dependency = "livekit-telemetry" }, +] +changelog = "livekit-telemetry/CHANGELOG.md" + [packages.livekit-common] versioned_files = [ "livekit-common/Cargo.toml", diff --git a/livekit-telemetry/CHANGELOG.md b/livekit-telemetry/CHANGELOG.md new file mode 100644 index 000000000..825c32f0d --- /dev/null +++ b/livekit-telemetry/CHANGELOG.md @@ -0,0 +1 @@ +# Changelog diff --git a/livekit-telemetry/Cargo.toml b/livekit-telemetry/Cargo.toml new file mode 100644 index 000000000..7364e25b3 --- /dev/null +++ b/livekit-telemetry/Cargo.toml @@ -0,0 +1,25 @@ +[package] +name = "livekit-telemetry" +description = "Client telemetry core for LiveKit: buffers events on-device and exports them as OTLP" +version = "0.1.0" +readme = "README.md" +license.workspace = true +edition.workspace = true +repository.workspace = true + +[dependencies] +log = { workspace = true } +prost = { workspace = true } +rand = { workspace = true } +# 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"] } +uniffi = { workspace = true, features = ["scaffolding-ffi-buffer-fns"], optional = true } + +[features] +uniffi = ["dep:uniffi"] + +# How CI checks this crate's features, read by +# `.github/workflows/feature-combinations-curated.yml` via `cargo metadata`. +[package.metadata.feature-combinations] +mode = "powerset" diff --git a/livekit-telemetry/README.md b/livekit-telemetry/README.md new file mode 100644 index 000000000..83ca1fa3a --- /dev/null +++ b/livekit-telemetry/README.md @@ -0,0 +1,19 @@ +# LiveKit Telemetry + +**Important**: +This is an internal crate for client telemetry in LiveKit client SDKs and is not meant to be +used directly. + +It currently holds the shared data model and its encoding: + +- **Records.** `TelemetryEvent`, `LogRecord`, `Attribute` / `AttributeValue`, `Severity` and + `SpanOutcome`: the events, log lines and span outcomes a platform SDK hands over. +- **Buffers.** Bounded, drop-oldest queues for records and spans, each filed under the session + (trace) it belongs to; every drop is counted. +- **Encoding.** [OTLP/HTTP](https://opentelemetry.io/docs/specs/otlp/) protobuf logs and traces, + using the `opentelemetry-proto` message types only (no gRPC stack). + +The pipeline that batches, caches and uploads records is not part of the crate yet. + +Event names, attributes and span rules are defined in [`SPEC.md`](SPEC.md). It is the contract +for the whole crate, including the parts not implemented yet. diff --git a/livekit-telemetry/SPEC.md b/livekit-telemetry/SPEC.md new file mode 100644 index 000000000..65679921c --- /dev/null +++ b/livekit-telemetry/SPEC.md @@ -0,0 +1,259 @@ +# Client telemetry spec + +Source of truth for event names, attributes and cadences emitted by LiveKit client SDKs. +Additive-only by convention; LiveKit-defined names carry the `lk.` prefix, everything else +follows [OpenTelemetry semantic conventions](https://github.com/open-telemetry/semantic-conventions). + +Every record carries a wall-clock timestamp: one without `timestamp_ns` is stamped when it is +captured (queued), never at export; an explicit `timestamp_ns` is kept. Attribute keys the SDK +owns are `lk.*`, `session.id` and `error.type`: an app's attributes and custom events cannot set +them, and on a span only the core writes `lk.outcome` and `error.type` (see [Spans](#spans)). + +## Resource attributes + +Set once per pipeline (`TelemetryConfig.resource`): + +| Key | Who sets it | Example | +|---|---|---| +| `service.name` | platform SDK | `livekit-client-swift` | +| `service.version` | platform SDK | `2.9.0` | +| `os.name`, `os.version` | platform SDK | `iOS`, `18.5` | +| `device.model.identifier` | platform SDK | `iPhone16,1` | +| `telemetry.sdk.name/language/version` | core | `livekit-telemetry`, `rust`, `0.1.0` | + +## Events + +An event with no `body` is exported with its name as the body as well as in `event_name`: log +viewers key their line on the body, and not every backend surfaces `event_name` yet. + +```yaml +event: lk.ping +area: sdk +severity: info +attributes: + lk.ping.seq: int # optional, monotonically increasing per pipeline +cadence: on demand — pipeline smoke test, never emitted in production paths +platforms: all +``` + +```yaml +event: lk.telemetry.report +area: sdk (self-telemetry) +severity: info +attributes: # counts since the previous report + lk.telemetry.uploads.sent: int # batches accepted + lk.telemetry.uploads.bytes: int # compressed bytes accepted — what telemetry cost the uplink + lk.telemetry.uploads.failed: int # attempts that failed transiently (no answer, 429, 5xx) + lk.telemetry.cache.batches: int # batches waiting in the cache right now (a gauge) + # the rest only when non-zero: + lk.telemetry.uploads.timeouts: int # attempts that hit export_timeout_ms + lk.telemetry.uploads.unauthorized: int # 401/403 answers: tokens the collector refused + lk.telemetry.holds.capped: int # soft holds that reached the 60 s cap + lk.telemetry.cache.write_errors: int # batches the disk refused (full, gone), kept in memory + lk.telemetry.dropped.queue_full: int # records evicted from the in-memory queue + lk.telemetry.dropped.cache_error: int # records no cache, not even memory, could take + lk.telemetry.dropped.cache_full: int # records evicted by the cache's size / file-count bound + lk.telemetry.dropped.expired: int # records in batches past the 24 h age limit + lk.telemetry.dropped.corrupt: int # records in cached batches that failed their CRC + lk.telemetry.dropped.invalid: int # custom events / attributes over the limits (rejected) + lk.telemetry.dropped.rejected: int # records the collector rejected (final 4xx/5xx, partial success) + lk.telemetry.dropped.oversized: int # single records larger than the collector accepts (413) + lk.telemetry.dropped.throttled: int # records evicted from the cache during a server-directed pause + lk.telemetry.dropped.rate_limited: int # discrete events dropped by the flood guard +cadence: appended to the next upload whenever a loss, a failure, a refusal, a capped hold or a + disk error happened since the previous report — never its own request, never persisted + on its own, so a broken uploader never reports through itself — and once at shutdown as + the session summary, so fleet-wide success rates have denominators. Losses by policy + (a project that receives nothing, the opt-out) are local only (`Telemetry::stats`). +platforms: all +``` + +```yaml +event: lk.device.thermal.changed +area: device +attributes: + lk.device.thermal.state: enum(nominal | fair | serious | critical) +cadence: on change (+ initial value on the first `set_device_state`); `unknown` (no thermal + source on the platform, the default) is no reading: no event, no stretch +platforms: ios, macos, android — optional elsewhere +``` + +```yaml +event: lk.device.low_power.changed +area: device +attributes: + lk.device.low_power.enabled: bool +cadence: on change (+ initial value); `None` (no source on the platform, the default) is no + reading: no event, no stretch +platforms: ios, macos, android — optional elsewhere +``` + +```yaml +event: lk.device.app_state.changed +area: device +attributes: + lk.device.app_state: enum(foreground | background) +cadence: on change (+ initial value); entering background also forces a flush +platforms: all +``` + +```yaml +event: lk.device.memory.changed +area: device +attributes: + lk.device.memory.pressure: enum(normal | warning | critical) + # Apple: DispatchSource memory-pressure levels; Android onTrimMemory: RUNNING_LOW / + # BACKGROUND → warning, RUNNING_CRITICAL / COMPLETE → critical +cadence: on change (+ initial value) +platforms: ios, macos, android — optional elsewhere +``` + +```yaml +event: lk.device.network.changed +area: device +attributes: + network.connection.type: enum(wifi | cell | wired | vpn | bluetooth | other | unavailable | unknown) # OTel semconv + lk.device.network.expensive: bool # cellular / hotspot (NWPath.isExpensive, metered) + lk.device.network.constrained: bool # Low Data Mode / Data Saver / navigator.connection.saveData +cadence: on change of any attribute (+ initial value) +platforms: ios, macos, android — web: Chromium only +``` + +```yaml +event: lk.device.battery.changed +area: device +attributes: + hw.battery.charge: double # 0.0–1.0 (OTel hardware semconv) + hw.battery.state: enum(charging | discharging) # OTel hardware semconv +cadence: on charging change and when the level crosses 20 % or 10 % unplugged — never per + percent; silent where the level is unknown (desktops, tvOS) +platforms: ios, android — optional elsewhere +``` + +```yaml +event: lk.device.audio_route.changed +area: device +attributes: + lk.device.audio_route.reason: enum(new_device | old_device_unavailable | category_change | override | wake_from_sleep | no_suitable_route | route_configuration_change | unknown) # AVAudioSession names; `unknown` where the platform gives none + lk.device.audio_route.outputs: string # comma-separated enum(speaker | receiver | wired_headset | bluetooth | car_audio | air_play | hdmi | usb | other) +cadence: on change +platforms: ios — android: audio device callbacks; optional elsewhere +``` + +```yaml +event: lk.device.audio.interruption +area: device +attributes: + lk.device.audio.interruption: enum(began | ended) +cadence: on change +platforms: ios — android: audio focus loss/gain; optional elsewhere +``` + +```yaml +event: lk.device.capture.failed +area: device +severity: warn +attributes: + lk.device.capture.device: enum(camera | microphone | screen_share) + lk.device.capture.reason: enum(permission_denied | not_found | in_use | disconnected | other) # the getUserMedia failure taxonomy +cadence: on failure +platforms: all — ios: authorization status, capture interruptions; android: permission checks, camera callbacks; web: DOMException names +``` + +## Log records + +A `TelemetryEvent` with an empty `name` is a plain log record (OTLP log without `event_name`): +`severity` + `body` (the message) + `code.function.name`, `code.file.path`, `code.line.number` +(semconv), `lk.log.source` (`sdk` | `ffi` | `webrtc`) and `lk.log.logger` (type, module or file). The +platform hands the core a typed `LogRecord` via `log(record)`; the core applies the floor: WebRTC only +at `error`, the SDK and the core at the configured `log_severity`, the core's own telemetry module +never. Only `warn` and `error` records +leave the device; `trace`/`debug`/`info` are dropped in `emit`. + +The core's own warnings and errors (Rust `log` records from `livekit*` targets, never +`livekit_telemetry*`) reach the pipeline by themselves through `livekit-uniffi`'s log forwarder: +it copies them as `ffi` records — every JWT-shaped substring (`eyJ…` `.` … `.` …) masked as +``, so a token quoted in an error never leaves the device — and forwards the console entry +unchanged. The copy is active wherever the platform calls `log_forward_bootstrap` (Swift's +default `OSLogger` with `ffi: true` does), whatever level it passes: that level filters only what +is forwarded to the platform's console (the global `log` level is kept at `warn` or looser) — Rust allows one logger per process, and without the +forwarder the core's records go nowhere. Platforms must not feed forwarded Rust log entries +(what `log_forward_receive` returns) to `telemetry_log` / `log(record)`: the core already copied +them, so they would be counted twice. `telemetry_log` is for the platform's own and WebRTC's +lines. + +## Spans + +A span is **one attempt** at an operation. The scope (one Room connection lifetime, across +reconnects) is the trace; its id is generated by the core when the pipeline starts and rides on +every span and log record. Spans are exported when they end — never a long-lived scope span. + +| Rule | Value | +|---|---| +| Names | `lk.connect`, `lk.reconnect`, `lk.publish`, `lk.subscribe` — verbs, never ids | +| Kind | `CLIENT` for connect/reconnect (a call to the SFU), `INTERNAL` otherwise | +| Status | OTel `Unset` on success **and** cancellation, `Error` (+ `error.type`, message) on failure | +| `lk.outcome` | exactly once on every span: `ok` \| `error` \| `cancelled` — rollups read this, never the status. The core writes it and `error.type` from how the span ended, dropping span attributes of either name; `error.type` appears at most once, only when the span ended with an error type | +| `error.type` | platform-defined, a type name (≤ 128 bytes), never a message: e.g. Swift sends `LiveKitError.`, `CancellationError` or the Swift error type; dashboards group by it per `service.name` | +| Checkpoints | span events in the span's envelope (`ws_open`, `join_recv`, `pc_connected`, `attempt 2 full`, …); real events stay log records pointing at the span via `span_id` | +| Limits | 128 events and 128 attributes per span (OTel defaults); 256 open spans per pipeline | + +```yaml +span: lk.connect +kind: client +attributes: + lk.connect.attempt: int # 1 for the user-initiated connect +checkpoints: + required: ws_open, signal, join_recv, pc_created # every platform, in this order + best-effort: engine, pc_connected, offer_sent, answer_sent, room_connected # where the SDK has the moment +outcome: ok | error (error.type) | cancelled +``` + +Dashboards compute connect phases only from the required checkpoints; best-effort ones refine +a platform's own view and may be missing or ordered differently (Android reports seven of the +nine today). + +```yaml +span: lk.reconnect +kind: client +attributes: + lk.reconnect.reason: enum(signal_disconnected | publisher_failed | subscriber_failed | transport_failed | switch_candidate | network_changed | debug | unknown) + lk.reconnect.mode: enum(quick | full) # mode of the last attempt + lk.reconnect.attempts: int +checkpoints: "attempt " per attempt +outcome: ok | error | cancelled # cancelled when disconnect() or a newer reconnect wins +``` + +```yaml +span: lk.publish +kind: internal +parent: the ambient span, when any; a pre-connect publish (a microphone published before the + connect completes) is its own span in the session's trace — `lk.connect` has ended by then +attributes: + lk.track.kind: enum(audio | video) + lk.track.source: enum(camera | microphone | screen_share | screen_share_audio | unknown) + lk.track.sid: string # on success +outcome: ok | error (error.type) | cancelled +``` + +```yaml +span: lk.subscribe +kind: internal +starts: when the intent to subscribe exists — a remote publish under autoSubscribe, or the + manual subscribe call. Tracks already in the room at join: platforms call + `subscribe_started` for each at connect (the join response lists them), so `lk.subscribe` + measures from join on every platform; a `subscribed` with no intent before it still opens + the span (a fallback, measured from the confirmation) +ends: at first media (the first inbound stats reading with bytes; the core sees it) → ok; + unsubscribe / unpublish before media → cancelled; + subscription failure → error; no media within 30 s → error (error.type = timed_out), + enforced on the core's own clock (no reading needed) and kept through a later disconnect +owner: the core (`Scope::subscribe_started / subscribed / track_ended / subscribe_failed`); + an SDK only reports the remote track's lifecycle +attributes: + lk.track.sid: string + lk.track.kind: enum(audio | video) + lk.track.source: enum(camera | microphone | screen_share | screen_share_audio | unknown) + lk.participant.remote_identity: string +checkpoints: subscribed, first_media +``` diff --git a/livekit-telemetry/src/event.rs b/livekit-telemetry/src/event.rs new file mode 100644 index 000000000..16bb47209 --- /dev/null +++ b/livekit-telemetry/src/event.rs @@ -0,0 +1,290 @@ +// 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::time::{SystemTime, UNIX_EPOCH}; + +/// A discrete telemetry event. +/// +/// Exported as one OTLP log record whose `event_name` is [`name`](Self::name), following the +/// OTel logs data model (events are log records with a top-level event name). +#[cfg_attr(feature = "uniffi", derive(uniffi::Record))] +#[derive(Debug, Clone, PartialEq)] +pub struct TelemetryEvent { + /// Event name. LiveKit-defined events use the `lk.` prefix (e.g. `lk.ping`); see `SPEC.md`. + pub name: String, + pub severity: Severity, + /// Optional human-readable message (the OTLP log record body). + #[cfg_attr(feature = "uniffi", uniffi(default))] + pub body: Option, + pub attributes: Vec, + /// Wall-clock time in nanoseconds since the Unix epoch. `None` is stamped when the record is + /// queued, never at export. + #[cfg_attr(feature = "uniffi", uniffi(default))] + pub timestamp_ns: Option, + /// The in-flight span this record belongs to (a handle from `begin_span`), if any. The trace + /// id is always the session's and is attached by the core. + #[cfg_attr(feature = "uniffi", uniffi(default))] + pub span_id: Option, +} + +impl TelemetryEvent { + /// An `Info` event without attributes, stamped when emitted. + pub fn new(name: impl Into) -> Self { + Self { + name: name.into(), + severity: Severity::Info, + body: None, + attributes: Vec::new(), + timestamp_ns: None, + span_id: None, + } + } + + /// Link this record to an in-flight span. + pub fn in_span(mut self, span: u64) -> Self { + self.span_id = Some(span); + self + } + + /// The same event with `severity`. + pub fn with_severity(mut self, severity: Severity) -> Self { + self.severity = severity; + self + } + + /// The same event with a display body. + pub fn with_body(mut self, body: impl Into) -> Self { + self.body = Some(body.into()); + self + } + + /// The same event with one more attribute. + pub fn with_attribute( + mut self, + key: impl Into, + value: impl Into, + ) -> Self { + self.attributes.push(Attribute::new(key, value)); + self + } + + /// A consumer's own event. Always namespaced under `custom.` so it can never be mistaken for + /// a LiveKit-defined `lk.*` event, and the backend can filter or quota it separately; + /// attributes keep the caller's namespace (`acme.checkout.step`). + pub fn custom(name: &str, attributes: Vec) -> Self { + let name = format!("custom.{}", name.trim_start_matches("custom.")); + Self { attributes, body: Some(name.clone()), ..Self::new(name) } + } + + /// Rough encoded size in bytes — strings plus a fixed overhead per field. Drives the byte + /// bounds on queue flushing and request size; cheaper than encoding and close enough for both. + pub fn size_hint(&self) -> usize { + 32 + self.name.len() + + self.body.as_ref().map_or(0, String::len) + + self.attributes.iter().map(|a| 4 + a.key.len() + a.value.size_hint()).sum::() + } +} + +impl Default for TelemetryEvent { + /// An unnamed `Info` record: a plain log record once it has a body (see `SPEC.md`). + fn default() -> Self { + Self::new("") + } +} + +/// Event severity, mapped onto the OTel severity numbers (`TRACE`=1 … `ERROR`=17). +#[cfg_attr(feature = "uniffi", derive(uniffi::Enum))] +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] +pub enum Severity { + Trace, + Debug, + Info, + Warn, + Error, +} + +/// Where a log line came from. WebRTC is chatty at warn, so only its errors become records; +/// the SDK and the core use the configured floor. +#[cfg_attr(feature = "uniffi", derive(uniffi::Enum))] +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum LogSource { + Sdk, + Ffi, + WebRtc, +} + +impl LogSource { + fn as_str(self) -> &'static str { + match self { + Self::Sdk => "sdk", + Self::Ffi => "ffi", + Self::WebRtc => "webrtc", + } + } +} + +/// A log line as the platform captured it, where it happened. The core turns it into a record: +/// semconv `code.*` attributes, `lk.log.source`, `lk.log.logger`, filed under the span's session. +/// Stamp `timestamp_ns` at capture; the record may cross an executor hop before it gets here. +#[cfg_attr(feature = "uniffi", derive(uniffi::Record))] +#[derive(Debug, Clone, PartialEq)] +pub struct LogRecord { + pub severity: Severity, + pub source: LogSource, + /// The line itself (the OTLP log body). Not `message`: AGENTS.md keeps that name off every + /// exported record. + pub body: String, + /// The logger: a type, module or file name (`Room`, `livekit::rtc_engine`, `sctp.cc`). + #[cfg_attr(feature = "uniffi", uniffi(default))] + pub logger: Option, + #[cfg_attr(feature = "uniffi", uniffi(default))] + pub function: Option, + #[cfg_attr(feature = "uniffi", uniffi(default))] + pub file: Option, + #[cfg_attr(feature = "uniffi", uniffi(default))] + pub line: Option, + #[cfg_attr(feature = "uniffi", uniffi(default))] + pub timestamp_ns: Option, + /// The in-flight span this line was logged under, if any. + #[cfg_attr(feature = "uniffi", uniffi(default))] + pub span_id: Option, +} + +impl From for TelemetryEvent { + fn from(record: LogRecord) -> Self { + let mut event = TelemetryEvent::default() + .with_severity(record.severity) + .with_body(record.body) + .with_attribute("lk.log.source", record.source.as_str()); + if let Some(logger) = record.logger.filter(|s| !s.is_empty()) { + event = event.with_attribute("lk.log.logger", logger); + } + if let Some(function) = record.function.filter(|s| !s.is_empty()) { + event = event.with_attribute("code.function.name", function); + } + if let Some(file) = record.file.filter(|s| !s.is_empty()) { + event = event.with_attribute("code.file.path", file); + } + if let Some(line) = record.line.filter(|l| *l > 0) { + event = event.with_attribute("code.line.number", line as i64); + } + event.timestamp_ns = record.timestamp_ns; + event.span_id = record.span_id; + event + } +} + +/// A key/value attribute on an event or on the resource. +#[cfg_attr(feature = "uniffi", derive(uniffi::Record))] +#[derive(Debug, Clone, PartialEq)] +pub struct Attribute { + pub key: String, + pub value: AttributeValue, +} + +impl Attribute { + /// A key/value pair. + pub fn new(key: impl Into, value: impl Into) -> Self { + Self { key: key.into(), value: value.into() } + } +} + +/// Attribute value: the scalar subset of OTLP `AnyValue`. +#[cfg_attr(feature = "uniffi", derive(uniffi::Enum))] +#[derive(Debug, Clone, PartialEq)] +pub enum AttributeValue { + Str(String), + Int(i64), + Double(f64), + Bool(bool), +} + +impl From<&str> for AttributeValue { + fn from(value: &str) -> Self { + Self::Str(value.to_owned()) + } +} + +impl From for AttributeValue { + fn from(value: String) -> Self { + Self::Str(value) + } +} + +impl From for AttributeValue { + fn from(value: i64) -> Self { + Self::Int(value) + } +} + +impl From for AttributeValue { + fn from(value: f64) -> Self { + Self::Double(value) + } +} + +impl From for AttributeValue { + fn from(value: bool) -> Self { + Self::Bool(value) + } +} + +/// Current wall-clock time in nanoseconds since the Unix epoch (0 if the clock is before 1970). +pub(crate) fn now_unix_nanos() -> u64 { + let now = + SystemTime::now().duration_since(UNIX_EPOCH).map(|d| d.as_nanos() as u64).unwrap_or(0); + #[cfg(test)] + let now = now.saturating_add_signed(CLOCK_JUMP_NS.with(|jump| jump.get())); + now +} + +#[cfg(test)] +thread_local! { + /// How far a test moved the wall clock, on this (current-thread runtime) thread. + pub(crate) static CLOCK_JUMP_NS: std::cell::Cell = const { std::cell::Cell::new(0) }; +} + +/// Limits on what an app hands over: long enough for any real identifier, short enough that one +/// app cannot bloat every record. Over-long input is rejected and counted, never truncated — a +/// truncated id silently collides with another. +pub(crate) const MAX_NAME_BYTES: usize = 128; +pub(crate) const MAX_KEY_BYTES: usize = 128; +pub(crate) const MAX_VALUE_BYTES: usize = 1024; +/// Custom attributes per room, and per custom event. +pub(crate) const MAX_CUSTOM_ATTRIBUTES: usize = 64; + +/// Keys the SDK owns: an app can neither set nor override them (`lk.*` — room, participant, +/// track, outcome —, `session.id` and `error.type`). +pub(crate) fn reserved(key: &str) -> bool { + key.starts_with("lk.") || key == "session.id" || key == "error.type" +} + +/// Whether an app-provided attribute is within the limits and outside the SDK's namespace. +pub(crate) fn valid_custom(key: &str, value: Option<&AttributeValue>) -> bool { + let value_ok = match value { + Some(AttributeValue::Str(s)) => s.len() <= MAX_VALUE_BYTES, + _ => true, + }; + !key.is_empty() && key.len() <= MAX_KEY_BYTES && !reserved(key) && value_ok +} + +impl AttributeValue { + /// Rough encoded size of the value. + pub(crate) fn size_hint(&self) -> usize { + match self { + AttributeValue::Str(s) => s.len(), + _ => 8, + } + } +} diff --git a/livekit-telemetry/src/lib.rs b/livekit-telemetry/src/lib.rs new file mode 100644 index 000000000..ac69c5014 --- /dev/null +++ b/livekit-telemetry/src/lib.rs @@ -0,0 +1,42 @@ +// 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. + +// Internals the pipeline (`Telemetry`, `Exporter`) consumes once it is in place. +#![allow(dead_code)] + +/// Event data model: what SDKs push in. +mod event; + +/// Bounded in-memory queue between `emit` and the exporter. +mod store; + +/// Pipeline health counters and the `lk.telemetry.report` event. +mod stats; + +/// Spans: one attempt at an operation, with explicit handles across the FFI. +mod scope; +mod span; + +/// OTLP/HTTP protobuf encoding of a batch. +mod otlp; + +/// OTLP protobuf types (re-exported from `opentelemetry-proto`). +mod proto; + +pub use event::*; +pub use span::SpanOutcome; +pub use stats::{TelemetryStats, TelemetryStatus}; + +#[cfg(feature = "uniffi")] +uniffi::setup_scaffolding!(); diff --git a/livekit-telemetry/src/otlp.rs b/livekit-telemetry/src/otlp.rs new file mode 100644 index 000000000..1a32b52ef --- /dev/null +++ b/livekit-telemetry/src/otlp.rs @@ -0,0 +1,323 @@ +// 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 prost::Message; + +use crate::span::SpanKind; +use crate::{ + event::now_unix_nanos, + proto::opentelemetry::proto::{ + collector::{logs::v1::ExportLogsServiceRequest, trace::v1::ExportTraceServiceRequest}, + common::v1::{any_value, AnyValue, InstrumentationScope, KeyValue}, + logs::v1::{LogRecord, ResourceLogs, ScopeLogs, SeverityNumber}, + resource::v1::Resource, + trace::v1::{span, status, ResourceSpans, ScopeSpans, Span, Status}, + }, + span::SpanRecord, + store::Queued, + Attribute, AttributeValue, Severity, SpanOutcome, +}; + +pub(crate) const CONTENT_TYPE: &str = "application/x-protobuf"; + +fn resource(attributes: &[Attribute]) -> Option { + Some(Resource { + attributes: attributes.iter().map(KeyValue::from).collect(), + ..Default::default() + }) +} + +fn scope() -> Option { + Some(InstrumentationScope { + name: env!("CARGO_PKG_NAME").to_owned(), + version: env!("CARGO_PKG_VERSION").to_owned(), + ..Default::default() + }) +} + +/// Encode one batch as an OTLP `ExportLogsServiceRequest`: one resource, one instrumentation +/// scope (this crate), one log record per event. Every record carries its session's trace id and +/// attributes; records emitted inside a span carry its span id too. +pub(crate) fn encode_logs( + resource_attributes: &[Attribute], + global: &[Attribute], + events: Vec, +) -> Vec { + ExportLogsServiceRequest { + resource_logs: vec![ResourceLogs { + resource: resource(resource_attributes), + scope_logs: vec![ScopeLogs { + scope: scope(), + log_records: events.into_iter().map(|e| log_record(e, global)).collect(), + ..Default::default() + }], + ..Default::default() + }], + } + .encode_to_vec() +} + +/// Encode finished spans as an OTLP `ExportTraceServiceRequest`, each under its session's trace id. +pub(crate) fn encode_spans( + resource_attributes: &[Attribute], + global: &[Attribute], + spans: Vec, +) -> Vec { + ExportTraceServiceRequest { + resource_spans: vec![ResourceSpans { + resource: resource(resource_attributes), + scope_spans: vec![ScopeSpans { + scope: scope(), + spans: spans.into_iter().map(|s| otlp_span(s, global)).collect(), + ..Default::default() + }], + ..Default::default() + }], + } + .encode_to_vec() +} + +/// An encoded batch cut in two halves (a 413 answer), each with its record count; `None` when it +/// holds a single record or is not the one-resource, one-scope shape this crate encodes. +pub(crate) type Halves = Option<[(Vec, u64); 2]>; + +pub(crate) fn split_logs(encoded: &[u8]) -> Halves { + let mut request = ExportLogsServiceRequest::decode(encoded).ok()?; + let [resource] = &mut request.resource_logs[..] else { return None }; + let [scope] = &mut resource.scope_logs[..] else { return None }; + let n = scope.log_records.len(); + if n < 2 { + return None; + } + let second = scope.log_records.split_off(n / 2); + let first = request.encode_to_vec(); + request.resource_logs[0].scope_logs[0].log_records = second; + Some([(first, (n / 2) as u64), (request.encode_to_vec(), (n - n / 2) as u64)]) +} + +pub(crate) fn split_spans(encoded: &[u8]) -> Halves { + let mut request = ExportTraceServiceRequest::decode(encoded).ok()?; + let [resource] = &mut request.resource_spans[..] else { return None }; + let [scope] = &mut resource.scope_spans[..] else { return None }; + let n = scope.spans.len(); + if n < 2 { + return None; + } + let second = scope.spans.split_off(n / 2); + let first = request.encode_to_vec(); + request.resource_spans[0].scope_spans[0].spans = second; + Some([(first, (n / 2) as u64), (request.encode_to_vec(), (n - n / 2) as u64)]) +} + +fn log_record(Queued { mut event, session, .. }: Queued, global: &[Attribute]) -> LogRecord { + session.decorate(&mut event.attributes, global); + // `Queued::new` stamps capture time; this fallback covers only a `Queued` built by hand. + let time_unix_nano = event.timestamp_ns.unwrap_or_else(now_unix_nanos); + // Events carry a display body (OTel: "a string display message of the event"); the name is + // the last resort so no event ever renders as an empty line. `otel.event.name` (semconv 1.39) + // duplicates `EventName` for backends that do not surface the field yet. + let body = event.body.or_else(|| (!event.name.is_empty()).then(|| event.name.clone())); + if !event.name.is_empty() { + event.attributes.push(Attribute::new("otel.event.name", event.name.clone())); + } + LogRecord { + time_unix_nano, + observed_time_unix_nano: time_unix_nano, + severity_number: SeverityNumber::from(event.severity) as i32, + severity_text: severity_text(event.severity).to_owned(), + body: body.map(|text| AnyValue { value: Some(any_value::Value::StringValue(text)) }), + attributes: event.attributes.iter().map(KeyValue::from).collect(), + event_name: event.name, + trace_id: session.trace_id.to_vec(), + span_id: event.span_id.map(|id| id.to_be_bytes().to_vec()).unwrap_or_default(), + ..Default::default() + } +} + +fn otlp_span(mut record: SpanRecord, global: &[Attribute]) -> Span { + let session = record.session.clone(); + session.decorate(&mut record.attributes, global); + // The outcome is the core's, as `session.id` is: a span's own copies never ship. + record.attributes.retain(|a| a.key != "lk.outcome" && a.key != "error.type"); + let mut attributes: Vec = record.attributes.iter().map(KeyValue::from).collect(); + attributes.extend(record.outcome_attributes().iter().map(KeyValue::from)); + Span { + trace_id: session.trace_id.to_vec(), + span_id: record.span_id.to_be_bytes().to_vec(), + parent_span_id: record.parent_span_id.map(|p| p.to_be_bytes().to_vec()).unwrap_or_default(), + name: record.name, + kind: span::SpanKind::from(record.kind) as i32, + start_time_unix_nano: record.start_ns, + end_time_unix_nano: record.end_ns, + attributes, + events: record + .events + .into_iter() + .map(|e| span::Event { + time_unix_nano: e.time_ns, + name: e.name, + attributes: e.attributes.iter().map(KeyValue::from).collect(), + ..Default::default() + }) + .collect(), + // OTel: instrumentation should not set `Ok`; success and cancellation stay `Unset` and + // are told apart by `lk.outcome`. + status: Some(Status { + code: match record.outcome { + SpanOutcome::Error => status::StatusCode::Error, + SpanOutcome::Ok | SpanOutcome::Cancelled => status::StatusCode::Unset, + } as i32, + message: record.error_type.unwrap_or_default(), + }), + ..Default::default() + } +} + +impl From for span::SpanKind { + /// The OTLP value of a span kind. + fn from(kind: SpanKind) -> Self { + match kind { + SpanKind::Internal => Self::Internal, + SpanKind::Client => Self::Client, + } + } +} + +impl From for SeverityNumber { + fn from(severity: Severity) -> Self { + match severity { + Severity::Trace => SeverityNumber::Trace, + Severity::Debug => SeverityNumber::Debug, + Severity::Info => SeverityNumber::Info, + Severity::Warn => SeverityNumber::Warn, + Severity::Error => SeverityNumber::Error, + } + } +} + +fn severity_text(severity: Severity) -> &'static str { + match severity { + Severity::Trace => "TRACE", + Severity::Debug => "DEBUG", + Severity::Info => "INFO", + Severity::Warn => "WARN", + Severity::Error => "ERROR", + } +} + +impl From<&Attribute> for KeyValue { + fn from(attribute: &Attribute) -> Self { + let value = match &attribute.value { + AttributeValue::Str(s) => any_value::Value::StringValue(s.clone()), + AttributeValue::Int(i) => any_value::Value::IntValue(*i), + AttributeValue::Double(d) => any_value::Value::DoubleValue(*d), + AttributeValue::Bool(b) => any_value::Value::BoolValue(*b), + }; + KeyValue { + key: attribute.key.clone(), + value: Some(AnyValue { value: Some(value) }), + ..Default::default() + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::TelemetryEvent; + + #[test] + fn encodes_events_as_otlp_log_records() { + let resource = [Attribute::new("service.name", "test")]; + let event = TelemetryEvent::new("lk.ping") + .with_severity(Severity::Warn) + .with_body("hi") + .with_attribute("lk.ping.seq", 7i64); + let session = crate::scope::ScopeState::with_trace_id([7u8; 16]); + let bytes = encode_logs(&resource, &[], vec![Queued::new(event, session)]); + + let decoded = ExportLogsServiceRequest::decode(&bytes[..]).expect("valid OTLP"); + let resource_logs = &decoded.resource_logs[0]; + let res_attr = &resource_logs.resource.as_ref().expect("resource").attributes[0]; + assert_eq!(res_attr.key, "service.name"); + let scope_logs = &resource_logs.scope_logs[0]; + assert_eq!(scope_logs.scope.as_ref().expect("scope").name, "livekit-telemetry"); + let record = &scope_logs.log_records[0]; + assert_eq!(record.event_name, "lk.ping"); + assert_eq!(record.trace_id, vec![7u8; 16]); + assert!(record.span_id.is_empty()); + assert_eq!(record.severity_number, SeverityNumber::Warn as i32); + assert_eq!(record.severity_text, "WARN"); + assert!(record.time_unix_nano > 0); + assert_eq!(record.attributes[0].key, "lk.ping.seq"); + assert_eq!( + record.attributes[0].value.as_ref().and_then(|v| v.value.clone()), + Some(any_value::Value::IntValue(7)) + ); + } + + #[test] + fn queued_records_carry_their_capture_time_not_their_export_time() { + const HOUR_NS: u64 = 3_600_000_000_000; + let store = crate::store::Store::new(10, usize::MAX, Default::default()); + let session = crate::scope::ScopeState::new(); + let captured = now_unix_nanos(); + store.push(Queued::new(TelemetryEvent::new("unstamped"), session.clone())); + let own = TelemetryEvent { timestamp_ns: Some(7), ..TelemetryEvent::new("stamped") }; + store.push(Queued::new(own, session)); + // An hour in the queue (an outage) before the exporter encodes the batch. + crate::event::CLOCK_JUMP_NS.with(|jump| jump.set(HOUR_NS as i64)); + let bytes = encode_logs(&[], &[], store.drain(10, usize::MAX)); + let decoded = ExportLogsServiceRequest::decode(&bytes[..]).expect("valid OTLP"); + let records = &decoded.resource_logs[0].scope_logs[0].log_records; + assert!((captured..captured + HOUR_NS).contains(&records[0].time_unix_nano)); + assert_eq!(records[1].time_unix_nano, 7, "a caller's own timestamp is kept"); + } + + #[test] + fn spans_carry_the_cores_outcome_and_error_type_once() { + let mut spans = crate::span::Spans::new(1); + let id = spans.begin("lk.connect", SpanKind::Client, None); + let own = vec![Attribute::new("lk.outcome", "ok"), Attribute::new("error.type", "spoof")]; + spans.end(id, SpanOutcome::Error, Some("timeout".into()), own); + let bytes = encode_spans(&[], &[], spans.drain(1, usize::MAX)); + let decoded = ExportTraceServiceRequest::decode(&bytes[..]).expect("valid OTLP"); + let exported = &decoded.resource_spans[0].scope_spans[0].spans[0]; + assert_eq!(exported.kind, span::SpanKind::Client as i32); + let values = |key: &str| -> Vec<_> { + exported + .attributes + .iter() + .filter(|kv| kv.key == key) + .filter_map(|kv| kv.value.clone()?.value) + .collect() + }; + assert_eq!(values("lk.outcome"), [any_value::Value::StringValue("error".into())]); + assert_eq!(values("error.type"), [any_value::Value::StringValue("timeout".into())]); + } + + #[test] + fn events_without_a_body_carry_their_name_as_body() { + let session = crate::scope::ScopeState::with_trace_id([7u8; 16]); + let event = TelemetryEvent::new("lk.rtc.stats.sample"); + let bytes = encode_logs(&[], &[], vec![Queued::new(event, session)]); + let decoded = ExportLogsServiceRequest::decode(&bytes[..]).expect("valid OTLP"); + let record = &decoded.resource_logs[0].scope_logs[0].log_records[0]; + assert_eq!(record.event_name, "lk.rtc.stats.sample"); + assert_eq!( + record.body.as_ref().and_then(|b| b.value.clone()), + Some(any_value::Value::StringValue("lk.rtc.stats.sample".into())) + ); + } +} diff --git a/livekit-telemetry/src/proto/mod.rs b/livekit-telemetry/src/proto/mod.rs new file mode 100644 index 000000000..195d5e2bb --- /dev/null +++ b/livekit-telemetry/src/proto/mod.rs @@ -0,0 +1,26 @@ +// 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. + +//! OTLP protobuf types from the upstream `opentelemetry-proto` crate (prost-generated, +//! `gen-tonic-messages` only — no tonic/gRPC). The crate tracks a newer proto revision than +//! the OTLP 1.x stable surface, so construct its messages with `..Default::default()`. +//! +//! Only the types are used; the `opentelemetry`/`opentelemetry_sdk` crates it depends on are +//! dead code here and LTO removes them from release binaries (measured: +8 bytes on iOS). + +pub mod opentelemetry { + pub mod proto { + pub use opentelemetry_proto::tonic::{collector, common, logs, resource, trace}; + } +} diff --git a/livekit-telemetry/src/scope.rs b/livekit-telemetry/src/scope.rs new file mode 100644 index 000000000..e839796eb --- /dev/null +++ b/livekit-telemetry/src/scope.rs @@ -0,0 +1,218 @@ +// 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::{ + fmt, + sync::{Arc, Mutex}, +}; + +use crate::{Attribute, AttributeValue}; + +/// One session's identity: the trace id every one of its records carries, and the attributes +/// attached to them at export time (`lk.room.sid`, `lk.participant.identity`, …). +pub(crate) struct ScopeState { + pub trace_id: [u8; 16], + /// SDK-owned attributes (`lk.room.*`, `lk.participant.*`), attached at export: late-known + /// identity (the room sid arrives after join) still reaches the records captured before. + attributes: Mutex>, + /// The app's correlation attributes, copied into each record when it is captured, so a + /// later change never rewrites what is already queued. + custom: Mutex>, + /// The project host this session's batches go to (`Scope::set_server`); `None` until then. + route: Mutex>, + /// The last `(url, token)` handed over, so handing the same pair again costs a comparison; + /// `None` again once disconnected (see [`ScopeState::in_call`]). + server: Mutex>, +} + +impl ScopeState { + /// A fresh session: random, non-zero trace id (OTLP treats all-zero as absent). + pub fn new() -> Arc { + Self::with_trace_id(rand::random::().max(1).to_be_bytes()) + } + + pub fn with_trace_id(trace_id: [u8; 16]) -> Arc { + Arc::new(Self { + trace_id, + attributes: Mutex::new(Vec::new()), + custom: Mutex::new(Vec::new()), + route: Mutex::new(None), + server: Mutex::new(None), + }) + } + + /// In a call: it has a server and has not disconnected since. + pub fn in_call(&self) -> bool { + self.server.lock().unwrap_or_else(|e| e.into_inner()).is_some() + } + + /// Whether `(url, token)` is what this session already has; remembers it if not. + pub fn same_server(&self, url: &str, token: &str) -> bool { + let mut server = self.server.lock().unwrap_or_else(|e| e.into_inner()); + if server.as_ref().is_some_and(|(u, t)| u == url && t == token) { + return true; + } + *server = Some((url.to_owned(), token.to_owned())); + false + } + + pub fn route(&self) -> Option { + self.route.lock().unwrap_or_else(|e| e.into_inner()).clone() + } + + pub fn set_route(&self, host: String) { + *self.route.lock().unwrap_or_else(|e| e.into_inner()) = Some(host); + } + + /// The trace id as 32 hex characters. + pub fn hex(&self) -> String { + format!("{:032x}", u128::from_be_bytes(self.trace_id)) + } + + pub fn set_attribute(&self, key: &str, value: Option) { + let mut attributes = self.attributes.lock().unwrap_or_else(|e| e.into_inner()); + attributes.retain(|a| a.key != key); + if let Some(value) = value { + attributes.push(Attribute::new(key, value)); + } + } + + /// Whether `set_custom(key, value)` would be accepted right now: within the limits, outside + /// the SDK's namespace, not one attribute too many. + pub fn accepts_custom(&self, key: &str, value: Option<&AttributeValue>) -> bool { + crate::event::valid_custom(key, value) + && fits(&self.custom.lock().unwrap_or_else(|e| e.into_inner()), key, value) + } + + /// Set or remove an app correlation attribute; `false` when rejected (see + /// [`accepts_custom`](Self::accepts_custom)). The capacity check and the write share one + /// lock, so concurrent writers cannot overshoot the limit. + pub fn set_custom(&self, key: &str, value: Option) -> bool { + if !crate::event::valid_custom(key, value.as_ref()) { + return false; + } + let mut custom = self.custom.lock().unwrap_or_else(|e| e.into_inner()); + if !fits(&custom, key, value.as_ref()) { + return false; + } + custom.retain(|a| a.key != key); + if let Some(value) = value { + custom.push(Attribute::new(key, value)); + } + true + } + + /// The app's correlation attributes right now. + pub fn custom_snapshot(&self) -> Vec { + self.custom.lock().unwrap_or_else(|e| e.into_inner()).clone() + } + + /// Copy the app's correlation attributes into a record being captured; the record's own + /// attributes win. + pub fn snapshot_custom(&self, own: &mut Vec) { + let custom = self.custom.lock().unwrap_or_else(|e| e.into_inner()); + for attribute in custom.iter() { + if !own.iter().any(|a| a.key == attribute.key) { + own.push(attribute.clone()); + } + } + } + + /// At export: the session's SDK-owned attributes win over anything the record carries + /// (an app cannot spoof them), then the pipeline-wide ones (`global`) fill in, and + /// `session.id` (OTel semconv) — the trace id, so a record can be joined to its session even + /// where a backend drops trace ids from logs. + pub fn decorate(&self, own: &mut Vec, global: &[Attribute]) { + let session = self.attributes.lock().unwrap_or_else(|e| e.into_inner()); + for attribute in session.iter() { + own.retain(|a| a.key != attribute.key); + own.push(attribute.clone()); + } + for attribute in global { + if !own.iter().any(|a| a.key == attribute.key) { + own.push(attribute.clone()); + } + } + own.retain(|a| a.key != "session.id"); + own.push(Attribute::new("session.id", self.hex())); + } +} + +/// Whether `custom` has room for `key`: a removal or a replacement always fits. +fn fits(custom: &[Attribute], key: &str, value: Option<&AttributeValue>) -> bool { + value.is_none() + || custom.iter().any(|a| a.key == key) + || custom.len() < crate::event::MAX_CUSTOM_ATTRIBUTES +} + +impl PartialEq for ScopeState { + fn eq(&self, other: &Self) -> bool { + self.trace_id == other.trace_id + } +} + +impl fmt::Debug for ScopeState { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!(f, "Scope({})", self.hex()) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::event::MAX_CUSTOM_ATTRIBUTES; + + #[test] + fn custom_attributes_reject_sdk_keys_and_fit_replacements_at_capacity() { + for key in ["lk.room.sid", "session.id", "error.type"] { + assert!(!crate::event::valid_custom(key, None), "{key} is the SDK's"); + } + let scope = ScopeState::new(); + for i in 0..MAX_CUSTOM_ATTRIBUTES { + assert!(scope.set_custom(&format!("key.{i}"), Some(1i64.into()))); + } + assert!(!scope.set_custom("extra", Some(1i64.into()))); + assert!(scope.set_custom("key.0", Some(2i64.into())), "replacing an existing key fits"); + assert!(scope.custom_snapshot().contains(&Attribute::new("key.0", 2i64))); + assert!(scope.set_custom("key.0", None)); + assert!(scope.set_custom("extra", Some(1i64.into()))); + assert_eq!(scope.custom_snapshot().len(), MAX_CUSTOM_ATTRIBUTES); + } + + #[test] + fn concurrent_custom_insertions_respect_capacity() { + const WRITERS: usize = 16; + for _ in 0..64 { + let scope = ScopeState::new(); + for i in 0..MAX_CUSTOM_ATTRIBUTES - 1 { + assert!(scope.set_custom(&format!("key.{i}"), Some(1i64.into()))); + } + let barrier = std::sync::Barrier::new(WRITERS); + let accepted = std::thread::scope(|threads| { + let handles: Vec<_> = (0..WRITERS) + .map(|i| { + let (scope, barrier) = (&scope, &barrier); + threads.spawn(move || { + barrier.wait(); + scope.set_custom(&format!("race.{i}"), Some(1i64.into())) + }) + }) + .collect(); + handles.into_iter().map(|h| h.join().unwrap()).filter(|a| *a).count() + }); + assert_eq!(accepted, 1, "one slot left, one writer wins"); + assert_eq!(scope.custom_snapshot().len(), MAX_CUSTOM_ATTRIBUTES); + } + } +} diff --git a/livekit-telemetry/src/span.rs b/livekit-telemetry/src/span.rs new file mode 100644 index 000000000..8c43235c0 --- /dev/null +++ b/livekit-telemetry/src/span.rs @@ -0,0 +1,416 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::collections::{HashMap, VecDeque}; +use std::sync::Arc; + +use crate::{event::now_unix_nanos, scope::ScopeState, Attribute, AttributeValue}; + +/// OTel span kind, restricted to what client operations need; implied by [`crate::SpanName`]. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum SpanKind { + /// An operation inside the SDK (publish, subscribe). + Internal, + /// A call to the SFU that waits for its answer (connect, reconnect). + Client, +} + +/// How an attempt ended. OTel status knows only `Unset`/`Ok`/`Error`, so `Cancelled` travels as +/// `status = Unset` plus the `lk.outcome` attribute — every span carries `lk.outcome` so rollups +/// never have to infer it (a user hanging up mid-connect is not a failure). +#[cfg_attr(feature = "uniffi", derive(uniffi::Enum))] +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum SpanOutcome { + Ok, + Error, + Cancelled, +} + +impl SpanOutcome { + pub(crate) fn as_str(self) -> &'static str { + match self { + SpanOutcome::Ok => "ok", + SpanOutcome::Error => "error", + SpanOutcome::Cancelled => "cancelled", + } + } +} + +/// A checkpoint inside a span (OTLP span event). Structural to one attempt — the connect +/// sequence's `ws_open → join_recv → pc_connected → …` — hence in the span's own envelope rather +/// than a standalone log record (OTEP 4430 keeps that legal). +#[derive(Debug, Clone, PartialEq)] +pub(crate) struct SpanEvent { + pub name: String, + pub time_ns: u64, + pub attributes: Vec, +} + +/// One attempt at an operation, from `begin_span` to `end_span`. +#[derive(Debug, Clone, PartialEq)] +pub(crate) struct SpanRecord { + pub span_id: u64, + pub parent_span_id: Option, + pub name: String, + pub kind: SpanKind, + pub start_ns: u64, + pub end_ns: u64, + pub outcome: SpanOutcome, + pub error_type: Option, + pub attributes: Vec, + pub events: Vec, + /// The session (trace) the span belongs to. + pub session: Arc, + /// The session's project when the span ended (see `Queued::route`). + pub route: Option, +} + +/// OTel default span limits. +const MAX_EVENTS_PER_SPAN: usize = 128; +const MAX_ATTRIBUTES_PER_SPAN: usize = 128; +/// Spans a host can leave open before the oldest is abandoned (counted as dropped). +const MAX_OPEN_SPANS: usize = 256; + +/// Open spans by handle, plus the finished ones waiting for the exporter. +/// +/// Handles are opaque `u64`s minted here; the host keeps them (a Swift `Span` object, a Kotlin +/// value) and never sees ambient context — that is the platform's job (task-locals, coroutine +/// context, zones), not the FFI's. +pub(crate) struct Spans { + open: HashMap, + /// Insertion order of `open`, to abandon the oldest when the cap is hit. + open_order: VecDeque, + finished: Vec, + finished_capacity: usize, + /// Which session every recent span belongs to — open, finished or already exported — so a + /// log record that arrives after its span ended (a warning logged right before a failing + /// publish ends, delivered a hop later) is still filed under the right session. + sessions: HashMap>, + /// Spans that ended or were abandoned, oldest first. Only these age out of `sessions`: an + /// open span keeps its session however many spans start after it. + session_order: VecDeque, + next_id: u64, + pub dropped: u64, + /// The first drop since the last drain logs a warning; the rest are only counted. + full_warned: bool, + /// The opt-out, checked under the lock that guards the registry (see `Store`). + revoked: Arc, +} + +/// Ended spans whose session stays resolvable. +// ponytail: a fixed ring; a time-based expiry if a long session ever ends more spans than this +// between a log and its export. +const REMEMBERED_SPANS: usize = 1024; + +impl Spans { + /// A registry keeping at most `finished_capacity` finished spans for the exporter (at least + /// one: zero is treated as one). + pub fn new(finished_capacity: usize) -> Self { + Self { + open: HashMap::new(), + open_order: VecDeque::new(), + finished: Vec::new(), + finished_capacity: finished_capacity.max(1), + sessions: HashMap::new(), + session_order: VecDeque::new(), + // Span ids must be non-zero (OTLP treats all-zero as absent); start at 1 and mix in + // randomness so ids from two pipelines in one process never collide. 63 bits: a + // platform whose integers are signed (Dart) must be able to hand an id back. + next_id: (rand::random::() >> 1) | 1, + dropped: 0, + full_warned: false, + revoked: Arc::default(), + } + } + + /// Tie the registry to a pipeline's opt-out. + pub fn with_consent(mut self, revoked: Arc) -> Self { + self.revoked = revoked; + self + } + + fn revoked(&self) -> bool { + self.revoked.load(std::sync::atomic::Ordering::SeqCst) + } + + /// Open a span in `session`'s trace. + pub fn begin_in( + &mut self, + name: &str, + kind: SpanKind, + parent: Option, + session: Arc, + ) -> u64 { + if self.revoked() { + return 0; + } + let id = self.next_id; + self.next_id = ((self.next_id + 1) & (u64::MAX >> 1)).max(1); + if self.open.len() >= MAX_OPEN_SPANS { + if let Some(oldest) = self.open_order.pop_front() { + self.open.remove(&oldest); + self.retire(oldest); + self.count_drop("open", MAX_OPEN_SPANS); + } + } + self.sessions.insert(id, session.clone()); + let record = SpanRecord { + span_id: id, + parent_span_id: parent.filter(|p| *p != 0), + name: name.to_owned(), + kind, + start_ns: now_unix_nanos(), + end_ns: 0, + outcome: SpanOutcome::Ok, + error_type: None, + attributes: Vec::new(), + events: Vec::new(), + session, + route: None, + }; + self.open.insert(id, record); + self.open_order.push_back(id); + id + } + + /// Whether a span with one of these names is still open (the exporter holds uploads while + /// `lk.connect` / `lk.reconnect` are). + /// The session a recent span belongs to — open, ended or exported (log records emitted inside + /// a span are filed there, and they may arrive after the span ended). + pub fn scope_of(&self, id: u64) -> Option> { + self.sessions.get(&id).cloned() + } + + #[cfg(test)] + pub fn begin(&mut self, name: &str, kind: SpanKind, parent: Option) -> u64 { + self.begin_in(name, kind, parent, ScopeState::new()) + } + + #[cfg(test)] + pub fn open_count(&self) -> usize { + self.open.len() + } + + pub fn any_open(&self, names: &[&str]) -> bool { + self.open.values().any(|span| names.contains(&span.name.as_str())) + } + + pub fn add_event(&mut self, id: u64, name: &str, attributes: Vec) { + if self.revoked() { + return; + } + let Some(span) = self.open.get_mut(&id) else { return }; + if span.events.len() >= MAX_EVENTS_PER_SPAN { + return; + } + span.events.push(SpanEvent { + name: name.to_owned(), + time_ns: now_unix_nanos(), + attributes, + }); + } + + /// Close a span; the finished record waits for the next export. Unknown ids are ignored + /// (double `end` is harmless, like OTel's). + pub fn end( + &mut self, + id: u64, + outcome: SpanOutcome, + error_type: Option, + mut attributes: Vec, + ) { + if self.revoked() { + return; + } + let Some(mut span) = self.open.remove(&id) else { return }; + self.open_order.retain(|open| *open != id); + span.end_ns = now_unix_nanos().max(span.start_ns); + span.outcome = outcome; + span.error_type = error_type; + span.session.snapshot_custom(&mut attributes); + span.route = span.session.route(); + attributes.truncate(MAX_ATTRIBUTES_PER_SPAN); + span.attributes = attributes; + self.retire(id); + if self.finished.len() >= self.finished_capacity { + self.finished.remove(0); + self.count_drop("finished", self.finished_capacity); + } + self.finished.push(span); + } + + /// A span ended or was abandoned: its session stays resolvable until [`REMEMBERED_SPANS`] + /// more have. + fn retire(&mut self, id: u64) { + self.session_order.push_back(id); + if self.session_order.len() > REMEMBERED_SPANS { + if let Some(old) = self.session_order.pop_front() { + self.sessions.remove(&old); + } + } + } + + /// Count a span dropped from a full buffer. + /// The first drop since the last drain logs a warning; the rest are only counted. + fn count_drop(&mut self, buffer: &str, capacity: usize) { + self.dropped += 1; + if !self.full_warned { + self.full_warned = true; + log::warn!("{buffer} spans full ({capacity}): dropping the oldest"); + } + } + + /// Take the finished spans, oldest first. + /// At most `max` spans and about `max_bytes` (always at least one, so an oversized span + /// still ships). + pub fn drain(&mut self, max: usize, max_bytes: usize) -> Vec { + self.full_warned = false; + let (mut n, mut bytes) = (0, 0); + for span in self.finished.iter().take(max) { + let size = span.size_hint(); + if n > 0 && bytes + size > max_bytes { + break; + } + bytes += size; + n += 1; + } + self.finished.drain(..n).collect() + } + + /// Forget every open and finished span (opt-out); returns how many went. + pub fn clear(&mut self) -> u64 { + let n = self.open.len() + self.finished.len(); + self.open.clear(); + self.open_order.clear(); + self.finished.clear(); + self.sessions.clear(); + self.session_order.clear(); + n as u64 + } + + pub fn take_dropped(&mut self) -> u64 { + std::mem::take(&mut self.dropped) + } +} + +impl SpanRecord { + /// Rough encoded size, like `TelemetryEvent::size_hint`. + pub(crate) fn size_hint(&self) -> usize { + let attributes = |attributes: &[Attribute]| { + attributes.iter().map(|a| 4 + a.key.len() + a.value.size_hint()).sum::() + }; + 64 + self.name.len() + + attributes(&self.attributes) + + self + .events + .iter() + .map(|e| 16 + e.name.len() + attributes(&e.attributes)) + .sum::() + } + + /// `lk.outcome` and `error.type`, the attributes every span carries beyond the caller's. + pub(crate) fn outcome_attributes(&self) -> Vec { + let mut attributes = vec![Attribute::new("lk.outcome", self.outcome.as_str())]; + if let Some(error_type) = &self.error_type { + attributes.push(Attribute::new("error.type", AttributeValue::Str(error_type.clone()))); + } + attributes + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn spans_open_record_events_and_finish_in_order() { + let mut spans = Spans::new(8); + let parent = spans.begin("lk.connect", SpanKind::Client, None); + let child = spans.begin("lk.publish", SpanKind::Internal, Some(parent)); + spans.add_event(parent, "ws_open", vec![]); + spans.end(child, SpanOutcome::Cancelled, None, vec![]); + spans.end( + parent, + SpanOutcome::Error, + Some("timeout".into()), + vec![Attribute::new("lk.connect.attempt", 1i64)], + ); + spans.end(parent, SpanOutcome::Ok, None, vec![]); // double end: ignored + + let finished = spans.drain(10, usize::MAX); + assert_eq!(finished.len(), 2); + assert_eq!(finished[0].name, "lk.publish"); + assert_eq!(finished[0].parent_span_id, Some(parent)); + assert_eq!(finished[0].outcome, SpanOutcome::Cancelled); + assert_eq!(finished[1].events[0].name, "ws_open"); + assert_eq!(finished[1].error_type.as_deref(), Some("timeout")); + assert!(finished[1].end_ns >= finished[1].start_ns); + assert_eq!(spans.take_dropped(), 0); + } + + #[test] + fn open_spans_keep_their_session_however_many_spans_end_after_them() { + let mut spans = Spans::new(REMEMBERED_SPANS + 1); + let room = ScopeState::new(); + let connect = spans.begin_in("lk.connect", SpanKind::Client, None, room.clone()); + let first = spans.begin("lk.publish", SpanKind::Internal, None); + spans.end(first, SpanOutcome::Ok, None, vec![]); + for _ in 0..REMEMBERED_SPANS { + let id = spans.begin("lk.publish", SpanKind::Internal, None); + spans.end(id, SpanOutcome::Ok, None, vec![]); + } + assert!(spans.scope_of(first).is_none(), "ended spans age out"); + assert_eq!(spans.scope_of(connect), Some(room.clone())); + spans.end(connect, SpanOutcome::Error, Some("timeout".into()), vec![]); + assert_eq!(spans.scope_of(connect), Some(room), "and stays resolvable once it ended"); + assert_eq!(spans.sessions.len(), REMEMBERED_SPANS); + } + + #[test] + fn abandoned_spans_are_counted_warned_once_and_age_out() { + let mut spans = Spans::new(8); + let ids: Vec<_> = (0..MAX_OPEN_SPANS + REMEMBERED_SPANS + 1) + .map(|_| spans.begin("lk.publish", SpanKind::Internal, None)) + .collect(); + assert_eq!(spans.open_count(), MAX_OPEN_SPANS); + assert_eq!(spans.take_dropped(), REMEMBERED_SPANS as u64 + 1); + assert!(spans.full_warned, "the first drop logged, the rest only counted"); + assert!(spans.scope_of(ids[0]).is_none()); + assert_eq!(spans.sessions.len(), MAX_OPEN_SPANS + REMEMBERED_SPANS); + spans.drain(10, usize::MAX); + assert!(!spans.full_warned, "a drain starts a new episode"); + } + + #[test] + fn zero_finished_capacity_keeps_one_span_instead_of_panicking() { + let mut spans = Spans::new(0); + for _ in 0..2 { + let id = spans.begin("lk.publish", SpanKind::Internal, None); + spans.end(id, SpanOutcome::Ok, None, vec![]); + } + assert_eq!(spans.drain(10, usize::MAX).len(), 1); + assert_eq!(spans.take_dropped(), 1); + } + + #[test] + fn finished_spans_are_bounded() { + let mut spans = Spans::new(1); + for _ in 0..2 { + let id = spans.begin("lk.publish", SpanKind::Internal, None); + spans.end(id, SpanOutcome::Ok, None, vec![]); + } + assert_eq!(spans.drain(10, usize::MAX).len(), 1); + assert_eq!(spans.take_dropped(), 1); + } +} diff --git a/livekit-telemetry/src/stats.rs b/livekit-telemetry/src/stats.rs new file mode 100644 index 000000000..75ab223ed --- /dev/null +++ b/livekit-telemetry/src/stats.rs @@ -0,0 +1,305 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::sync::atomic::{AtomicU64, Ordering}; + +use crate::TelemetryEvent; + +/// Declares the pipeline's health counters once: the shared atomics ([`Counters`]), their +/// point-in-time copy ([`Snapshot`]) and the delta between two copies. +macro_rules! counters { + ($($(#[$doc:meta])* $name:ident,)*) => { + /// Pipeline health counters, shared by the store, the exporter and + /// [`Telemetry::stats`](crate::Telemetry::stats). Loss reasons follow the OpenTelemetry + /// SDK self-metrics conventions (`queue_full`, `rejected`, `timeout`) where one exists. + #[derive(Default)] + pub(crate) struct Counters { + $($(#[$doc])* pub $name: AtomicU64,)* + } + + /// A point-in-time copy of [`Counters`]. + #[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] + pub(crate) struct Snapshot { + $(pub $name: u64,)* + } + + impl Counters { + pub fn snapshot(&self) -> Snapshot { + Snapshot { $($name: self.$name.load(Ordering::Relaxed),)* } + } + } + + impl Snapshot { + /// Counts accumulated since `earlier`. + pub fn since(&self, earlier: &Snapshot) -> Snapshot { + Snapshot { $($name: self.$name.saturating_sub(earlier.$name),)* } + } + } + }; +} + +counters! { + /// Records evicted from the in-memory queue (`max_queue_size`). + queue_full, + /// Records lost because no cache, not even memory, could take their batch. + cache_error, + /// Records evicted from the cache by its size or file-count bound. + cache_full, + /// Records in cached batches past the 24 h age limit. + expired, + /// Records in cached batches that failed their integrity check (truncated, corrupt). + corrupt, + /// Records the core refused at the door: a custom event or attribute over the limits. + invalid, + /// Records the collector rejected: a final 4xx/5xx, or refused in a partial success. + rejected, + /// Single records larger than the collector accepts (413 down to one record). + oversized, + /// Records evicted from the cache while the collector held uploads off (`Retry-After`). + throttled, + /// Records for a project that receives nothing (self-hosted, no grant, disabled, 404). + disabled, + /// Discrete events dropped by the flood guard (`max_events_per_10min`). + rate_limited, + /// Records deleted by the opt-out. + purged, + /// Batches the collector accepted. + uploads_sent, + /// Compressed bytes the collector accepted — what telemetry actually cost the uplink. + upload_bytes, + /// Upload attempts that failed transiently (no answer, 429, 5xx). + upload_failures, + /// Upload attempts that hit `export_timeout_ms` (a slow network, or a stalled collector). + upload_timeouts, + /// 401/403 answers: a token the collector refused (the batches wait for the next one). + auth_denied, + /// Soft holds that reached the cap and let one batch through: the policy was starving + /// telemetry, and data arrived late. + hold_cap_hits, + /// Batches the disk cache could not store (disk full, directory gone), kept in memory + /// instead: they survive a failed upload, not the process. + cache_write_errors, +} + +impl Counters { + pub fn add(counter: &AtomicU64, n: u64) { + counter.fetch_add(n, Ordering::Relaxed); + } +} + +impl Snapshot { + /// Everything that was lost, by any reason. + pub fn dropped(&self) -> u64 { + self.queue_full + + self.cache_error + + self.cache_full + + self.expired + + self.corrupt + + self.invalid + + self.rejected + + self.oversized + + self.throttled + + self.disabled + + self.rate_limited + + self.purged + } + + /// Anything worth telling the backend about: data lost (other than by policy: a project + /// that receives nothing, an opt-out), uploads failing or refused, holds starving uploads, + /// the disk refusing writes. + pub fn has_problems(&self) -> bool { + self.dropped() - self.disabled - self.purged + + self.upload_failures + + self.upload_timeouts + + self.auth_denied + + self.hold_cap_hits + + self.cache_write_errors + > 0 + } + + /// The `lk.telemetry.report` event: what this pipeline sent, dropped or failed to upload + /// since the previous report — deltas by reason, riding along with the next batch, never + /// persisted on its own, never an extra request — plus one at shutdown, so every session + /// leaves a summary the fleet's success rates can be computed from. + pub fn report(&self, cached_batches: u64) -> TelemetryEvent { + let mut event = TelemetryEvent::new("lk.telemetry.report") + .with_body(format!( + "telemetry: {} batches sent ({} B), {} failed, {} dropped, {} cached", + self.uploads_sent, + self.upload_bytes, + self.upload_failures + self.upload_timeouts, + self.dropped() - self.disabled - self.purged, + cached_batches + )) + .with_attribute("lk.telemetry.uploads.sent", self.uploads_sent as i64) + .with_attribute("lk.telemetry.uploads.bytes", self.upload_bytes as i64) + .with_attribute("lk.telemetry.uploads.failed", self.upload_failures as i64) + .with_attribute("lk.telemetry.cache.batches", cached_batches as i64); + for (key, value) in [ + ("lk.telemetry.uploads.timeouts", self.upload_timeouts), + ("lk.telemetry.uploads.unauthorized", self.auth_denied), + ("lk.telemetry.holds.capped", self.hold_cap_hits), + ("lk.telemetry.cache.write_errors", self.cache_write_errors), + ("lk.telemetry.dropped.queue_full", self.queue_full), + ("lk.telemetry.dropped.cache_error", self.cache_error), + ("lk.telemetry.dropped.cache_full", self.cache_full), + ("lk.telemetry.dropped.expired", self.expired), + ("lk.telemetry.dropped.corrupt", self.corrupt), + ("lk.telemetry.dropped.invalid", self.invalid), + ("lk.telemetry.dropped.rejected", self.rejected), + ("lk.telemetry.dropped.oversized", self.oversized), + ("lk.telemetry.dropped.throttled", self.throttled), + ("lk.telemetry.dropped.rate_limited", self.rate_limited), + ] { + if value > 0 { + event = event.with_attribute(key, value as i64); + } + } + event + } +} + +/// What the upload policy is doing right now. +#[cfg_attr(feature = "uniffi", derive(uniffi::Enum))] +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum TelemetryStatus { + /// Uploading as data arrives. + Ok, + /// Uploads wait for a connect/reconnect to finish or for the device (Low Data Mode, low + /// battery); capped at 60 s. + Held, + /// The last upload failed; retrying after a jittered exponential backoff. + Paused, + /// The collector asked for a pause (`Retry-After`). + Throttled, + /// No destination or no usable token yet (no room connected, the token expired, lacks the + /// grant or was refused); everything waits in the cache for the next token. + Waiting, + /// No project receives telemetry: every server this process talked to is self-hosted, never + /// granted observability, has no ingest, or disabled data recording. + Off, +} + +impl std::fmt::Display for TelemetryStatus { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.write_str(match self { + Self::Ok => "ok", + Self::Held => "held", + Self::Paused => "paused", + Self::Throttled => "throttled", + Self::Waiting => "waiting", + Self::Off => "off", + }) + } +} + +/// Pipeline health as seen by the SDK: [`Telemetry::stats`](crate::Telemetry::stats). +#[cfg_attr(feature = "uniffi", derive(uniffi::Record))] +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct TelemetryStats { + /// What the upload policy is doing right now. + pub status: TelemetryStatus, + /// Records lost for any reason (sum of the `dropped_*` fields). + pub dropped: u64, + pub dropped_queue_full: u64, + pub dropped_cache_error: u64, + pub dropped_cache_full: u64, + pub dropped_expired: u64, + pub dropped_corrupt: u64, + pub dropped_invalid: u64, + pub dropped_rejected: u64, + pub dropped_oversized: u64, + pub dropped_throttled: u64, + pub dropped_disabled: u64, + pub dropped_rate_limited: u64, + pub dropped_purged: u64, + /// Batches the collector accepted. + pub uploads_sent: u64, + /// Compressed bytes the collector accepted. + pub upload_bytes: u64, + /// Upload attempts that failed transiently (no answer, 429, 5xx). + pub upload_failures: u64, + /// Upload attempts that timed out. + pub upload_timeouts: u64, + /// Tokens the collector refused (401/403). + pub uploads_unauthorized: u64, + /// Soft holds that reached the cap. + pub holds_capped: u64, + /// Batches the disk cache could not store, kept in memory instead. + pub cache_write_errors: u64, + /// Batches currently waiting in the cache. + pub cached_batches: u64, +} + +/// One line for a debug console: status, throughput, one backlog number, one loss number, then +/// the loss breakdown. +impl std::fmt::Display for TelemetryStats { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!( + f, + "{}, sent {} ({} B), backlog {}, failed {}, lost {} (queue {}, cache {}, expired {}, \ + corrupt {}, invalid {}, rejected {}, oversized {}, throttled {}, rate-limited {}, \ + disabled {}, purged {}; unauthorized {}, holds capped {}, write errors {})", + self.status, + self.uploads_sent, + self.upload_bytes, + self.cached_batches, + self.upload_failures + self.upload_timeouts, + self.dropped, + self.dropped_queue_full, + self.dropped_cache_full + self.dropped_cache_error, + self.dropped_expired, + self.dropped_corrupt, + self.dropped_invalid, + self.dropped_rejected, + self.dropped_oversized, + self.dropped_throttled, + self.dropped_rate_limited, + self.dropped_disabled, + self.dropped_purged, + self.uploads_unauthorized, + self.holds_capped, + self.cache_write_errors, + ) + } +} + +impl TelemetryStats { + pub(crate) fn new(snapshot: Snapshot, cached_batches: u64, status: TelemetryStatus) -> Self { + Self { + status, + dropped: snapshot.dropped(), + dropped_queue_full: snapshot.queue_full, + dropped_cache_error: snapshot.cache_error, + dropped_cache_full: snapshot.cache_full, + dropped_expired: snapshot.expired, + dropped_corrupt: snapshot.corrupt, + dropped_invalid: snapshot.invalid, + dropped_rejected: snapshot.rejected, + dropped_oversized: snapshot.oversized, + dropped_throttled: snapshot.throttled, + dropped_disabled: snapshot.disabled, + dropped_rate_limited: snapshot.rate_limited, + dropped_purged: snapshot.purged, + uploads_sent: snapshot.uploads_sent, + upload_bytes: snapshot.upload_bytes, + upload_failures: snapshot.upload_failures, + upload_timeouts: snapshot.upload_timeouts, + uploads_unauthorized: snapshot.auth_denied, + holds_capped: snapshot.hold_cap_hits, + cache_write_errors: snapshot.cache_write_errors, + cached_batches, + } + } +} diff --git a/livekit-telemetry/src/store.rs b/livekit-telemetry/src/store.rs new file mode 100644 index 000000000..87ebf864b --- /dev/null +++ b/livekit-telemetry/src/store.rs @@ -0,0 +1,192 @@ +// 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::VecDeque, + sync::{ + atomic::{AtomicBool, Ordering}, + Arc, Mutex, + }, +}; + +use crate::{event::now_unix_nanos, scope::ScopeState, stats::Counters, TelemetryEvent}; + +/// An event waiting for export, filed under the session whose trace id and attributes it +/// will carry. +pub(crate) struct Queued { + pub event: TelemetryEvent, + pub session: Arc, + /// The project the session was routed to when the record was captured: immutable, so a Room + /// that later reconnects to another project never takes queued records along. `None`: the + /// session had no server yet. + pub route: Option, +} + +impl Queued { + /// Capture a record for `session`, with its owner as of now: the session's project and its + /// correlation attributes (the record's own win). An unstamped record is stamped here, so + /// time spent queued never shifts it to its export time. + pub fn new(mut event: TelemetryEvent, session: Arc) -> Self { + event.timestamp_ns.get_or_insert_with(now_unix_nanos); + let route = session.route(); + session.snapshot_custom(&mut event.attributes); + Self { event, session, route } + } +} + +/// Bounded FIFO of events waiting for export. +/// +/// When full, the *oldest* event is dropped so the freshest context survives a burst +/// (the queue role of OTel's `BatchLogRecordProcessor`, with drop-oldest instead of +/// drop-newest); every eviction is counted as `queue_full`. Tracks its approximate size in +/// bytes so the exporter can flush early (design doc: every tick *or* at 256 KB) and bound a +/// request's size. +// ponytail: one mutex around a VecDeque; a lock-free ring only if `emit` shows up in a profile. +pub(crate) struct Store { + queue: Mutex, + /// The opt-out, checked under the queue lock `clear` also takes: nothing lands after a purge. + revoked: Arc, + /// Test-only: runs at the start of `push`, before the lock (to hold a producer there). + #[cfg(test)] + pub(crate) pause: Mutex>>, + capacity: usize, + flush_threshold: usize, + counters: Arc, +} + +#[derive(Default)] +struct Queue { + events: VecDeque, + bytes: usize, + /// The first drop since the last drain logs a warning; the rest are only counted. + full_warned: bool, +} + +impl Store { + pub fn new(capacity: usize, flush_threshold: usize, counters: Arc) -> Self { + Self { + queue: Mutex::new(Queue::default()), + revoked: Arc::default(), + #[cfg(test)] + pause: Mutex::new(None), + capacity, + flush_threshold, + counters, + } + } + + /// Tie the queue to a pipeline's opt-out. + pub fn with_consent(mut self, revoked: Arc) -> Self { + self.revoked = revoked; + self + } + + /// Queue an event. Returns `true` when this push carried the queue across + /// `flush_threshold` bytes — the caller should wake the exporter. + pub fn push(&self, queued: Queued) -> bool { + #[cfg(test)] + { + let pause = self.pause.lock().unwrap_or_else(|e| e.into_inner()).clone(); + if let Some(pause) = pause { + pause(); + } + } + let mut queue = self.queue.lock().unwrap_or_else(|e| e.into_inner()); + if self.revoked.load(Ordering::SeqCst) { + return false; + } + if queue.events.len() >= self.capacity { + if let Some(oldest) = queue.events.pop_front() { + queue.bytes = queue.bytes.saturating_sub(oldest.event.size_hint()); + } + Counters::add(&self.counters.queue_full, 1); + if !queue.full_warned { + queue.full_warned = true; + log::warn!( + "queue full ({} records): dropping oldest until the exporter drains", + self.capacity + ); + } + } + let before = queue.bytes; + queue.bytes += queued.event.size_hint(); + queue.events.push_back(queued); + before < self.flush_threshold && queue.bytes >= self.flush_threshold + } + + /// Remove and return the oldest events: at most `max` of them and about `max_bytes` in total + /// (always at least one, so an oversized event still ships). + pub fn drain(&self, max: usize, max_bytes: usize) -> Vec { + let mut queue = self.queue.lock().unwrap_or_else(|e| e.into_inner()); + queue.full_warned = false; + let mut out = Vec::new(); + let mut bytes = 0; + while out.len() < max { + let Some(next) = queue.events.front() else { break }; + let size = next.event.size_hint(); + if !out.is_empty() && bytes + size > max_bytes { + break; + } + bytes += size; + queue.bytes = queue.bytes.saturating_sub(size); + out.extend(queue.events.pop_front()); + } + out + } + + pub fn is_empty(&self) -> bool { + self.queue.lock().unwrap_or_else(|e| e.into_inner()).events.is_empty() + } + + /// Drop everything queued (opt-out); returns how many records went. + pub fn clear(&self) -> u64 { + let mut queue = self.queue.lock().unwrap_or_else(|e| e.into_inner()); + queue.bytes = 0; + queue.events.drain(..).count() as u64 + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn queued(event: TelemetryEvent) -> Queued { + Queued::new(event, ScopeState::new()) + } + + #[test] + fn drops_oldest_when_full() { + let counters = Arc::new(Counters::default()); + let store = Store::new(2, usize::MAX, counters.clone()); + for name in ["a", "b", "c"] { + store.push(queued(TelemetryEvent::new(name))); + } + let names: Vec<_> = store.drain(10, usize::MAX).into_iter().map(|q| q.event.name).collect(); + assert_eq!(names, ["b", "c"]); + assert_eq!(counters.snapshot().queue_full, 1); + assert!(store.drain(10, usize::MAX).is_empty()); + } + + #[test] + fn reports_the_threshold_crossing_once_and_drains_by_bytes() { + let store = Store::new(100, 100, Arc::default()); + let event = || queued(TelemetryEvent::new("e").with_body("x".repeat(30))); // 63 bytes + assert!(!store.push(event()), "63 < 100"); + assert!(store.push(event()), "126 crosses 100"); + assert!(!store.push(event()), "already above: no second wake-up"); + assert_eq!(store.drain(10, 130).len(), 2, "two fit in 130 bytes"); + assert_eq!(store.drain(10, 1).len(), 1, "an oversized event still ships alone"); + assert!(store.drain(10, usize::MAX).is_empty()); + } +} diff --git a/livekit-telemetry/uniffi.toml b/livekit-telemetry/uniffi.toml new file mode 100644 index 000000000..bb41bd4cc --- /dev/null +++ b/livekit-telemetry/uniffi.toml @@ -0,0 +1,7 @@ +[bindings.swift] +ffi_module_name = "RustLiveKitTelemetry" + +[bindings.kotlin] +# The Kotlin checksum test is broken on ARM in every UniFFI release this workspace can use; see +# livekit-uniffi/uniffi.toml for the full explanation. +omit_checksums = true