websocket_service.dart 4.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164
  1. import 'dart:async';
  2. import 'dart:convert';
  3. import 'dart:io';
  4. import 'dart:math' as math;
  5. import 'package:web_socket_channel/web_socket_channel.dart';
  6. import 'package:web_socket_channel/io.dart';
  7. import 'package:asr_client/common/constants.dart';
  8. import 'package:asr_client/models/websocket_message.dart';
  9. export 'package:asr_client/models/websocket_message.dart' show ServerMessage;
  10. enum ConnectionState { disconnected, connecting, connected, error }
  11. abstract interface class WebSocketService {
  12. Stream<ServerMessage> get messageStream;
  13. Stream<ConnectionState> get connectionStateStream;
  14. ConnectionState get currentState;
  15. Future<void> connect(String url);
  16. void send(WebSocketMessage message);
  17. void sendBinary(List<int> bytes);
  18. void disconnect();
  19. }
  20. final class WebSocketServiceImpl implements WebSocketService {
  21. WebSocketServiceImpl();
  22. WebSocketChannel? _channel;
  23. final _messageController = StreamController<ServerMessage>.broadcast();
  24. final _stateController = StreamController<ConnectionState>.broadcast();
  25. Timer? _pingTimer;
  26. Timer? _reconnectTimer;
  27. String? _currentUrl;
  28. bool _shouldReconnect = false;
  29. var _reconnectAttempt = 0;
  30. @override
  31. Stream<ServerMessage> get messageStream => _messageController.stream;
  32. @override
  33. Stream<ConnectionState> get connectionStateStream => _stateController.stream;
  34. @override
  35. ConnectionState get currentState => _lastState;
  36. ConnectionState _lastState = ConnectionState.disconnected;
  37. void _setState(ConnectionState state) {
  38. _lastState = state;
  39. if (!_stateController.isClosed) {
  40. _stateController.add(state);
  41. }
  42. }
  43. @override
  44. Future<void> connect(String url) async {
  45. if (_lastState == ConnectionState.connecting ||
  46. _lastState == ConnectionState.connected) {
  47. return;
  48. }
  49. _currentUrl = url;
  50. _shouldReconnect = true;
  51. await _connectInternal(url, rethrowError: true);
  52. }
  53. Future<void> _connectInternal(String url, {bool rethrowError = false}) async {
  54. _setState(ConnectionState.connecting);
  55. try {
  56. final socket = await WebSocket.connect(
  57. url,
  58. ).timeout(const Duration(seconds: 10));
  59. _channel = IOWebSocketChannel(socket);
  60. _channel!.stream.listen(
  61. (event) {
  62. if (event is String) {
  63. final message = ServerMessage.fromJsonString(event);
  64. _messageController.add(message);
  65. }
  66. },
  67. onError: (Object error) {
  68. _setState(ConnectionState.error);
  69. _messageController.add(
  70. ErrorMessage(message: 'WebSocket error: $error'),
  71. );
  72. _scheduleReconnect();
  73. },
  74. onDone: () {
  75. _setState(ConnectionState.disconnected);
  76. _scheduleReconnect();
  77. },
  78. );
  79. _setState(ConnectionState.connected);
  80. _reconnectAttempt = 0;
  81. _startPingTimer();
  82. } on Exception catch (e) {
  83. _setState(ConnectionState.error);
  84. _messageController.add(ErrorMessage(message: 'Connection failed: $e'));
  85. if (rethrowError) rethrow;
  86. _scheduleReconnect();
  87. }
  88. }
  89. void _startPingTimer() {
  90. _pingTimer?.cancel();
  91. _pingTimer = Timer.periodic(AppConstants.pingInterval, (_) {
  92. _channel?.sink.add(jsonEncode({'type': 'ping'}));
  93. });
  94. }
  95. void _scheduleReconnect() {
  96. _pingTimer?.cancel();
  97. _channel = null;
  98. if (!_shouldReconnect || _currentUrl == null) return;
  99. final delay = Duration(
  100. seconds: math.min(
  101. AppConstants.reconnectBaseDelay.inSeconds *
  102. math.pow(2, _reconnectAttempt).toInt(),
  103. AppConstants.reconnectMaxDelay.inSeconds,
  104. ),
  105. );
  106. _reconnectAttempt++;
  107. _reconnectTimer?.cancel();
  108. _reconnectTimer = Timer(delay, () {
  109. if (_shouldReconnect && _currentUrl != null) {
  110. unawaited(_connectInternal(_currentUrl!, rethrowError: false));
  111. }
  112. });
  113. }
  114. @override
  115. void send(WebSocketMessage message) {
  116. if (_channel != null && _lastState == ConnectionState.connected) {
  117. _channel!.sink.add(message.toJsonString());
  118. }
  119. }
  120. @override
  121. void sendBinary(List<int> bytes) {
  122. if (_channel != null && _lastState == ConnectionState.connected) {
  123. _channel!.sink.add(bytes);
  124. }
  125. }
  126. @override
  127. void disconnect() {
  128. _shouldReconnect = false;
  129. _reconnectTimer?.cancel();
  130. _pingTimer?.cancel();
  131. _channel?.sink.close();
  132. _channel = null;
  133. _setState(ConnectionState.disconnected);
  134. }
  135. void dispose() {
  136. disconnect();
  137. _messageController.close();
  138. _stateController.close();
  139. }
  140. }