| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134 |
- import 'dart:async';
- import 'dart:convert';
- import 'package:dart_jsonwebtoken/dart_jsonwebtoken.dart';
- import 'package:flutter/foundation.dart';
- import 'package:livekit_client/livekit_client.dart';
- import 'package:asr_client/models/websocket_message.dart';
- /// Wraps LiveKit Room connection and relays transcription data messages.
- abstract interface class LiveKitService {
- Stream<ServerMessage> get messageStream;
- bool get isConnected;
- Future<void> connect({required String url, required String room, required String identity, String token = ''});
- Future<void> disconnect();
- }
- final class LiveKitServiceImpl implements LiveKitService {
- LiveKitServiceImpl({Room? room}) : _room = room ?? Room();
- Room _room;
- final _messageController = StreamController<ServerMessage>.broadcast();
- var _connected = false;
- EventsListener<RoomEvent>? _listener;
- @override
- Stream<ServerMessage> get messageStream => _messageController.stream;
- @override
- bool get isConnected => _connected;
- @override
- Future<void> connect({
- required String url,
- required String room,
- required String identity,
- String token = '',
- }) async {
- if (_connected) return;
- final effectiveToken = token.isNotEmpty ? token : _makeToken(room, identity);
- debugPrint('[lk] connecting to $url room=$room');
- // Listen for room events (including connection failures)
- _listener = _room.createListener()
- ..on<DataReceivedEvent>((e) {
- debugPrint('[lk] data received topic=${e.topic} bytes=${e.data.length}');
- if (e.topic == 'transcription') {
- try {
- final json = jsonDecode(utf8.decode(e.data)) as Map<String, dynamic>;
- final message = ServerMessage.fromJson(json);
- _messageController.add(message);
- } catch (ex) {
- debugPrint('[lk] bad data: $ex');
- }
- }
- })
- ..on<RoomDisconnectedEvent>((e) {
- debugPrint('[lk] disconnected: ${e.reason}');
- _connected = false;
- })
- ..on<RoomReconnectingEvent>((e) {
- debugPrint('[lk] reconnecting');
- })
- ..on<RoomReconnectedEvent>((e) {
- debugPrint('[lk] reconnected');
- });
- try {
- await _room.connect(
- url,
- effectiveToken,
- );
- } catch (e, st) {
- debugPrint('[lk] connect failed: $e');
- debugPrint('[lk] stack: $st');
- rethrow;
- }
- debugPrint('[lk] connected, publishing mic...');
- try {
- await _room.localParticipant!.setMicrophoneEnabled(true);
- debugPrint('[lk] mic published');
- } catch (e, st) {
- debugPrint('[lk] mic FAILED: $e\n$st');
- }
- _connected = true;
- _messageController.add(const ConnectedMessage(model: {}));
- }
- @override
- Future<void> disconnect() async {
- if (!_connected) return;
- _connected = false;
- _listener?.dispose();
- _listener = null;
- await _room.disconnect();
- // Room instances are single-use; create a fresh one for the next connect.
- _room = Room();
- }
- String _makeToken(String room, String identity) {
- // Dev token: for localhost LiveKit Server with key "devkey"
- // In production, tokens should be generated server-side.
- const key = 'devkey';
- const secret = 'secretsecretsecretsecretsecret12';
- final jwt = JWT(
- {
- 'iss': key,
- 'sub': identity,
- 'name': identity,
- 'nbf': DateTime.now()
- .subtract(const Duration(minutes: 1))
- .millisecondsSinceEpoch ~/
- 1000,
- 'exp': DateTime.now()
- .add(const Duration(hours: 6))
- .millisecondsSinceEpoch ~/
- 1000,
- 'video': {
- 'roomJoin': true,
- 'room': room,
- 'canPublish': true,
- 'canSubscribe': true,
- 'canPublishData': true,
- },
- },
- );
- return jwt.sign(SecretKey(secret));
- }
- }
|