Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions livekit-telemetry/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,8 @@ base64 = "0.22"
serde_json = { workspace = true }
# Server URLs are parsed, never string-matched, before a token may follow them.
url = "2.3"
# gzip for request bodies; the deflate backend (miniz_oxide) is already linked by data streams.
flate2 = "1"
# OTLP message types only (`gen-tonic-messages` = prost structs, no tonic). Same prost as
# livekit-protocol so a single prost is linked.
opentelemetry-proto = { version = "0.32", default-features = false, features = ["logs", "trace", "gen-tonic-messages"] }
Expand Down
81 changes: 81 additions & 0 deletions livekit-telemetry/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,81 @@
# LiveKit Telemetry

**Important**:
This is an internal crate that powers client telemetry in LiveKit client SDKs (through
`livekit-uniffi`: Swift, Kotlin, Dart) and is not meant to be used directly.

The core buffers records on the device and ships them out-of-band as standard
[OTLP/HTTP](https://opentelemetry.io/docs/specs/otlp/) logs and traces to the LiveKit Cloud
project a room belongs to. Everything hard lives here once — destination and credentials,
batching, encoding, retry, persistence, holds, loss accounting — while platforms provide only
what they are uniquely placed to: OS signals (instruments) and the byte-moving transport.

```text
SDK / instruments ──emit()──▶ Telemetry ──▶ Store ──drain──▶ Exporter ──push──▶ BatchCache ──upload oldest-first──▶ TelemetryTransport
(source) (records) (tick · encode) (MemoryCache | FileCache) (NetTransport | host HTTP)
```

| Layer | OTel equivalent | Type |
|---|---|---|
| source | `Logger.emit`, `Tracer.start` | [`Telemetry`], [`Scope`] (one per room), [`Span`] |
| store | `BatchLogRecordProcessor` queue | `Store` (in-memory, drop-oldest, bounded) |
| exporter | batch processor timer + exporter | [`Exporter`] (actor, spawn `run()`) |
| cache | disk buffering | [`BatchCache`]: [`MemoryCache`], [`FileCache`] (`storage_dir`) |
| transport | exporter's HTTP client | [`TelemetryTransport`]; `NetTransport` over `livekit-net` (feature `net`) |

## Usage

```rust
# use std::sync::Arc;
# use livekit_telemetry::*;
# struct Discard;
# #[async_trait::async_trait]
# impl TelemetryTransport for Discard {
# async fn send(&self, _: ExportRequest) -> Result<ExportResponse, ExportError> {
# Ok(ExportResponse::accepted())
# }
# }
# #[tokio::main(flavor = "current_thread")] async fn main() {
let config = TelemetryConfig { storage_dir: None, ..Default::default() };
let (telemetry, exporter) = Telemetry::new(config, Arc::new(Discard));
tokio::spawn(exporter.run());

// One scope per room. Its server URL and token decide where its records go; hand the token
// over again on every refresh.
let room = telemetry.begin_scope();
room.set_server("wss://my-project.livekit.cloud", "<participant token>");
room.set_attribute("app.call_id", Some("c-42".into()));
room.emit_custom("checkout", vec![Attribute::new("plan", "pro")]);

let connect = room.start(SpanName::Connect, None);
connect.step(SpanStep::WsOpen);
connect.end(SpanOutcome::Ok, None);

telemetry.set_device_state(DeviceState { thermal: ThermalState::Serious, ..Default::default() });
telemetry.shutdown().await; // cache, then upload what the network allows
println!("{}", telemetry.stats()); // drops by reason, uploads, cached batches
# }
```

Event names, attributes, cadences and the upload policy are defined in [`SPEC.md`](SPEC.md).

## Design notes

- **The room decides the destination.** A platform passes only the server URL and the token
([`Scope::set_server`]). The core derives the ingest URL (LiveKit Cloud hosts only), reads
the token's unverified claims (observability grant, expiry) and routes each batch by its
owner — project and session — so two rooms on two projects never share a token.
- **Write-ahead cache.** Every batch is encoded, gzipped and cached *before* the network is
tried, then uploaded oldest-first and removed on an answer. Tokens are attached at upload
time and never persisted.
- **The core reads every answer.** Transports return status, headers and body; the core alone
classifies OTLP partial success, 4xx, 413 (split), 429/503 `Retry-After` and `RetryInfo`,
and Cloud's disabled and credential answers.
- **Loss is counted, never silent.** Every way a record can be lost has a counter; the deltas
ride along with the next batch as `lk.telemetry.report` and are readable locally through
[`Telemetry::stats`].
- **Device state comes from the host; the policy lives here.** Thermal, power, memory,
network and lifecycle are OS APIs the host watches and pushes through
[`Telemetry::set_device_state`]; the core stretches its cadence and holds uploads.
- **Size.** OTLP types come from `opentelemetry-proto` (`gen-tonic-messages`, no tonic); JWT
claims are read with `base64` and `serde_json`, both already linked by the UniFFI library.
Loading
Loading