diff --git a/packages/stream_chat/CHANGELOG.md b/packages/stream_chat/CHANGELOG.md index f7094a53..f28d7e65 100644 --- a/packages/stream_chat/CHANGELOG.md +++ b/packages/stream_chat/CHANGELOG.md @@ -8,6 +8,8 @@ ✅ Added - Expose `ChannelMute` class. [#1473](https://github.com/GetStream/stream-chat-flutter/issues/1473) +- Added synchronization to the `StreamChatClient.sync` + api. [#1392](https://github.com/GetStream/stream-chat-flutter/issues/1392) ## 6.0.0 diff --git a/packages/stream_chat/lib/src/client/client.dart b/packages/stream_chat/lib/src/client/client.dart index 26d229b4..72dfe989 100644 --- a/packages/stream_chat/lib/src/client/client.dart +++ b/packages/stream_chat/lib/src/client/client.dart @@ -32,6 +32,7 @@ import 'package:stream_chat/src/event_type.dart'; import 'package:stream_chat/src/ws/connection_status.dart'; import 'package:stream_chat/src/ws/websocket.dart'; import 'package:stream_chat/version.dart'; +import 'package:synchronized/extension.dart'; /// Handler function used for logging records. Function requires a single /// [LogRecord] as the only parameter. @@ -488,37 +489,40 @@ class StreamChatClient { /// Get the events missed while offline to sync the offline storage /// Will automatically fetch [cids] and [lastSyncedAt] if [persistenceEnabled] - Future sync({List? cids, DateTime? lastSyncAt}) async { - cids ??= await _chatPersistenceClient?.getChannelCids(); - if (cids == null || cids.isEmpty) { - return; - } - - lastSyncAt ??= await _chatPersistenceClient?.getLastSyncAt(); - if (lastSyncAt == null) { - return; - } - - try { - final res = await _chatApi.general.sync(cids, lastSyncAt); - final events = res.events - ..sort((a, b) => a.createdAt.compareTo(b.createdAt)); - - for (final event in events) { - logger.fine('event.type: ${event.type}'); - final messageText = event.message?.text; - if (messageText != null) { - logger.fine('event.message.text: $messageText'); - } - handleEvent(event); + Future sync({List? cids, DateTime? lastSyncAt}) { + return synchronized(() async { + final channels = cids ?? await _chatPersistenceClient?.getChannelCids(); + if (channels == null || channels.isEmpty) { + return; } - final now = DateTime.now(); - _lastSyncedAt = now; - _chatPersistenceClient?.updateLastSyncAt(now); - } catch (e, stk) { - logger.severe('Error during sync', e, stk); - } + final syncAt = + lastSyncAt ?? await _chatPersistenceClient?.getLastSyncAt(); + if (syncAt == null) { + return; + } + + try { + final res = await _chatApi.general.sync(channels, syncAt); + final events = res.events + ..sort((a, b) => a.createdAt.compareTo(b.createdAt)); + + for (final event in events) { + logger.fine('event.type: ${event.type}'); + final messageText = event.message?.text; + if (messageText != null) { + logger.fine('event.message.text: $messageText'); + } + handleEvent(event); + } + + final now = DateTime.now(); + _lastSyncedAt = now; + _chatPersistenceClient?.updateLastSyncAt(now); + } catch (e, stk) { + logger.severe('Error during sync', e, stk); + } + }); } final _queryChannelsStreams = >>{}; diff --git a/packages/stream_chat/pubspec.yaml b/packages/stream_chat/pubspec.yaml index f0bbc763..61c1ff97 100644 --- a/packages/stream_chat/pubspec.yaml +++ b/packages/stream_chat/pubspec.yaml @@ -22,6 +22,7 @@ dependencies: mime: ^1.0.4 rate_limiter: ^1.0.0 rxdart: ^0.27.7 + synchronized: ^3.0.0 uuid: ^3.0.7 web_socket_channel: ^2.3.0