Skip to content

feat(uniffi): expose the telemetry core - #1396

Open
pblazej wants to merge 3 commits into
blaze/telemetry-stack/7-process-pipelinefrom
blaze/telemetry
Open

pblazej wants to merge 3 commits into
blaze/telemetry-stack/7-process-pipelinefrom
blaze/telemetry

Conversation

@pblazej

@pblazej pblazej commented Sep 3, 2026 •

Copy link
Copy Markdown
Contributor

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 state
  • log_forward.rs: the core's own warnings and errors copied into telemetry
  • Dart tests
  • The telemetry_ping example and the README's local testing guide
What a platform calls, and when
When Call
SDK init telemetry_configure(config, transport, instruments) — storage and tuning only; Dart: telemetry_configure_pulled(config, []) + queue.try_next()
Room created telemetry_scope()
Connect, and every token refresh scope.set_server(url, token) — cheap, idempotent
Room lifecycle scope.start(SpanName, parent) → step / set_track / end; subscribe_started / subscribed (tracks already in the room at join: subscribe_started at connect, so lk.subscribe measures from join; a lone subscribed still opens the span as a fallback) / subscribe_failed; track_ended(sid); set_room; disconnected(reason)
Stats every scope.stats_poll_interval_ms(): scope.record_peer_stats(report, {trackId → sid}, ts) per peer connection
OS signals telemetry_set_device_state (thermal Unknown / low power None without a source), telemetry_device_event, telemetry_log
App API (the only public surface) scope.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 each getStats or submit

No 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.

# 1. Local OTel backend (OTLP/HTTP on :4318, UI on :3000)
docker run -d --name lk-lgtm -p 3000:3000 -p 4318:4318 grafana/otel-lgtm

# 2. Record a small session (lk.connect with its checkpoints, lk.publish, lk.subscribe ended by
#    first media, a custom event) and export it to the local backend; it prints the trace id
LK_TELEMETRY_ENDPOINT=http://localhost:4318 cargo run -p telemetry_ping

# 3. Optional: the same through the Dart bindings (pull queue)
(cd livekit-uniffi && cargo make dart-package) && cd livekit-uniffi/packages/dart && dart pub get \
  && LK_TELEMETRY_ENDPOINT=http://localhost:4318 dart test test/live_collector_test.dart

# 4. Optional: the file cache — stop the backend, run step 2 with LK_TELEMETRY_DIR=/tmp/lk-telemetry,
#    start the backend again and run once more to watch the cached batches replay

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 set LK_TELEMETRY_ENDPOINT the same way (with livekit-server --dev for a real Room).

Verification

cargo test -p livekit-telemetry --all-features (178 unit, 2 doc), cargo test -p livekit-uniffi (23), clippy clean for livekit-telemetry; livekit-uniffi has 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 from 6aba1b68, the identical tree.)

Stack, at a9ec8d52 from a clean checkout:

  • cargo fmt -- --check
  • cargo clippy -p livekit-telemetry --all-targets --all-features -- -D warnings
  • cargo check -p livekit-telemetry --all-targets --no-default-features with features [], [net], [uniffi], [net,uniffi]
  • cargo test -p livekit-telemetry: 175 unit, 2 doc; --all-features: 178 unit, 2 doc
  • cargo build -p livekit-uniffi; cargo test -p livekit-uniffi: 23
  • cargo clippy -p livekit-uniffi --all-targets --no-deps: nothing in telemetry.rs, log_forward.rs or lib.rs (the crate's other findings, in data_track/data_stream tests, are on main too)
  • cargo build -p telemetry_ping; cargo clippy -p telemetry_ping --no-deps -- -D warnings
  • git diff --quiet 6aba1b68 a9ec8d52: identical trees (117f88c7)

Stages 1–7 pass the same livekit-telemetry gates on their own tips (see each PR).

`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.
@pblazej pblazej added the internal to tag changes that don't require changelog documentation label Oct 1, 2026
@pblazej pblazej changed the title Telemetry feat(uniffi): expose the telemetry core Oct 1, 2026
@pblazej
pblazej changed the base branch from main to blaze/telemetry-stack/7-process-pipeline October 1, 2026 13:53
@pblazej
pblazej added this pull request to stack #1486 October 1, 2026 14:23
@pblazej
pblazej marked this pull request as ready for review October 1, 2026 14:33
@pblazej
pblazej requested a review from ladvoc as a code owner October 1, 2026 14:33
@pblazej
pblazej requested review from 1egoman, davidliu, hiroshihorie, lukasIO and xianshijing-lk and removed request for ladvoc October 1, 2026 14:33

@devin-ai-integration devin-ai-integration Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Devin Review found 4 potential issues.

1 flag not posted on this PR by your GitHub settings — view it in Devin Review. (Configure)

Devin Review

Comment on lines +89 to +92
crate::runtime::runtime().spawn(exporter.run());
if let Some(previous) = previous {
crate::runtime::runtime().spawn(async move { previous.shutdown().await });
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔴 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.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

Comment on lines +516 to +523
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)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔴 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.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

Comment on lines +73 to +76
let queue = TelemetryExportQueue::new();
queue.process.store(true, Ordering::SeqCst);
install(livekit_telemetry::Telemetry::new(config, queue.clone()), instruments);
queue

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 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.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

Comment on lines +279 to +286
_ => (
skip,
value
.find(|c: char| {
c.is_whitespace() || matches!(c, '&' | '"' | '\'' | ',' | ';' | ')' | '<' | '>')
})
.unwrap_or(value.len()),
),

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟨 Incomplete credential masking in exported logs

When telemetry_copy captures a warning containing an unquoted credential with delimiter characters, mask_jwts redacts only its prefix. The remaining credential bytes enter telemetry.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

internal to tag changes that don't require changelog documentation

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant