From c0cb91b94001172b7efdcc6990c4b7b022d307bd Mon Sep 17 00:00:00 2001 From: Damodar Lohani Date: Thu, 7 Sep 2023 07:45:07 +0545 Subject: [PATCH 1/6] fix for realtime multiple subscription --- .../flutter/lib/src/realtime_mixin.dart.twig | 82 +++++++++++-------- .../lib/src/realtime_subscription.dart.twig | 11 ++- .../src/realtime_subscription_test.dart.twig | 20 ++--- 3 files changed, 66 insertions(+), 47 deletions(-) diff --git a/templates/flutter/lib/src/realtime_mixin.dart.twig b/templates/flutter/lib/src/realtime_mixin.dart.twig index accf973c68..8ee6394f72 100644 --- a/templates/flutter/lib/src/realtime_mixin.dart.twig +++ b/templates/flutter/lib/src/realtime_mixin.dart.twig @@ -2,7 +2,7 @@ import 'dart:async'; import 'dart:convert'; import 'package:flutter/foundation.dart'; import 'package:web_socket_channel/web_socket_channel.dart'; -import 'package:web_socket_channel/status.dart'; +import 'package:web_socket_channel/status.dart' as status; import 'exception.dart'; import 'realtime_subscription.dart'; import 'client.dart'; @@ -15,15 +15,20 @@ typedef GetFallbackCookie = String? Function(); mixin RealtimeMixin { late Client client; - final Map>> _channels = {}; + final Set _channels = {}; WebSocketChannel? _websok; String? _lastUrl; late WebSocketFactory getWebSocket; GetFallbackCookie? getFallbackCookie; int? get closeCode => _websok?.closeCode; + int _subscriptionsCounter = 0; + Map _subscriptions = {}; + bool _notifyDone = true; + StreamSubscription? _websocketSubscription; Future _closeConnection() async { - await _websok?.sink.close(normalClosure); + await _websocketSubscription?.cancel(); + await _websok?.sink.close(status.normalClosure, 'Ending session'); _lastUrl = null; } @@ -36,14 +41,16 @@ mixin RealtimeMixin { if (_lastUrl == uri.toString() && _websok?.closeCode == null) { return; } + _notifyDone = false; await _closeConnection(); _lastUrl = uri.toString(); _websok = await getWebSocket(uri); + _notifyDone = true; } debugPrint('subscription: $_lastUrl'); try { - _websok?.stream.listen((response) { + _websocketSubscription = _websok?.stream.listen((response) { final data = RealtimeResponse.fromJson(response); switch (data.type) { case 'error': @@ -67,28 +74,25 @@ mixin RealtimeMixin { break; case 'event': final message = RealtimeMessage.fromMap(data.data); - for(var channel in message.channels) { - if (_channels[channel] != null) { - for( var stream in _channels[channel]!) { - stream.sink.add(message); + for (var subscription in _subscriptions.values) { + for (var channel in message.channels) { + if (subscription.channels.contains(channel)) { + subscription.controller.add(message); } } } break; } }, onDone: () { - for (var list in _channels.values) { - for (var stream in list) { - stream.close(); - } + if (!_notifyDone) return; + for (var subscription in _subscriptions.values) { + subscription.close(); } _channels.clear(); _closeConnection(); }, onError: (err, stack) { - for (var list in _channels.values) { - for (var stream in list) { - stream.sink.addError(err, stack); - } + for (var subscription in _subscriptions.values) { + subscription.controller.addError(err, stack); } if (_websok?.closeCode != null && _websok?.closeCode != 1008) { debugPrint("Reconnecting in one second."); @@ -96,19 +100,19 @@ mixin RealtimeMixin { } }); } catch (e) { - if (e is {{spec.title | caseUcfirst}}Exception) { + if (e is AppwriteException) { rethrow; } if (e is WebSocketChannelException) { - throw {{spec.title | caseUcfirst}}Exception(e.message); + throw AppwriteException(e.message); } - throw {{spec.title | caseUcfirst}}Exception(e.toString()); + throw AppwriteException(e.toString()); } } Uri _prepareUri() { if (client.endPointRealtime == null) { - throw {{spec.title | caseUcfirst}}Exception( + throw AppwriteException( "Please set endPointRealtime to connect to realtime server"); } var uri = Uri.parse(client.endPointRealtime!); @@ -118,7 +122,7 @@ mixin RealtimeMixin { port: uri.port, queryParameters: { "project": client.config['project'], - "channels[]": _channels.keys.toList(), + "channels[]": _channels.toList(), }, path: uri.path + "/realtime", ); @@ -126,35 +130,41 @@ mixin RealtimeMixin { RealtimeSubscription subscribeTo(List channels) { StreamController controller = StreamController.broadcast(); - for(var channel in channels) { - if (!_channels.containsKey(channel)) { - _channels[channel] = []; - } - _channels[channel]!.add(controller); - } + _channels.addAll(channels); Future.delayed(Duration.zero, () => _createSocket()); + int counter = _subscriptionsCounter++; RealtimeSubscription subscription = RealtimeSubscription( - stream: controller.stream, + controller: controller, + channels: channels, close: () async { + _subscriptions.remove(counter); + _subscriptionsCounter--; controller.close(); - for(var channel in channels) { - _channels[channel]!.remove(controller); - if (_channels[channel]!.isEmpty) { - _channels.remove(channel); - } - } - if(_channels.isNotEmpty) { + _cleanup(channels); + + if (_channels.isNotEmpty) { await Future.delayed(Duration.zero, () => _createSocket()); } else { await _closeConnection(); } }); + _subscriptions[counter] = subscription; return subscription; } + void _cleanup(List channels) { + for (var channel in channels) { + bool found = _subscriptions.values + .any((subscription) => subscription.channels.contains(channel)); + if (!found) { + _channels.remove(channel); + } + } + } + void handleError(RealtimeResponse response) { if (response.data['code'] == 1008) { - throw {{spec.title | caseUcfirst}}Exception(response.data["message"], response.data["code"]); + throw AppwriteException(response.data["message"], response.data["code"]); } else { debugPrint("Reconnecting in one second."); Future.delayed(const Duration(seconds: 1), () { diff --git a/templates/flutter/lib/src/realtime_subscription.dart.twig b/templates/flutter/lib/src/realtime_subscription.dart.twig index e45d3b4191..1707691691 100644 --- a/templates/flutter/lib/src/realtime_subscription.dart.twig +++ b/templates/flutter/lib/src/realtime_subscription.dart.twig @@ -1,3 +1,5 @@ +import 'dart:async'; + import 'realtime_message.dart'; /// Realtime Subscription @@ -5,9 +7,16 @@ class RealtimeSubscription { /// Stream of [RealtimeMessage]s final Stream stream; + final StreamController controller; + + /// List of channels + List channels; + /// Closes the subscription final Future Function() close; /// Initializes a [RealtimeSubscription] - RealtimeSubscription({required this.stream, required this.close}); + RealtimeSubscription( + {required this.close, required this.channels, required this.controller}) + : stream = controller.stream; } diff --git a/templates/flutter/test/src/realtime_subscription_test.dart.twig b/templates/flutter/test/src/realtime_subscription_test.dart.twig index 1cf82e1630..f74ff25347 100644 --- a/templates/flutter/test/src/realtime_subscription_test.dart.twig +++ b/templates/flutter/test/src/realtime_subscription_test.dart.twig @@ -1,20 +1,20 @@ -import 'package:mockito/mockito.dart'; -import 'package:{{language.params.packageName}}/src/realtime_message.dart'; -import 'package:{{language.params.packageName}}/src/realtime_subscription.dart'; +import 'package:appwrite/src/realtime_message.dart'; +import 'package:appwrite/src/realtime_subscription.dart'; import 'package:flutter_test/flutter_test.dart'; - -class MockStream extends Mock implements Stream {} - - +import 'dart:async'; void main() { group('RealtimeSubscription', () { - final mockStream = MockStream(); + final mockStream = StreamController.broadcast(); final mockCloseFunction = () async {}; - final subscription = RealtimeSubscription(stream: mockStream, close: mockCloseFunction); + final subscription = RealtimeSubscription( + controller: mockStream, + close: mockCloseFunction, + channels: ['documents']); test('should have the correct stream and close function', () { - expect(subscription.stream, equals(mockStream)); + expect(subscription.controller, equals(mockStream)); + expect(subscription.stream, equals(mockStream.stream)); expect(subscription.close, equals(mockCloseFunction)); }); }); From d56ecef0ebe06087771e1805552cbe2bac2b2320 Mon Sep 17 00:00:00 2001 From: Damodar Lohani Date: Sun, 17 Dec 2023 00:50:11 +0000 Subject: [PATCH 2/6] use time mircosecond instead of counter for subscription id --- templates/flutter/lib/src/realtime_mixin.dart.twig | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/templates/flutter/lib/src/realtime_mixin.dart.twig b/templates/flutter/lib/src/realtime_mixin.dart.twig index 8ee6394f72..ec3727aabb 100644 --- a/templates/flutter/lib/src/realtime_mixin.dart.twig +++ b/templates/flutter/lib/src/realtime_mixin.dart.twig @@ -132,12 +132,12 @@ mixin RealtimeMixin { StreamController controller = StreamController.broadcast(); _channels.addAll(channels); Future.delayed(Duration.zero, () => _createSocket()); - int counter = _subscriptionsCounter++; + int id = DateTime.now().microsecondsSinceEpoch; RealtimeSubscription subscription = RealtimeSubscription( controller: controller, channels: channels, close: () async { - _subscriptions.remove(counter); + _subscriptions.remove(id); _subscriptionsCounter--; controller.close(); _cleanup(channels); @@ -148,7 +148,7 @@ mixin RealtimeMixin { await _closeConnection(); } }); - _subscriptions[counter] = subscription; + _subscriptions[id] = subscription; return subscription; } From bba375a2f6ed88b293bb4c3cc4e2cb4da87b1e8b Mon Sep 17 00:00:00 2001 From: Damodar Lohani Date: Mon, 8 Apr 2024 12:21:26 +0545 Subject: [PATCH 3/6] Update realtime_mixin.dart.twig --- templates/flutter/lib/src/realtime_mixin.dart.twig | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/templates/flutter/lib/src/realtime_mixin.dart.twig b/templates/flutter/lib/src/realtime_mixin.dart.twig index ec3727aabb..1e99f40a26 100644 --- a/templates/flutter/lib/src/realtime_mixin.dart.twig +++ b/templates/flutter/lib/src/realtime_mixin.dart.twig @@ -100,19 +100,19 @@ mixin RealtimeMixin { } }); } catch (e) { - if (e is AppwriteException) { + if (e is {{spec.title | caseUcfirst}}Exception) { rethrow; } if (e is WebSocketChannelException) { - throw AppwriteException(e.message); + throw {{spec.title | caseUcfirst}}Exception(e.message); } - throw AppwriteException(e.toString()); + throw {{spec.title | caseUcfirst}}Exception(e.toString()); } } Uri _prepareUri() { if (client.endPointRealtime == null) { - throw AppwriteException( + throw {{spec.title | caseUcfirst}}Exception( "Please set endPointRealtime to connect to realtime server"); } var uri = Uri.parse(client.endPointRealtime!); From 18ad858ea12af4717263b59a8bc18b663f03ad75 Mon Sep 17 00:00:00 2001 From: Damodar Lohani Date: Mon, 15 Apr 2024 06:27:22 +0545 Subject: [PATCH 4/6] Update realtime_mixin.dart.twig --- templates/flutter/lib/src/realtime_mixin.dart.twig | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/templates/flutter/lib/src/realtime_mixin.dart.twig b/templates/flutter/lib/src/realtime_mixin.dart.twig index 1e99f40a26..f70ac2e1b4 100644 --- a/templates/flutter/lib/src/realtime_mixin.dart.twig +++ b/templates/flutter/lib/src/realtime_mixin.dart.twig @@ -25,6 +25,7 @@ mixin RealtimeMixin { Map _subscriptions = {}; bool _notifyDone = true; StreamSubscription? _websocketSubscription; + bool _creatingSocket = false; Future _closeConnection() async { await _websocketSubscription?.cancel(); @@ -33,12 +34,15 @@ mixin RealtimeMixin { } _createSocket() async { + if(_creatingSocket) return; + _creatingSocket = true; final uri = _prepareUri(); if (_websok == null) { _websok = await getWebSocket(uri); _lastUrl = uri.toString(); } else { if (_lastUrl == uri.toString() && _websok?.closeCode == null) { + _creatingSocket = false; return; } _notifyDone = false; @@ -107,6 +111,8 @@ mixin RealtimeMixin { throw {{spec.title | caseUcfirst}}Exception(e.message); } throw {{spec.title | caseUcfirst}}Exception(e.toString()); + } finally { + _creatingSocket = false; } } From b4b06a2b86e948005e9360b9344161c5c3e53f13 Mon Sep 17 00:00:00 2001 From: Damodar Lohani Date: Wed, 17 Apr 2024 15:27:10 +0545 Subject: [PATCH 5/6] do not close subscription if creating socket is true do not close subscription if creating socket is true to prevent concurrent modification error --- templates/flutter/lib/src/realtime_mixin.dart.twig | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/templates/flutter/lib/src/realtime_mixin.dart.twig b/templates/flutter/lib/src/realtime_mixin.dart.twig index f70ac2e1b4..69da4fca5d 100644 --- a/templates/flutter/lib/src/realtime_mixin.dart.twig +++ b/templates/flutter/lib/src/realtime_mixin.dart.twig @@ -88,7 +88,7 @@ mixin RealtimeMixin { break; } }, onDone: () { - if (!_notifyDone) return; + if (!_notifyDone || _creatingSocket) return; for (var subscription in _subscriptions.values) { subscription.close(); } From 2fc1e1f5b285c413d04c95b04ad9d886327e3dc2 Mon Sep 17 00:00:00 2001 From: Damodar Lohani Date: Tue, 23 Apr 2024 06:45:03 +0545 Subject: [PATCH 6/6] check channels while creating socket --- templates/flutter/lib/src/realtime_mixin.dart.twig | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/templates/flutter/lib/src/realtime_mixin.dart.twig b/templates/flutter/lib/src/realtime_mixin.dart.twig index 69da4fca5d..210de4974f 100644 --- a/templates/flutter/lib/src/realtime_mixin.dart.twig +++ b/templates/flutter/lib/src/realtime_mixin.dart.twig @@ -34,7 +34,7 @@ mixin RealtimeMixin { } _createSocket() async { - if(_creatingSocket) return; + if(_creatingSocket || _channels.isEmpty) return; _creatingSocket = true; final uri = _prepareUri(); if (_websok == null) {