diff --git a/.changeset/add-livekit-telemetry-crate.md b/.changeset/add-livekit-telemetry-crate.md new file mode 100644 index 000000000..4e6eb7d95 --- /dev/null +++ b/.changeset/add-livekit-telemetry-crate.md @@ -0,0 +1,5 @@ +--- +livekit-telemetry: minor +--- + +Add `livekit-telemetry`, the shared client telemetry core: records (events, warn/error logs, spans, RTC stats windows) are buffered on the device, committed to a write-ahead cache (memory, or fsynced files with `storage_dir`) and exported as OTLP/HTTP logs and traces to the LiveKit Cloud project each Room belongs to. A Room hands over only its server URL and token (`Scope::set_server`, again on every refresh); the core validates the URL, reads the token's grant and expiry, and uploads each record with the token of the Room and project that captured it. Every collector answer is classified (partial success, 4xx, 413 split, 401/403 disabled vs credential, 404, 429/503 `Retry-After`/`RetryInfo`, 5xx, no answer) with jittered backoff per destination, soft and hard device holds, loss counters by reason (`lk.telemetry.report`, `Telemetry::stats`), Room-scoped custom events and correlation attributes, and an opt-out that purges everything unsent. Defaults export and window RTC stats once a minute; `LK_TELEMETRY_ENDPOINT` points uploads at a local collector for tests. 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/SPEC.md b/livekit-telemetry/SPEC.md new file mode 100644 index 000000000..07a350f87 --- /dev/null +++ b/livekit-telemetry/SPEC.md @@ -0,0 +1,254 @@ +# 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). + +## 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` | always present: `ok` \| `error` \| `cancelled` — rollups read this, never the status | +| `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..957de0027 --- /dev/null +++ b/livekit-telemetry/src/event.rs @@ -0,0 +1,282 @@ +// 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` stamps the event at emit time. + #[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 — 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::() + } +} + +/// 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::new("") + .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 — and the session id). +pub(crate) fn reserved(key: &str) -> bool { + key.starts_with("lk.") || key == "session.id" +} + +/// 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..1276000d0 --- /dev/null +++ b/livekit-telemetry/src/otlp.rs @@ -0,0 +1,273 @@ +// 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); + 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); + 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: match record.kind { + SpanKind::Internal => span::SpanKind::Internal, + SpanKind::Client => span::SpanKind::Client, + } 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 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 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..bf282f015 --- /dev/null +++ b/livekit-telemetry/src/scope.rs @@ -0,0 +1,165 @@ +// 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)); + } + } + + /// Set or remove an app correlation attribute. `false` when rejected: over the limits, in + /// the SDK's namespace, or one attribute too many. + /// Whether `set_custom(key, value)` would be accepted: within the limits, outside the + /// SDK's namespace, not one attribute too many. + pub fn accepts_custom(&self, key: &str, value: Option<&AttributeValue>) -> bool { + if !crate::event::valid_custom(key, value) { + return false; + } + let custom = self.custom.lock().unwrap_or_else(|e| e.into_inner()); + value.is_none() + || custom.iter().any(|a| a.key == key) + || custom.len() < crate::event::MAX_CUSTOM_ATTRIBUTES + } + + /// Set or remove an app correlation attribute; `false` when rejected (see + /// [`accepts_custom`](Self::accepts_custom)). + pub fn set_custom(&self, key: &str, value: Option) -> bool { + if !self.accepts_custom(key, value.as_ref()) { + return false; + } + let mut custom = self.custom.lock().unwrap_or_else(|e| e.into_inner()); + 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())); + } +} + +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()) + } +} diff --git a/livekit-telemetry/src/span.rs b/livekit-telemetry/src/span.rs new file mode 100644 index 000000000..6a553f05f --- /dev/null +++ b/livekit-telemetry/src/span.rs @@ -0,0 +1,347 @@ +// 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>, + session_order: VecDeque, + next_id: u64, + pub dropped: u64, + /// The opt-out, checked under the lock that guards the registry (see `Store`). + revoked: Arc, +} + +/// Spans whose session stays resolvable after they ended. +// ponytail: a fixed ring; a time-based expiry if a long session ever opens more spans than this +// between a log and its export. +const REMEMBERED_SPANS: usize = 1024; + +impl Spans { + pub fn new(finished_capacity: usize) -> Self { + Self { + open: HashMap::new(), + open_order: VecDeque::new(), + finished: Vec::new(), + finished_capacity, + 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, + 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.dropped += 1; + } + } + self.sessions.insert(id, session.clone()); + 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); + } + } + 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; + if self.finished.len() >= self.finished_capacity { + self.finished.remove(0); + self.dropped += 1; + } + self.finished.push(span); + } + + /// 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 { + 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 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..88977b227 --- /dev/null +++ b/livekit-telemetry/src/store.rs @@ -0,0 +1,190 @@ +// 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::{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). + pub fn new(mut event: TelemetryEvent, session: Arc) -> Self { + 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 eviction of an episode logs; the rest are 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