websocket_service.dart 4.1 KB

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