Skip to content

feat(telemetry): add the process-wide pipeline and opt-out - #1485

Open
pblazej wants to merge 2 commits into
blaze/telemetry-stack/6-pipeline-testsfrom
blaze/telemetry-stack/7-process-pipeline
Open

pblazej wants to merge 2 commits into
blaze/telemetry-stack/6-pipeline-testsfrom
blaze/telemetry-stack/7-process-pipeline

Conversation

@pblazej

@pblazej pblazej commented Oct 1, 2026 •

Copy link
Copy Markdown
Contributor

Summary

global, the process-wide pipeline, and the opt-out. From here livekit-telemetry is exactly what #1396 ships. Every link below points at the stage that adds the test (the UniFFI ones at #1396).

Changes

  • global: one pipeline per process, installed once and reachable from anywhere, with the platform's instruments started on it
  • global::disable: the opt-out, in effect before it returns
  • device_tests.rs: the device contract, one test per row
  • The last temporary dead_code allowances removed
Device contract: every failure scenario → behaviour → max lost → why → test (32 rows)

What can go wrong on a device, and the most it can cost. OTel specifies no client persistence or device policy, so every row is custom. U = records not yet committed (≤ one flush interval, ≤ 2 048 queued) + open spans (≤ 256) + open RTC windows; B = one batch (≤ 512 records); E = records evicted or expired, always counted. N = records committed to the cache and not yet sent (≤ 4 MiB compressed, ≤ 512 batches, ≤ 24 h old). A batch is committed when the cache's push returns (file fsynced, renamed, directory fsynced on Unix); crash points are driven deterministically by test-only fault hooks around every write, fsync, rename, delete and publish.

Max lost: 0 = nothing; 0 until limits = nothing until the cache's bounds (4 MiB, 512 files, 24 h) evict, then E; U, B, E, N as defined above. Every row is custom (no OTel spec covers client persistence or device policy); Why names the platform behaviour, spec or reason behind the row.

