Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .changes/client-telemetry
Original file line number Diff line number Diff line change
@@ -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"
1 change: 1 addition & 0 deletions .changes/uniffi-rust-core-dev-wiring
Original file line number Diff line number Diff line change
@@ -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"
13 changes: 13 additions & 0 deletions .github/workflows/build.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
3 changes: 3 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -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
22 changes: 22 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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-<triple>.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.
Expand Down
4 changes: 3 additions & 1 deletion analysis_options.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
3 changes: 3 additions & 0 deletions example/lib/main.dart
Original file line number Diff line number Diff line change
Expand Up @@ -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');
Expand Down Expand Up @@ -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());
}

Expand Down
8 changes: 8 additions & 0 deletions example/lib/pages/connect.dart
Original file line number Diff line number Diff line change
Expand Up @@ -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 }

Expand Down Expand Up @@ -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),
),
],
);
}
Expand Down
18 changes: 18 additions & 0 deletions example/lib/utils.dart
Original file line number Diff line number Diff line change
@@ -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<void> 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';
}
}
1 change: 1 addition & 0 deletions example/pubspec.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ dependencies:
livekit_client:
path: ../


dev_dependencies:
flutter_test:
sdk: flutter
Expand Down
1 change: 1 addition & 0 deletions lib/livekit_client.dart
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
40 changes: 31 additions & 9 deletions lib/src/core/engine.dart
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -145,6 +146,13 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {

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;
Expand Down Expand Up @@ -258,6 +266,7 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
connectOptions: this.connectOptions,
roomOptions: this.roomOptions,
);
telemetry?.step(ConnectStep.signal);

// wait for join response
await events.waitFor<EngineJoinResponseEvent>(
Expand All @@ -278,6 +287,9 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
'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');
Expand Down Expand Up @@ -321,6 +333,8 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
_isReconnecting = false;

_clearPendingReconnect();
_reconnectSpan?.cancel();
_reconnectSpan = null;
}

@internal
Expand Down Expand Up @@ -696,6 +710,7 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
publisher?.onOffer = (offer) {
logger.fine('publisher onOffer');
signalClient.sendOffer(offer);
telemetry?.step(ConnectStep.offerSent);
};

// in subscriber primary mode, server side opens sub data channels.
Expand Down Expand Up @@ -1094,10 +1109,13 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {

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();

Expand Down Expand Up @@ -1165,6 +1183,7 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
final fullReconnect = fullReconnectOnNext;
fullReconnectOnNext = false;
_attemptIsFullReconnect = fullReconnect;
_reconnectSpan?.attempt(_reconnectAttempts + 1, full: fullReconnect);

var succeeded = false;
try {
Expand All @@ -1182,18 +1201,15 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
);
}

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');
Expand All @@ -1214,6 +1230,8 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
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
Expand Down Expand Up @@ -1420,6 +1438,7 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {

void _setUpSignalListeners() => _signalListener
..on<SignalJoinResponseEvent>((event) async {
telemetry?.step(ConnectStep.joinRecv);
// create peer connections
_subscriberPrimary = event.response.subscriberPrimary;
_serverInfo = event.response.serverInfo;
Expand All @@ -1445,6 +1464,7 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {

if (publisher == null && subscriber == null) {
await _createPeerConnections(rtcConfiguration);
telemetry?.step(ConnectStep.pcCreated);
}

if (!_subscriberPrimary || event.response.fastPublish) {
Expand Down Expand Up @@ -1495,6 +1515,7 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
})
..on<SignalConnectedEvent>((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
Expand Down Expand Up @@ -1544,6 +1565,7 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
logger.finer('sdp: ${answer.sdp}');
await subscriber!.pc.setLocalDescription(answer);
signalClient.sendAnswer(answer);
telemetry?.step(ConnectStep.answerSent);
} catch (_) {
logger.severe('[$objectId] Failed to createAnswer()');
}
Expand Down Expand Up @@ -1620,7 +1642,7 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
DisconnectReason reason = DisconnectReason.clientInitiated,
}) async {
_isClosed = true;
events.emit(EngineClosingEvent());
events.emit(EngineClosingEvent(reason: reason));
if (connectionState == ConnectionState.connected) {
await signalClient.sendLeave();
} else {
Expand Down
61 changes: 55 additions & 6 deletions lib/src/core/room.dart
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -122,6 +123,33 @@ class Room extends DisposableChangeNotifier with EventsEmittable<RoomEvent> {

@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.<name>` 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<String, String> 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<EngineEvent> _engineListener;
//
Expand Down Expand Up @@ -179,12 +207,16 @@ class Room extends DisposableChangeNotifier with EventsEmittable<RoomEvent> {
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);
Expand Down Expand Up @@ -273,6 +305,23 @@ class Room extends DisposableChangeNotifier with EventsEmittable<RoomEvent> {
ConnectOptions? connectOptions,
@Deprecated('deprecated, please use roomOptions in Room constructor') RoomOptions? roomOptions,
FastConnectOptions? fastConnectOptions,
}) {
Future<void> connect() => _connect(
url,
token,
connectOptions: connectOptions,
roomOptions: roomOptions,
fastConnectOptions: fastConnectOptions,
);
return telemetry?.connect(url, token, connect) ?? connect();
}

Future<void> _connect(
String url,
String token, {
ConnectOptions? connectOptions,
RoomOptions? roomOptions,
FastConnectOptions? fastConnectOptions,
}) async {
var effectiveRoomOptions = roomOptions ?? this.roomOptions;
if (lkPlatformIs(PlatformType.web) &&
Expand Down
Loading
Loading