From abab2e0a0606b9ef7c6fef9af12e173a2b9773f6 Mon Sep 17 00:00:00 2001 From: Sahil Kumar Date: Thu, 4 May 2023 16:54:12 +0530 Subject: [PATCH 1/6] feat(llc): added synchronization to the `StreamChatClient.sync` api. Signed-off-by: xsahil03x --- packages/stream_chat/CHANGELOG.md | 2 + .../stream_chat/lib/src/client/client.dart | 62 ++++++++++--------- packages/stream_chat/pubspec.yaml | 1 + 3 files changed, 36 insertions(+), 29 deletions(-) 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 From 91a378f720c0030c8b8af958e8d54b1fb4968e5c Mon Sep 17 00:00:00 2001 From: Sahil Kumar Date: Fri, 5 May 2023 01:26:40 +0530 Subject: [PATCH 2/6] refactor: remove mutex as drift already uses one internally. Signed-off-by: xsahil03x --- .../src/stream_chat_persistence_client.dart | 287 ++++++++---------- packages/stream_chat_persistence/pubspec.yaml | 1 - 2 files changed, 128 insertions(+), 160 deletions(-) diff --git a/packages/stream_chat_persistence/lib/src/stream_chat_persistence_client.dart b/packages/stream_chat_persistence/lib/src/stream_chat_persistence_client.dart index 48e563dc..9d2cb955 100644 --- a/packages/stream_chat_persistence/lib/src/stream_chat_persistence_client.dart +++ b/packages/stream_chat_persistence/lib/src/stream_chat_persistence_client.dart @@ -1,7 +1,5 @@ import 'package:flutter/foundation.dart'; -import 'package:mutex/mutex.dart'; import 'package:stream_chat/stream_chat.dart'; - import 'package:stream_chat_persistence/src/db/drift_chat_database.dart'; /// Various connection modes on which [StreamChatPersistenceClient] can work @@ -48,7 +46,6 @@ class StreamChatPersistenceClient extends ChatPersistenceClient { final Logger _logger; final ConnectionMode _connectionMode; final bool _webUseIndexedDbIfSupported; - final _mutex = ReadWriteMutex(); void _defaultLogHandler(LogRecord record) { print( @@ -59,9 +56,6 @@ class StreamChatPersistenceClient extends ChatPersistenceClient { if (record.stackTrace != null) print(record.stackTrace); } - Future _readProtected(AsyncValueGetter func) => - _mutex.protectRead(func); - bool get _debugIsConnected { assert(() { if (db == null) { @@ -104,90 +98,84 @@ class StreamChatPersistenceClient extends ChatPersistenceClient { Future getConnectionInfo() { assert(_debugIsConnected, ''); _logger.info('getConnectionInfo'); - return _readProtected(() => db!.connectionEventDao.connectionEvent); + return db!.connectionEventDao.connectionEvent; } @override Future updateConnectionInfo(Event event) { assert(_debugIsConnected, ''); _logger.info('updateConnectionInfo'); - return _readProtected( - () => db!.connectionEventDao.updateConnectionEvent(event), - ); + return db!.connectionEventDao.updateConnectionEvent(event); } @override Future updateLastSyncAt(DateTime lastSyncAt) { assert(_debugIsConnected, ''); _logger.info('updateLastSyncAt'); - return _readProtected( - () => db!.connectionEventDao.updateLastSyncAt(lastSyncAt), - ); + return db!.connectionEventDao.updateLastSyncAt(lastSyncAt); } @override Future getLastSyncAt() { assert(_debugIsConnected, ''); _logger.info('getLastSyncAt'); - return _readProtected(() => db!.connectionEventDao.lastSyncAt); + return db!.connectionEventDao.lastSyncAt; } @override Future deleteChannels(List cids) { assert(_debugIsConnected, ''); _logger.info('deleteChannels'); - return _readProtected(() => db!.channelDao.deleteChannelByCids(cids)); + return db!.channelDao.deleteChannelByCids(cids); } @override Future> getChannelCids() { assert(_debugIsConnected, ''); _logger.info('getChannelCids'); - return _readProtected(() => db!.channelDao.cids); + return db!.channelDao.cids; } @override Future deleteMessageByIds(List messageIds) { assert(_debugIsConnected, ''); _logger.info('deleteMessageByIds'); - return _readProtected(() => db!.messageDao.deleteMessageByIds(messageIds)); + return db!.messageDao.deleteMessageByIds(messageIds); } @override Future deletePinnedMessageByIds(List messageIds) { assert(_debugIsConnected, ''); _logger.info('deletePinnedMessageByIds'); - return _readProtected( - () => db!.pinnedMessageDao.deleteMessageByIds(messageIds), - ); + return db!.pinnedMessageDao.deleteMessageByIds(messageIds); } @override Future deleteMessageByCids(List cids) { assert(_debugIsConnected, ''); _logger.info('deleteMessageByCids'); - return _readProtected(() => db!.messageDao.deleteMessageByCids(cids)); + return db!.messageDao.deleteMessageByCids(cids); } @override Future deletePinnedMessageByCids(List cids) { assert(_debugIsConnected, ''); _logger.info('deletePinnedMessageByCids'); - return _readProtected(() => db!.pinnedMessageDao.deleteMessageByCids(cids)); + return db!.pinnedMessageDao.deleteMessageByCids(cids); } @override Future> getMembersByCid(String cid) { assert(_debugIsConnected, ''); _logger.info('getMembersByCid'); - return _readProtected(() => db!.memberDao.getMembersByCid(cid)); + return db!.memberDao.getMembersByCid(cid); } @override Future getChannelByCid(String cid) { assert(_debugIsConnected, ''); _logger.info('getChannelByCid'); - return _readProtected(() => db!.channelDao.getChannelByCid(cid)); + return db!.channelDao.getChannelByCid(cid); } @override @@ -197,11 +185,9 @@ class StreamChatPersistenceClient extends ChatPersistenceClient { }) { assert(_debugIsConnected, ''); _logger.info('getMessagesByCid'); - return _readProtected( - () => db!.messageDao.getMessagesByCid( - cid, - messagePagination: messagePagination, - ), + return db!.messageDao.getMessagesByCid( + cid, + messagePagination: messagePagination, ); } @@ -212,37 +198,34 @@ class StreamChatPersistenceClient extends ChatPersistenceClient { }) { assert(_debugIsConnected, ''); _logger.info('getPinnedMessagesByCid'); - return _readProtected( - () => db!.pinnedMessageDao.getMessagesByCid( - cid, - messagePagination: messagePagination, - ), + return db!.pinnedMessageDao.getMessagesByCid( + cid, + messagePagination: messagePagination, ); } @override - Future> getReadsByCid(String cid) { + Future> getReadsByCid(String cid) async { assert(_debugIsConnected, ''); _logger.info('getReadsByCid'); - return _readProtected(() => db!.readDao.getReadsByCid(cid)); + return db!.readDao.getReadsByCid(cid); } @override - Future>> getChannelThreads(String cid) { + Future>> getChannelThreads(String cid) async { assert(_debugIsConnected, ''); _logger.info('getChannelThreads'); - return _readProtected(() async { - final messages = await db!.messageDao.getThreadMessages(cid); - final messageByParentIdDictionary = >{}; - for (final message in messages) { - final parentId = message.parentId!; - messageByParentIdDictionary[parentId] = [ - ...messageByParentIdDictionary[parentId] ?? [], - message, - ]; - } - return messageByParentIdDictionary; - }); + final messages = await db!.messageDao.getThreadMessages(cid); + final messageByParentIdDictionary = >{}; + for (final message in messages) { + final parentId = message.parentId!; + messageByParentIdDictionary[parentId] = [ + ...messageByParentIdDictionary[parentId] ?? [], + message, + ]; + } + + return messageByParentIdDictionary; } @override @@ -252,11 +235,9 @@ class StreamChatPersistenceClient extends ChatPersistenceClient { }) { assert(_debugIsConnected, ''); _logger.info('getReplies'); - return _readProtected( - () => db!.messageDao.getThreadMessagesByParentId( - parentId, - options: options, - ), + return db!.messageDao.getThreadMessagesByParentId( + parentId, + options: options, ); } @@ -268,73 +249,69 @@ class StreamChatPersistenceClient extends ChatPersistenceClient { Please use channelStateSort instead.''') List>? sort, List>? channelStateSort, PaginationParams? paginationParams, - }) { + }) async { assert(_debugIsConnected, ''); assert( sort == null || channelStateSort == null, 'sort and channelStateSort cannot be used together', ); _logger.info('getChannelStates'); - return _readProtected( - () async { - final channels = await db!.channelQueryDao.getChannels( - filter: filter, - sort: sort, - ); - final channelStates = await Future.wait( - channels.map((e) => getChannelStateByCid(e.cid)), - ); - - // Only sort the channel states if the channels are not already sorted. - if (sort == null) { - var chainedComparator = (ChannelState a, ChannelState b) { - final dateA = a.channel?.lastMessageAt ?? a.channel?.createdAt; - final dateB = b.channel?.lastMessageAt ?? b.channel?.createdAt; - - if (dateA == null && dateB == null) { - return 0; - } else if (dateA == null) { - return 1; - } else if (dateB == null) { - return -1; - } else { - return dateB.compareTo(dateA); - } - }; - - if (channelStateSort != null && channelStateSort.isNotEmpty) { - chainedComparator = (a, b) { - int result; - for (final comparator in channelStateSort - .map((it) => it.comparator) - .withNullifyer) { - try { - result = comparator(a, b); - } catch (e) { - result = 0; - } - if (result != 0) return result; - } - return 0; - }; - } - - channelStates.sort(chainedComparator); - } - - final offset = paginationParams?.offset; - if (offset != null && offset > 0 && channelStates.isNotEmpty) { - channelStates.removeRange(0, offset); - } - - if (paginationParams?.limit != null) { - return channelStates.take(paginationParams!.limit).toList(); - } - - return channelStates; - }, + final channels = await db!.channelQueryDao.getChannels( + filter: filter, + sort: sort, ); + + final channelStates = await Future.wait( + channels.map((e) => getChannelStateByCid(e.cid)), + ); + + // Only sort the channel states if the channels are not already sorted. + if (sort == null) { + var chainedComparator = (ChannelState a, ChannelState b) { + final dateA = a.channel?.lastMessageAt ?? a.channel?.createdAt; + final dateB = b.channel?.lastMessageAt ?? b.channel?.createdAt; + + if (dateA == null && dateB == null) { + return 0; + } else if (dateA == null) { + return 1; + } else if (dateB == null) { + return -1; + } else { + return dateB.compareTo(dateA); + } + }; + + if (channelStateSort != null && channelStateSort.isNotEmpty) { + chainedComparator = (a, b) { + int result; + for (final comparator + in channelStateSort.map((it) => it.comparator).withNullifyer) { + try { + result = comparator(a, b); + } catch (e) { + result = 0; + } + if (result != 0) return result; + } + return 0; + }; + } + + channelStates.sort(chainedComparator); + } + + final offset = paginationParams?.offset; + if (offset != null && offset > 0 && channelStates.isNotEmpty) { + channelStates.removeRange(0, offset); + } + + if (paginationParams?.limit != null) { + return channelStates.take(paginationParams!.limit).toList(); + } + + return channelStates; } @override @@ -345,12 +322,10 @@ class StreamChatPersistenceClient extends ChatPersistenceClient { }) { assert(_debugIsConnected, ''); _logger.info('updateChannelQueries'); - return _readProtected( - () => db!.channelQueryDao.updateChannelQueries( - filter, - cids, - clearQueryCache: clearQueryCache, - ), + return db!.channelQueryDao.updateChannelQueries( + filter, + cids, + clearQueryCache: clearQueryCache, ); } @@ -358,60 +333,56 @@ class StreamChatPersistenceClient extends ChatPersistenceClient { Future updateChannels(List channels) { assert(_debugIsConnected, ''); _logger.info('updateChannels'); - return _readProtected(() => db!.channelDao.updateChannels(channels)); + return db!.channelDao.updateChannels(channels); } @override Future bulkUpdateMembers(Map?> members) { assert(_debugIsConnected, ''); _logger.info('bulkUpdateMembers'); - return _readProtected(() => db!.memberDao.bulkUpdateMembers(members)); + return db!.memberDao.bulkUpdateMembers(members); } @override Future bulkUpdateMessages(Map?> messages) { assert(_debugIsConnected, ''); _logger.info('bulkUpdateMessages'); - return _readProtected(() => db!.messageDao.bulkUpdateMessages(messages)); + return db!.messageDao.bulkUpdateMessages(messages); } @override Future bulkUpdatePinnedMessages(Map?> messages) { assert(_debugIsConnected, ''); _logger.info('bulkUpdatePinnedMessages'); - return _readProtected( - () => db!.pinnedMessageDao.bulkUpdateMessages(messages), - ); + return db!.pinnedMessageDao.bulkUpdateMessages(messages); } @override Future updatePinnedMessageReactions(List reactions) { assert(_debugIsConnected, ''); _logger.info('updatePinnedMessageReactions'); - return _readProtected( - () => db!.pinnedMessageReactionDao.updateReactions(reactions), - ); + return db!.pinnedMessageReactionDao.updateReactions(reactions); } @override Future updateReactions(List reactions) { assert(_debugIsConnected, ''); _logger.info('updateReactions'); - return _readProtected(() => db!.reactionDao.updateReactions(reactions)); + return db!.reactionDao.updateReactions(reactions); } @override Future bulkUpdateReads(Map?> reads) { assert(_debugIsConnected, ''); _logger.info('bulkUpdateReads'); - return _readProtected(() => db!.readDao.bulkUpdateReads(reads)); + return db!.readDao.bulkUpdateReads(reads); } @override Future updateUsers(List users) { assert(_debugIsConnected, ''); _logger.info('updateUsers'); - return _readProtected(() => db!.userDao.updateUsers(users)); + return db!.userDao.updateUsers(users); } @override @@ -420,53 +391,51 @@ class StreamChatPersistenceClient extends ChatPersistenceClient { ) { assert(_debugIsConnected, ''); _logger.info('deletePinnedMessageReactionsByMessageId'); - return _readProtected( - () => - db!.pinnedMessageReactionDao.deleteReactionsByMessageIds(messageIds), - ); + return db!.pinnedMessageReactionDao.deleteReactionsByMessageIds(messageIds); } @override Future deleteReactionsByMessageId(List messageIds) { assert(_debugIsConnected, ''); _logger.info('deleteReactionsByMessageId'); - return _readProtected( - () => db!.reactionDao.deleteReactionsByMessageIds(messageIds), - ); + return db!.reactionDao.deleteReactionsByMessageIds(messageIds); } @override Future deleteMembersByCids(List cids) { assert(_debugIsConnected, ''); _logger.info('deleteMembersByCids'); - return _readProtected(() => db!.memberDao.deleteMemberByCids(cids)); + return db!.memberDao.deleteMemberByCids(cids); + } + + @override + Future updateChannelThreads( + String cid, + Map> threads, + ) { + assert(_debugIsConnected, ''); + _logger.info('updateChannelThreads'); + return db!.transaction(() => super.updateChannelThreads(cid, threads)); } @override Future updateChannelStates(List channelStates) { assert(_debugIsConnected, ''); _logger.info('updateChannelStates'); - return _readProtected( - () async => db!.transaction( - () async { - await super.updateChannelStates(channelStates); - }, - ), - ); + return db!.transaction(() => super.updateChannelStates(channelStates)); } @override - Future disconnect({bool flush = false}) async => - _mutex.protectWrite(() async { - _logger.info('disconnect'); - if (db != null) { - _logger.info('Disconnecting'); - if (flush) { - _logger.info('Flushing'); - await db!.flush(); - } - await db!.disconnect(); - db = null; - } - }); + Future disconnect({bool flush = false}) async { + _logger.info('disconnect'); + if (db != null) { + _logger.info('Disconnecting'); + if (flush) { + _logger.info('Flushing'); + await db!.flush(); + } + await db!.disconnect(); + db = null; + } + } } diff --git a/packages/stream_chat_persistence/pubspec.yaml b/packages/stream_chat_persistence/pubspec.yaml index 23275396..83fa7a11 100644 --- a/packages/stream_chat_persistence/pubspec.yaml +++ b/packages/stream_chat_persistence/pubspec.yaml @@ -15,7 +15,6 @@ dependencies: sdk: flutter logging: ^1.0.1 meta: ^1.8.0 - mutex: ^3.0.0 path: ^1.8.2 path_provider: ^2.0.1 sqlite3_flutter_libs: ^0.5.0 From 167d60dcacae417573de40f0ab98b46ae6c58ed7 Mon Sep 17 00:00:00 2001 From: Sahil Kumar Date: Fri, 5 May 2023 01:27:06 +0530 Subject: [PATCH 3/6] doc: advice to use `regular` as the default connectionMode Signed-off-by: xsahil03x --- packages/stream_chat_persistence/README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/stream_chat_persistence/README.md b/packages/stream_chat_persistence/README.md index 0851335b..1d7f8baf 100644 --- a/packages/stream_chat_persistence/README.md +++ b/packages/stream_chat_persistence/README.md @@ -37,7 +37,7 @@ The usage is pretty simple. ```dart final chatPersistentClient = StreamChatPersistenceClient( logLevel: Level.INFO, - connectionMode: ConnectionMode.background, + connectionMode: ConnectionMode.regular, ); ``` 2. Pass the instance to the official Stream chat client. From 9e9d2dfc18f97783b91f4b9257986f03ab73cab5 Mon Sep 17 00:00:00 2001 From: Sahil Kumar Date: Fri, 5 May 2023 01:55:03 +0530 Subject: [PATCH 4/6] chore: improve comparator. Signed-off-by: xsahil03x --- .../src/stream_chat_persistence_client.dart | 66 +++++++++++-------- 1 file changed, 37 insertions(+), 29 deletions(-) diff --git a/packages/stream_chat_persistence/lib/src/stream_chat_persistence_client.dart b/packages/stream_chat_persistence/lib/src/stream_chat_persistence_client.dart index 9d2cb955..b19977db 100644 --- a/packages/stream_chat_persistence/lib/src/stream_chat_persistence_client.dart +++ b/packages/stream_chat_persistence/lib/src/stream_chat_persistence_client.dart @@ -268,38 +268,14 @@ class StreamChatPersistenceClient extends ChatPersistenceClient { // Only sort the channel states if the channels are not already sorted. if (sort == null) { - var chainedComparator = (ChannelState a, ChannelState b) { - final dateA = a.channel?.lastMessageAt ?? a.channel?.createdAt; - final dateB = b.channel?.lastMessageAt ?? b.channel?.createdAt; - - if (dateA == null && dateB == null) { - return 0; - } else if (dateA == null) { - return 1; - } else if (dateB == null) { - return -1; - } else { - return dateB.compareTo(dateA); - } - }; - + var comparator = _defaultChannelStateComparator; if (channelStateSort != null && channelStateSort.isNotEmpty) { - chainedComparator = (a, b) { - int result; - for (final comparator - in channelStateSort.map((it) => it.comparator).withNullifyer) { - try { - result = comparator(a, b); - } catch (e) { - result = 0; - } - if (result != 0) return result; - } - return 0; - }; + comparator = _combineComparators( + channelStateSort.map((it) => it.comparator).withNullifyer, + ); } - channelStates.sort(chainedComparator); + channelStates.sort(comparator); } final offset = paginationParams?.offset; @@ -439,3 +415,35 @@ class StreamChatPersistenceClient extends ChatPersistenceClient { } } } + +// Creates a new combined [Comparator] which sorts items +// by the given [comparators]. +Comparator _combineComparators(Iterable> comparators) { + return (T a, T b) { + for (final comparator in comparators) { + try { + final result = comparator(a, b); + if (result != 0) return result; + } catch (e) { + // If the comparator throws an exception, we ignore it and + // continue with the next comparator. + continue; + } + } + return 0; + }; +} + +// The default [Comparator] used to sort [ChannelState]s. +int _defaultChannelStateComparator(ChannelState a, ChannelState b) { + final dateA = a.channel?.lastMessageAt ?? a.channel?.createdAt; + final dateB = b.channel?.lastMessageAt ?? b.channel?.createdAt; + + if (dateA == null && dateB == null) return 0; + if (dateA == null) return 1; + if (dateB == null) { + return -1; + } else { + return dateB.compareTo(dateA); + } +} From 7cc8e9ea5a8a4914e6668df56527d6cb9e976e58 Mon Sep 17 00:00:00 2001 From: Sahil Kumar Date: Fri, 5 May 2023 01:55:26 +0530 Subject: [PATCH 5/6] ci: run action when draft status changes. Signed-off-by: xsahil03x --- .github/workflows/stream_flutter_workflow.yml | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/.github/workflows/stream_flutter_workflow.yml b/.github/workflows/stream_flutter_workflow.yml index af57f49f..b63933bb 100644 --- a/.github/workflows/stream_flutter_workflow.yml +++ b/.github/workflows/stream_flutter_workflow.yml @@ -8,6 +8,11 @@ on: pull_request: paths: - 'packages/**' + types: + - opened + - reopened + - synchronize + - ready_for_review push: branches: - master From ee62abdbc40c8bd6368f5456589276c6ae6710d7 Mon Sep 17 00:00:00 2001 From: Sahil Kumar Date: Fri, 5 May 2023 02:21:06 +0530 Subject: [PATCH 6/6] test: add test. Signed-off-by: xsahil03x --- .../lib/src/stream_chat_persistence_client.dart | 1 + .../test/stream_chat_persistence_client_test.dart | 13 +++++++++++++ 2 files changed, 14 insertions(+) diff --git a/packages/stream_chat_persistence/lib/src/stream_chat_persistence_client.dart b/packages/stream_chat_persistence/lib/src/stream_chat_persistence_client.dart index b19977db..3d5a5bdb 100644 --- a/packages/stream_chat_persistence/lib/src/stream_chat_persistence_client.dart +++ b/packages/stream_chat_persistence/lib/src/stream_chat_persistence_client.dart @@ -90,6 +90,7 @@ class StreamChatPersistenceClient extends ChatPersistenceClient { 'disconnect the previous instance before connecting again.', ); } + _logger.info('connect'); db = databaseProvider?.call(userId, _connectionMode) ?? await _defaultDatabaseProvider(userId, _connectionMode); } diff --git a/packages/stream_chat_persistence/test/stream_chat_persistence_client_test.dart b/packages/stream_chat_persistence/test/stream_chat_persistence_client_test.dart index 8d7ab8be..c99668fd 100644 --- a/packages/stream_chat_persistence/test/stream_chat_persistence_client_test.dart +++ b/packages/stream_chat_persistence/test/stream_chat_persistence_client_test.dart @@ -55,6 +55,15 @@ void main() { expect(client.db, isNull); }); + test('client function throws stateError if db is not yet connected', () { + final client = StreamChatPersistenceClient(logLevel: Level.ALL); + expect( + // Running a function that requires db connection. + () => client.getReplies('testParentId'), + throwsA(isA()), + ); + }); + group('client functions', () { const userId = 'testUserId'; final mockDatabase = MockChatDatabase(); @@ -66,6 +75,10 @@ void main() { await client.connect(userId, databaseProvider: _mockDatabaseProvider); }); + tearDown(() async { + await client.disconnect(); + }); + test('getReplies', () async { const parentId = 'testParentId'; final replies = List.generate(3, (index) => Message(id: 'testId$index'));