Conversation
a951a3b to
6f17b80
Compare
`telemetry_configure` installs the process pipeline with a host-implemented `TelemetryTransport` or the `livekit-net` HTTP client, and `telemetry_configure_pulled` hands bindings with thread-bound callbacks (Dart) a bounded pull queue instead. `telemetry_scope` gives each Room its `TelemetryScope`: `set_server` at connect and on every token refresh, spans, the subscribe lifecycle, `record_peer_stats` paced by `stats_poll_interval_ms`, `track_ended`, `disconnected`, and the app API (`emit_custom`, `set_attribute`). `telemetry_disable` is the synchronous opt-out and `telemetry_is_disabled` reads it process-wide; device state, device events and log records come from the host. Dart tests cover the pull queue, and macOS links reserve header padding for Flutter's install name rewrite.
Once a platform installs the log forwarder, warnings and errors from `livekit*` targets also reach the telemetry pipeline as log records, with any compact JWT and bearer or token credential values masked; the telemetry crate's own records are never copied. The platform's level still filters only what is forwarded to it, so the console output is unchanged.
`telemetry_ping` records a small session (an `lk.connect` span with its checkpoints, an `lk.publish`, an `lk.subscribe` ended by first media, a custom event) and exports it to the collector `LK_TELEMETRY_ENDPOINT` names; the README shows how to run it against a local OpenTelemetry backend, and a Dart test does the same through the pull queue.
6aba1b6 to
a9ec8d5
Compare
There was a problem hiding this comment.
Devin Review found 4 potential issues.
1 flag not posted on this PR by your GitHub settings — view it in Devin Review. (Configure)
| crate::runtime::runtime().spawn(exporter.run()); | ||
| if let Some(previous) = previous { | ||
| crate::runtime::runtime().spawn(async move { previous.shutdown().await }); | ||
| } |
There was a problem hiding this comment.
🔴 Reconfiguration duplicates cached telemetry batches
Reconfiguring with the same storage directory starts a new exporter before the old one finishes draining. Both exporters can upload the same cached batch, duplicating logs and spans.
Learn more
A pipeline writes encoded batches into its configured on-disk cache before uploading. Reconfiguration creates another pipeline using that directory, starts its exporter, then asynchronously shuts down the previous pipeline. Exporter::begin_round reads the directory independently for each exporter, so an old in-flight batch is also eligible for the new exporter. Both requests can reach the collector before either exporter removes the batch.
Example: The old exporter posts batch A and waits for the HTTP response. Reconfigure using the same storage directory. The new exporter reads A and posts it too; both requests succeed, yielding two copies of A.
Recommended fix: Serialize replacement so the previous exporter stops before the new exporter begins replaying the shared cache. Preserve records captured during the handover, and test replacement with an old request held in flight.
Was this helpful? React with 👍 or 👎 to provide feedback.
| let opted_out = self.process.load(Ordering::SeqCst) && global::is_disabled(); | ||
| if opted_out || self.finished.load(Ordering::SeqCst) || pending.done.is_closed() { | ||
| return None; | ||
| } | ||
| let mut inflight = self.inflight.lock().unwrap_or_else(|e| e.into_inner()); | ||
| inflight.retain(|_, done| !done.is_closed()); | ||
| inflight.insert(pending.export.id, pending.done); | ||
| Some(pending.export) |
There was a problem hiding this comment.
🔴 Finished queue still serves exports
If finish runs after serve checks the flag, serve can return a request after finish completes. The host can upload an export it already cancelled.
Learn more
A host gets requests through next or try_next, and calls finish to stop serving them. finish sets finished and clears inflight, but the flag test here is not synchronized with insertion into inflight. A concurrent serve can pass the test, wait for the mutex, then insert and return after finish has cleared the map. The host receives a cancelled request and can still perform its HTTP call.
Example: Thread A passes the finished check for request 7. Thread B calls finish and returns. Thread A inserts request 7 and next returns it for upload, instead of None.
Recommended fix: Synchronize the finished check and in-flight insertion with finish under one lock or equivalent lifecycle guard. Test the cancellation race with a blocked serve between the flag check and insertion.
Was this helpful? React with 👍 or 👎 to provide feedback.
| let queue = TelemetryExportQueue::new(); | ||
| queue.process.store(true, Ordering::SeqCst); | ||
| install(livekit_telemetry::Telemetry::new(config, queue.clone()), instruments); | ||
| queue |
There was a problem hiding this comment.
🟡 Disabled pipeline leaves pull queue waiting
After telemetry_disable, telemetry_configure_pulled returns an open queue without starting an exporter. A host awaiting next blocks indefinitely unless it calls finish itself.
Learn more
The process-wide opt-out is permanent. global::install refuses a pipeline after it, and the local install skips spawning its exporter. The queue still holds its sender, so its next method waits for a request that no exporter can generate. try_next returns None, but cannot tell an empty live queue from this permanently inactive one.
Example: A client calls telemetry_disable, then initializes an SDK using telemetry_configure_pulled and starts an async next loop. The first next never returns, rather than ending the loop.
Recommended fix: Mark the returned queue finished when installation is refused. Expose installation success to telemetry_configure_pulled so it can close this queue without affecting a live one.
Was this helpful? React with 👍 or 👎 to provide feedback.
| _ => ( | ||
| skip, | ||
| value | ||
| .find(|c: char| { | ||
| c.is_whitespace() || matches!(c, '&' | '"' | '\'' | ',' | ';' | ')' | '<' | '>') | ||
| }) | ||
| .unwrap_or(value.len()), | ||
| ), |
There was a problem hiding this comment.
Summary
livekit-telemetry— the shared client-telemetry core — and its UniFFI surface for the Swift, Kotlin and Dart SDKs. Wire spec (events, attributes, cadence, upload policy):livekit-telemetry/SPEC.md.Changes
livekit-uniffi/src/telemetry.rs: configure and configure-pulled, scopes, spans, stats, the app API, the opt-out, device statelog_forward.rs: the core's own warnings and errors copied into telemetrytelemetry_pingexample and the README's local testing guideWhat a platform calls, and when
telemetry_configure(config, transport, instruments)— storage and tuning only; Dart:telemetry_configure_pulled(config, [])+queue.try_next()telemetry_scope()scope.set_server(url, token)— cheap, idempotentscope.start(SpanName, parent)→step/set_track/end;subscribe_started/subscribed(tracks already in the room at join:subscribe_startedat connect, solk.subscribemeasures from join; a lonesubscribedstill opens the span as a fallback) /subscribe_failed;track_ended(sid);set_room;disconnected(reason)scope.stats_poll_interval_ms():scope.record_peer_stats(report, {trackId → sid}, ts)per peer connectiontelemetry_set_device_state(thermalUnknown/ low powerNonewithout a source),telemetry_device_event,telemetry_logscope.emit_custom(name, {k: v}),scope.set_attribute(key, value?),telemetry_disable()— synchronous, in effect when it returns (TODO on platforms: token discussion);telemetry_is_disabled()reads that opt-out process-wide (every thread and Dart isolate), for platforms to check right before eachgetStatsor submitNo endpoint, header or sink exists in any client-language API. Local end-to-end tests set
LK_TELEMETRY_ENDPOINT(read by the core; see the README).Local testing against a local OpenTelemetry backend (LGTM)
Everything can be tried against a local OpenTelemetry backend (LGTM) in a couple of minutes; no LiveKit server is needed for the Rust example.
Then open http://localhost:3000 → Explore: Loki
{service_name="telemetry_ping"}for the log records (otel.event.name:custom.ping,lk.rtc.stats.sample,lk.telemetry.report), Tempo with the printed trace id for the spans (lk.connect,lk.publish,lk.subscribe). Platform e2e tests setLK_TELEMETRY_ENDPOINTthe same way (withlivekit-server --devfor a real Room).Verification
cargo test -p livekit-telemetry --all-features(178 unit, 2 doc),cargo test -p livekit-uniffi(23), clippy clean forlivekit-telemetry;livekit-uniffihas no findings in the files this PR touches (pre-existing ones elsewhere, see Verification); Swift (SPM_PLATFORMS="macos ios" cargo make swift-package-debug), Kotlin (cargo make android-package-local, generated Kotlin compiled) and Dart (cargo make dart-package+dart test) packages built from this branch. (Swift, Kotlin and Dart packages built from6aba1b68, the identical tree.)Stack, at
a9ec8d52from a clean checkout:cargo fmt -- --checkcargo clippy -p livekit-telemetry --all-targets --all-features -- -D warningscargo check -p livekit-telemetry --all-targets --no-default-featureswith features[],[net],[uniffi],[net,uniffi]cargo test -p livekit-telemetry: 175 unit, 2 doc;--all-features: 178 unit, 2 doccargo build -p livekit-uniffi;cargo test -p livekit-uniffi: 23cargo clippy -p livekit-uniffi --all-targets --no-deps: nothing intelemetry.rs,log_forward.rsorlib.rs(the crate's other findings, indata_track/data_streamtests, are onmaintoo)cargo build -p telemetry_ping;cargo clippy -p telemetry_ping --no-deps -- -D warningsgit diff --quiet 6aba1b68 a9ec8d52: identical trees (117f88c7)Stages 1–7 pass the same
livekit-telemetrygates on their own tips (see each PR).