Repository navigation
Conversation
Changeset ✓This PR includes a changeset covering all affected packages:
|
Add `livekit-telemetry`: the records SDKs push in (events, log records, attributes), the bounded in-memory queue, health counters and the `lk.telemetry.report` event, the span store, a session's identity (trace id, SDK and app attributes, route), and the OTLP/HTTP protobuf encoding of logs and traces on `opentelemetry-proto` message types.
ddc7bc2 to
fc0ef58
Compare
| --- | ||
| 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. |
There was a problem hiding this comment.
Just as a reminder here @pblazej, make sure you create https://crates.io/crates/livekit-telemetry and publish an initial version prior to this merging, otherwise knope publishing will break.
There was a problem hiding this comment.
Thanks. 0.1.0 will be on crates.io before this merges, so the release job's check passes. The changeset is now one line that's true for this PR alone.
|
|
||
| impl From<LogRecord> for TelemetryEvent { | ||
| fn from(record: LogRecord) -> Self { | ||
| let mut event = TelemetryEvent::new("") |
There was a problem hiding this comment.
thought: Given ::new("") is a thing, should TelemetryEvent maybe implement Default?
There was a problem hiding this comment.
Added Default (an unnamed Info record); From<LogRecord> uses it (0424f50).
There was a problem hiding this comment.
Most likely you don't need impl Default for TelemetryEvent, you could do #[derive(Default)] instead, but this is a nitpick:.
A record queued without a timestamp was stamped when the exporter encoded it, so time spent queued (an outage, a slow drain) shifted it to its export time and could reorder it against later records. `Queued::new` now stamps it when it is captured; the encoder keeps its fallback only for a `Queued` built by hand. The span registry aged session mappings out by start order, so an open span lost its session after 1,024 newer spans began. Only ended or abandoned spans now enter the ring that ages out; an open span keeps its session until it ends. `Spans::new(0)` panicked on the first ended span (`remove(0)` on an empty buffer); the capacity is now treated as at least one. Dropping an abandoned open span or the oldest finished span is still counted, and now logs a warning on the first drop since the last drain, the same rule `Store` applies to its queue.
… span outcomes `set_custom` checked the 64-attribute limit under one lock and inserted under another, so concurrent writers adding distinct keys could all pass the check and overshoot the limit (four writers for one free slot in the new test). The check and the write now share one lock; `accepts_custom` uses the same rule. A span ended with its own `lk.outcome` or `error.type` attribute was exported with those copies next to the core's, so a backend reading the first value could count a failed span as successful. The encoder now drops them before appending the core's, the same rule it applies to `session.id`, and SPEC.md says so. `error.type` joins `lk.*` and `session.id` as an SDK-owned key, so an app's correlation attributes and custom events cannot set it either: one rule for app input and for the span encoder.
…ersion
Review nits, no behaviour change:
- `TelemetryEvent` implements `Default` as an unnamed `Info` record,
`new("")`, which `From<LogRecord>` now uses.
- The span kind's OTLP mapping moves from an inline match in the encoder
to `From<SpanKind>`, like the severity mapping beside it.
- `size_hint` says its estimate is in bytes.
Cargo.toml declares `readme = "README.md"`, but the file was missing, so `cargo package` and `cargo publish` failed before producing an archive. The README describes what the crate holds at this point (the data model, its bounded buffers and the OTLP encoding) and says SPEC.md is the contract for the whole crate, including parts not implemented yet. SPEC.md now states the timestamp rule (a record without one is stamped when it is captured, an explicit one is kept), lists the SDK-owned keys in one place (`lk.*`, `session.id`, `error.type`), and says each span carries `lk.outcome` exactly once and `error.type` at most once. The changeset summary is one line that holds for this PR on its own, rather than describing the pipeline that lands in later PRs.
| /// 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 { |
There was a problem hiding this comment.
curiously, why do we need LogSource ?
There was a problem hiding this comment.
It's recorded as lk.log.source; the pipeline uses it to keep WebRTC lines at error only (chatty at warn) and to never feed the core's own telemetry logs back in.
| #[cfg_attr(feature = "uniffi", uniffi(default))] | ||
| pub function: Option<String>, | ||
| #[cfg_attr(feature = "uniffi", uniffi(default))] | ||
| pub file: Option<String>, |
There was a problem hiding this comment.
I would assume we want to keep the telemetry logs overhead small, and not sure if function + file + line are necessary here.
I think line is cheap, wIll it make sense to do things like
line + short file/module identifier
?
There was a problem hiding this comment.
Optional, mapped to OTel's code.*, and the pipeline only exports warning and error lines, so the cost is small. function stays stable across releases, unlike line. The core's own lines will drop the redundant module path from function in the UniFFI stage.
| 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<String>, |
There was a problem hiding this comment.
what is the key use case for logger ?
I'd think we want to keep things simple from the beginning
There was a problem hiding this comment.
An optional grouping key (type, module or WebRTC tag), sent as lk.log.logger. The pipeline also uses it to recognise and skip its own telemetry logs.
| } | ||
| LogRecord { | ||
| time_unix_nano, | ||
| observed_time_unix_nano: time_unix_nano, |
There was a problem hiding this comment.
are we logging both observed_time_unix_nano and time_unix_nano ?
and why observed_time_unix_nano is needed ?
There was a problem hiding this comment.
Both, same value. observed_time_unix_nano is when the collecting side saw it; for in-process records OTel SDKs set it equal to the timestamp. 9 bytes per record.
| // limitations under the License. | ||
|
|
||
| // Internals the pipeline (`Telemetry`, `Exporter`) consumes once it is in place. | ||
| #![allow(dead_code)] |
There was a problem hiding this comment.
issue: Instead of disabling lint globally for the crate, consider annotating individual declarations as needed.
There was a problem hiding this comment.
Agreed. The crate-wide allow is gone; each dormant module carries #[expect(dead_code, reason = …)], which warns once the module has no dead code left. Fixed in b417592.
|
|
||
| /// 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 { |
There was a problem hiding this comment.
suggestion: Consider adding a unit test to ensure this is within range of the actual encoded size.
There was a problem hiding this comment.
Good call: measured, the hint was 1.2–3.8× under the encoded size. Overheads are now named constants calibrated against the encoder, and a test keeps the hint within 0.8–1.25× for typical records. Fixed in b559dab.
|
|
||
| /// An event waiting for export, filed under the session whose trace id and attributes it | ||
| /// will carry. | ||
| pub(crate) struct Queued { |
There was a problem hiding this comment.
nitpick: Consider renaming QueuedEvent for clarity of what is being queued.
| outcome: SpanOutcome, | ||
| error_type: Option<String>, | ||
| mut attributes: Vec<Attribute>, | ||
| ) { |
There was a problem hiding this comment.
how would it handle the out of order problem ? for instance, we are expecting things like
begin span 42
event A
event B
end span 42
but it comes as
begin span 42
event A
end span 42
event B
Do we still log event B ? so that we can track what is going on ?
There was a problem hiding this comment.
B is dropped, as OTel ignores calls on an ended span. It only happens if the platform steps after end; a log written then still lands with the span id, under the right session. Documented in 94107d2.
| /// Upload attempts that failed transiently (no answer, 429, 5xx). | ||
| pub upload_failures: u64, | ||
| /// Upload attempts that timed out. | ||
| pub upload_timeouts: u64, |
There was a problem hiding this comment.
curiously, what are the key difference between upload_failures and upload_timeouts? do we have ways to distinguish them in real world ? or they will be flaky anyway ?
There was a problem hiding this comment.
A timeout is our export_timeout_ms expiring with no answer; a failure is a retryable answer (429, 5xx) or a network error before it. Different fixes, so counted apart; the summary adds them. Documented in 94107d2.
| pub fn push(&self, queued: Queued) -> bool { | ||
| #[cfg(test)] | ||
| { | ||
| let pause = self.pause.lock().unwrap_or_else(|e| e.into_inner()).clone(); |
There was a problem hiding this comment.
I am not a Rust guy, but wonder if it is a good idea to inject a test hook in production code.
I am open to it if that is a common technique in our rust code base
There was a problem hiding this comment.
Compiled only into tests, like TwirpClient's test hook. It lets the later pipeline's opt-out race test park a producer just before the queue lock.
There was a problem hiding this comment.
(FWIW, for better or worse, this is being done in some other places...)
| let mut out = Vec::new(); | ||
| let mut bytes = 0; | ||
| while out.len() < max { | ||
| let Some(next) = queue.events.front() else { break }; |
There was a problem hiding this comment.
any chance this can get into an infinitely loop ?
There was a problem hiding this comment.
No. Each pass stops or removes one event, so at most min(max, len) events plus one stopping pass. Documented in 94107d2.
…, clamp store capacity - Custom events no longer copy their name into the body; the encoder's name fallback puts the same body on the wire. - Reserve the otel.* and code.* namespaces for the SDK, so an app cannot duplicate or spoof otel.event.name or a log line's code.* attributes. - Spans set no status message, and carry error.type only when they ended in error. - Size hints use named overheads calibrated against the encoder; a test keeps encoded/hint within 0.8..=1.25 for typical records and a span. - Store::new treats a zero capacity or flush threshold as one, and clearing the queue or the spans starts a new warning episode.
…ename QueuedEvent, scope dead_code - Span events carry only a name and a time; no caller set attributes. - Snapshot::problem_drops() lists the non-policy loss reasons, shared by has_problems() and the report's dropped count. - ScopeState::decorate releases the session lock before the global attributes are merged. - Rename Queued to QueuedEvent. - Replace the crate-wide dead_code allow with an expectation on each module the pipeline consumes later, and on MAX_NAME_BYTES.
- Store::drain terminates: each pass stops or removes one event, so it removes at most min(max, len) events with at most one extra, stopping pass; max == 0 returns nothing. - A checkpoint after its span ended is ignored, as in OTel; log records carrying the span id still land in its session. - Upload failures are errors before the timeout; timeouts are counted apart.
| RECORD_OVERHEAD_BYTES | ||
| + 2 * self.name.len() | ||
| + self.body.as_ref().map_or(self.name.len(), String::len) | ||
| + attributes_size_hint(&self.attributes) |
There was a problem hiding this comment.
Curiously, does the attributes include the global / session attributes already ?
or those are negligible in size ?
There was a problem hiding this comment.
Session attributes are now attached at capture, so size_hint counts them (and records keep their room after a disconnect). Fixed in 775b105.
| span.error_type = error_type; | ||
| span.session.snapshot_custom(&mut attributes); | ||
| span.route = span.session.route(); | ||
| attributes.truncate(MAX_ATTRIBUTES_PER_SPAN); |
There was a problem hiding this comment.
should we print out a warning if a truncation happens ?
There was a problem hiding this comment.
Over-cap attributes and checkpoints are now counted on the span and warn once, like the buffer limits. Fixed in 87f13dd.
… export - ScopeState::decorate now runs when a record is queued (QueuedEvent::new) and when a span ends, right after the app's correlation attributes, instead of in the encoder. A record captured before the platform clears the room (a disconnect) keeps its lk.room.* / lk.participant.* attributes, and size hints, queue bytes and request byte bounds count them and session.id. - Pipeline-wide attributes still fill in at encode, since the queues never see them. Precedence is unchanged: the session's keys win over the record's own, then the app's correlation attributes, then the pipeline's. - Recalibrate RECORD_OVERHEAD_BYTES (128 -> 78) and SPAN_OVERHEAD_BYTES (144 -> 94), which no longer include session.id; the calibration test now measures records decorated with a room. - SPEC: a record's attributes, like its owner and timestamp, are taken at capture.
…ts exceed the cap - Attributes past MAX_ATTRIBUTES_PER_SPAN (128) are counted on the span as OTLP dropped_attributes_count, and checkpoints past MAX_EVENTS_PER_SPAN (128) as dropped_events_count. Neither adds to the span drop count: no record was lost. - Each limit (a full span buffer, attributes, checkpoints) logs its first hit since the last drain and only counts the rest, with a flag of its own, so one limit never silences another, even within one end(). A drain or a clear starts a new episode. - The attribute cap covers the span's own attributes and the app's correlation attributes; the app's are appended last, so they go first. The session's attributes, session.id, lk.outcome and error.type are added past the cap and always ship. Documented on Spans::end and Spans::add_event, and in SPEC. - Tests capture warnings with a test-only log::Log (installed once per process, lines filtered by the thread that logged them) and assert the warnings themselves, for these limits and the existing span and queue buffers.
| 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: |
There was a problem hiding this comment.
I'm sure you've considered this, but what's the downside of having this data model declared in protocol as shared protos?
Summary
Adds
livekit-telemetry, the shared client-telemetry core, starting with its data model and the OTLP encoding.Changes
lk.telemetry.reporteventScopeState(trace id, SDK and app attributes, route)opentelemetry-prototypesREADME.md,SPEC.md(resource attributes, events, log records, spans; the contract for the whole crate) and the crate changeset — the only PR the changeset check seesTemporary
Until #1483: a crate-level
#![allow(dead_code)]for internals the pipeline consumes, andScopeStatewithout subscribe tracking (#1483 completesscope.rs).Architecture
lk.outcomecarriesok | error | cancelled.getStats()report per peer connection) and polling cadence.TelemetryTransportforeign trait (Swift/Kotlin), a bounded pull queue (Dart, whose callbacks are isolate-bound: polled withtry_next()from a timer, no instruments passed, so Rust never calls into Dart), orNetTransportover the unmodifiedlivekit-netregistry.Verification
At
87f13dd2, from a clean checkout (CI's test workflow runs only for PRs intomain, so these were run locally; there is no clippy job in CI):cargo fmt -- --checkcargo clippy -p livekit-telemetry --all-targets --all-features -- -D warningscargo check -p livekit-telemetry --all-targets --no-default-featureswith features[],[uniffi]cargo test -p livekit-telemetry: 23 unit;--all-features: 23 unitcargo publish --dry-run -p livekit-telemetrycargo doc -D warningsreports one unresolved link, toTelemetry::stats, which #1483 adds.