Scenario Behaviour Max lost Why Test
App to background RTC windows close, everything committed and the whole cache flushed at once, whatever is left of the interval's allowance 0 once committed iOS suspends a backgrounded app within seconds (background execution limits); Android may freeze or kill a cached process entering_background_flushes_immediately
entering_the_background_flushes_even_with_the_allowance_spent
Killed / crashed (no callback) nothing waits for a termination callback; committed batches replay next launch U neither OS guarantees a termination callback (memory-pressure kills, crashes, force-quit) a_crash_loses_only_what_was_not_yet_cached
batch_is_written_before_upload_and_replayed_on_next_start
pre_connect_records_replay_after_a_restart
Killed after an answer, before the delete the batch is sent again at the next launch 0 lost; that batch duplicated an answer and a file delete cannot be atomic: at-least-once over at-most-once a_crash_between_the_answer_and_the_delete_resends_the_batch_once
Killed while a request is out the batch is sent again at the next launch 0 lost; ≤ 1 batch duplicated a request whose answer never arrived may have been ingested: at-least-once a_batch_in_flight_at_a_crash_is_sent_again_once
Killed mid-write .tmp + fsync + rename; a torn .tmp removed at open the uncommitted batch POSIX rename(2) is atomic; fsync(2) first, so a crash never exposes a torn file a_write_killed_half_way_leaves_no_partial_batch
file_cache_evicts_oldest_beyond_max_bytes_and_drops_stray_tmp
Killed at any step of a split (413) journaled; finished or rolled back at open, before any eviction, at the byte and at the file-count limit; a failure after the commit point is a success (halves published by the next listing); a reconfigure's second cache on the directory never breaks a split in progress (one process-wide split lock); binding a pre-connect batch to its project is the same journaled replace, rolled forward at open once its bound copy is journaled (only that copy names the project) 0 lost, 0 duplicated RFC 9110 §15.5.14: the body must shrink; journaling keeps the replace as crash-safe as a push a_split_crashed_at_every_step_loses_and_duplicates_nothing
a_split_failing_after_its_commit_point_is_kept_and_purgeable
a_second_cache_on_the_directory_never_breaks_a_split_in_progress
a_stuck_journal_is_retried_after_a_write_not_every_listing
a_crash_while_binding_still_replays_to_the_rooms_project
Corrupt / truncated cache file full gzip + CRC check before sending; dropped, counted corrupt that batch (≤ B) storage can truncate or damage files; the gzip CRC-32 (RFC 1952) detects it before upload corrupt_and_stray_cache_files_are_dropped_and_counted
File there but unreadable (file protection, permissions) kept, retried later; not corrupt 0 iOS Data Protection makes files unreadable while the device is locked: not corruption an_inaccessible_batch_is_kept_not_counted_corrupt
Delete fails after acceptance pending deletion: not sent again this launch (never rebound to a new id either), deleted when possible 0 lost; each such batch sent once more after a restart a delete that did not happen must not look like one that did; no resend within the launch a_failed_delete_after_acceptance_is_not_sent_again
an_accepted_batch_awaiting_its_delete_is_not_rebound_and_resent
Disk full, directory gone, fsync fails never claimed committed: the file is removed and the batch kept in memory (bounded by max_cache_bytes, oldest evicted and counted), still uploaded; cache.write_errors; a split that cannot commit keeps the parent spill evictions (E) while alive; the spill on a crash ENOSPC is ordinary on phones, iOS may purge Caches while the app runs, a failed fsync(2) leaves durability unknown a_full_disk_keeps_batches_in_memory_and_says_so
a_failed_directory_sync_keeps_the_batch_in_memory_not_half_on_disk
durability_failures_are_errors_not_silent_successes
Several processes sharing one storage dir not supported: each platform passes its own per-app dir; there is no cross-process lock (the split lock is per process) not tested; ≤ B per split in progress, uncounted (another process's recovery can delete its journal); each cached batch may be sent once per process app sandboxes give each app its own container, and an app's extensions must be given a directory of their own not tested
Cache full (4 MiB / 512 files) oldest evicted, counted (throttled during a server pause) E a bounded on-device footprint; the OS may purge cache directories anyway cache_eviction_is_counted_as_a_drop
file_cache_caps_the_number_of_batches
a_hold_longer_than_the_cache_reports_what_it_cost
24 h age expired at start and while running E stale telemetry loses its value; bounds how long a user's data stays on the device batches_older_than_a_day_are_dropped_and_counted_at_start
batches_expire_while_the_app_runs
Clock jump this launch's batches need the monotonic clock to agree 0 users and networks change the wall clock; the monotonic clock cannot jump a_clock_jump_does_not_expire_a_fresh_backlog
Offline hard hold: nothing attempted, no failure counted; checked before every request, also mid-pass; the soft cap is not scheduled meanwhile 0 until limits attempts while offline only cost battery and radio; the OS reports reachability device_holds_record_without_attempting
a_device_change_during_a_pass_stops_its_remaining_requests
a_soft_hold_turning_hard_does_not_spin
Low Data Mode, battery ≤ 10 % unplugged soft hold; one batch per 60 s on its own clock 0 until limits the user asked for less traffic (Low Data Mode, Data Saver); a nearly empty battery device_holds_record_without_attempting
the_soft_hold_cap_is_scheduled_and_splits_count_against_it
Thermal, memory, low power, background, CPU-limited encoder cadence ×2 … ×4; relief brings the pending tick forward at once; an unknown thermal / low-power reading is neutral and makes no change event 0 the OS asks apps to do less under thermal, memory and power pressure; WebRTC reports a CPU-limited encoder device_state_emits_change_events_and_stretches_cadence
cpu_limitation_counter_drives_cadence_pressure
relief_brings_the_pending_tick_forward
unknown_thermal_and_low_power_are_silent_and_neutral
Connect / reconnect soft hold 0 signaling and ICE/DTLS own the uplink while a call is set up: telemetry never competes with media uploads_hold_while_connecting_but_never_beyond_the_cap
typed_spans_hold_uploads_while_connecting_and_export_when_ended
A slow collector the exporter stays responsive (deadlines, commands) while a request is out 0 an actor blocked on I/O would miss its deadlines and commands deadlines_are_served_while_a_request_is_out
Shutdown drains within export_timeout_ms, then cancels the request on the wire and stops; the rest stays committed; the device state pushed at install always ships by flush + shutdown 0 committed an SDK must never hold up the app's teardown; bounded by export_timeout_ms the_initial_device_state_always_ships_by_flush_and_shutdown
shutdown_during_a_slow_backlog_is_bounded_and_joins_the_exporter
shutdown_drains_every_finished_span_batch
shutdown_leaves_a_session_summary
Queue overflow, event flood drop-oldest (2 048), 300 events / 10 min; counted exactly the dropped bounded memory, like the OTel batch processor's maxQueueSize; drop-oldest keeps the freshest context in_memory_bounds_count_every_drop
queue_overflow_is_counted_by_reason
flood_guard_caps_events_but_not_stats_windows
Notification storm (subscribes, device changes, tokens) wake-ups only re-evaluate; a new batch only at the tick, the queue threshold or backgrounding; exactly one request allowance per interval while a Room is in a call 0 the collector's per-project quota is shared by the fleet: wake-ups must not multiply requests notifications_within_an_interval_never_exceed_one_allowance
Backlog with no call not metered: once no Room is in a call (none connected, or all disconnected) each pass sends the whole cache 0 no media to protect: drain while the app is running the_allowance_applies_only_while_a_room_is_in_a_call
The core's own warnings and errors copied by the log forwarder (livekit* targets, never livekit_telemetry*) wherever the platform calls log_forward_bootstrap (Swift's default OSLogger with ffi: true does), whatever level it passes — that level filters only the console forwarding; platforms must not feed forwarded Rust entries to telemetry_log; any compact JWT (three base64url segments decoding to JSON, whatever the header's encoding, also inside punctuation such as (...<jwt>)) masked as <jwt>, bearer and token credential values (token=, "token": "…", quoted or spaced) as <redacted>; console entry unchanged — the core's own failures are otherwise invisible; Rust allows one global logger per process only_the_cores_own_warnings_and_errors_are_copied (added in #1396)
the_platforms_level_filters_the_console_not_the_telemetry_copy (added in #1396)
tokens_in_copied_records_are_masked (added in #1396)
the_console_entry_is_what_it_always_was (added in #1396) (lands in #1396)
Server-requested delay across a restart a Retry-After / RetryInfo pause is held in memory only: after a relaunch each project gets one request (once its token is handed over) before a new 429/503 pauses it again 0 lost (the throttled batch stays cached); ≤ 1 early request per project per relaunch; not tested RFC 9110 §10.2.3 and RFC 6585 §4 ask the client to wait; one probe per relaunch is the cost of not persisting the deadline not tested
Retries exhausted pause, never delete 0 until limits OTLP/HTTP: retryable failures are retried with backoff; giving up would lose data silently a_failing_server_never_costs_a_batch
Opt-out (also mid-upload, racing a capture or a configure) in effect when telemetry_disable() returns (no scope, no capture; a later configure only purges its storage dir and starts no exporter), the purge then runs on the core's runtime (a repeated opt-out joins the one already running); telemetry_is_disabled() is true process-wide before the call returns, so a platform with per-isolate state checks it before each getStats or submit; lifecycle serialized with instruments; every generation revoked under the locks that commit data, its exporter cancelled and awaited (draining ones too), queue/spans/windows/cache purged, later configures refused; only actual deletions counted; incomplete deletes reported N + U by design; ≤ B per generation already on the wire or handed to a pull-queue host withdrawing consent must take effect at once and erase unsent data (GDPR Art. 7(3), Art. 17) the_opt_out_is_in_effect_before_its_purge_runs
an_opt_out_counts_the_halves_of_an_unpublished_split
opting_out_stops_collection_and_purges_everything
opting_out_during_an_upload_resurrects_nothing
a_capture_racing_the_opt_out_leaves_nothing_behind
an_instrument_never_starts_after_the_opt_out_stopped_it
the_opt_out_cancels_a_draining_generation_before_returning
the_opt_out_withdraws_a_replaced_generations_pulled_request (added in #1396)
a_repeated_opt_out_waits_for_the_first_purge (added in #1396)
after_the_opt_out_nothing_is_captured
an_incomplete_purge_is_reported
purges_count_only_actual_deletions_once
clearing_an_unreadable_directory_fails
the_process_pipeline_no_ops_until_installed (lands in #1396)
Dart host stops serving the pull queue ≤ 1 request queued; a request nobody polls for times out as usual; requests the exporter gave up on, or pulled after the opt-out, are never served, by next() or try_next(); unknown ids ignored 0 served stale Dart FFI callbacks are bound to their isolate: Rust must never call into, or wake, a dead one try_next_serves_without_waiting_and_never_a_stale_request (added in #1396)
a_suspended_host_is_never_served_stale_requests (added in #1396)
a_cancelled_request_is_never_served (added in #1396)
finishing_discards_what_is_queued (added in #1396) (lands in #1396)
Invalid app input rejected, never truncated; counted invalid exactly the rejected a truncated id collides; attribute limits as in the OTel attribute-limits spec custom_data_is_room_scoped_validated_and_snapshotted
Correlation attributes or project change mid-window windows capture owner + attributes at open; a change closes them first, under the lock readings record under 0 a window describes one owner and one set of attributes; emitted records are immutable rtc_windows_carry_the_attributes_of_their_readings
a_window_closed_before_a_change_keeps_its_snapshot_whenever_it_is_queued
a_project_change_splits_the_rtc_windows_by_project
Subscribe never gets media timed_out on the core's clock, whoever looks first, kept through disconnect 0 no first media is the failure to measure: nothing else reports it the_subscribe_span_is_owned_by_the_core
late_first_media_is_timed_out_whoever_looks_first
a_subscribe_past_its_deadline_stays_timed_out_through_disconnect
Long-lived process spans hold the pipeline weakly; per-track state retired; span state ≤ 128 steps / attributes; retained caller strings bounded (an oversized span name is kept only as invalid, detached spans included) 0 apps run for days across many Rooms and tracks: nothing may grow without bound a_pending_subscribe_never_keeps_the_pipeline_alive
per_track_state_is_retired_with_the_track
the_same_track_in_two_sessions_never_merges
retained_span_state_is_bounded
final_limits_hold_after_decoration_and_caller_strings_are_bounded

Uninstall or an OS cache purge is beyond the SDK's reach (N + U). Defaults, polling (fast after a subscribe or a publish) and explicit flush (drains without the per-pass budget, within holds and pauses): defaults_export_and_window_once_a_minute
the_core_paces_stats_polling
publishing_polls_fast_until_the_first_outbound_reading
device_holds_uploads_and_a_backlog_replays_within_the_budget.

Verification

At 172d65f8, from a clean checkout (CI's test workflow runs only for PRs into main, so these were run locally; there is no clippy job in CI):

  • 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 doc -D warnings reports exactly what it reports on 6aba1b68 (private-item links, ExportError::from_response).

`global::install` makes one pipeline the process's and starts the
platform's instruments on it; `scope`, `emit`, `log`, `device_event`,
`set_device_state` and `stats` reach it from anywhere and no-op until it is
installed. `global::disable` is the opt-out: before it returns the
pipeline is removed, its instruments stopped and every generation revoked,
so nothing is captured, cached or sent again and later installs are
refused; the future it returns deletes everything still held.
One test per row of the device contract: crashes and kills at every
storage step, full or vanished storage, offline and constrained networks,
background and suspension, clock changes, device pressure, the opt-out at
every point of a capture and upload, and the lifecycle of the process
pipeline and its instruments.
@pblazej pblazej added the internal to tag changes that don't require changelog documentation label Oct 1, 2026
@pblazej
pblazej force-pushed the blaze/telemetry-stack/7-process-pipeline branch from bd2245f to 172d65f Compare 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

@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 2 potential issues.

2 flags not posted on this PR by your GitHub settings — view them in Devin Review. (Configure)

Devin Review

Comment on lines +140 to +144
let generations: Vec<Generation> =
GENERATIONS.lock().unwrap_or_else(|e| e.into_inner()).drain(..).collect();
generations
.into_iter()
.filter_map(|g| Some((g.shared.upgrade()?, g.commands)))

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.

🔴 Repeated opt-out bypasses pending purge

When disable() runs twice before the first purge finishes, the second call returns success without awaiting it. GENERATIONS was drained by the first call, so cached batches and exporters can remain active after the second returns.

Learn more

The opt-out has a synchronous phase that revokes all generations and an asynchronous phase that clears storage and stops exporters. The first disable() moves all entries out of GENERATIONS before its returned future is polled. A second disable() therefore returns a future with no work, even when the first future is pending. Awaiting that second future does not establish that the purge has finished.

Example: Call let first = disable() while a disk-backed pipeline has cached records; call disable().await before polling first. The second call returns true with the cache still on disk, rather than waiting for its deletion.

Recommended fix: Track the active purge as shared lifecycle state. Make subsequent disable() futures join its completion and return its final result; ensure this also works when a prior future is dropped.

Devin Review


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

Comment on lines +97 to +105
let previous = SHARED
.write()
.unwrap_or_else(|e| e.into_inner())
.replace(Installed { telemetry, instruments: instruments.clone() });
if let Some(previous) = &previous {
for instrument in &previous.instruments {
instrument.stop();
}
}

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.

🟡 Old instrument records reach new pipeline

When install() replaces a pipeline, an outgoing instrument's stop() can emit into the replacement. SHARED already points to the new pipeline, so shutdown records acquire the wrong session and destination.

Learn more

Global capture functions look up the currently installed telemetry through current. install switches that pointer before calling the outgoing instruments' stop methods. A stop method that forwards a final event using the global API therefore sends it to a different generation, whose project and session may differ.

Example: A device instrument calls global::device_event during stop() as project A is replaced with project B. Its final event goes into B's pipeline rather than A's.

Recommended fix: Stop the outgoing instruments while the outgoing pipeline is still current, then publish the replacement and start its instruments, retaining the lifecycle lock across that transition.

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