From d65b3152f499cc25d3262189aa8e84edd2dd1698 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?B=C5=82az=CC=87ej=20Pankowski?= <86720177+pblazej@users.noreply.github.com> Date: Thu, 1 Oct 2026 13:19:41 +0200 Subject: [PATCH 1/8] feat(telemetry): add UniFFI bridge and example setup Exposes the telemetry core's public API through UniFFI. Example app displays the Rust core build version as an FFI smoke test. --- example/lib/main.dart | 3 ++ example/lib/pages/connect.dart | 8 ++++ example/lib/utils.dart | 18 ++++++++ example/pubspec.yaml | 1 + lib/src/uniffi/uniffi.dart | 48 +++++++++++++++++++++ lib/src/uniffi/uniffi_io.dart | 79 ++++++++++++++++++++++++++++++++++ lib/src/uniffi/uniffi_web.dart | 27 ++++++++++++ pubspec.lock | 24 +++++++++++ pubspec.yaml | 4 ++ 9 files changed, 212 insertions(+) create mode 100644 lib/src/uniffi/uniffi.dart create mode 100644 lib/src/uniffi/uniffi_io.dart create mode 100644 lib/src/uniffi/uniffi_web.dart diff --git a/example/lib/main.dart b/example/lib/main.dart index fd056f6d5..dff140022 100644 --- a/example/lib/main.dart +++ b/example/lib/main.dart @@ -4,6 +4,7 @@ import 'package:livekit_example/theme.dart'; import 'package:logging/logging.dart'; import 'package:intl/intl.dart'; import 'pages/connect.dart'; +import 'utils.dart'; void main() async { final format = DateFormat('HH:mm:ss'); @@ -31,6 +32,8 @@ void main() async { zeroPlayoutDelay: true, enableWARP: true, ); + // Smoke test the Rust core with one synchronous FFI call at startup. + Logger('LiveKitExample').info(rustCoreVersionLabel()); runApp(const LiveKitExampleApp()); } diff --git a/example/lib/pages/connect.dart b/example/lib/pages/connect.dart index 74ef623b1..0b3f0e476 100644 --- a/example/lib/pages/connect.dart +++ b/example/lib/pages/connect.dart @@ -12,6 +12,7 @@ import 'package:permission_handler/permission_handler.dart'; import 'package:shared_preferences/shared_preferences.dart'; import '../exts.dart'; +import '../utils.dart'; enum _ConnectOption { autoSubscribe, e2ee } @@ -436,6 +437,13 @@ class _ConnectIntro extends StatelessWidget { context, ).textTheme.titleMedium?.copyWith(color: LKColors.textSecondary), ), + const SizedBox(height: 4), + Text( + rustCoreVersionLabel(), + style: Theme.of( + context, + ).textTheme.titleSmall?.copyWith(color: LKColors.textSecondary), + ), ], ); } diff --git a/example/lib/utils.dart b/example/lib/utils.dart index fc74277dc..2142b23cf 100644 --- a/example/lib/utils.dart +++ b/example/lib/utils.dart @@ -1,3 +1,21 @@ +// The Rust core facade is experimental, the example is its smoke test. +// ignore_for_file: experimental_member_use + import 'dart:async'; +import 'package:livekit_client/livekit_client.dart'; FutureOr Function()? onWindowShouldClose; + +/// One synchronous call into the Rust core through the uniffi facade. The +/// example shows the result at startup and on the connect page as a smoke +/// test that the native library was bundled and loads on this platform. +String rustCoreVersionLabel() { + if (!LiveKitUniffi.isAvailable) { + return 'Rust core not available on this platform'; + } + try { + return 'Rust core ${LiveKitUniffi.buildVersion}'; + } catch (error) { + return 'Rust core failed to load: $error'; + } +} diff --git a/example/pubspec.yaml b/example/pubspec.yaml index 38fd761a1..e6f089a44 100644 --- a/example/pubspec.yaml +++ b/example/pubspec.yaml @@ -25,6 +25,7 @@ dependencies: livekit_client: path: ../ + dev_dependencies: flutter_test: sdk: flutter diff --git a/lib/src/uniffi/uniffi.dart b/lib/src/uniffi/uniffi.dart new file mode 100644 index 000000000..c25367087 --- /dev/null +++ b/lib/src/uniffi/uniffi.dart @@ -0,0 +1,48 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +import 'package:meta/meta.dart'; + +import 'uniffi_io.dart' if (dart.library.js_interop) 'uniffi_web.dart' as impl; + +/// Facade over the Rust core exposed by the `livekit_uniffi` package. +/// +/// `livekit_uniffi` reaches Rust through Dart's Native Assets: its build hook +/// bundles a `cdylib` into the host app and the generated bindings call into it +/// with `@Native`. None of that exists on the web, where there is no dynamic +/// library to load, so every entry point here is split native/web through the +/// same conditional-import pattern the rest of the SDK uses (see +/// `support/platform.dart`). Web builds must never reach the generated +/// bindings -- importing them at all would break `dart compile js`/`wasm`. +/// +/// Callers get [isAvailable] to branch on, and platform-specific code paths +/// stay out of the public API surface. +@experimental +abstract final class LiveKitUniffi { + /// Whether the Rust core can be called on this platform. + /// + /// False on web. Every other member throws [UnsupportedError] when this is + /// false, rather than returning a silently wrong value. + static bool get isAvailable => impl.isAvailable; + + /// Version string reported by the Rust core. + /// + /// The simplest possible round trip -- a synchronous, argument-free call + /// returning a string -- so it doubles as the smoke test that the whole + /// chain is wired up: build hook resolved the library, `@Native` bound the + /// symbol, and a value came back across the FFI boundary. + /// + /// Throws [UnsupportedError] on web. + static String get buildVersion => impl.buildVersion(); +} diff --git a/lib/src/uniffi/uniffi_io.dart b/lib/src/uniffi/uniffi_io.dart new file mode 100644 index 000000000..64a9cad87 --- /dev/null +++ b/lib/src/uniffi/uniffi_io.dart @@ -0,0 +1,79 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +import 'package:livekit_uniffi/livekit_uniffi.dart' as uniffi; + +// The telemetry integration (`telemetry/telemetry_io.dart`) reaches the +// bindings through this file, so this stays the SDK's single import site. +export 'package:livekit_uniffi/livekit_telemetry.dart' + show + AnswerSentSpanStep, + AppState, + AttemptSpanStep, + AttributeValue, + AudioOutput, + AudioRouteChangedDeviceEvent, + AudioRouteReason, + BoolAttributeValue, + CaptureDevice, + CaptureFailedDeviceEvent, + CaptureFailure, + ConnectSpanName, + DeviceState, + DisconnectReason, + DoubleAttributeValue, + EngineSpanStep, + ExportRequest, + ExportResponse, + IntAttributeValue, + JoinRecvSpanStep, + LogRecord, + LogSource, + MemoryPressure, + NetworkType, + OfferSentSpanStep, + PcConnectedSpanStep, + PcCreatedSpanStep, + PublishSpanName, + ReconnectReason, + ReconnectSpanName, + RejectedExportException, + RetryableExportException, + RoomConnectedSpanStep, + RoomIdentity, + RtcStat, + Sdk, + Severity, + SignalSpanStep, + SpanName, + SpanOutcome, + SpanStep, + SpanTrack, + StrAttributeValue, + TelemetryConfig, + TelemetryInstrument, + TelemetryResource, + ThermalState, + TrackKind, + TrackSource, + WsOpenSpanStep; +export 'package:livekit_uniffi/livekit_uniffi.dart'; + +/// Native implementation of [LiveKitUniffi]. See `uniffi.dart`. +/// +/// This is the only file in the SDK that may import the generated bindings: +/// the conditional import in `uniffi.dart` keeps it out of web builds. +const bool isAvailable = true; + +String buildVersion() => uniffi.buildVersion(); diff --git a/lib/src/uniffi/uniffi_web.dart b/lib/src/uniffi/uniffi_web.dart new file mode 100644 index 000000000..db2ab1968 --- /dev/null +++ b/lib/src/uniffi/uniffi_web.dart @@ -0,0 +1,27 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +/// Web implementation of [LiveKitUniffi]. See `uniffi.dart`. +/// +/// Native Assets bundles a `cdylib`, which the web has no way to load, so the +/// Rust core is simply absent here. This file deliberately does not import +/// `package:livekit_uniffi/...` -- doing so would pull `dart:ffi` into a web +/// compile and fail the build. +const bool isAvailable = false; + +Never buildVersion() => throw UnsupportedError( + 'LiveKitUniffi.buildVersion is not available on web: the Rust core is ' + 'delivered as a native library. Guard calls with ' + 'LiveKitUniffi.isAvailable.', +); diff --git a/pubspec.lock b/pubspec.lock index 2ef21c52c..d9570a489 100644 --- a/pubspec.lock +++ b/pubspec.lock @@ -17,6 +17,14 @@ packages: url: "https://pub.dev" source: hosted version: "13.0.0" + archive: + dependency: transitive + description: + name: archive + sha256: "6c5bcd986e06b94e3c40244af471750840a3d2341d1f9763a1100a14add517b4" + url: "https://pub.dev" + source: hosted + version: "4.3.0" args: dependency: transitive description: @@ -432,6 +440,14 @@ packages: url: "https://pub.dev" source: hosted version: "6.1.0" + livekit_uniffi: + dependency: "direct main" + description: + name: livekit_uniffi + sha256: "2771936d7a799410ab9f9db29ae596648df36322222d5c928f3fa4612b8edce4" + url: "https://pub.dev" + source: hosted + version: "0.1.12" logger: dependency: transitive description: @@ -616,6 +632,14 @@ packages: url: "https://pub.dev" source: hosted version: "1.5.2" + posix: + dependency: transitive + description: + name: posix + sha256: bc1bad54ad2b735816e31f8d4600cfde6c7839975085ddfbca48b6c9f7c4044e + url: "https://pub.dev" + source: hosted + version: "6.5.2" protobuf: dependency: "direct main" description: diff --git a/pubspec.yaml b/pubspec.yaml index 28ef2bc60..1ef4289a7 100644 --- a/pubspec.yaml +++ b/pubspec.yaml @@ -57,6 +57,10 @@ dependencies: flutter_webrtc: 1.6.2+hotfix.3 dart_webrtc: ^1.8.0 + # Rust core (livekit-uniffi), delivered as a bundled cdylib via Native Assets. + # Native platforms only, see lib/src/uniffi/. + livekit_uniffi: ^0.1.12 + dev_dependencies: flutter_test: sdk: flutter From 5996ab4dd0ba4495b2865a7ed744fc011e2ebbcc Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?B=C5=82az=CC=87ej=20Pankowski?= <86720177+pblazej@users.noreply.github.com> Date: Thu, 1 Oct 2026 13:19:41 +0200 Subject: [PATCH 2/8] feat(telemetry): add core OTLP pipeline Process-wide export queue with device state monitoring and RTC stats collection. Serves OTLP/HTTP pulled export with persistent storage and fail-open behavior. --- lib/src/logger.dart | 4 +- lib/src/support/sdk_logger.dart | 97 ++++ lib/src/telemetry/telemetry.dart | 95 ++++ lib/src/telemetry/telemetry_io.dart | 813 +++++++++++++++++++++++++++ lib/src/telemetry/telemetry_web.dart | 25 + 5 files changed, 1033 insertions(+), 1 deletion(-) create mode 100644 lib/src/support/sdk_logger.dart create mode 100644 lib/src/telemetry/telemetry.dart create mode 100644 lib/src/telemetry/telemetry_io.dart create mode 100644 lib/src/telemetry/telemetry_web.dart diff --git a/lib/src/logger.dart b/lib/src/logger.dart index d235a55d3..7468935a2 100644 --- a/lib/src/logger.dart +++ b/lib/src/logger.dart @@ -14,6 +14,8 @@ import 'package:logging/logging.dart'; +import 'support/sdk_logger.dart'; + enum LoggerLevel { kALL, kFINEST, @@ -27,7 +29,7 @@ enum LoggerLevel { kOFF, } -final logger = Logger('livekit'); +final Logger logger = SdkLogger(Logger('livekit')); /// disable logging void disableLogging() { diff --git a/lib/src/support/sdk_logger.dart b/lib/src/support/sdk_logger.dart new file mode 100644 index 000000000..6749dce5d --- /dev/null +++ b/lib/src/support/sdk_logger.dart @@ -0,0 +1,97 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +import 'dart:async'; + +import 'package:logging/logging.dart'; + +/// Receives the SDK warnings and errors the logger's level filters out (telemetry): those it +/// emits reach its `onRecord` listeners as usual. +void Function(LogRecord record)? sdkFilteredWarningCapture; + +/// The SDK's `logger`: the `livekit` [Logger] itself for everything the app sees (level, records, +/// listeners, how a message is evaluated), plus [sdkFilteredWarningCapture] for a warning or +/// error its level drops, so the console level never decides what telemetry gets. +class SdkLogger implements Logger { + SdkLogger(this._logger); + + final Logger _logger; + + @override + void log(Level logLevel, Object? message, [Object? error, StackTrace? stackTrace, Zone? zone]) { + final capture = logLevel >= Level.WARNING ? sdkFilteredWarningCapture : null; + if (capture == null || _logger.isLoggable(logLevel)) { + return _logger.log(logLevel, message, error, stackTrace, zone); // exactly as without telemetry + } + // Filtered out: the console sees nothing, telemetry gets the record built the way the logger + // would have. A message that fails to evaluate is dropped, never thrown to the caller. + try { + if (message is Function) message = (message as Object? Function())(); + final text = message is String ? message : message.toString(); + capture( + LogRecord( + logLevel, + text, + fullName, + error, + stackTrace, + zone ?? Zone.current, + message is String ? null : message, + ), + ); + } catch (_) {} + } + + @override + void finest(Object? message, [Object? error, StackTrace? stackTrace]) => + log(Level.FINEST, message, error, stackTrace); + @override + void finer(Object? message, [Object? error, StackTrace? stackTrace]) => log(Level.FINER, message, error, stackTrace); + @override + void fine(Object? message, [Object? error, StackTrace? stackTrace]) => log(Level.FINE, message, error, stackTrace); + @override + void config(Object? message, [Object? error, StackTrace? stackTrace]) => + log(Level.CONFIG, message, error, stackTrace); + @override + void info(Object? message, [Object? error, StackTrace? stackTrace]) => log(Level.INFO, message, error, stackTrace); + @override + void warning(Object? message, [Object? error, StackTrace? stackTrace]) => + log(Level.WARNING, message, error, stackTrace); + @override + void severe(Object? message, [Object? error, StackTrace? stackTrace]) => + log(Level.SEVERE, message, error, stackTrace); + @override + void shout(Object? message, [Object? error, StackTrace? stackTrace]) => log(Level.SHOUT, message, error, stackTrace); + + @override + String get name => _logger.name; + @override + String get fullName => _logger.fullName; + @override + Logger? get parent => _logger.parent; + @override + Map get children => _logger.children; + @override + Level get level => _logger.level; + @override + set level(Level? value) => _logger.level = value; + @override + Stream get onLevelChanged => _logger.onLevelChanged; + @override + Stream get onRecord => _logger.onRecord; + @override + void clearListeners() => _logger.clearListeners(); + @override + bool isLoggable(Level value) => _logger.isLoggable(value); +} diff --git a/lib/src/telemetry/telemetry.dart b/lib/src/telemetry/telemetry.dart new file mode 100644 index 000000000..7174bda93 --- /dev/null +++ b/lib/src/telemetry/telemetry.dart @@ -0,0 +1,95 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +import 'dart:async'; + +import '../core/room.dart'; +import '../proto/livekit_models.pb.dart' as lk_models; +import '../publication/remote.dart'; +import '../publication/track_publication.dart'; +import '../track/local/local.dart'; +import '../track/options.dart'; +import '../types/internal.dart'; +import 'telemetry_io.dart' if (dart.library.js_interop) 'telemetry_web.dart' as impl; + +// Client telemetry, internal to the SDK. The pipeline — destination, token, batching, retries, +// cache, holds, stats mapping, span state — lives in the Rust core, one per process; Dart installs +// it with the first Room, feeds it OS signals and moves its bytes. A compiled no-op on web. + +/// Checkpoints of the `lk.connect` span, in the core's vocabulary. +enum ConnectStep { wsOpen, signal, joinRecv, pcCreated, offerSent, answerSent, engine, pcConnected, roomConnected } + +/// One Room's telemetry session: its trace, spans, RTC instrument and destination. +abstract interface class RoomTelemetry { + /// A new Room's session, installing the process pipeline first if this is the first Room; + /// null on web and after [disable]. + static RoomTelemetry? create(Room room) => impl.createRoomTelemetry(room); + + /// Opt-out for the rest of the process; see `LiveKitClient.disableTelemetry`. + static Future disable() => impl.disableTelemetry(); + + /// A capture device failed to start (getUserMedia / getDisplayMedia threw). + static void captureFailed(LocalTrackOptions options, Object error) => impl.captureFailed(options, error); + + /// One `Room.connect` attempt as the `lk.connect` span; hands the core the server and token. + Future connect(String url, String token, Future Function() body); + + /// A checkpoint of the open `lk.connect` span, if any. + void step(ConnectStep step); + + /// One reconnect cycle as the `lk.reconnect` span; attempts are its checkpoints. + TraceSpan? reconnect(ClientDisconnectReason reason, lk_models.ReconnectReason? protoReason); + + /// One publish attempt as the `lk.publish` span, under the ambient span or the open connect span. + Future publish(LocalTrack track, Future Function() body); + + /// Intent to subscribe to a track manually: the core opens its `lk.subscribe` span. + void subscribeStarted(RemoteTrackPublication publication); + + /// See `Room.emitTelemetryEvent`. + void emitCustom(String name, Map attributes); + + /// See `Room.setTelemetryAttribute`. + void setAttribute(String key, String? value); +} + +/// A span the SDK drives across several steps (a reconnect cycle). +abstract interface class TraceSpan { + /// `attempt N quick|full`. + void attempt(int number, {required bool full}); + + void end(); + + /// End on an error; its type becomes `error.type`. + void fail(Object error); + + void cancel(); +} + +/// Zone keys of the ambient span and Room: warn/error records logged inside point at that span or +/// land in that Room's session (`LogRecord.zone`); a publish nests under the ambient span. +const Symbol telemetrySpanKey = #livekitTelemetrySpan; +const Symbol telemetryRoomKey = #livekitTelemetryRoom; + +extension TraceSpanZone on TraceSpan? { + /// Run [body] with this span as the ambient one (a no-op when null). + Future run(Future Function() body) => + this == null ? body() : runZoned(body, zoneValues: {telemetrySpanKey: this}); +} + +extension RoomTelemetryZone on RoomTelemetry? { + /// Run [body] with this Room as the ambient one (a no-op when null); stream listeners + /// subscribed inside keep running in it. + T run(T Function() body) => this == null ? body() : runZoned(body, zoneValues: {telemetryRoomKey: this}); +} diff --git a/lib/src/telemetry/telemetry_io.dart b/lib/src/telemetry/telemetry_io.dart new file mode 100644 index 000000000..6d426ec37 --- /dev/null +++ b/lib/src/telemetry/telemetry_io.dart @@ -0,0 +1,813 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +import 'dart:async'; +import 'dart:io'; + +import 'package:flutter/widgets.dart' show AppLifecycleState, WidgetsBinding, WidgetsBindingObserver; + +import 'package:connectivity_plus/connectivity_plus.dart'; +import 'package:http/http.dart' as http; +import 'package:logging/logging.dart'; +import 'package:meta/meta.dart'; +import 'package:path/path.dart' as p; + +import '../core/room.dart'; +import '../events.dart'; +import '../extensions.dart'; +import '../hardware/hardware.dart'; +import '../internal/events.dart'; +import '../livekit.dart'; +import '../logger.dart'; +import '../participant/participant.dart'; +import '../proto/livekit_models.pb.dart' as lk_models; +import '../publication/remote.dart'; +import '../publication/track_publication.dart'; +import '../support/platform.dart'; +import '../support/sdk_logger.dart'; +import '../track/local/local.dart'; +import '../track/options.dart'; +import '../types/internal.dart'; +import '../types/other.dart'; +import '../uniffi/uniffi_io.dart' as ffi; +import 'telemetry.dart'; + +// MARK: - Pipeline + +bool _installed = false; + +/// Opted out: Dart collects nothing more, whatever Room is still connected. +bool _disabled = false; + +/// Telemetry never fails the SDK or the app: a core error (a Rust panic surfaces as an exception) +/// is logged and dropped, never thrown into a connect, a publish, a teardown or a log call. +T? _quiet(T Function() body) { + try { + return body(); + } catch (error) { + logger.fine('telemetry: $error'); + return null; + } +} + +/// Stops serving the installed pipeline's export queue. +void Function()? _stopServing; + +/// Started and stopped from Dart, never handed to the core: a callback the core holds would +/// point into this isolate after it is gone (hot restart, a recreated engine). +final _instruments = [_DeviceTelemetry(), _LogCapture()]; + +/// Where a new Room's session comes from: a test can stand in a core that fails. +@visibleForTesting +ffi.TelemetryScope? Function() telemetryScopeFactory = ffi.telemetryScope; + +RoomTelemetry? createRoomTelemetry(Room room) { + try { + if (!_installed) { + _installed = true; + _install(); + } + final scope = telemetryScopeFactory(); + return scope == null ? null : _RoomTelemetry(room, scope); + } catch (error) { + // Fail-open: the app runs without telemetry rather than not at all (e.g. no native library). + logger.fine('telemetry unavailable: $error'); + return null; + } +} + +/// Installed synchronously with the first Room, so a Room never misses its session. Every tuning +/// value is the core's default. +void _install() { + // The core creates the cache directory but not its parents (`~/.cache` may not exist yet). + var storageDir = storageDirectory; + try { + if (storageDir != null) Directory(storageDir).parent.createSync(recursive: true); + } on FileSystemException catch (_) { + storageDir = null; // batches in memory only + } + final queue = ffi.telemetryConfigurePulled( + config: ffi.TelemetryConfig( + sdk: ffi.TelemetryResource( + sdk: ffi.Sdk.flutter, + sdkVersion: LiveKitClient.version, + osName: Platform.operatingSystem, + osVersion: Platform.operatingSystemVersion, + ), + storageDir: storageDir, + ), + instruments: const [], + ); + if (ffi.telemetryStats() == null) return; // refused: this process opted out + _stopServing = _serve(queue); + for (final instrument in _instruments) { + instrument.start(); + } +} + +/// The app's own temporary directory on mobile (purgeable, never backed up); on desktop a per-user +/// directory named after the app: the user's cache directory on Linux, where the temporary one is +/// shared by every user. Null (batches in memory only) on Linux without an absolute cache home. +@visibleForTesting +final String? storageDirectory = () { + final base = Platform.isLinux ? _linuxCacheHome() : Directory.systemTemp.path; + return base == null + ? null + : p.join(base, 'livekit-telemetry-${p.basenameWithoutExtension(Platform.resolvedExecutable)}'); +}(); + +/// `$XDG_CACHE_HOME`, else `$HOME/.cache`; relative or empty values are ignored (XDG Base Directory). +String? _linuxCacheHome() { + final xdg = Platform.environment['XDG_CACHE_HOME']; + if (xdg != null && p.isAbsolute(xdg)) return xdg; + final home = Platform.environment['HOME']; + return home != null && p.isAbsolute(home) ? p.join(home, '.cache') : null; +} + +Future disableTelemetry() async { + // Nothing installs after an opt-out; before the first Room, the cache a previous launch left is + // this SDK's to delete (the core purges only a pipeline it installed). + final installed = _installed; + _installed = true; + _stopCapture(); + // In effect now; what was not sent is purged in the background. Not awaited through + // `telemetryFlush`: a pending Rust future would call into this isolate even after it is gone. + _quiet(ffi.telemetryDisable); + final storageDir = storageDirectory; + if (installed || storageDir == null) return; + try { + await Directory(storageDir).delete(recursive: true); + } on FileSystemException catch (_) {} // nothing cached +} + +/// Stops everything this isolate collects and serves. +void _stopCapture() { + _disabled = true; + // Each step on its own: one that throws never skips the rest. + for (final instrument in _instruments) { + _quiet(instrument.stop); + } + for (final session in _inCall) { + session._stats?.cancel(); + } + final stopServing = _stopServing; + _stopServing = null; + _quiet(() => stopServing?.call()); +} + +/// Opted out, here or in any other isolate of the process (the core keeps the process-wide flag); +/// a core that cannot answer counts as opted out. Learning it from the core stops this isolate's +/// capture too. Synchronous: callers check it right before a read or a submission. +bool _optedOut() { + if (_disabled) return true; + if (!(_quiet(ffi.telemetryIsDisabled) ?? true)) return false; + _stopCapture(); + return true; +} + +void captureFailed(LocalTrackOptions options, Object error) { + final text = '$error'; + try { + ffi.telemetryDeviceEvent( + event: ffi.CaptureFailedDeviceEvent( + device: options is ScreenShareCaptureOptions + ? ffi.CaptureDevice.screenShare + : options is VideoCaptureOptions + ? ffi.CaptureDevice.camera + : ffi.CaptureDevice.microphone, + // getUserMedia error names (DOMException on web, the plugin's messages natively). + reason: text.contains('NotAllowed') || text.toLowerCase().contains('permission') + ? ffi.CaptureFailure.permissionDenied + : text.contains('NotFound') + ? ffi.CaptureFailure.notFound + : text.contains('NotReadable') || text.contains('in use') + ? ffi.CaptureFailure.inUse + : ffi.CaptureFailure.other, + ), + ); + } catch (_) {} // no native library: nothing to report to +} + +// MARK: - Transport + +/// Serves the core's pull queue by polling it: uniffi-dart callbacks cannot run on the core's +/// threads, and a pending Rust future would hold a continuation into this isolate that aborts the +/// VM once the isolate is gone (hot restart, a recreated engine). In the root zone, so no caller's +/// zone (a test's fake clock, a Room's) owns the timer. Every second while a Room is in a call or +/// requests keep coming, else every 5 s: an export nobody picks up within the core's 10 s export +/// timeout is withdrawn and retried later. +// ponytail: polling; a Dart port the core signals would remove the wake-ups. +void Function() _serve(ffi.TelemetryExportQueue queue) { + final client = http.Client(); + Timer? timer; + var stopped = false; + Future poll() async { + try { + for (var pending = queue.tryNext(); pending != null; pending = queue.tryNext()) { + _pollFast(); + try { + queue.complete(id: pending.id, response: await sendExport(client, pending.request)); + } on FormatException catch (error) { + queue.fail(id: pending.id, error: ffi.RejectedExportException('invalid url: $error')); + } catch (error) { + queue.fail( + id: pending.id, + error: ffi.RetryableExportException(reason: '$error', retryAfterMs: null), + ); + } + } + } catch (error) { + logger.fine('telemetry export not served: $error'); // the next tick tries again + } + if (stopped) return; + final fast = _inCall.isNotEmpty || _clock.elapsed < _fastUntil; + timer = Timer(Duration(seconds: fast ? 1 : 5), poll); + } + + timer = Zone.root.run(() => Timer(Duration.zero, poll)); + return () { + stopped = true; + timer?.cancel(); + _quiet(queue.finish); + _quiet(client.close); + }; +} + +/// Rooms in a call: their exports are polled for every second. +final _inCall = <_RoomTelemetry>{}; +final _clock = Stopwatch()..start(); // monotonic: wall-clock changes cannot stretch the fast window +var _fastUntil = Duration.zero; + +/// Poll every second for a while: a call ended, or requests are coming. +void _pollFast() => _fastUntil = _clock.elapsed + const Duration(seconds: 15); + +/// The bound on one request, the core's export timeout. +@visibleForTesting +Duration exportTimeout = const Duration(seconds: 10); + +/// One request: status, headers and body come back untouched and the core decides what they +/// mean; only a missing answer throws. Never follows a redirect — the request carries the +/// participant token — so a 3xx is the answer. +@visibleForTesting +Future sendExport(http.Client client, ffi.ExportRequest request) async { + // Aborting closes the socket: a collector that accepts and never answers cannot hold the queue. + final post = http.AbortableRequest('POST', Uri.parse(request.url), abortTrigger: Future.delayed(exportTimeout)) + ..followRedirects = false + ..headers.addAll(request.headers) + ..bodyBytes = request.body; + final response = await http.Response.fromStream(await client.send(post)); + return ffi.ExportResponse(status: response.statusCode, headers: response.headers, body: response.bodyBytes); +} + +// MARK: - Instruments + +/// Feeds the pipeline while it runs. +abstract interface class _Instrument { + void start(); + void stop(); +} + +/// SDK warnings and errors, next to (never instead of) the app's own log handling: filed under +/// the ambient span, else the Room whose handler logged them, else the process. +/// The Rust core copies its own warnings itself once a log forwarder is installed (the SDK installs +/// none); forwarded Rust entries must never be fed here, or they would count twice. +class _LogCapture implements _Instrument { + StreamSubscription? _records; + + // What the logger emits, from its records; what its level filters out (`disableLogging()`, a + // level above WARNING), from the SDK logger: the console level never decides what telemetry gets. + @override + void start() { + // `onRecord` is the root logger's stream unless logging is hierarchical: keep the SDK's own. + _records ??= logger.onRecord.where((r) => r.loggerName == logger.name && r.level >= Level.WARNING).listen(_log); + sdkFilteredWarningCapture = _log; + } + + @override + void stop() { + sdkFilteredWarningCapture = null; + unawaited(_records?.cancel()); + _records = null; + } + + static void _log(LogRecord record) { + // Errors are dropped silently, never logged: this runs inside a log call or the logger's own + // delivery, which a record logged from here would re-enter. + try { + _capture(record); + } catch (_) {} + } + + static void _capture(LogRecord record) { + // A zone outlives its span (listeners and timers created inside keep it): an ended span + // files the record under its Room instead. + final ambient = record.zone?[telemetrySpanKey]; + final span = ambient is _TraceSpan && ambient._span?.isEnded() == false ? ambient : null; // no _quiet here + final room = record.zone?[telemetryRoomKey] ?? (ambient is _TraceSpan ? ambient._session?.target : null); + final lowered = ffi.LogRecord( + severity: record.level >= Level.SEVERE ? ffi.Severity.error : ffi.Severity.warn, + source: ffi.LogSource.sdk, + body: record.error == null ? record.message : '${record.message} ${record.error}', + logger: record.loggerName, + timestampNs: record.time.microsecondsSinceEpoch * 1000, + spanId: span?._span?.context()?.spanId, + ); + if (span == null && room is _RoomTelemetry) { + room._scope.log(record: lowered); + } else { + ffi.telemetryLog(record: lowered); + } + } +} + +/// The device instrument: app lifecycle, memory pressure and network type as one `DeviceState`, +/// audio output changes as events. Every OS callback feeds one change stream, applied in order. +/// Thermal state, low-power mode, battery and audio interruptions need plugins the SDK does not +/// have: reported as unknown. +class _DeviceTelemetry with WidgetsBindingObserver implements _Instrument { + var _appState = ffi.AppState.foreground; + var _memory = ffi.MemoryPressure.normal; + var _network = ffi.NetworkType.unknown; + List? _outputs; + StreamController? _changes; + final _subscriptions = >[]; + Timer? _memoryRelief; + + @override + void start() { + final changes = _changes = StreamController(sync: true); + changes.stream.listen( + (apply) => _quiet(() { + apply(); + ffi.telemetrySetDeviceState( + state: ffi.DeviceState( + thermal: ffi.ThermalState.unknown, // no source without a platform plugin; low power too + appState: _appState, + memory: _memory, + network: _network, + networkExpensive: _network == ffi.NetworkType.cell, + ), + ); + }), + ); + try { + final state = WidgetsBinding.instance.lifecycleState; // already in the background, say + if (state != null) _appState = _appStateOf(state); + } catch (_) {} // no binding (plain Dart tests): foreground + changes.add(() {}); // the initial state + try { + WidgetsBinding.instance.addObserver(this); + } catch (_) {} // no binding (plain Dart tests): no lifecycle to observe + if (lkPlatformIsTest()) return; // no plugins under `flutter test` + unawaited(Connectivity().checkConnectivity().then(_networkChanged, onError: (_) {})); + _subscriptions + ..add(Connectivity().onConnectivityChanged.listen(_networkChanged)) + ..add(Hardware.instance.onDeviceChange.stream.listen(_devicesChanged)); + } + + @override + void stop() { + try { + WidgetsBinding.instance.removeObserver(this); + } catch (_) {} + for (final subscription in _subscriptions) { + unawaited(subscription.cancel()); + } + _subscriptions.clear(); + _memoryRelief?.cancel(); + unawaited(_changes?.close()); + _changes = null; + } + + void _change(void Function() apply) { + final changes = _changes; + if (changes != null && !changes.isClosed) changes.add(apply); + } + + @override + void didChangeAppLifecycleState(AppLifecycleState state) => _change(() => _appState = _appStateOf(state)); + + static ffi.AppState _appStateOf(AppLifecycleState state) => + state == AppLifecycleState.resumed || state == AppLifecycleState.inactive + ? ffi.AppState.foreground + : ffi.AppState.background; + + /// The OS warns but never says the pressure is over: back to normal after a quiet minute. + // ponytail: fixed relief delay, tune if the core's cadence stretching needs a better estimate. + @override + void didHaveMemoryPressure() { + _change(() => _memory = ffi.MemoryPressure.warning); + _memoryRelief?.cancel(); + _memoryRelief = Timer(const Duration(minutes: 1), () => _change(() => _memory = ffi.MemoryPressure.normal)); + } + + void _networkChanged(List result) => _change( + () => _network = result.contains(ConnectivityResult.none) + ? ffi.NetworkType.unavailable + : result.contains(ConnectivityResult.mobile) + ? ffi.NetworkType.cell + : result.contains(ConnectivityResult.wifi) + ? ffi.NetworkType.wifi + : result.contains(ConnectivityResult.ethernet) + ? ffi.NetworkType.wired + : result.contains(ConnectivityResult.bluetooth) + ? ffi.NetworkType.bluetooth + : result.contains(ConnectivityResult.vpn) + ? ffi.NetworkType.vpn + : result.contains(ConnectivityResult.other) + ? ffi.NetworkType.other + : ffi.NetworkType.unknown, + ); + + /// Any device change surfaces here; only a change of the audio outputs is a route change. The + /// platform gives no reason. + void _devicesChanged(List devices) => _change(() { + final outputs = [for (final device in devices.where((d) => d.kind == 'audiooutput')) _output(device.label)]; + if (_outputs != null && outputs.join(',') == _outputs!.join(',')) return; + _outputs = outputs; + ffi.telemetryDeviceEvent( + event: ffi.AudioRouteChangedDeviceEvent(outputs: outputs, reason: ffi.AudioRouteReason.unknown), + ); + }); + + static ffi.AudioOutput _output(String label) { + final l = label.toLowerCase(); + if (l.contains('bluetooth')) return ffi.AudioOutput.bluetooth; + if (l.contains('headphone') || l.contains('headset') || l.contains('wired')) return ffi.AudioOutput.wiredHeadset; + if (l.contains('speaker')) return ffi.AudioOutput.speaker; + if (l.contains('earpiece') || l.contains('receiver')) return ffi.AudioOutput.receiver; + if (l.contains('hdmi')) return ffi.AudioOutput.hdmi; + if (l.contains('usb')) return ffi.AudioOutput.usb; + if (l.contains('airplay')) return ffi.AudioOutput.airPlay; + if (l.contains('carplay') || l.contains('car audio')) return ffi.AudioOutput.carAudio; + return ffi.AudioOutput.other; + } +} + +// MARK: - Room session + +/// One Room's session, created with the Room. Also its RTC instrument: reports the tracks' +/// lifecycle, from which the core runs `lk.subscribe` (intent → first media), and hands the core +/// one raw `getStats()` report per peer connection as often as it asks. +class _RoomTelemetry implements RoomTelemetry { + _RoomTelemetry(this._room, this._scope) { + // A refreshed token (server refresh, room move): uploads keep this Room's latest. Cancelled on + // dispose: the SignalClient can outlive the Room (injected, or held by its connectivity + // subscription). + final signal = _room.engine.signalClient.events; + final cancelTokenListener = signal.on((event) { + if (_url != null) _quiet(() => _scope.setServer(url: _url!, token: event.token)); + }); + // A late or revised room sid / name. + final cancelRoomListener = signal.on((event) => _setRoom(event.room, null)); + // `dispose()` without `disconnect()` emits neither closing nor disconnected: the call ends here. + _room.onDispose(() async { + await cancelTokenListener(); + await cancelRoomListener(); + _stats?.cancel(); + _connect?.cancel(); + _connect = null; + if (_inCall.remove(this)) _quiet(() => _scope.disconnected(reason: ffi.DisconnectReason.clientInitiated)); + }); + _room.engine.events + ..on((event) => _setRoom(event.response.room, event.response.participant)) + ..on((event) => _setRoom(event.response.room, event.response.participant)) + // The app's `disconnect()` while connecting cancels the attempt; a server-initiated close + // fails it (through the error the connect then throws). + ..on((event) { + // Out of the call even when `disconnect()` times out and no RoomDisconnectedEvent follows. + _inCall.remove(this); + if (event.reason != DisconnectReason.clientInitiated) return; + _connect?.cancel(); + _connect = null; + }); + _room.events + ..on((_) { + _inCall.add(this); + step(ConnectStep.roomConnected); + // Tracks published before this Room joined: with autoSubscribe the intent is the join. + if (_room.connectOptions.autoSubscribe) { + for (final participant in _room.remoteParticipants.values.toList()) { + for (final publication in participant.trackPublications.values.toList()) { + if (publication.track == null) subscribeStarted(publication); + } + } + } + _poll(); + }) + ..on((event) { + _stats?.cancel(); + _inCall.remove(this); + _pollFast(); // the session's last records ship now + _quiet(() => _scope.disconnected(reason: _disconnectReason(event.reason))); + }) + // A new outbound track gets its first reading soon (the core asks for it), not a whole + // poll interval later. + ..on((_) => _poll()) + ..on((event) => _quiet(() => _scope.trackEnded(sid: event.publication.sid))) + ..on((event) { + // With autoSubscribe the intent exists the moment the track is known. + if (_room.connectOptions.autoSubscribe) subscribeStarted(event.publication); + }) + ..on((event) { + _quiet(() => _scope.subscribed(track: _spanTrack(event.publication))); + _poll(); // the core polls every second until first media + }) + ..on((event) => _quiet(() => _scope.trackEnded(sid: event.publication.sid))) + ..on((event) => _quiet(() => _scope.trackEnded(sid: event.publication.sid))) + ..on((event) { + if (event.sid != null) _quiet(() => _scope.subscribeFailed(sid: event.sid!, errorType: event.reason.name)); + }); + } + + final Room _room; + final ffi.TelemetryScope _scope; + + /// The session's room and participant; empty values (a room update without a sid) keep what + /// the session had. + void _setRoom(lk_models.Room room, lk_models.ParticipantInfo? participant) { + String? pick(String? value, String? old) => value == null || value.isEmpty ? old : value; + final was = _identity; + final now = _identity = ( + sid: pick(room.sid, was?.sid), + name: pick(room.name, was?.name), + participantSid: pick(participant?.sid, was?.participantSid), + participantIdentity: pick(participant?.identity, was?.participantIdentity), + ); + if (now == was) return; + _quiet( + () => _scope.setRoom( + room: ffi.RoomIdentity( + sid: now.sid, + name: now.name, + participantSid: now.participantSid, + participantIdentity: now.participantIdentity, + ), + ), + ); + } + + ({String? sid, String? name, String? participantSid, String? participantIdentity})? _identity; + String? _url; + Timer? _stats; + + // The `lk.connect` span ends once both halves are done: the engine's primary peer connection + // connected (`pc_connected`) and the join response was applied (`room_connected`); they arrive + // in either order. + _TraceSpan? _connect; + final _halves = {}; + + @override + Future connect(String url, String token, Future Function() body) async { + _url = url; + _quiet(() => _scope.setServer(url: url, token: token)); + final span = _connect = _start(ffi.ConnectSpanName()); + _halves.clear(); + try { + await span.run(body); + } catch (error) { + span?.fail(error); + if (identical(_connect, span)) _connect = null; + rethrow; + } + } + + @override + void step(ConnectStep step) { + final span = _connect; + if (span == null) return; + _quiet(() => span._span?.step(step: _step(step))); + if (step == ConnectStep.pcConnected || step == ConnectStep.roomConnected) _halves.add(step); + if (_halves.length < 2) return; + span.end(); + _connect = null; + } + + static ffi.SpanStep _step(ConnectStep step) => switch (step) { + ConnectStep.wsOpen => ffi.WsOpenSpanStep(), + ConnectStep.signal => ffi.SignalSpanStep(), + ConnectStep.joinRecv => ffi.JoinRecvSpanStep(), + ConnectStep.pcCreated => ffi.PcCreatedSpanStep(), + ConnectStep.offerSent => ffi.OfferSentSpanStep(), + ConnectStep.answerSent => ffi.AnswerSentSpanStep(), + ConnectStep.engine => ffi.EngineSpanStep(), + ConnectStep.pcConnected => ffi.PcConnectedSpanStep(), + ConnectStep.roomConnected => ffi.RoomConnectedSpanStep(), + }; + + @override + TraceSpan? reconnect(ClientDisconnectReason reason, lk_models.ReconnectReason? protoReason) => _quiet( + () => _start( + ffi.ReconnectSpanName( + protoReason != null && protoReason != lk_models.ReconnectReason.RR_UNKNOWN + ? ffi.telemetryReconnectReason(proto: protoReason.value) + : switch (reason) { + ClientDisconnectReason.signal => ffi.ReconnectReason.signalDisconnected, + ClientDisconnectReason.peerConnectionClosed || + ClientDisconnectReason.peerConnectionFailed || + ClientDisconnectReason.negotiationFailed => ffi.ReconnectReason.transportFailed, + _ => ffi.ReconnectReason.unknown, + }, + ), + ), + ); + + @override + Future publish(LocalTrack track, Future Function() body) async { + final ambient = Zone.current[telemetrySpanKey]; + final span = _start(ffi.PublishSpanName(), parent: ambient is _TraceSpan && ambient._live ? ambient : _connect); + span?._setTrack(ffi.SpanTrack(kind: _kind(track.kind), source: _source(track.source))); + try { + final publication = await span.run(body); + // The sid lets the core poll fast until this track's first outbound reading. + span?._setTrack(ffi.SpanTrack(sid: publication.sid, kind: _kind(track.kind), source: _source(track.source))); + span?.end(); + return publication; + } catch (error) { + span?.fail(error); + rethrow; + } + } + + @override + void subscribeStarted(RemoteTrackPublication publication) { + _quiet(() => _scope.subscribeStarted(track: _spanTrack(publication))); + _poll(); // the core polls every second until first media + } + + @override + void emitCustom(String name, Map attributes) => + _quiet(() => _scope.emitCustom(name: name, attributes: attributes)); + + @override + void setAttribute(String key, String? value) => _quiet(() => _scope.setAttribute(key: key, value: value)); + + _TraceSpan? _start(ffi.SpanName name, {_TraceSpan? parent}) => + _quiet(() => _TraceSpan(this, _scope.start(name: name, parent: parent?._span))); + + ffi.SpanTrack _spanTrack(RemoteTrackPublication publication) => ffi.SpanTrack( + sid: publication.sid, + kind: _kind(publication.kind), + source: _source(publication.source), + remoteIdentity: publication.participant.identity, + ); + + /// Arms the stats timer at the core's current interval, or brings an armed one forward: a track + /// signal never postpones a poll, and a read in flight re-arms when it is done (one poll at a + /// time). The loop ends with the session. + void _poll() { + if (_reading) return; + final interval = _disabled ? null : _quiet(_scope.statsPollIntervalMs); + if (interval == null) return; + final due = _clock.elapsed + Duration(milliseconds: interval); + if (_stats?.isActive ?? false) { + if (_due <= due) return; + _stats!.cancel(); + } + _due = due; + _stats = Timer(due - _clock.elapsed, () async { + if (_disabled || _room.connectionState == ConnectionState.disconnected) return; + _reading = true; + try { + await _recordPeerStats(); + } finally { + _reading = false; + } + if (_room.connectionState != ConnectionState.disconnected) _poll(); + }); + } + + Duration _due = Duration.zero; + bool _reading = false; + + /// One `getStats()` report per peer connection, with every track this Room sends or receives. + Future _recordPeerStats() async { + final tracks = { + for (final participant in [_room.localParticipant, ..._room.remoteParticipants.values]) + for (final publication in participant?.trackPublications.values ?? []) + ?publication.track?.mediaStreamTrack.id: publication.sid, + }; + for (final transport in [_room.engine.publisher, _room.engine.subscriber]) { + if (transport == null) continue; + if (_optedOut()) return; // every request is gated: the previous one may have failed meanwhile + try { + // Bounded: a peer connection that never answers holds no poll (a late answer is dropped). + final report = await transport.pc.getStats().timeout(const Duration(seconds: 5)); + if (_optedOut()) return; // opted out while reading + _scope.recordPeerStats( + report: [ + for (final stat in report) ffi.RtcStat(kind: stat.type, id: stat.id, members: _members(stat.values)), + ], + tracks: tracks, + timestampNs: null, + ); + } catch (error) { + logger.fine('telemetry stats not read: $error'); // a peer connection closing meanwhile + } + } + } +} + +/// The core's span; names, timing, attributes, outcome and export live in the core. +class _TraceSpan implements TraceSpan { + _TraceSpan(_RoomTelemetry session, ffi.TelemetrySpan span) : _session = WeakReference(session), _span = span { + liveSpanHandles++; + } + + /// Weak, and both released when the span ends: a zone (and every long-lived listener created in + /// it, like the signal client's connectivity subscription) can outlive the span and its Room, + /// but then holds only this empty object. + WeakReference<_RoomTelemetry>? _session; + ffi.TelemetrySpan? _span; + + /// Not ended yet. + bool get _live => _span != null && (_quiet(() => !_span!.isEnded()) ?? false); + + void _setTrack(ffi.SpanTrack track) => _quiet(() => _span?.setTrack(track: track)); + + /// Ends the core's span with [finish] and lets go of its handle. + void _finish(void Function(ffi.TelemetrySpan span) finish) { + final span = _span; + if (span == null) return; + _span = _session = null; + liveSpanHandles--; + _quiet(() => finish(span)); + _quiet(span.dispose); + } + + @override + void attempt(int number, {required bool full}) => _quiet( + () => _span?.step( + step: ffi.AttemptSpanStep(number: number, full: full), + ), + ); + + @override + void end() => _finish((span) => span.end(outcome: ffi.SpanOutcome.ok, error: null)); + + @override + void fail(Object error) => _finish((span) => span.fail(error: error is String ? error : '${error.runtimeType}')); + + @override + void cancel() => _finish((span) => span.cancel()); +} + +/// Spans whose native handle this isolate still holds (open ones). +@visibleForTesting +int liveSpanHandles = 0; + +// MARK: - Vocabulary + +ffi.DisconnectReason _disconnectReason(DisconnectReason? reason) => switch (reason) { + null => ffi.DisconnectReason.unknown, + DisconnectReason.disconnected => ffi.DisconnectReason.signalClose, + DisconnectReason.signalingConnectionFailure => ffi.DisconnectReason.joinFailure, + DisconnectReason.reconnectAttemptsExceeded => ffi.DisconnectReason.reconnectFailed, + _ => ffi.telemetryDisconnectReason( + proto: lk_models.DisconnectReason.values + .firstWhere((p) => p.toSDKType() == reason, orElse: () => lk_models.DisconnectReason.UNKNOWN_REASON) + .value, + ), +}; + +ffi.TrackKind _kind(TrackType kind) => kind == TrackType.AUDIO ? ffi.TrackKind.audio : ffi.TrackKind.video; + +ffi.TrackSource _source(TrackSource source) => switch (source) { + TrackSource.camera => ffi.TrackSource.camera, + TrackSource.microphone => ffi.TrackSource.microphone, + TrackSource.screenShareVideo => ffi.TrackSource.screenShare, + TrackSource.screenShareAudio => ffi.TrackSource.screenShareAudio, + TrackSource.unknown => ffi.TrackSource.unknown, +}; + +/// A stats entry's members as the core takes them; nested maps flattened with a dot +/// (`qualityLimitationDurations.cpu`), sequences dropped. No member names are known here. +Map _members(Map values, [String prefix = '']) { + final members = {}; + values.forEach((key, value) { + final name = prefix.isEmpty ? '$key' : '$prefix.$key'; + if (value is Map) { + members.addAll(_members(value, name)); + } else if (value is bool) { + members[name] = ffi.BoolAttributeValue(value); + } else if (value is int) { + members[name] = ffi.IntAttributeValue(value); + } else if (value is num) { + members[name] = ffi.DoubleAttributeValue(value.toDouble()); + } else if (value is String) { + members[name] = ffi.StrAttributeValue(value); + } + }); + return members; +} diff --git a/lib/src/telemetry/telemetry_web.dart b/lib/src/telemetry/telemetry_web.dart new file mode 100644 index 000000000..fa99d194f --- /dev/null +++ b/lib/src/telemetry/telemetry_web.dart @@ -0,0 +1,25 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +import '../core/room.dart'; +import '../track/options.dart'; +import 'telemetry.dart'; + +// The Rust core is a native library the web cannot load: no Room ever gets a session. + +RoomTelemetry? createRoomTelemetry(Room room) => null; + +Future disableTelemetry() async {} + +void captureFailed(LocalTrackOptions options, Object error) {} From 9eb2deb080cedc8d560e8d15f490f5a2df79220e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?B=C5=82az=CC=87ej=20Pankowski?= <86720177+pblazej@users.noreply.github.com> Date: Thu, 1 Oct 2026 13:19:41 +0200 Subject: [PATCH 3/8] feat(telemetry): instrument Engine connection lifecycle Record Engine connect checkpoints on the lk.connect span that Room wiring starts, add reconnect spans, and propagate the mapped disconnect reason. --- lib/src/core/engine.dart | 40 ++++++++++++++++++++++++++++-------- lib/src/internal/events.dart | 6 ++++-- 2 files changed, 35 insertions(+), 11 deletions(-) diff --git a/lib/src/core/engine.dart b/lib/src/core/engine.dart index ba82dfefb..a8a66fe75 100644 --- a/lib/src/core/engine.dart +++ b/lib/src/core/engine.dart @@ -42,6 +42,7 @@ import '../support/disposable.dart'; import '../support/platform.dart' show lkPlatformIsTest, lkPlatformIs, PlatformType; import '../support/region_url_provider.dart'; import '../support/websocket.dart'; +import '../telemetry/telemetry.dart'; import '../track/local/local.dart'; import '../track/local/video.dart'; import '../types/internal.dart'; @@ -145,6 +146,13 @@ class Engine extends Disposable with EventsEmittable { bool _attemptingReconnect = false; + /// The Room's telemetry session, set by the Room. + @internal + RoomTelemetry? telemetry; + + /// One reconnect cycle = one `lk.reconnect` span; attempts are its checkpoints. + TraceSpan? _reconnectSpan; + RegionUrlProvider? _regionUrlProvider; lk_models.ServerInfo? _serverInfo; @@ -258,6 +266,7 @@ class Engine extends Disposable with EventsEmittable { connectOptions: this.connectOptions, roomOptions: this.roomOptions, ); + telemetry?.step(ConnectStep.signal); // wait for join response await events.waitFor( @@ -278,6 +287,9 @@ class Engine extends Disposable with EventsEmittable { 'Timed out waiting for PeerConnection to connect, please check your network for ice connectivity', ), ); + telemetry + ?..step(ConnectStep.engine) + ..step(ConnectStep.pcConnected); events.emit(const EngineConnectedEvent()); } catch (error) { logger.fine('Connect Error $error'); @@ -321,6 +333,8 @@ class Engine extends Disposable with EventsEmittable { _isReconnecting = false; _clearPendingReconnect(); + _reconnectSpan?.cancel(); + _reconnectSpan = null; } @internal @@ -696,6 +710,7 @@ class Engine extends Disposable with EventsEmittable { publisher?.onOffer = (offer) { logger.fine('publisher onOffer'); signalClient.sendOffer(offer); + telemetry?.step(ConnectStep.offerSent); }; // in subscriber primary mode, server side opens sub data channels. @@ -1094,10 +1109,13 @@ class Engine extends Disposable with EventsEmittable { if (_reconnectAttempts == 0) { _reconnectStart = DateTime.timestamp(); + _reconnectSpan ??= telemetry?.reconnect(reason, reconnectReason); } if (_reconnectAttempts >= _reconnectCount) { logger.fine('reconnectAttempts exceeded, disconnecting...'); + _reconnectSpan?.fail('reconnectAttemptsExceeded'); + _reconnectSpan = null; _isClosed = true; await cleanUp(); @@ -1165,6 +1183,7 @@ class Engine extends Disposable with EventsEmittable { final fullReconnect = fullReconnectOnNext; fullReconnectOnNext = false; _attemptIsFullReconnect = fullReconnect; + _reconnectSpan?.attempt(_reconnectAttempts + 1, full: fullReconnect); var succeeded = false; try { @@ -1182,18 +1201,15 @@ class Engine extends Disposable with EventsEmittable { ); } - if (fullReconnect) { - await restartConnection(); - } else { - await resumeConnection( - reason, - reconnectReason: reconnectReason, - ); - } + await _reconnectSpan.run( + () => fullReconnect ? restartConnection() : resumeConnection(reason, reconnectReason: reconnectReason), + ); _clearPendingReconnect(); _attemptingReconnect = false; _isReconnecting = false; succeeded = true; + _reconnectSpan?.end(); + _reconnectSpan = null; } catch (e) { _reconnectAttempts = _reconnectAttempts + 1; logger.fine('attemptReconnect: ${fullReconnect ? 'full reconnect' : 'resume'} failed: $e'); @@ -1214,6 +1230,8 @@ class Engine extends Disposable with EventsEmittable { unawaited(handleReconnect(ClientDisconnectReason.reconnectRetry)); } else { logger.fine('attemptReconnect: disconnecting...'); + _reconnectSpan?.fail(e); + _reconnectSpan = null; // clean up before emitting, room's EngineDisconnectedEvent handler // drops the event while fullReconnectOnNext is still true and // cleanUp() is what resets it @@ -1420,6 +1438,7 @@ class Engine extends Disposable with EventsEmittable { void _setUpSignalListeners() => _signalListener ..on((event) async { + telemetry?.step(ConnectStep.joinRecv); // create peer connections _subscriberPrimary = event.response.subscriberPrimary; _serverInfo = event.response.serverInfo; @@ -1445,6 +1464,7 @@ class Engine extends Disposable with EventsEmittable { if (publisher == null && subscriber == null) { await _createPeerConnections(rtcConfiguration); + telemetry?.step(ConnectStep.pcCreated); } if (!_subscriberPrimary || event.response.fastPublish) { @@ -1495,6 +1515,7 @@ class Engine extends Disposable with EventsEmittable { }) ..on((event) async { logger.fine('Signal connected'); + telemetry?.step(ConnectStep.wsOpen); // The attempt counter is not reset here. A resume opens its socket before // the peer connections are restored, so a reset on socket connect would // let an attempt that fails afterwards start again from zero and never @@ -1544,6 +1565,7 @@ class Engine extends Disposable with EventsEmittable { logger.finer('sdp: ${answer.sdp}'); await subscriber!.pc.setLocalDescription(answer); signalClient.sendAnswer(answer); + telemetry?.step(ConnectStep.answerSent); } catch (_) { logger.severe('[$objectId] Failed to createAnswer()'); } @@ -1620,7 +1642,7 @@ class Engine extends Disposable with EventsEmittable { DisconnectReason reason = DisconnectReason.clientInitiated, }) async { _isClosed = true; - events.emit(EngineClosingEvent()); + events.emit(EngineClosingEvent(reason: reason)); if (connectionState == ConnectionState.connected) { await signalClient.sendLeave(); } else { diff --git a/lib/src/internal/events.dart b/lib/src/internal/events.dart index 63c54a831..75e2cd5ab 100644 --- a/lib/src/internal/events.dart +++ b/lib/src/internal/events.dart @@ -210,10 +210,12 @@ class EngineDisconnectedEvent with InternalEvent, EngineEvent { @internal class EngineClosingEvent with InternalEvent, EngineEvent { - const EngineClosingEvent(); + /// Why the engine closes: [DisconnectReason.clientInitiated] for the app's own `disconnect()`. + final DisconnectReason reason; + const EngineClosingEvent({this.reason = DisconnectReason.clientInitiated}); @override - String toString() => '${runtimeType}()'; + String toString() => '${runtimeType}(reason: $reason)'; } @internal From 9c6bc38e210f9373c4cc560599496e020e15ad4d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?B=C5=82az=CC=87ej=20Pankowski?= <86720177+pblazej@users.noreply.github.com> Date: Thu, 1 Oct 2026 13:19:41 +0200 Subject: [PATCH 4/8] feat(telemetry): instrument participants and tracks Track lk.publish spans for local participants, lk.subscribe for remote. Capture getUserMedia failures with mapped reason. --- lib/src/participant/local.dart | 17 +++++++++++++---- lib/src/publication/remote.dart | 1 + lib/src/track/local/local.dart | 12 +++++++++++- 3 files changed, 25 insertions(+), 5 deletions(-) diff --git a/lib/src/participant/local.dart b/lib/src/participant/local.dart index 61cd38434..6c41def49 100644 --- a/lib/src/participant/local.dart +++ b/lib/src/participant/local.dart @@ -157,8 +157,13 @@ class LocalParticipant extends Participant { LocalAudioTrack track, { AudioPublishOptions? publishOptions, }) async { - final result = await _publishRunner.run(() => _publishAudioTrack(track, publishOptions: publishOptions)); - return result! as LocalTrackPublication; + // One publish attempt = one `lk.publish` span. + Future> publish() async { + final result = await _publishRunner.run(() => _publishAudioTrack(track, publishOptions: publishOptions)); + return result! as LocalTrackPublication; + } + + return room.telemetry?.publish(track, publish) ?? publish(); } Future?> _publishAudioTrack( @@ -273,8 +278,12 @@ class LocalParticipant extends Participant { LocalVideoTrack track, { VideoPublishOptions? publishOptions, }) async { - final result = await _publishRunner.run(() => _publishVideoTrack(track, publishOptions: publishOptions)); - return result! as LocalTrackPublication; + Future> publish() async { + final result = await _publishRunner.run(() => _publishVideoTrack(track, publishOptions: publishOptions)); + return result! as LocalTrackPublication; + } + + return room.telemetry?.publish(track, publish) ?? publish(); } Future?> _publishVideoTrack( diff --git a/lib/src/publication/remote.dart b/lib/src/publication/remote.dart index a90bf3c5c..acf467d86 100644 --- a/lib/src/publication/remote.dart +++ b/lib/src/publication/remote.dart @@ -349,6 +349,7 @@ class RemoteTrackPublication extends TrackPublication logger.fine('ignoring subscribe() request...'); return; } + participant.room.telemetry?.subscribeStarted(this); _sendUpdateSubscription(subscribed: true); } diff --git a/lib/src/track/local/local.dart b/lib/src/track/local/local.dart index 7252cb045..83918b935 100644 --- a/lib/src/track/local/local.dart +++ b/lib/src/track/local/local.dart @@ -31,6 +31,7 @@ import '../../logger.dart'; import '../../participant/remote.dart'; import '../../support/native.dart'; import '../../support/platform.dart'; +import '../../telemetry/telemetry.dart'; import '../../types/other.dart'; import '../options.dart'; import '../processor.dart'; @@ -245,7 +246,16 @@ abstract class LocalTrack extends Track { /// Creates a [rtc.MediaStream] from [LocalTrackOptions]. @internal - static Future createStream( + static Future createStream(LocalTrackOptions options) async { + try { + return await _createStream(options); + } catch (error) { + RoomTelemetry.captureFailed(options, error); + rethrow; + } + } + + static Future _createStream( LocalTrackOptions options, ) async { final constraints = { From 18526276b974df0bf5770386b8694cd91d9b02f5 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?B=C5=82az=CC=87ej=20Pankowski?= <86720177+pblazej@users.noreply.github.com> Date: Thu, 1 Oct 2026 13:19:41 +0200 Subject: [PATCH 5/8] feat(telemetry): expose public API Add LiveKitClient.disableTelemetry() to opt out process-wide. Room.emitTelemetryEvent() sends custom. events with attributes. Room.setTelemetryAttribute() sets correlation attributes per session. --- lib/livekit_client.dart | 1 + lib/src/core/room.dart | 61 +++++++++++++++++++++++++++++++++++++---- lib/src/livekit.dart | 9 ++++++ 3 files changed, 65 insertions(+), 6 deletions(-) diff --git a/lib/livekit_client.dart b/lib/livekit_client.dart index f8527e14b..e529ed165 100644 --- a/lib/livekit_client.dart +++ b/lib/livekit_client.dart @@ -90,6 +90,7 @@ export 'src/token_source/custom.dart'; export 'src/token_source/caching.dart'; export 'src/token_source/development.dart'; export 'src/token_source/jwt.dart'; +export 'src/uniffi/uniffi.dart'; /// Misspelled alias for [kRpcVersion]. Kept for backwards compatibility with code /// that referenced the original typo. diff --git a/lib/src/core/room.dart b/lib/src/core/room.dart index 4ef28bbf0..5afcc025c 100644 --- a/lib/src/core/room.dart +++ b/lib/src/core/room.dart @@ -44,6 +44,7 @@ import '../support/disposable.dart'; import '../support/http_client.dart'; import '../support/platform.dart'; import '../support/region_url_provider.dart'; +import '../telemetry/telemetry.dart'; import '../track/audio_management.dart'; import '../track/local/audio.dart'; import '../track/local/video.dart'; @@ -122,6 +123,33 @@ class Room extends DisposableChangeNotifier with EventsEmittable { @internal final Engine engine; + + /// This Room's telemetry session, one trace for the Room's lifetime; null on web and after + /// `LiveKitClient.disableTelemetry`. + @internal + late final RoomTelemetry? telemetry = RoomTelemetry.create(this); + + /// Records an app event in this Room's telemetry, exported as `custom.` next to the SDK's + /// own records, with this Room's correlation attributes. A no-op on web. + /// + /// ```dart + /// room.emitTelemetryEvent('checkout.started', attributes: {'cart.items': '3'}); + /// ``` + /// + /// Names and keys up to 128 bytes, values up to 1024 bytes, at most 64 attributes and no `lk.` + /// keys; anything else is dropped, never truncated. + void emitTelemetryEvent(String name, {Map attributes = const {}}) => + telemetry?.emitCustom(name, attributes); + + /// Sets a correlation attribute on every telemetry record this Room captures from now on, to + /// match them with your own data (an order id, a tenant); `null` removes it. Same limits as + /// [emitTelemetryEvent], at most 64 per Room. A no-op on web. + /// + /// ```dart + /// room.setTelemetryAttribute('app.order_id', orderId); + /// ``` + void setTelemetryAttribute(String key, String? value) => telemetry?.setAttribute(key, value); + // suppport for multiple event listeners late final EventsListener _engineListener; // @@ -179,12 +207,16 @@ class Room extends DisposableChangeNotifier with EventsEmittable { connectOptions: connectOptions, roomOptions: roomOptions, ) { - // - _engineListener = this.engine.createListener(); - _setUpEngineListeners(); - - _signalListener = this.engine.signalClient.createListener(); - _setUpSignalListeners(); + this.engine.telemetry = telemetry; + // The Room's handlers run in its telemetry zone: a warning one of them logs lands in this + // Room's session even with no span open. + telemetry.run(() { + _engineListener = this.engine.createListener(); + _setUpEngineListeners(); + + _signalListener = this.engine.signalClient.createListener(); + _setUpSignalListeners(); + }); _rpcClientManager = RpcClientManager(this); _rpcServerManager = RpcServerManager(this); @@ -273,6 +305,23 @@ class Room extends DisposableChangeNotifier with EventsEmittable { ConnectOptions? connectOptions, @Deprecated('deprecated, please use roomOptions in Room constructor') RoomOptions? roomOptions, FastConnectOptions? fastConnectOptions, + }) { + Future connect() => _connect( + url, + token, + connectOptions: connectOptions, + roomOptions: roomOptions, + fastConnectOptions: fastConnectOptions, + ); + return telemetry?.connect(url, token, connect) ?? connect(); + } + + Future _connect( + String url, + String token, { + ConnectOptions? connectOptions, + RoomOptions? roomOptions, + FastConnectOptions? fastConnectOptions, }) async { var effectiveRoomOptions = roomOptions ?? this.roomOptions; if (lkPlatformIs(PlatformType.web) && diff --git a/lib/src/livekit.dart b/lib/src/livekit.dart index b4a79f236..4f4141f37 100644 --- a/lib/src/livekit.dart +++ b/lib/src/livekit.dart @@ -19,12 +19,21 @@ import 'audio/audio_session.dart'; import 'support/native.dart'; import 'support/platform.dart' show PlatformType, lkPlatformIs, lkPlatformIsMobile, lkPlatformIsDesktop; import 'support/webrtc_initialize_options.dart'; +import 'telemetry/telemetry.dart'; /// Main entry point to connect to a room. /// {@category Room} class LiveKitClient { static const version = '2.13.0'; + /// Opts this process out of client telemetry, in effect as soon as it is called: collection + /// stops and Rooms created afterwards collect nothing. Everything not yet sent — queued, open or + /// cached on disk — is deleted in the background. Call it at every launch, before creating a + /// Room, to collect nothing at all. A no-op on web. + /// + /// TODO: final shape pending the token/consent discussion. + static Future disableTelemetry() => RoomTelemetry.disable(); + /// Initialize the WebRTC plugin. /// /// Optional: call once at startup to enable [bypassVoiceProcessing] before From caf4a31818750526bebedba59a9a6a5d51133b7d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?B=C5=82az=CC=87ej=20Pankowski?= <86720177+pblazej@users.noreply.github.com> Date: Thu, 1 Oct 2026 13:19:41 +0200 Subject: [PATCH 6/8] test(telemetry): add comprehensive test suite Unit tests for SDK logger, telemetry core errors, isolate safety, opt-out. E2E test against local collector validates the full pipeline. --- test/mock/e2e_container.dart | 19 +- test/mock/media_stream_mock.dart | 119 +++++ test/mock/peerconnection_mock.dart | 95 +++- test/support/sdk_logger_test.dart | 87 ++++ test/telemetry/otelcol.yaml | 22 + .../telemetry/telemetry_core_errors_test.dart | 161 ++++++ test/telemetry/telemetry_e2e_test.dart | 490 ++++++++++++++++++ test/telemetry/telemetry_isolate_test.dart | 81 +++ .../telemetry_opt_out_submit_test.dart | 85 +++ test/telemetry/telemetry_opt_out_test.dart | 69 +++ test/telemetry/telemetry_session_test.dart | 165 ++++++ test/uniffi/uniffi_test.dart | 40 ++ 12 files changed, 1424 insertions(+), 9 deletions(-) create mode 100644 test/mock/media_stream_mock.dart create mode 100644 test/support/sdk_logger_test.dart create mode 100644 test/telemetry/otelcol.yaml create mode 100644 test/telemetry/telemetry_core_errors_test.dart create mode 100644 test/telemetry/telemetry_e2e_test.dart create mode 100644 test/telemetry/telemetry_isolate_test.dart create mode 100644 test/telemetry/telemetry_opt_out_submit_test.dart create mode 100644 test/telemetry/telemetry_opt_out_test.dart create mode 100644 test/telemetry/telemetry_session_test.dart create mode 100644 test/uniffi/uniffi_test.dart diff --git a/test/mock/e2e_container.dart b/test/mock/e2e_container.dart index 3c9527213..244ae994d 100644 --- a/test/mock/e2e_container.dart +++ b/test/mock/e2e_container.dart @@ -58,13 +58,15 @@ class E2EContainer { /// is rewritten so the local participant's [Participant.clientProtocol] takes /// that value (used to exercise v1 vs v2 caller paths in self-loop tests). /// When [captureOutbound] is true, all DataPackets sent over the reliable - /// data channel are recorded in [capturedDataPackets]. + /// data channel are recorded in [capturedDataPackets]. [otherParticipants] + /// are already in the room when it joins. Future connectRoom({ int? localClientProtocol, bool captureOutbound = false, ConnectOptions? connectOptions, @Deprecated('mirrors the deprecated Room.connect parameter') RoomOptions? roomOptions, lk_models.ClientConfiguration? clientConfiguration, + List otherParticipants = const [], }) async { final connectFuture = room.connect( exampleUri, @@ -73,7 +75,13 @@ class E2EContainer { // ignore: deprecated_member_use_from_same_package roomOptions: roomOptions, ); - unawaited(answerJoin(localClientProtocol: localClientProtocol, clientConfiguration: clientConfiguration)); + unawaited( + answerJoin( + localClientProtocol: localClientProtocol, + clientConfiguration: clientConfiguration, + otherParticipants: otherParticipants, + ), + ); await connectFuture; @@ -115,10 +123,15 @@ class E2EContainer { Future answerJoin({ int? localClientProtocol, lk_models.ClientConfiguration? clientConfiguration, + List otherParticipants = const [], }) async { // Give the SDK a tick to start waiting for the join response. await Future.delayed(const Duration(milliseconds: 1)); - final resp = _buildJoinResponse(localClientProtocol, clientConfiguration); + var resp = _buildJoinResponse(localClientProtocol, clientConfiguration); + if (otherParticipants.isNotEmpty) { + // A copy: the default response is shared by every test. + resp = lk_rtc.SignalResponse(join: resp.join.deepCopy()..otherParticipants.addAll(otherParticipants)); + } wsConnector.onData(resp.writeToBuffer()); wsConnector.onData(offerResponse.writeToBuffer()); } diff --git a/test/mock/media_stream_mock.dart b/test/mock/media_stream_mock.dart new file mode 100644 index 000000000..355c5765e --- /dev/null +++ b/test/mock/media_stream_mock.dart @@ -0,0 +1,119 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +import 'dart:typed_data'; + +import 'package:flutter_webrtc/flutter_webrtc.dart' as rtc; + +/// A media stream with no native counterpart, for tracks built in unit tests. +class FakeMediaStream extends rtc.MediaStream { + final List _tracks = []; + + FakeMediaStream(String id) : super(id, 'fake-owner'); + + @override + bool? get active => true; + + @override + Future addTrack(rtc.MediaStreamTrack track, {bool addToNative = true}) async { + _tracks.add(track); + } + + @override + Future clone() async => FakeMediaStream('${id}_clone'); + + @override + List getAudioTracks() => _tracks.where((t) => t.kind == 'audio').toList(); + + @override + Future getMediaTracks() async {} + + @override + List getTracks() => List.from(_tracks); + + @override + List getVideoTracks() => _tracks.where((t) => t.kind == 'video').toList(); + + @override + Future removeTrack(rtc.MediaStreamTrack track, {bool removeFromNative = true}) async { + _tracks.remove(track); + } +} + +class FakeMediaStreamTrack implements rtc.MediaStreamTrack { + @override + rtc.StreamTrackCallback? onEnded; + + @override + rtc.StreamTrackCallback? onMute; + + @override + rtc.StreamTrackCallback? onUnMute; + + @override + bool enabled; + + @override + final String id; + + @override + final String kind; + + @override + String? get label => '$kind-track'; + + @override + bool? get muted => false; + + FakeMediaStreamTrack({required this.id, required this.kind, this.enabled = true}); + + @override + Future applyConstraints([Map? constraints]) async {} + + @override + Future clone() async => FakeMediaStreamTrack(id: id, kind: kind, enabled: enabled); + + @override + Future dispose() async {} + + @override + Future adaptRes(int width, int height) async {} + + @override + Map getConstraints() => const {}; + + @override + Map getSettings() => const {}; + + @override + Future stop() async {} + + @override + void enableSpeakerphone(bool enable) {} + + @override + Future captureFrame() => throw UnimplementedError(); + + @override + Future hasTorch() async => false; + + @override + Future setTorch(bool torch) async {} + + @override + Future switchCamera() async => false; + + @override + String toString() => 'FakeMediaStreamTrack(id: $id, kind: $kind, enabled: $enabled)'; +} diff --git a/test/mock/peerconnection_mock.dart b/test/mock/peerconnection_mock.dart index 1dc3c8755..e96167db3 100644 --- a/test/mock/peerconnection_mock.dart +++ b/test/mock/peerconnection_mock.dart @@ -43,11 +43,61 @@ void resetMockDataChannels() { _dataChannels.clear(); } +/// Remote MediaStreamTrack ids (→ kind) every mock peer connection reports a growing +/// `inbound-rtp` stream for, as if media arrived. +final mockInboundTracks = {}; + +/// A sender for [MockPeerConnection.addTransceiver]. +class MockRtpSender extends RTCRtpSender { + MockRtpSender(this._track); + + final MediaStreamTrack? _track; + + @override + MediaStreamTrack? get track => _track; + + @override + String get senderId => 'mock-sender'; + + @override + Future> getStats() async => []; + + @override + Future replaceTrack(MediaStreamTrack? track) async {} + + @override + Future dispose() async {} + + @override + dynamic noSuchMethod(Invocation invocation) => throw UnimplementedError('${invocation.memberName}'); +} + +class MockRtpTransceiver extends RTCRtpTransceiver { + MockRtpTransceiver(this.sender); + + @override + final RTCRtpSender sender; + + @override + String get mid => '0'; + + @override + Future setCodecPreferences(List codecs) async {} + + @override + Future stop() async {} + + @override + dynamic noSuchMethod(Invocation invocation) => throw UnimplementedError('${invocation.memberName}'); +} + class MockPeerConnection extends RTCPeerConnection { static const _offerType = 'offer'; static const _answerType = 'answer'; bool closed = false; + final _sentTracks = []; + int _polls = 0; RTCSessionDescription? _localDescription; RTCSessionDescription? _remoteDescription; @@ -171,9 +221,9 @@ class MockPeerConnection extends RTCPeerConnection { MediaStreamTrack? track, RTCRtpMediaType? kind, RTCRtpTransceiverInit? init, - }) { - // TODO: implement addTransceiver - throw UnimplementedError(); + }) async { + if (track != null) _sentTracks.add(track); + return MockRtpTransceiver(MockRtpSender(track)); } @override @@ -273,8 +323,42 @@ a=rtpmap:32 MPV/90000 @override Future> getSenders() async => List.empty(); + /// `getStats()` calls on every mock peer connection in this process. + static int statsCalls = 0; + + /// Runs inside every `getStats()` before it answers (a slow or failing read). + static Future Function()? onGetStats; + @override - Future> getStats([MediaStreamTrack? track]) async => List.empty(); + Future> getStats([MediaStreamTrack? track]) async { + statsCalls++; + await onGetStats?.call(); + // Growing counters, like a peer connection with media flowing: an `outbound-rtp` (and its + // `media-source`) per sent track, an `inbound-rtp` per [mockInboundTracks] entry. + final bytes = 4000 * ++_polls; + final now = DateTime.now().microsecondsSinceEpoch.toDouble(); + return [ + StatsReport('CIT01', 'codec', now, {'mimeType': 'audio/opus'}), + for (final (i, sent) in _sentTracks.indexed) ...[ + StatsReport('MS$i', 'media-source', now, {'trackIdentifier': sent.id, 'kind': sent.kind}), + StatsReport('OT$i', 'outbound-rtp', now, { + 'kind': sent.kind, + 'mediaSourceId': 'MS$i', + 'bytesSent': bytes, + 'packetsSent': bytes ~/ 100, + 'codecId': 'CIT01', + }), + ], + for (final (i, MapEntry(key: id, value: kind)) in mockInboundTracks.entries.indexed) + StatsReport('IT$i', 'inbound-rtp', now, { + 'kind': kind, + 'trackIdentifier': id, + 'bytesReceived': bytes, + 'packetsReceived': bytes ~/ 100, + 'codecId': 'CIT01', + }), + ]; + } @override Future> getTransceivers() async => List.empty(); @@ -300,6 +384,5 @@ a=rtpmap:32 MPV/90000 ]) async => MockPeerConnection(); @override - // TODO: implement restartIce - Future restartIce() => throw UnimplementedError(); + Future restartIce() async {} } diff --git a/test/support/sdk_logger_test.dart b/test/support/sdk_logger_test.dart new file mode 100644 index 000000000..3de7a4b1f --- /dev/null +++ b/test/support/sdk_logger_test.dart @@ -0,0 +1,87 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +import 'package:flutter_test/flutter_test.dart'; +import 'package:logging/logging.dart'; + +import 'package:livekit_client/src/support/sdk_logger.dart'; + +void main() { + late SdkLogger logger; + late List emitted, captured; + + setUp(() { + logger = SdkLogger(Logger.detached('sdk_logger_test')); + emitted = []; + captured = []; + final listening = logger.onRecord.listen(emitted.add); + sdkFilteredWarningCapture = captured.add; // as with telemetry installed + addTearDown(() { + sdkFilteredWarningCapture = null; + return listening.cancel(); + }); + }); + + test('an emitted record is the logger\'s own, message converted and evaluated once', () { + logger.level = Level.ALL; + final message = _Counting(); + var evaluations = 0; + var inner = 0; + logger + ..warning(message) + ..severe(() => ++evaluations) + ..warning( + () => + () => ++inner, + ); // a lazy message whose value is a Function: not called again + expect(message.conversions, 1); + expect(evaluations, 1); + expect(inner, 0); + expect(emitted.map((r) => r.message).take(2), ['counted', '1']); + expect(captured, isEmpty, reason: 'telemetry takes emitted records from onRecord'); + }); + + test('a filtered warning goes to telemetry only, evaluated once', () { + logger.level = Level.OFF; + var evaluations = 0; + logger + ..warning(() => ++evaluations) + ..info(() => ++evaluations); // below WARNING: never evaluated, as before + expect(evaluations, 1); + expect(captured.single.message, '1'); + expect(emitted, isEmpty); + }); + + test('a filtered message that throws never reaches the caller', () { + logger.level = Level.OFF; + expect(() => logger.warning(() => throw StateError('message')), returnsNormally); + expect(captured, isEmpty); + }); + + test('without telemetry a filtered message is never evaluated', () { + sdkFilteredWarningCapture = null; + logger.level = Level.OFF; + expect(() => logger.warning(() => throw StateError('message')), returnsNormally); + }); +} + +class _Counting { + var conversions = 0; + + @override + String toString() { + conversions++; + return 'counted'; + } +} diff --git a/test/telemetry/otelcol.yaml b/test/telemetry/otelcol.yaml new file mode 100644 index 000000000..61f37e5e3 --- /dev/null +++ b/test/telemetry/otelcol.yaml @@ -0,0 +1,22 @@ +# Collector for the telemetry end-to-end test: accepts OTLP/HTTP on :4319 and writes every batch +# as one JSON line the test reads back. Run: otelcol-contrib --config test/telemetry/otelcol.yaml +# then: LK_TELEMETRY_ENDPOINT=http://127.0.0.1:4319 flutter test test/telemetry/ +receivers: + otlp: + protocols: + http: + endpoint: 127.0.0.1:4319 +exporters: + file: + path: /tmp/livekit-telemetry-otlp.jsonl +service: + telemetry: + metrics: + level: none + pipelines: + logs: + receivers: [otlp] + exporters: [file] + traces: + receivers: [otlp] + exporters: [file] diff --git a/test/telemetry/telemetry_core_errors_test.dart b/test/telemetry/telemetry_core_errors_test.dart new file mode 100644 index 000000000..ad8127a0e --- /dev/null +++ b/test/telemetry/telemetry_core_errors_test.dart @@ -0,0 +1,161 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +/// Telemetry never fails or outlives what it observes. A core that fails on every call never fails +/// the SDK: connect, publish (and its error), app events, logging at every level, a reconnect and +/// teardown run as without telemetry (an exception escaping a telemetry listener would fail the +/// test as an uncaught error). A disposed Room leaves no telemetry listener behind. +@TestOn('vm') +library; + +import 'package:flutter_test/flutter_test.dart'; +import 'package:logging/logging.dart'; + +import 'package:livekit_client/livekit_client.dart'; +import 'package:livekit_client/src/core/signal_client.dart'; +import 'package:livekit_client/src/proto/livekit_models.pb.dart' as lk_models; +import 'package:livekit_client/src/proto/livekit_rtc.pb.dart' as lk_rtc; +import 'package:livekit_client/src/telemetry/telemetry_io.dart' show liveSpanHandles, telemetryScopeFactory; +import 'package:livekit_client/src/uniffi/uniffi_io.dart' as ffi; +import '../mock/e2e_container.dart'; +import '../mock/media_stream_mock.dart'; +import '../mock/peerconnection_mock.dart'; +import '../mock/test_data.dart'; +import '../mock/websocket_mock.dart'; + +void main() { + test('a failing core never fails a call', () async { + final real = telemetryScopeFactory; + final level = Logger.root.level; + addTearDown(() { + telemetryScopeFactory = real; + Logger.root.level = level; + }); + telemetryScopeFactory = () => real() == null ? null : _BrokenScope(); + // Every level: a diagnostic logged while the logger delivers a record must not re-enter it. + Logger.root.level = Level.ALL; + final disposeErrors = []; + final records = Logger.root.onRecord + .where((r) => r.message.contains('error during dispose')) + .listen((r) => disposeErrors.add(r.message)); + addTearDown(records.cancel); + + final container = E2EContainer(); + final room = container.room; + final ws = container.wsConnector; + if (room.telemetry == null) { + markTestSkipped('no native library'); + return; + } + await container.connectRoom(); + room + ..emitTelemetryEvent('app.event') + ..setTelemetryAttribute('app.key', 'value'); + + final track = LocalAudioTrack( + TrackSource.microphone, + FakeMediaStream('local_stream'), + FakeMediaStreamTrack(id: 'mic-1', kind: 'audio'), + const AudioCaptureOptions(), + ); + final publishing = room.localParticipant!.publishAudioTrack(track); + await Future.delayed(const Duration(milliseconds: 50)); + ws.onData( + lk_rtc.SignalResponse( + trackPublished: lk_rtc.TrackPublishedResponse(cid: track.getCid(), track: localAudioTrack), + ).writeToBuffer(), + ); + expect((await publishing).sid, localAudioTrack.sid); + await expectLater( + room.localParticipant!.publishAudioTrack(track), + throwsA(isA()), + reason: 'the publish error, not the core\'s', + ); + + // A warning from one of the Room's handlers goes to its (failing) session. + ws.onData( + lk_rtc.SignalResponse( + streamStateUpdate: lk_rtc.StreamStateUpdate( + streamStates: [lk_rtc.StreamStateInfo(participantSid: 'nobody', trackSid: 'TR_nobody')], + ), + ).writeToBuffer(), + ); + + // A quick reconnect, then the client hangs up. + final handlers = ws.handlers; + ws.onDispose(); + for (var i = 0; i < 200 && identical(ws.handlers, handlers); i++) { + await Future.delayed(const Duration(milliseconds: 10)); + } + ws.onData(lk_rtc.SignalResponse(reconnect: lk_rtc.ReconnectResponse()).writeToBuffer()); + await room.events.waitFor(duration: const Duration(seconds: 5)); + final disconnecting = room.disconnect(); + await Future.delayed(const Duration(milliseconds: 50)); + ws.onData( + lk_rtc.SignalResponse( + leave: lk_rtc.LeaveRequest( + reason: lk_models.DisconnectReason.CLIENT_INITIATED, + action: lk_rtc.LeaveRequest_Action.DISCONNECT, + ), + ).writeToBuffer(), + ); + await disconnecting; + expect(room.connectionState, ConnectionState.disconnected); + await container.dispose(); + + // A Room disposed while still connected. + resetMockDataChannels(); + final connected = E2EContainer(); + await connected.connectRoom(); + await connected.dispose(); + expect(disposeErrors, isEmpty, reason: 'telemetry fails no dispose step'); + }); + + test('a disposed Room leaves no telemetry listener and no span handle behind', () async { + // The SignalClient can outlive its Room; a listener left on it would keep the Room alive, and + // the connect zone its subscriptions keep would hold the connect span. + final baseline = SignalClient(MockWebSocketConnector().connect).events.listeners.length; + final spans = liveSpanHandles; + for (var i = 0; i < 3; i++) { + resetMockDataChannels(); + final container = E2EContainer(); + if (container.room.telemetry == null) { + markTestSkipped('no native library'); + return; + } + final connecting = container.connectRoom(); + expect(liveSpanHandles, spans + 1, reason: 'the open connect span holds its handle'); + await connecting; + expect(liveSpanHandles, spans, reason: 'the ended connect span released it'); + expect(container.client.events.listeners.length, greaterThan(baseline)); + await container.dispose(); + expect(container.client.events.listeners.length, baseline, reason: 'round $i'); + expect(liveSpanHandles, spans, reason: 'round $i'); + } + }); +} + +/// Every call fails, as a Rust panic surfaces in Dart; spans it starts fail the same way. +class _BrokenScope implements ffi.TelemetryScope { + @override + ffi.TelemetrySpan start({required ffi.SpanName name, required ffi.TelemetrySpan? parent}) => _BrokenSpan(); + + @override + dynamic noSuchMethod(Invocation invocation) => throw StateError('core panicked'); +} + +class _BrokenSpan implements ffi.TelemetrySpan { + @override + dynamic noSuchMethod(Invocation invocation) => throw StateError('core panicked'); +} diff --git a/test/telemetry/telemetry_e2e_test.dart b/test/telemetry/telemetry_e2e_test.dart new file mode 100644 index 000000000..c0c21d710 --- /dev/null +++ b/test/telemetry/telemetry_e2e_test.dart @@ -0,0 +1,490 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +/// End to end through the Rust core into a local OpenTelemetry collector: the SDK's mock signal +/// and peer connection, and `otelcol-contrib --config test/telemetry/otelcol.yaml`, which writes +/// every OTLP request as a JSON line. The core reads `LK_TELEMETRY_ENDPOINT` when the pipeline +/// starts, so the story runs only with it set (`LK_TELEMETRY_ENDPOINT=http://127.0.0.1:4319 +/// flutter test test/telemetry/`) and skips otherwise. The pipeline is process-wide, hence one story. +@TestOn('vm') +@Timeout(Duration(seconds: 90)) +library; + +import 'dart:async'; +import 'dart:convert'; +import 'dart:io'; +import 'dart:typed_data'; + +import 'package:flutter/widgets.dart' show AppLifecycleState; + +import 'package:flutter_test/flutter_test.dart'; +import 'package:http/http.dart' as http; + +import 'package:livekit_client/livekit_client.dart'; +import 'package:livekit_client/src/internal/events.dart'; +import 'package:livekit_client/src/proto/livekit_models.pb.dart' as lk_models; +import 'package:livekit_client/src/proto/livekit_rtc.pb.dart' as lk_rtc; +import 'package:livekit_client/src/telemetry/telemetry_io.dart' show exportTimeout, sendExport, storageDirectory; +import 'package:livekit_client/src/uniffi/uniffi_io.dart' as ffi; +import '../core/signal_client_test.dart'; +import '../mock/e2e_container.dart'; +import '../mock/media_stream_mock.dart'; +import '../mock/peerconnection_mock.dart'; +import '../mock/test_data.dart'; + +/// What the file collector writes: test/telemetry/otelcol.yaml's path, or `LK_TELEMETRY_OTLP_FILE` +/// for a collector of your own (another port, another file). +final collectorOutput = Platform.environment['LK_TELEMETRY_OTLP_FILE'] ?? '/tmp/livekit-telemetry-otlp.jsonl'; + +Future main() async { + final binding = TestWidgetsFlutterBinding.ensureInitialized(); + // The test binding answers every HttpClient request with a 400; the transport needs the real + // collector (and the transport tests their local servers). + HttpOverrides.global = null; + + final endpoint = Platform.environment['LK_TELEMETRY_ENDPOINT']; + final skip = endpoint == null + ? 'LK_TELEMETRY_ENDPOINT is not set' + : await _reachable(Uri.parse(endpoint)) + ? null + : 'no collector at $endpoint'; + + test('a call, from connect to opt-out, reaches the collector', skip: skip, () async { + final start = DateTime.now().microsecondsSinceEpoch * 1000; + final marker = 'e2e-${DateTime.now().microsecondsSinceEpoch}'; + resetMockDataChannels(); + + final container = E2EContainer(); + final room = container.room; + final ws = container.wsConnector; + expect(room.telemetry, isNotNull, reason: 'the native library is bundled and the pipeline installed'); + // Someone published before this Room joined: with autoSubscribe the join is the intent. + final early = lk_models.ParticipantInfo( + sid: 'PA_early', + identity: 'early', + state: lk_models.ParticipantInfo_State.ACTIVE, + tracks: [lk_models.TrackInfo(sid: 'TR_early', type: lk_models.TrackType.AUDIO)], + ); + await container.connectRoom(otherParticipants: [early]); + await Future.delayed(const Duration(milliseconds: 500)); + mockInboundTracks['TR_early'] = 'audio'; + container.engine.events.emit( + EngineTrackAddedEvent( + track: FakeMediaStreamTrack(id: 'TR_early', kind: 'audio'), + stream: FakeMediaStream('PA_early|early_stream'), + receiver: null, + ), + ); + await room.events.waitFor(duration: const Duration(seconds: 2)); + + // A track unpublished before any media arrived: its subscribe is cancelled. + lk_rtc.SignalResponse update(List tracks) => lk_rtc.SignalResponse( + update: lk_rtc.ParticipantUpdate(participants: [early.deepCopy()..tracks.addAll(tracks)]), + ); + ws.onData(update([lk_models.TrackInfo(sid: 'TR_gone', type: lk_models.TrackType.VIDEO)]).writeToBuffer()); + await room.events.waitFor(duration: const Duration(seconds: 2)); + ws.onData(update(const []).writeToBuffer()); + await room.events.waitFor(duration: const Duration(seconds: 2)); + + // App data: a correlation attribute on everything from now on (one set, one removed). + room.setTelemetryAttribute('app.call_id', marker); + room.setTelemetryAttribute('app.removed', marker); + room.setTelemetryAttribute('app.removed', null); + + // Publish a microphone through the mock peer connection: the SDK sends AddTrack and waits for + // the server's TrackPublished answer. Publishing it again fails: an `lk.publish` error. + final track = LocalAudioTrack( + TrackSource.microphone, + FakeMediaStream('local_stream'), + FakeMediaStreamTrack(id: 'mic-1', kind: 'audio'), + const AudioCaptureOptions(), + ); + final publishing = room.localParticipant!.publishAudioTrack(track); + await Future.delayed(const Duration(milliseconds: 50)); + ws.onData( + lk_rtc.SignalResponse( + trackPublished: lk_rtc.TrackPublishedResponse(cid: track.getCid(), track: localAudioTrack), + ).writeToBuffer(), + ); + expect((await publishing).sid, localAudioTrack.sid); + await expectLater(room.localParticipant!.publishAudioTrack(track), throwsA(isA())); + // …and a microphone that cannot be opened (no capture device under `flutter test`). + await expectLater(LocalAudioTrack.create(), throwsA(anything)); + + // Subscribe: a remote participant joins and publishes (autoSubscribe: the intent), its track arrives, + // and its media flows (the mock reports growing inbound bytes): first media. + mockInboundTracks[remoteAudioTrack.sid] = 'audio'; + ws.onData(participantJoinResponse.writeToBuffer()); + container.engine.events.emit( + EngineTrackAddedEvent( + track: FakeMediaStreamTrack(id: remoteAudioTrack.sid, kind: 'audio'), + stream: FakeMediaStream('${remoteParticipantData.sid}|remote_stream'), + receiver: null, + ), + ); + await room.events.waitFor(duration: const Duration(seconds: 2)); + await Future.delayed(const Duration(seconds: 3)); // first media, at the core's 1 s polls + + // A warning from one of the Room's own handlers, with no span open: the Room's session. + ws.onData( + lk_rtc.SignalResponse( + streamStateUpdate: lk_rtc.StreamStateUpdate( + streamStates: [ + lk_rtc.StreamStateInfo(participantSid: marker, trackSid: 'TR_nobody', state: lk_rtc.StreamState.ACTIVE), + ], + ), + ).writeToBuffer(), + ); + room.emitTelemetryEvent('e2e.checkpoint', attributes: {'e2e.marker': marker}); + + // The server moves the participant to another room: what follows carries the new identity. + ws.onData( + lk_rtc.SignalResponse( + roomMoved: lk_rtc.RoomMovedResponse( + room: lk_models.Room(sid: 'RM_moved', name: 'moved_room'), + participant: localParticipantData, + token: 'moved-$token', + ), + ).writeToBuffer(), + ); + await room.events.waitFor(duration: const Duration(seconds: 2)); + room.emitTelemetryEvent('e2e.moved', attributes: {'e2e.marker': marker}); + + // Device changes a host never makes on its own. + binding + ..handleAppLifecycleStateChanged(AppLifecycleState.paused) + ..handleAppLifecycleStateChanged(AppLifecycleState.resumed) + ..handleMemoryPressure(); + + // A quick reconnect: the socket drops, the SDK resumes, the server answers. + final handlers = ws.handlers; + ws.onDispose(); + for (var i = 0; i < 200 && identical(ws.handlers, handlers); i++) { + await Future.delayed(const Duration(milliseconds: 10)); + } + expect(ws.uri!.queryParameters['reconnect'], '1', reason: 'a resume, not a re-join'); + ws.onData(lk_rtc.SignalResponse(reconnect: lk_rtc.ReconnectResponse()).writeToBuffer()); + await room.events.waitFor(duration: const Duration(seconds: 5)); + + // A server-refreshed token is taken over without a hiccup. + ws.onData(lk_rtc.SignalResponse(refreshToken: 'refreshed-$token').writeToBuffer()); + + // The client hangs up: the server acknowledges with a Leave. + final disconnecting = room.disconnect(); + await Future.delayed(const Duration(milliseconds: 50)); + ws.onData( + lk_rtc.SignalResponse( + leave: lk_rtc.LeaveRequest( + reason: lk_models.DisconnectReason.CLIENT_INITIATED, + action: lk_rtc.LeaveRequest_Action.DISCONNECT, + ), + ).writeToBuffer(), + ); + await disconnecting; + await container.dispose(); + mockInboundTracks.clear(); + + await _flush(); + final stats = ffi.telemetryStats()!; + expect(stats.cachedBatches, 0, reason: 'the whole call shipped: ${ffi.telemetryDiagnostics()}'); + expect(stats.dropped, 0, reason: ffi.telemetryDiagnostics()); + + // Opt-out: what was not yet sent is deleted, nothing is collected afterwards. + Room().emitTelemetryEvent('e2e.pending', attributes: {'e2e.marker': marker}); + final purged = LiveKitClient.disableTelemetry(); + expect(Room().telemetry, isNull, reason: 'in effect at once: a Room created afterwards collects nothing'); + await purged; + await ffi.telemetryFlush(); // waits for the background purge + final cache = Directory(storageDirectory!); + expect(cache.existsSync() ? cache.listSync() : const [], isEmpty, reason: 'the on-disk cache is purged'); + + // What reached the backend: asserted when it is the file collector (test/telemetry/otelcol.yaml); + // any other backend (a local LGTM stack) is looked up by the call's marker instead. + if (Uri.parse(endpoint!).port != 4319 && !Platform.environment.containsKey('LK_TELEMETRY_OTLP_FILE')) { + print('telemetry e2e: find app.call_id=$marker in the backend'); + return; + } + await Future.delayed(const Duration(seconds: 2)); // the collector's file write + final otlp = OtlpFile(collectorOutput, since: start); + expect(otlp.logs.where((l) => l.eventName == 'custom.e2e.pending'), isEmpty, reason: 'deleted, never uploaded'); + final checkpoint = otlp.logs.singleWhere( + (l) => l.eventName == 'custom.e2e.checkpoint' && l.attributes['e2e.marker'] == marker, + ); + final trace = checkpoint.traceId; + expect(trace, hasLength(32), reason: 'the Room has its own trace'); + final spans = otlp.spans.where((s) => s.traceId == trace).toList(); + final logs = otlp.logs.where((l) => l.traceId == trace).toList(); + + // Connect: one span with this platform's checkpoints; the reconnect is its own span. + final connect = spans.singleWhere((s) => s.name == 'lk.connect'); + expect( + connect.events, + containsAll(['ws_open', 'signal', 'join_recv', 'pc_created', 'answer_sent', 'pc_connected', 'room_connected']), + ); + expect(connect.attributes, containsPair('lk.outcome', 'ok')); + expect(connect.attributes, containsPair('lk.connect.attempt', '1')); + final reconnect = spans.singleWhere((s) => s.name == 'lk.reconnect'); + expect(reconnect.attributes, containsPair('lk.reconnect.reason', 'signal_disconnected')); + expect(reconnect.attributes, containsPair('lk.outcome', 'ok')); + expect(reconnect.events, contains('attempt 1 quick')); + + // Publish: one microphone, one duplicate refused; subscribe: intent → first media. + final publishes = spans.where((s) => s.name == 'lk.publish').toList(); + expect( + publishes.where((s) => s.attributes['lk.outcome'] == 'ok' && s.attributes['lk.track.sid'] == localAudioTrack.sid), + hasLength(1), + ); + expect(publishes.where((s) => s.attributes['error.type'] == 'TrackPublishException'), hasLength(1)); + final subscribes = {for (final s in spans.where((s) => s.name == 'lk.subscribe')) s.attributes['lk.track.sid']: s}; + for (final sid in ['TR_early', remoteAudioTrack.sid]) { + expect(subscribes[sid]?.attributes, containsPair('lk.outcome', 'ok'), reason: sid); + expect(subscribes[sid]?.events, containsAll(['subscribed', 'first_media']), reason: sid); + } + expect(subscribes['TR_gone']?.attributes, containsPair('lk.outcome', 'cancelled'), reason: 'ended before media'); + final joined = subscribes['TR_early']!; + expect( + joined.eventNanos['subscribed']! - joined.startNanos, + greaterThan(400 * 1000 * 1000), + reason: 'a track published before the join is wanted from the join on, not from its arrival', + ); + + // RTC windows from one report per peer connection, each track in its direction. + final windows = logs.where((l) => l.eventName == 'lk.rtc.stats.sample').toList(); + for (final direction in ['outbound', 'inbound']) { + expect( + windows.where( + (w) => w.attributes['lk.track.kind'] == 'audio' && w.attributes['lk.track.direction'] == direction, + ), + isNotEmpty, + reason: '$direction audio window', + ); + } + expect(windows.every((w) => w.attributes['app.call_id'] == marker), isTrue, reason: 'windows carry app attributes'); + + // App data and SDK records. + expect(checkpoint.attributes, containsPair('app.call_id', marker)); + expect(otlp.logs.where((l) => l.attributes.containsKey('app.removed')), isEmpty); + final moved = logs.singleWhere((l) => l.eventName == 'custom.e2e.moved'); + expect(moved.attributes, containsPair('lk.room.name', 'moved_room')); + expect(moved.attributes, containsPair('lk.room.sid', 'RM_moved')); + final warning = logs.singleWhere((l) => l.body == 'Participant not found for sid $marker'); + expect(warning.spanId, isEmpty, reason: 'no span open: filed under the Room\'s session'); + // `lk.room.*` is the scope's current identity when the record ships: after the move. + expect(warning.attributes, containsPair('lk.room.name', 'moved_room')); + expect( + otlp.logs.where((l) => l.eventName.isEmpty).every((l) => l.severity >= 13), + isTrue, + reason: 'log records are warnings and errors only', + ); + + // The session ends once, never on the reconnect. + final ended = logs.singleWhere((l) => l.eventName == 'lk.room.disconnected'); + expect(ended.attributes, containsPair('lk.disconnect.reason', 'client_initiated')); + + // Device: the state's initial values and what this platform can post. + for (final event in ['lk.device.memory.changed', 'lk.device.network.changed', 'lk.device.app_state.changed']) { + expect(otlp.logs.where((l) => l.eventName == event), isNotEmpty, reason: event); + } + for (final event in ['lk.device.thermal.changed', 'lk.device.low_power.changed']) { + expect(otlp.logs.where((l) => l.eventName == event), isEmpty, reason: '$event: no source, reported unknown'); + } + expect(otlp.logs.where((l) => l.attributes['lk.device.app_state'] == 'background'), isNotEmpty); + expect(otlp.logs.where((l) => l.attributes['lk.device.memory.pressure'] == 'warning'), isNotEmpty); + expect( + otlp.logs.where( + (l) => l.eventName == 'lk.device.capture.failed' && l.attributes['lk.device.capture.device'] == 'microphone', + ), + isNotEmpty, + ); + }); + + group('transport', () { + late HttpServer collector; + final client = http.Client(); + + setUp(() async => collector = await HttpServer.bind(InternetAddress.loopbackIPv4, 0)); + tearDown(() => collector.close(force: true)); + + ffi.ExportRequest requestTo(int port) => ffi.ExportRequest( + url: 'http://127.0.0.1:$port/v1/logs', + headers: {'Authorization': 'Bearer token'}, + body: Uint8List.fromList([1, 2, 3]), + ); + + test('passes the collector\'s answer through untouched', () async { + collector.listen((request) async { + final body = await request.fold>([], (all, chunk) => all..addAll(chunk)); + request.response + ..statusCode = body.length == 3 && request.headers.value('authorization') == 'Bearer token' ? 429 : 400 + ..headers.set('Retry-After', '7') + ..write('quota'); + await request.response.close(); + }); + final response = await sendExport(client, requestTo(collector.port)); + expect(response.status, 429); + expect(response.headers['retry-after'], '7'); + expect(utf8.decode(response.body), 'quota'); + }); + + test('never follows a redirect with the token', () async { + final elsewhere = await HttpServer.bind(InternetAddress.loopbackIPv4, 0); + addTearDown(() => elsewhere.close(force: true)); + var followed = false; + elsewhere.listen((request) { + followed = true; + unawaited(request.response.close()); + }); + collector.listen((request) async { + request.response + ..statusCode = 307 + ..headers.set('Location', 'http://127.0.0.1:${elsewhere.port}/v1/logs'); + await request.response.close(); + }); + final response = await sendExport(client, requestTo(collector.port)); + expect(response.status, 307, reason: 'the redirect is the answer'); + expect(followed, isFalse); + }); + + test('a collector that never answers cannot hold the queue', () async { + final unanswered = []; + collector.listen(unanswered.add); // accepts, never replies + exportTimeout = const Duration(milliseconds: 300); + addTearDown(() => exportTimeout = const Duration(seconds: 10)); + await expectLater( + sendExport(client, requestTo(collector.port)).timeout(const Duration(seconds: 5)), + throwsA(isA()), + ); + expect(unanswered, hasLength(1)); + }); + + test('no answer is its only error', () async { + final port = collector.port; + await collector.close(force: true); + await expectLater(sendExport(client, requestTo(port)), throwsA(isA())); + }); + }); +} + +Future _reachable(Uri endpoint) async { + try { + final socket = await Socket.connect(endpoint.host, endpoint.port, timeout: const Duration(seconds: 1)); + socket.destroy(); + return true; + } catch (_) { + return false; + } +} + +/// Ships everything the core holds. +Future _flush() async { + for (var i = 0; i < 10; i++) { + await ffi.telemetryFlush(); // drains the whole cache the network allows + if ((ffi.telemetryStats()?.cachedBatches ?? 0) == 0) break; + } +} + +/// What the collector wrote: OTLP/JSON, one export request per line; only the records stamped at +/// or after `since` (unix nanoseconds). +class OtlpFile { + final logs = []; + final spans = []; + + OtlpFile(String path, {required int since}) { + for (final line in File(path).readAsLinesSync()) { + if (!line.startsWith('{')) continue; + final request = jsonDecode(line) as Map; + for (final scope in _children(request, 'resourceLogs', 'scopeLogs')) { + for (final record in scope['logRecords'] as List? ?? []) { + if (_nanos(record['timeUnixNano']) < since) continue; + logs.add( + OtlpLog( + eventName: record['eventName'] as String? ?? '', + body: (record['body'] as Map?)?['stringValue'] as String?, + traceId: record['traceId'] as String? ?? '', + spanId: record['spanId'] as String? ?? '', + severity: record['severityNumber'] as int? ?? 0, + attributes: _attributes(record['attributes']), + ), + ); + } + } + for (final scope in _children(request, 'resourceSpans', 'scopeSpans')) { + for (final span in scope['spans'] as List? ?? []) { + if (_nanos(span['startTimeUnixNano']) < since) continue; + final events = span['events'] as List? ?? []; + spans.add( + OtlpSpan( + name: span['name'] as String? ?? '', + traceId: span['traceId'] as String? ?? '', + startNanos: _nanos(span['startTimeUnixNano']), + attributes: _attributes(span['attributes']), + events: [for (final event in events) event['name'] as String], + eventNanos: {for (final event in events) event['name'] as String: _nanos(event['timeUnixNano'])}, + ), + ); + } + } + } + } + + static Iterable> _children(Map request, String resources, String scopes) => [ + for (final resource in request[resources] as List? ?? []) + for (final scope in resource[scopes] as List? ?? []) scope as Map, + ]; + + /// OTLP/JSON writes uint64 as a decimal string. + static int _nanos(Object? value) => value is int ? value : int.tryParse('$value') ?? 0; + + /// OTLP/JSON attributes (`[{key, value: {stringValue | intValue | boolValue | doubleValue}}]`) as strings. + static Map _attributes(Object? value) => { + for (final pair in value as List? ?? []) + if ((pair['value'] as Map).isNotEmpty) pair['key'] as String: '${(pair['value'] as Map).values.first}', + }; +} + +class OtlpLog { + final String eventName; + final String? body; + final String traceId; + final String spanId; + final int severity; + final Map attributes; + OtlpLog({ + required this.eventName, + required this.body, + required this.traceId, + required this.spanId, + required this.severity, + required this.attributes, + }); +} + +class OtlpSpan { + final String name; + final String traceId; + final int startNanos; + final Map attributes; + + /// Span event names: the checkpoints (`ws_open`, `first_media`, `attempt 1 quick`, …). + final List events; + final Map eventNanos; + OtlpSpan({ + required this.name, + required this.traceId, + required this.startNanos, + required this.attributes, + required this.events, + required this.eventNanos, + }); +} diff --git a/test/telemetry/telemetry_isolate_test.dart b/test/telemetry/telemetry_isolate_test.dart new file mode 100644 index 000000000..1173e4f6b --- /dev/null +++ b/test/telemetry/telemetry_isolate_test.dart @@ -0,0 +1,81 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +/// The core outlives the isolate that installed its pipeline (hot restart, a recreated engine): +/// it must hold nothing that calls back into a dead isolate. Its own process: the opt-out is +/// process-wide. +@TestOn('vm') +library; + +import 'dart:convert'; +import 'dart:io'; +import 'dart:isolate'; + +import 'package:flutter_test/flutter_test.dart'; + +import 'package:livekit_client/livekit_client.dart'; +import 'package:livekit_client/src/telemetry/telemetry_io.dart' show storageDirectory; +import 'package:livekit_client/src/uniffi/uniffi_io.dart' as ffi; + +void main() { + test('pipelines installed by isolates that are gone are replaced and disabled safely', () async { + // Hot restart twice: each isolate installs the pipeline (replacing its predecessor's), gives it + // a destination and records to ship, sends the app to the background (which uploads them), and + // dies before serving the request. + for (var i = 0; i < 2; i++) { + if (!await Isolate.run(_installAndQueue)) { + markTestSkipped('no native library'); + return; + } + await Future.delayed(const Duration(milliseconds: 500)); // the request is queued + } + // Their successor opts out before its first Room. The core stops the installed pipeline's + // instruments on replace and on opt-out, and withdraws the queued requests: an instrument or a + // pending Rust future it held would call into a dead isolate and abort the VM. + final leftover = File('${storageDirectory!}/leftover')..createSync(recursive: true); + final purged = LiveKitClient.disableTelemetry(); + expect(Room().telemetry, isNull, reason: 'in effect at once'); + await purged; + expect(leftover.existsSync(), isFalse, reason: 'the cache is deleted before any Room installed'); + }); +} + +bool _installAndQueue() { + if (Room().telemetry == null) return false; + // A Cloud project and a token with the observability grant (claims only: the core reads, the + // server would verify). Nothing is sent: the isolate is gone before its first poll. + final scope = ffi.telemetryScope()!..setServer(url: 'wss://isolate-probe.livekit.cloud', token: _grantedToken()); + for (var i = 0; i < 5; i++) { + scope.emitCustom(name: 'isolate.probe', attributes: {'i': '$i'}); + } + // Going to the background uploads the whole cache; no Rust future is left waiting on this isolate. + ffi.telemetrySetDeviceState( + state: ffi.DeviceState( + thermal: ffi.ThermalState.unknown, + appState: ffi.AppState.background, + memory: ffi.MemoryPressure.normal, + network: ffi.NetworkType.wifi, + ), + ); + return true; +} + +String _grantedToken() { + String part(Object json) => base64Url.encode(utf8.encode(jsonEncode(json))).replaceAll('=', ''); + final exp = DateTime.now().add(const Duration(hours: 1)).millisecondsSinceEpoch ~/ 1000; + return '${part({'alg': 'HS256', 'typ': 'JWT'})}.${part({ + 'exp': exp, + 'observability': {'write': true}, + })}.signature'; +} diff --git a/test/telemetry/telemetry_opt_out_submit_test.dart b/test/telemetry/telemetry_opt_out_submit_test.dart new file mode 100644 index 000000000..f879c953d --- /dev/null +++ b/test/telemetry/telemetry_opt_out_submit_test.dart @@ -0,0 +1,85 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +/// A stats read that succeeds after another isolate opted the process out is never submitted, and +/// nothing more is read (telemetry_opt_out_test.dart covers an opt-out in the reading isolate, +/// during a read that fails). Its own process: the opt-out is process-wide. +@TestOn('vm') +library; + +import 'dart:isolate'; + +import 'package:flutter_test/flutter_test.dart'; + +import 'package:livekit_client/src/telemetry/telemetry_io.dart' show telemetryScopeFactory; +import 'package:livekit_client/src/uniffi/uniffi_io.dart' as ffi; +import '../mock/e2e_container.dart'; +import '../mock/peerconnection_mock.dart'; + +void main() { + test('a read that completes after another isolate opted out is not submitted', () async { + final spy = _SpyScope(); + final real = telemetryScopeFactory; + addTearDown(() => telemetryScopeFactory = real); + telemetryScopeFactory = () => real() == null ? null : spy; + + final container = E2EContainer(); + if (container.room.telemetry == null) { + markTestSkipped('no native library'); + return; + } + expect(ffi.telemetryIsDisabled(), isFalse); + await container.connectRoom(); + for (var i = 0; i < 50 && spy.submissions == 0; i++) { + await Future.delayed(const Duration(milliseconds: 100)); + } + expect(spy.submissions, greaterThan(0), reason: 'submitting before the opt-out'); + + // Another isolate (another Flutter engine) opts out while this one's poll waits for a read that + // then succeeds. This isolate's own state never hears of it; only the core does. + var reads = -1, submissions = -1; + MockPeerConnection.onGetStats = () async { + MockPeerConnection.onGetStats = null; + await Isolate.run(ffi.telemetryDisable); + (reads, submissions) = (MockPeerConnection.statsCalls, spy.submissions); + await Future.delayed(const Duration(milliseconds: 200)); + }; + for (var i = 0; i < 50 && reads < 0; i++) { + await Future.delayed(const Duration(milliseconds: 100)); + } + await Future.delayed(const Duration(seconds: 3)); + expect(spy.submissions, submissions, reason: 'the late answer is not submitted'); + expect(MockPeerConnection.statsCalls, reads, reason: 'and nothing more is read'); + await container.dispose(); + }); +} + +/// Counts stats submissions and asks for a poll every second; everything else is accepted and +/// ignored (a span it cannot start is skipped). +class _SpyScope implements ffi.TelemetryScope { + var submissions = 0; + + @override + int statsPollIntervalMs() => 1000; + + @override + void recordPeerStats({ + required List report, + required Map tracks, + required int? timestampNs, + }) => submissions++; + + @override + dynamic noSuchMethod(Invocation invocation) => null; +} diff --git a/test/telemetry/telemetry_opt_out_test.dart b/test/telemetry/telemetry_opt_out_test.dart new file mode 100644 index 000000000..3893772bb --- /dev/null +++ b/test/telemetry/telemetry_opt_out_test.dart @@ -0,0 +1,69 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +/// The opt-out stops what Dart collects for a Room that stays connected. Its own process: the +/// opt-out is process-wide. +@TestOn('vm') +library; + +import 'package:flutter_test/flutter_test.dart'; + +import 'package:livekit_client/livekit_client.dart'; +import 'package:livekit_client/src/proto/livekit_models.pb.dart' as lk_models; +import '../mock/e2e_container.dart'; +import '../mock/peerconnection_mock.dart'; + +void main() { + test('a connected Room collects no stats after the opt-out', () async { + final container = E2EContainer(); + if (container.room.telemetry == null) { + markTestSkipped('no native library'); + return; + } + // A track already in the room waits for media: the core asks for a poll every second. + await container.connectRoom( + otherParticipants: [ + lk_models.ParticipantInfo( + sid: 'PA_other', + identity: 'other', + state: lk_models.ParticipantInfo_State.ACTIVE, + tracks: [lk_models.TrackInfo(sid: 'TR_other', type: lk_models.TrackType.AUDIO)], + ), + ], + ); + for (var i = 0; i < 50 && MockPeerConnection.statsCalls == 0; i++) { + await Future.delayed(const Duration(milliseconds: 100)); + } + expect(MockPeerConnection.statsCalls, greaterThan(0), reason: 'polling before the opt-out'); + + // The app opts out while a poll reads one peer connection, and that read fails: the poll + // must not go on to read the other one. + late Future purged; + var before = -1; + MockPeerConnection.onGetStats = () async { + MockPeerConnection.onGetStats = null; + purged = LiveKitClient.disableTelemetry(); + before = MockPeerConnection.statsCalls; + throw StateError('peer connection closing'); + }; + for (var i = 0; i < 50 && before < 0; i++) { + await Future.delayed(const Duration(milliseconds: 100)); + } + await purged; + await Future.delayed(const Duration(seconds: 3)); + expect(container.room.connectionState, ConnectionState.connected, reason: 'the Room is retained'); + expect(MockPeerConnection.statsCalls, before, reason: 'no stats collected after the opt-out'); + await container.dispose(); + }); +} diff --git a/test/telemetry/telemetry_session_test.dart b/test/telemetry/telemetry_session_test.dart new file mode 100644 index 000000000..0da2109d7 --- /dev/null +++ b/test/telemetry/telemetry_session_test.dart @@ -0,0 +1,165 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +/// What a Room's session hands the core, against a spy session: stats polling under track churn +/// and a stuck peer connection, room updates, and SDK warnings whatever the console level. +@TestOn('vm') +library; + +import 'dart:async'; + +import 'package:flutter_test/flutter_test.dart'; +import 'package:logging/logging.dart'; + +import 'package:livekit_client/livekit_client.dart'; +import 'package:livekit_client/src/proto/livekit_models.pb.dart' as lk_models; +import 'package:livekit_client/src/proto/livekit_rtc.pb.dart' as lk_rtc; +import 'package:livekit_client/src/telemetry/telemetry_io.dart' show telemetryScopeFactory; +import 'package:livekit_client/src/uniffi/uniffi_io.dart' as ffi; +import '../mock/e2e_container.dart'; +import '../mock/peerconnection_mock.dart'; + +void main() { + late _SpyScope spy; + late E2EContainer container; + + setUp(() async { + spy = _SpyScope(); + final real = telemetryScopeFactory; + addTearDown(() => telemetryScopeFactory = real); + telemetryScopeFactory = () => real() == null ? null : spy; + resetMockDataChannels(); + container = E2EContainer(); + if (container.room.telemetry == null) { + markTestSkipped('no native library'); + return; + } + await container.connectRoom( + otherParticipants: [ + lk_models.ParticipantInfo( + sid: 'PA_other', + identity: 'other', + state: lk_models.ParticipantInfo_State.ACTIVE, + tracks: [lk_models.TrackInfo(sid: 'TR_other', type: lk_models.TrackType.AUDIO)], + ), + ], + ); + addTearDown(() { + MockPeerConnection.onGetStats = null; + return container.dispose(); + }); + }); + + /// Track signals (a subscribe intent) every 100 ms for [duration]. + Future churn(Duration duration) async { + final publication = container.room.remoteParticipants.values.first.trackPublications.values.first; + for (final end = DateTime.now().add(duration); DateTime.now().isBefore(end);) { + container.room.telemetry?.subscribeStarted(publication); + await Future.delayed(const Duration(milliseconds: 100)); + } + } + + test('track signals every 100 ms never postpone a 1 s poll', () async { + if (container.room.telemetry == null) return; + final before = spy.submissions; + await churn(const Duration(seconds: 4)); + expect(spy.submissions - before, greaterThanOrEqualTo(4), reason: 'one poll per second, two peer connections'); + }); + + test('one read at a time, whatever signals arrive meanwhile', () async { + if (container.room.telemetry == null) return; + var inFlight = 0, most = 0; + MockPeerConnection.onGetStats = () async { + most = ++inFlight > most ? inFlight : most; + await Future.delayed(const Duration(milliseconds: 1500)); + inFlight--; + }; + await churn(const Duration(seconds: 4)); + expect(most, 1); + }); + + test('a peer connection that never answers holds no poll', () async { + if (container.room.telemetry == null) return; + MockPeerConnection.onGetStats = () { + MockPeerConnection.onGetStats = null; + return Completer().future; // never answers + }; + await Future.delayed(const Duration(seconds: 2)); + final stuck = spy.submissions; + await Future.delayed(const Duration(seconds: 6)); // past the 5 s bound + expect(spy.submissions, greaterThan(stuck), reason: 'the next reads go on'); + }); + + test('a room update refreshes the identity', () async { + if (container.room.telemetry == null) return; + container.wsConnector.onData( + lk_rtc.SignalResponse( + roomUpdate: lk_rtc.RoomUpdate(room: lk_models.Room(name: 'renamed_room')), + ).writeToBuffer(), + ); + await Future.delayed(const Duration(milliseconds: 100)); + expect(spy.room?.name, 'renamed_room'); + expect(spy.room?.sid, isNotNull, reason: 'an update without a sid keeps the one from the join'); + expect(spy.room?.participantIdentity, isNotNull); + }); + + test('SDK warnings reach telemetry whatever the console level', () async { + if (container.room.telemetry == null) return; + hierarchicalLoggingEnabled = true; + final level = logger.level; + addTearDown(() => logger.level = level); + disableLogging(); + final console = []; + final listening = logger.onRecord.listen(console.add); + addTearDown(listening.cancel); + // A warning from one of the Room's handlers: filed under its session. + container.wsConnector.onData( + lk_rtc.SignalResponse( + streamStateUpdate: lk_rtc.StreamStateUpdate( + streamStates: [lk_rtc.StreamStateInfo(participantSid: 'nobody', trackSid: 'TR_nobody')], + ), + ).writeToBuffer(), + ); + await Future.delayed(const Duration(milliseconds: 100)); + expect(spy.logs.map((r) => r.body), contains('Participant not found for sid nobody')); + expect(console, isEmpty, reason: 'the console stays silent'); + }); +} + +/// Records what the session hands the core and asks for a poll every second; everything else is +/// accepted and ignored (a span it cannot start is skipped). +class _SpyScope implements ffi.TelemetryScope { + var submissions = 0; + ffi.RoomIdentity? room; + final logs = []; + + @override + int statsPollIntervalMs() => 1000; + + @override + void recordPeerStats({ + required List report, + required Map tracks, + required int? timestampNs, + }) => submissions++; + + @override + void setRoom({required ffi.RoomIdentity room}) => this.room = room; + + @override + void log({required ffi.LogRecord record}) => logs.add(record); + + @override + dynamic noSuchMethod(Invocation invocation) => null; +} diff --git a/test/uniffi/uniffi_test.dart b/test/uniffi/uniffi_test.dart new file mode 100644 index 000000000..b0f6df841 --- /dev/null +++ b/test/uniffi/uniffi_test.dart @@ -0,0 +1,40 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +@TestOn('vm') +library; + +import 'package:flutter_test/flutter_test.dart'; + +import 'package:livekit_client/src/uniffi/uniffi.dart'; + +void main() { + // Exercises the whole delivery chain rather than any particular API: the + // build hook resolved a cdylib for this target, Native Assets bundled it, + // `@Native` bound the symbol, and a value crossed back from Rust. If the + // bindgen or the hook regresses, this is what fails first. + group('livekit_uniffi', () { + test('is available on native platforms', () { + expect(LiveKitUniffi.isAvailable, isTrue); + }); + + test('buildVersion returns the Rust core version', () { + final version = LiveKitUniffi.buildVersion; + expect(version, isNotEmpty); + // The crate stamps its own semver, so assert the shape rather than a + // literal that would need bumping on every livekit-uniffi release. + expect(version, matches(RegExp(r'^\d+\.\d+\.\d+'))); + }); + }); +} From ada7f668b6ad201ea62f291d06086e2610e33b54 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?B=C5=82az=CC=87ej=20Pankowski?= <86720177+pblazej@users.noreply.github.com> Date: Thu, 1 Oct 2026 13:19:41 +0200 Subject: [PATCH 7/8] ci(telemetry): run local collector alongside tests OTEL collector service provides HTTP endpoint for E2E validation. --- .github/workflows/build.yaml | 13 +++++++++++++ analysis_options.yaml | 4 +++- 2 files changed, 16 insertions(+), 1 deletion(-) diff --git a/.github/workflows/build.yaml b/.github/workflows/build.yaml index 61f8f2616..38ef8882e 100644 --- a/.github/workflows/build.yaml +++ b/.github/workflows/build.yaml @@ -83,6 +83,19 @@ jobs: - uses: ./.github/actions/setup-flutter - name: Dart Test Check run: flutter test + # The telemetry end-to-end test posts OTLP to 127.0.0.1:4319 and reads back what the + # collector wrote (test/telemetry/otelcol.yaml). Only that test gets the endpoint, and a + # collector that did not start fails the job instead of skipping the test. + - name: Telemetry E2E Test + run: | + curl -sSfL -o otelcol.tar.gz "https://github.com/open-telemetry/opentelemetry-collector-releases/releases/download/v0.162.0/otelcol-contrib_0.162.0_linux_amd64.tar.gz" + echo "fcc063749f730f8c21fe29f2d340ff174f5f1c5885bd3156fb6c985a3036fcc3 otelcol.tar.gz" | sha256sum -c - + tar -xzf otelcol.tar.gz otelcol-contrib + nohup ./otelcol-contrib --config test/telemetry/otelcol.yaml > "$RUNNER_TEMP/otelcol.log" 2>&1 & + timeout 30 bash -c 'until curl -s -o /dev/null http://127.0.0.1:4319; do sleep 1; done' || { cat "$RUNNER_TEMP/otelcol.log"; exit 1; } + flutter test test/telemetry/ + env: + LK_TELEMETRY_ENDPOINT: http://127.0.0.1:4319 build-for-android: name: Android diff --git a/analysis_options.yaml b/analysis_options.yaml index 1c0c39cf8..3f11b2aa7 100644 --- a/analysis_options.yaml +++ b/analysis_options.yaml @@ -22,12 +22,14 @@ analyzer: avoid_print: ignore deprecated_member_use_from_same_package: ignore - # Exclude protobuf files + # Exclude generated files: protobuf, and json_serializable output. Neither is + # hand-edited, so lint hits there can only be fixed by changing the generator. exclude: - "**/*.pb.dart" - "**/*.pbenum.dart" - "**/*.pbjson.dart" - "**/*.pbserver.dart" + - "**/*.g.dart" # - 'web/*.dart' # Xcode vendors Swift package checkouts under build when Swift Package # Manager is enabled and this package ships Package.swift. From 8755fbeedcd3ac11058232f82837f8537041c234 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?B=C5=82az=CC=87ej=20Pankowski?= <86720177+pblazej@users.noreply.github.com> Date: Thu, 1 Oct 2026 13:19:42 +0200 Subject: [PATCH 8/8] chore(telemetry): add changesets and UniFFI development guidance Changesets for client telemetry and uniffi-rust-core dev wiring. AGENTS.md documents UniFFI setup and local development: Native Assets, overrides and build hooks. --- .changes/client-telemetry | 1 + .changes/uniffi-rust-core-dev-wiring | 1 + .gitignore | 3 +++ AGENTS.md | 22 ++++++++++++++++++++++ 4 files changed, 27 insertions(+) create mode 100644 .changes/client-telemetry create mode 100644 .changes/uniffi-rust-core-dev-wiring diff --git a/.changes/client-telemetry b/.changes/client-telemetry new file mode 100644 index 000000000..3dbddef56 --- /dev/null +++ b/.changes/client-telemetry @@ -0,0 +1 @@ +minor type="added" "Client telemetry: every Room reports its connect, reconnect, publish and subscribe spans, RTC statistics, SDK warnings and device state to its LiveKit Cloud project through the shared Rust core (native platforms only), with LiveKitClient.disableTelemetry() to opt out and Room.emitTelemetryEvent / Room.setTelemetryAttribute for app events and correlation ids" diff --git a/.changes/uniffi-rust-core-dev-wiring b/.changes/uniffi-rust-core-dev-wiring new file mode 100644 index 000000000..bc9ac8dbc --- /dev/null +++ b/.changes/uniffi-rust-core-dev-wiring @@ -0,0 +1 @@ +minor type="added" "Wire in the livekit_uniffi Rust core behind a native-only facade, delivered as a bundled cdylib via Native Assets" diff --git a/.gitignore b/.gitignore index 1f1d37a82..076834508 100644 --- a/.gitignore +++ b/.gitignore @@ -89,3 +89,6 @@ lib/generated_plugin_registrant.dart # Test files - ignore any binary files in testfiles directory testfiles/*.bin + +# Local livekit_uniffi wiring (see AGENTS.md) +pubspec_overrides.yaml diff --git a/AGENTS.md b/AGENTS.md index f77196e0e..c36d9c336 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -32,6 +32,28 @@ CI (`build.yaml`) runs all of the above plus example-app builds for every platfo Web/native divergence is handled with conditional imports (e.g. `track/processor_native.dart` vs `processor_web.dart`) — new platform-specific code should follow that pattern. +## The Rust core (`livekit_uniffi`) + +`lib/src/uniffi/` wraps `livekit_uniffi`, a Dart package generated from the `livekit-uniffi` crate in the sibling `rust-sdks` repo. It reaches Rust through Dart's Native Assets: the package's `hook/build.dart` bundles a `cdylib` into the host app and the generated bindings call into it with `@Native`. This is why the SDK requires Flutter >= 3.38 / Dart >= 3.10. + +There is no dynamic library to load on the web, so `uniffi.dart` splits native/web the same way the rest of the SDK does. **`uniffi_io.dart` is the only file allowed to import `package:livekit_uniffi/...`** — importing it from anywhere reachable on web pulls `dart:ffi` into a web compile and breaks `flutter build web`/`--wasm`. Guard calls with `LiveKitUniffi.isAvailable`. + +### Local development loop + +`livekit_uniffi` is published on pub.dev and the normal dependency resolves it; its build hook downloads the matching prebuilt library from the `livekit-uniffi` GitHub release, so no Rust toolchain is needed. To develop against unreleased crate changes, point the package and the example at a `rust-sdks` checkout with a `pubspec_overrides.yaml` next to each `pubspec.yaml` (git-ignored, so nothing machine-specific is committed; without one the published package is used). Each path is relative to its own file, and overrides do not propagate from a dependency, hence both: + +```sh +printf 'dependency_overrides:\n livekit_uniffi:\n path: ../rust-sdks/livekit-uniffi/packages/dart\n' > pubspec_overrides.yaml +printf 'dependency_overrides:\n livekit_uniffi:\n path: ../../rust-sdks/livekit-uniffi/packages/dart\n' > example/pubspec_overrides.yaml +(cd ../rust-sdks/livekit-uniffi && cargo make dart-package) # generates packages/dart/: bindings, pubspec, build hook, host cdylib +flutter pub get +flutter test test/uniffi/ # smoke test: calls buildVersion() across the FFI boundary +``` + +Requires `cargo-make`, `protoc` and `tera`. The override rewrites `pubspec.lock` to the path source: keep that change out of commits. Re-run `cargo make dart-package` whenever the crate's exported surface changes — the build hook tracks the copied library, so a stale one won't be silently reused. + +Two things to know about that hook: it picks the locally built library purely on target *OS*, not architecture, so a host build can be bundled into an iOS or Android build by mistake — verify the desktop target first when debugging. And its download mode (used when no local library is present) fetches `build-.zip` from the `livekit-uniffi` GitHub release matching the package version and verifies its SHA-256, so a version whose release lacks assets fails at build time rather than at runtime. + ## Common pitfalls (from issue history) - `flutter_webrtc` is pinned to an exact version on purpose: livekit_client and flutter_webrtc must agree on the same WebRTC-SDK native pods, and mismatches break user builds (CocoaPods conflicts). Bump it only in sync with a matching WebRTC-SDK version.