From 99522651093cdda6efc0574a2aeec8d6cefe4b46 Mon Sep 17 00:00:00 2001 From: Salvatore Giordano Date: Fri, 12 Mar 2021 14:16:41 +0100 Subject: [PATCH] fix: persistence (#329) * fix(llc): Save original passed persistence for client reconnection. Signed-off-by: Sahil Kumar * fix(persistence): Check for empty channels before applying offset Signed-off-by: Sahil Kumar * fix(llc): Remove debug logs Signed-off-by: Sahil Kumar * fix(persistence): Reset isolate on successful disconnect Signed-off-by: Sahil Kumar * use mutexes in persistence client * fix(persistence): use `super.updateChannelStates` Signed-off-by: Sahil Kumar * fix(persistence): fix database multiple times creation. Signed-off-by: Sahil Kumar Co-authored-by: Sahil Kumar --- packages/stream_chat/lib/src/api/channel.dart | 8 +- packages/stream_chat/lib/src/client.dart | 82 ++--- .../stream_chat_flutter/example/pubspec.yaml | 8 + .../lib/src/dao/channel_query_dao.dart | 33 +- .../lib/src/dao/connection_event_dao.dart | 32 +- .../lib/src/db/moor_chat_database.dart | 8 +- .../lib/src/db/shared/native_db.dart | 18 +- .../lib/src/db/shared/unsupported_db.dart | 6 +- .../lib/src/db/shared/web_db.dart | 4 +- .../src/stream_chat_persistence_client.dart | 286 +++++++++++++----- packages/stream_chat_persistence/pubspec.yaml | 1 + 11 files changed, 332 insertions(+), 154 deletions(-) diff --git a/packages/stream_chat/lib/src/api/channel.dart b/packages/stream_chat/lib/src/api/channel.dart index 76209bfd..c63fe0f3 100644 --- a/packages/stream_chat/lib/src/api/channel.dart +++ b/packages/stream_chat/lib/src/api/channel.dart @@ -7,11 +7,11 @@ import 'package:logging/logging.dart'; import 'package:rxdart/rxdart.dart'; import 'package:stream_chat/src/api/retry_queue.dart'; import 'package:stream_chat/src/event_type.dart'; +import 'package:stream_chat/src/extensions/rate_limit.dart'; import 'package:stream_chat/src/models/attachment_file.dart'; import 'package:stream_chat/src/models/channel_state.dart'; import 'package:stream_chat/src/models/user.dart'; import 'package:stream_chat/stream_chat.dart'; -import 'package:stream_chat/src/extensions/rate_limit.dart'; /// This a the class that manages a specific channel. class Channel { @@ -1209,9 +1209,9 @@ class ChannelClientState { ChannelClientState( this._channel, ChannelState channelState, - ) : _debouncedUpdatePersistenceChannelState = _channel - ?._client?.chatPersistenceClient?.updateChannelState - ?.debounced(const Duration(seconds: 1)) { + ) : _debouncedUpdatePersistenceChannelState = ((ChannelState state) { + _channel?._client?.chatPersistenceClient?.updateChannelState(state); + }).debounced(const Duration(seconds: 1)) { retryQueue = RetryQueue( channel: _channel, logger: Logger('RETRY QUEUE ${_channel.cid}'), diff --git a/packages/stream_chat/lib/src/client.dart b/packages/stream_chat/lib/src/client.dart index d422dba1..11619b35 100644 --- a/packages/stream_chat/lib/src/client.dart +++ b/packages/stream_chat/lib/src/client.dart @@ -1,6 +1,7 @@ +// ignore_for_file: unnecessary_getters_setters + import 'dart:async'; import 'dart:convert'; -import 'package:stream_chat/src/extensions/map_extension.dart'; import 'package:dio/dio.dart'; import 'package:logging/logging.dart'; @@ -16,6 +17,7 @@ import 'package:stream_chat/src/attachment_file_uploader.dart'; import 'package:stream_chat/src/db/chat_persistence_client.dart'; import 'package:stream_chat/src/event_type.dart'; import 'package:stream_chat/src/exceptions.dart'; +import 'package:stream_chat/src/extensions/map_extension.dart'; import 'package:stream_chat/src/models/attachment_file.dart'; import 'package:stream_chat/src/models/channel_model.dart'; import 'package:stream_chat/src/models/channel_state.dart'; @@ -106,14 +108,22 @@ class StreamChatClient { logger.info('instantiating new client'); } + set chatPersistenceClient(ChatPersistenceClient value) { + _originalChatPersistenceClient = value; + } + + ChatPersistenceClient _originalChatPersistenceClient; + /// Chat persistence client - ChatPersistenceClient chatPersistenceClient; + ChatPersistenceClient get chatPersistenceClient => _chatPersistenceClient; + + ChatPersistenceClient _chatPersistenceClient; /// Attachment uploader AttachmentFileUploader attachmentFileUploader; /// Whether the chat persistence is available or not - bool get persistenceEnabled => chatPersistenceClient != null; + bool get persistenceEnabled => _chatPersistenceClient != null; RetryPolicy _retryPolicy; @@ -357,7 +367,7 @@ class StreamChatClient { /// Call this function to dispose the client void dispose() async { - await chatPersistenceClient?.disconnect(); + await _chatPersistenceClient?.disconnect(); await _disconnect(); httpClient.close(); await _controller.close(); @@ -446,8 +456,8 @@ class StreamChatClient { if (!event.isLocal) { if (_synced && event.createdAt != null) { - await chatPersistenceClient?.updateConnectionInfo(event); - await chatPersistenceClient?.updateLastSyncAt(event.createdAt); + await _chatPersistenceClient?.updateConnectionInfo(event); + await _chatPersistenceClient?.updateLastSyncAt(event.createdAt); } } @@ -478,8 +488,9 @@ class StreamChatClient { _wsConnectionStatus = ConnectionStatus.connecting; - if (persistenceEnabled) { - await chatPersistenceClient.connect(state.user.id); + if (_originalChatPersistenceClient != null) { + _chatPersistenceClient = _originalChatPersistenceClient; + await _chatPersistenceClient.connect(state.user.id); } _ws = WebSocket( @@ -508,34 +519,35 @@ class StreamChatClient { ), ); - if (status == ConnectionStatus.connected && - state.channels?.isNotEmpty == true) { - // ignore: unawaited_futures - queryChannelsOnline(filter: { - 'cid': { - '\$in': state.channels.keys.toList(), - }, - }).then( - (_) async { - await resync(); - handleEvent(Event( - type: EventType.connectionRecovered, - online: true, - )); - }, - ); - } else { - _synced = false; + if (status == ConnectionStatus.connected) { + handleEvent(Event( + type: EventType.connectionRecovered, + online: true, + )); + if (state.channels?.isNotEmpty == true) { + // ignore: unawaited_futures + queryChannelsOnline(filter: { + 'cid': { + '\$in': state.channels.keys.toList(), + }, + }).then( + (_) async { + await resync(); + }, + ); + } else { + _synced = false; + } } }; _connectionStatusSubscription = _ws.connectionStatusStream.listen(_connectionStatusHandler); - var event = await chatPersistenceClient?.getConnectionInfo(); + var event = await _chatPersistenceClient?.getConnectionInfo(); await _ws.connect().then((e) async { - await chatPersistenceClient?.updateConnectionInfo(e); + await _chatPersistenceClient?.updateConnectionInfo(e); event = e; await resync(); }).catchError((err, stacktrace) { @@ -551,14 +563,14 @@ class StreamChatClient { /// Get the events missed while offline to sync the offline storage Future resync([List cids]) async { - final lastSyncAt = await chatPersistenceClient?.getLastSyncAt(); + final lastSyncAt = await _chatPersistenceClient?.getLastSyncAt(); if (lastSyncAt == null) { _synced = true; return; } - cids ??= await chatPersistenceClient?.getChannelCids(); + cids ??= await _chatPersistenceClient?.getChannelCids(); if (cids?.isEmpty == true) { return; @@ -586,7 +598,7 @@ class StreamChatClient { res.events.forEach(handleEvent); - await chatPersistenceClient?.updateLastSyncAt(DateTime.now()); + await _chatPersistenceClient?.updateLastSyncAt(DateTime.now()); _synced = true; } catch (error) { logger.severe('Error during resync $error'); @@ -723,7 +735,7 @@ class StreamChatClient { final updateData = _mapChannelStateToChannel(channels); - await chatPersistenceClient?.updateChannelQueries( + await _chatPersistenceClient?.updateChannelQueries( filter, channels.map((c) => c.channel.cid).toList(), paginationParams?.offset == null || paginationParams.offset == 0, @@ -739,7 +751,7 @@ class StreamChatClient { @required List> sort, PaginationParams paginationParams = const PaginationParams(), }) async { - final offlineChannels = await chatPersistenceClient?.getChannelStates( + final offlineChannels = await _chatPersistenceClient?.getChannelStates( filter: filter, sort: sort, paginationParams: paginationParams, @@ -960,8 +972,8 @@ class StreamChatClient { logger.info('Disconnecting flushOfflineStorage: $flushChatPersistence; ' 'clearUser: $clearUser'); - await chatPersistenceClient?.disconnect(flush: flushChatPersistence); - chatPersistenceClient = null; + await _chatPersistenceClient?.disconnect(flush: flushChatPersistence); + _chatPersistenceClient = null; _connectCompleter = null; diff --git a/packages/stream_chat_flutter/example/pubspec.yaml b/packages/stream_chat_flutter/example/pubspec.yaml index 0c7ac24d..8a5270df 100644 --- a/packages/stream_chat_flutter/example/pubspec.yaml +++ b/packages/stream_chat_flutter/example/pubspec.yaml @@ -35,6 +35,14 @@ dependencies: # Use with the CupertinoIcons class for iOS style icons. cupertino_icons: ^1.0.0 +dependency_overrides: + stream_chat: + path: ../../stream_chat + stream_chat_flutter_core: + path: ../../stream_chat_flutter_core + stream_chat_persistence: + path: ../../stream_chat_persistence + dev_dependencies: flutter_test: sdk: flutter diff --git a/packages/stream_chat_persistence/lib/src/dao/channel_query_dao.dart b/packages/stream_chat_persistence/lib/src/dao/channel_query_dao.dart index 3f27cfd5..b5f6fe59 100644 --- a/packages/stream_chat_persistence/lib/src/dao/channel_query_dao.dart +++ b/packages/stream_chat_persistence/lib/src/dao/channel_query_dao.dart @@ -6,6 +6,7 @@ import 'package:stream_chat_persistence/src/db/moor_chat_database.dart'; import 'package:stream_chat_persistence/src/entity/channel_queries.dart'; import 'package:stream_chat_persistence/src/entity/channels.dart'; import 'package:stream_chat_persistence/src/entity/users.dart'; + import '../mapper/mapper.dart'; part 'channel_query_dao.g.dart'; @@ -33,24 +34,26 @@ class ChannelQueryDao extends DatabaseAccessor List cids, bool clearQueryCache, ) async { - final hash = _computeHash(filter); - if (clearQueryCache) { + return transaction(() async { + final hash = _computeHash(filter); + if (clearQueryCache) { + await batch((it) { + it.deleteWhere( + channelQueries, + (c) => c.queryHash.equals(hash), + ); + }); + } + await batch((it) { - it.deleteWhere( + it.insertAll( channelQueries, - (c) => c.queryHash.equals(hash), + cids.map((cid) { + return ChannelQueryEntity(queryHash: hash, channelCid: cid); + }).toList(), + mode: InsertMode.insertOrReplace, ); }); - } - - return batch((it) { - it.insertAll( - channelQueries, - cids.map((cid) { - return ChannelQueryEntity(queryHash: hash, channelCid: cid); - }).toList(), - mode: InsertMode.insertOrReplace, - ); }); } @@ -117,7 +120,7 @@ class ChannelQueryDao extends DatabaseAccessor cachedChannels.sort(chainedComparator); - if (paginationParams?.offset != null) { + if (paginationParams?.offset != null && cachedChannels.isNotEmpty) { cachedChannels.removeRange(0, paginationParams.offset); } diff --git a/packages/stream_chat_persistence/lib/src/dao/connection_event_dao.dart b/packages/stream_chat_persistence/lib/src/dao/connection_event_dao.dart index 8d2033b2..92764b10 100644 --- a/packages/stream_chat_persistence/lib/src/dao/connection_event_dao.dart +++ b/packages/stream_chat_persistence/lib/src/dao/connection_event_dao.dart @@ -3,6 +3,7 @@ import 'package:stream_chat/stream_chat.dart'; import 'package:stream_chat_persistence/src/db/moor_chat_database.dart'; import 'package:stream_chat_persistence/src/entity/connection_events.dart'; import 'package:stream_chat_persistence/src/entity/users.dart'; + import '../mapper/mapper.dart'; part 'connection_event_dao.g.dart'; @@ -27,20 +28,23 @@ class ConnectionEventDao extends DatabaseAccessor } /// Update stored connection event with latest data - Future updateConnectionEvent(Event event) async { - final connectionInfo = await select(connectionEvents).getSingle(); - return into(connectionEvents).insert( - ConnectionEventEntity( - id: 1, - lastSyncAt: connectionInfo?.lastSyncAt, - lastEventAt: event.createdAt ?? connectionInfo?.lastEventAt, - totalUnreadCount: - event.totalUnreadCount ?? connectionInfo?.totalUnreadCount, - ownUser: event.me?.toJson() ?? connectionInfo?.ownUser, - unreadChannels: event.unreadChannels ?? connectionInfo?.unreadChannels, - ), - mode: InsertMode.insertOrReplace, - ); + Future updateConnectionEvent(Event event) async { + return transaction(() async { + final connectionInfo = await select(connectionEvents).getSingle(); + await into(connectionEvents).insert( + ConnectionEventEntity( + id: 1, + lastSyncAt: connectionInfo?.lastSyncAt, + lastEventAt: event.createdAt ?? connectionInfo?.lastEventAt, + totalUnreadCount: + event.totalUnreadCount ?? connectionInfo?.totalUnreadCount, + ownUser: event.me?.toJson() ?? connectionInfo?.ownUser, + unreadChannels: + event.unreadChannels ?? connectionInfo?.unreadChannels, + ), + mode: InsertMode.insertOrReplace, + ); + }); } /// Update stored lastSyncAt with latest data diff --git a/packages/stream_chat_persistence/lib/src/db/moor_chat_database.dart b/packages/stream_chat_persistence/lib/src/db/moor_chat_database.dart index bdffeee7..57348330 100644 --- a/packages/stream_chat_persistence/lib/src/db/moor_chat_database.dart +++ b/packages/stream_chat_persistence/lib/src/db/moor_chat_database.dart @@ -60,7 +60,6 @@ class MoorChatDatabase extends _$MoorChatDatabase { /// Instantiate a new database instance MoorChatDatabase.connect( this._userId, - this._isolate, DatabaseConnection connection, ) : super.connect(connection); @@ -69,8 +68,6 @@ class MoorChatDatabase extends _$MoorChatDatabase { /// User id to which the database is connected String get userId => _userId; - MoorIsolate _isolate; - // you should bump this number whenever you change or add a table definition. @override int get schemaVersion => 2; @@ -89,8 +86,5 @@ class MoorChatDatabase extends _$MoorChatDatabase { ); /// Closes the database instance - Future disconnect() async { - await _isolate?.shutdownAll(); - await close(); - } + Future disconnect() => close(); } diff --git a/packages/stream_chat_persistence/lib/src/db/shared/native_db.dart b/packages/stream_chat_persistence/lib/src/db/shared/native_db.dart index c52e4a8f..9580ae50 100644 --- a/packages/stream_chat_persistence/lib/src/db/shared/native_db.dart +++ b/packages/stream_chat_persistence/lib/src/db/shared/native_db.dart @@ -76,17 +76,21 @@ class SharedDB { /// [MoorChatDatabase.connect] created on a background isolate. /// /// Generally used with [ConnectionMode.background]. - static Future constructMoorChatDatabase( + static MoorChatDatabase constructMoorChatDatabase( String userId, { bool logStatements = false, - }) async { + }) { final dbName = 'db_$userId'; - final isolate = await _createMoorIsolate( - dbName, - logStatements: logStatements, + return MoorChatDatabase.connect( + userId, + DatabaseConnection.delayed(Future(() async { + MoorIsolate isolate = await _createMoorIsolate( + dbName, + logStatements: logStatements, + ); + return isolate.connect(); + })), ); - final connection = await isolate.connect(); - return MoorChatDatabase.connect(userId, isolate, connection); } } diff --git a/packages/stream_chat_persistence/lib/src/db/shared/unsupported_db.dart b/packages/stream_chat_persistence/lib/src/db/shared/unsupported_db.dart index b597993d..811e5cc7 100644 --- a/packages/stream_chat_persistence/lib/src/db/shared/unsupported_db.dart +++ b/packages/stream_chat_persistence/lib/src/db/shared/unsupported_db.dart @@ -1,3 +1,5 @@ +import 'package:moor/backends.dart'; +import 'package:stream_chat_persistence/src/db/moor_chat_database.dart'; import 'package:stream_chat_persistence/stream_chat_persistence.dart'; /// A Helper class to construct new instances of [MoorChatDatabase] @@ -5,7 +7,7 @@ class SharedDB { /// Returns a new instance of database. /// /// Generally used with [ConnectionMode.regular]. - static dynamic constructDatabase( + static Future constructDatabase( String userId, { bool logStatements = false, bool persistOnDisk = true, @@ -16,7 +18,7 @@ class SharedDB { /// Return a new instance of moor chat database. /// /// Generally used with [ConnectionMode.background]. - static dynamic constructMoorChatDatabase( + static MoorChatDatabase constructMoorChatDatabase( String userId, { bool logStatements = false, }) { diff --git a/packages/stream_chat_persistence/lib/src/db/shared/web_db.dart b/packages/stream_chat_persistence/lib/src/db/shared/web_db.dart index a23e42ee..f44a41ec 100644 --- a/packages/stream_chat_persistence/lib/src/db/shared/web_db.dart +++ b/packages/stream_chat_persistence/lib/src/db/shared/web_db.dart @@ -22,10 +22,10 @@ class SharedDB { /// default constructor. /// /// Generally used with [ConnectionMode.background]. - static Future constructMoorChatDatabase( + static MoorChatDatabase constructMoorChatDatabase( String userId, { bool logStatements = false, - }) async { + }) { final dbName = 'db_$userId'; return MoorChatDatabase(dbName, logStatements: logStatements); } 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 44007b8d..832eabbd 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,4 +1,6 @@ +import 'package:logging/logging.dart' show LogRecord; import 'package:meta/meta.dart'; +import 'package:mutex/mutex.dart'; import 'package:stream_chat/stream_chat.dart'; import 'db/moor_chat_database.dart'; @@ -13,6 +15,12 @@ enum ConnectionMode { background, } +final levelEmojiMapper = { + Level.INFO: 'ℹ️', + Level.WARNING: '⚠️', + Level.SEVERE: '🚨', +}; + /// A [MoorChatDatabase] based implementation of the [ChatPersistenceClient] class StreamChatPersistenceClient extends ChatPersistenceClient { /// Creates a new instance of the stream chat persistence client @@ -20,15 +28,59 @@ class StreamChatPersistenceClient extends ChatPersistenceClient { /// Connection mode on which the client will work ConnectionMode connectionMode = ConnectionMode.regular, Level logLevel = Level.WARNING, + LogHandlerFunction logHandlerFunction, }) : assert(connectionMode != null), assert(logLevel != null), _connectionMode = connectionMode, - _logger = Logger.detached('💽')..level = logLevel; + _logger = Logger.detached('💽')..level = logLevel { + _logger.onRecord.listen(logHandlerFunction ?? _defaultLogHandler); + } + + /// A function that has a parameter of type [LogRecord]. + /// This is called on every new log record. + /// By default the client will use the handler returned by + /// [_getDefaultLogHandler]. + /// Setting it you can handle the log messages directly instead of have them + /// written to stdout, + /// this is very convenient if you use an error tracking tool or if you want + /// to centralize your logs into one facility. + /// + /// ```dart + /// myLogHandlerFunction = (LogRecord record) { + /// // do something with the record (ie. send it to Sentry or Fabric) + /// } + /// + /// final client = StreamChatPersistenceClient( + /// logHandlerFunction: myLogHandlerFunction, + /// ); + ///``` + LogHandlerFunction logHandlerFunction; @visibleForTesting MoorChatDatabase db; final Logger _logger; final ConnectionMode _connectionMode; + final _mutex = ReadWriteMutex(); + + void _defaultLogHandler(LogRecord record) { + print( + '(${record.time}) ' + '${levelEmojiMapper[record.level] ?? record.level.name} ' + '${record.loggerName} ${record.message}', + ); + if (record.stackTrace != null) print(record.stackTrace); + } + + Future readProtected(Future Function() f) async { + T ret; + await _mutex.protectRead(() async { + if (db == null) { + return; + } + ret = await f(); + }); + return ret; + } @override Future connect(String userId) async { @@ -45,67 +97,105 @@ class StreamChatPersistenceClient extends ChatPersistenceClient { return; case ConnectionMode.background: _logger.info('Connecting on background isolate'); - db = await SharedDB.constructMoorChatDatabase(userId); + db = SharedDB.constructMoorChatDatabase(userId); return; } } @override Future getConnectionInfo() { - return db.connectionEventDao.connectionEvent; + return readProtected(() { + _logger.info('getConnectionInfo'); + return db.connectionEventDao.connectionEvent; + }); } @override Future updateConnectionInfo(Event event) { - return db.connectionEventDao.updateConnectionEvent(event); + return readProtected(() { + _logger.info('updateConnectionInfo'); + return db.connectionEventDao.updateConnectionEvent(event); + }); } @override Future updateLastSyncAt(DateTime lastSyncAt) { - return db.connectionEventDao.updateLastSyncAt(lastSyncAt); + return readProtected(() { + _logger.info('updateLastSyncAt'); + return db.connectionEventDao.updateLastSyncAt(lastSyncAt); + }); } @override Future getLastSyncAt() { - return db.connectionEventDao.lastSyncAt; + return readProtected(() { + _logger.info('getLastSyncAt'); + return db.connectionEventDao.lastSyncAt; + }); } @override Future deleteChannels(List cids) { - return db.channelDao.deleteChannelByCids(cids); + return readProtected(() { + _logger.info('deleteChannels'); + return db.channelDao.deleteChannelByCids(cids); + }); } @override - Future> getChannelCids() => db.channelDao.cids; + Future> getChannelCids() { + return readProtected(() { + _logger.info('getChannelCids'); + return db.channelDao.cids; + }); + } @override Future deleteMessageByIds(List messageIds) { - return db.messageDao.deleteMessageByIds(messageIds); + return readProtected(() { + _logger.info('deleteMessageByIds'); + return db.messageDao.deleteMessageByIds(messageIds); + }); } @override Future deletePinnedMessageByIds(List messageIds) { - return db.pinnedMessageDao.deleteMessageByIds(messageIds); + return readProtected(() { + _logger.info('deletePinnedMessageByIds'); + return db.pinnedMessageDao.deleteMessageByIds(messageIds); + }); } @override Future deleteMessageByCids(List cids) { - return db.messageDao.deleteMessageByCids(cids); + return readProtected(() { + _logger.info('deleteMessageByCids'); + return db.messageDao.deleteMessageByCids(cids); + }); } @override Future deletePinnedMessageByCids(List cids) { - return db.pinnedMessageDao.deleteMessageByCids(cids); + return readProtected(() { + _logger.info('deletePinnedMessageByCids'); + return db.pinnedMessageDao.deleteMessageByCids(cids); + }); } @override Future> getMembersByCid(String cid) { - return db.memberDao.getMembersByCid(cid); + return readProtected(() { + _logger.info('getMembersByCid'); + return db.memberDao.getMembersByCid(cid); + }); } @override Future getChannelByCid(String cid) { - return db.channelDao.getChannelByCid(cid); + return readProtected(() { + _logger.info('getChannelByCid'); + return db.channelDao.getChannelByCid(cid); + }); } @override @@ -113,10 +203,13 @@ class StreamChatPersistenceClient extends ChatPersistenceClient { String cid, { PaginationParams messagePagination, }) { - return db.messageDao.getMessagesByCid( - cid, - messagePagination: messagePagination, - ); + return readProtected(() { + _logger.info('getMessagesByCid'); + return db.messageDao.getMessagesByCid( + cid, + messagePagination: messagePagination, + ); + }); } @override @@ -124,29 +217,38 @@ class StreamChatPersistenceClient extends ChatPersistenceClient { String cid, { PaginationParams messagePagination, }) { - return db.pinnedMessageDao.getMessagesByCid( - cid, - messagePagination: messagePagination, - ); + return readProtected(() { + _logger.info('getPinnedMessagesByCid'); + return db.pinnedMessageDao.getMessagesByCid( + cid, + messagePagination: messagePagination, + ); + }); } @override Future> getReadsByCid(String cid) { - return db.readDao.getReadsByCid(cid); + return readProtected(() { + _logger.info('getReadsByCid'); + return db.readDao.getReadsByCid(cid); + }); } @override Future>> getChannelThreads(String cid) 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; + return readProtected(() async { + _logger.info('getChannelThreads'); + 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 @@ -154,10 +256,13 @@ class StreamChatPersistenceClient extends ChatPersistenceClient { String parentId, { PaginationParams options, }) { - return db.messageDao.getThreadMessagesByParentId( - parentId, - options: options, - ); + return readProtected(() async { + _logger.info('getReplies'); + return db.messageDao.getThreadMessagesByParentId( + parentId, + options: options, + ); + }); } @override @@ -166,12 +271,15 @@ class StreamChatPersistenceClient extends ChatPersistenceClient { List> sort = const [], PaginationParams paginationParams, }) async { - final channels = await db.channelQueryDao.getChannels( - filter: filter, - sort: sort, - paginationParams: paginationParams, - ); - return Future.wait(channels.map((e) => getChannelStateByCid(e.cid))); + return readProtected(() async { + _logger.info('getChannelStates'); + final channels = await db.channelQueryDao.getChannels( + filter: filter, + sort: sort, + paginationParams: paginationParams, + ); + return Future.wait(channels.map((e) => getChannelStateByCid(e.cid))); + }); } @override @@ -180,72 +288,114 @@ class StreamChatPersistenceClient extends ChatPersistenceClient { List cids, bool clearQueryCache, ) { - return db.channelQueryDao.updateChannelQueries( - filter, - cids, - clearQueryCache, - ); + return readProtected(() async { + _logger.info('updateChannelQueries'); + return db.channelQueryDao.updateChannelQueries( + filter, + cids, + clearQueryCache, + ); + }); } @override Future updateChannels(List channels) { - return db.channelDao.updateChannels(channels); + return readProtected(() async { + _logger.info('updateChannels'); + return db.channelDao.updateChannels(channels); + }); } @override Future updateMembers(String cid, List members) { - return db.memberDao.updateMembers(cid, members); + return readProtected(() async { + _logger.info('updateMembers'); + return db.memberDao.updateMembers(cid, members); + }); } @override Future updateMessages(String cid, List messages) { - return db.messageDao.updateMessages(cid, messages); + return readProtected(() async { + _logger.info('updateMessages'); + return db.messageDao.updateMessages(cid, messages); + }); } @override Future updatePinnedMessages(String cid, List messages) { - return db.pinnedMessageDao.updateMessages(cid, messages); + return readProtected(() async { + _logger.info('updatePinnedMessages'); + return db.pinnedMessageDao.updateMessages(cid, messages); + }); } @override Future updateReactions(List reactions) { - return db.reactionDao.updateReactions(reactions); + return readProtected(() async { + _logger.info('updateReactions'); + return db.reactionDao.updateReactions(reactions); + }); } @override Future updateReads(String cid, List reads) { - return db.readDao.updateReads(cid, reads); + return readProtected(() async { + _logger.info('updateReads'); + return db.readDao.updateReads(cid, reads); + }); } @override Future updateUsers(List users) { - return db.userDao.updateUsers(users); + return readProtected(() async { + _logger.info('updateUsers'); + return db.userDao.updateUsers(users); + }); } @override Future deleteReactionsByMessageId(List messageIds) { - return db.reactionDao.deleteReactionsByMessageIds(messageIds); + return readProtected(() async { + _logger.info('deleteReactionsByMessageId'); + return db.reactionDao.deleteReactionsByMessageIds(messageIds); + }); } @override Future deleteMembersByCids(List cids) { - return db.memberDao.deleteMemberByCids(cids); + return readProtected(() async { + _logger.info('deleteMembersByCids'); + return db.memberDao.deleteMemberByCids(cids); + }); + } + + @override + Future updateChannelStates(List channelStates) { + return readProtected(() async { + return db.transaction(() async { + await super.updateChannelStates(channelStates); + }); + }); } @override Future disconnect({bool flush = false}) async { - if (db != null) { - _logger.info('Disconnecting'); - if (flush) { - _logger.info('Flushing'); - await db.batch((batch) { - db.allTables.forEach((table) { - db.delete(table).go(); + return _mutex.protectWrite(() async { + _logger.info('disconnect'); + if (db != null) { + _logger.info('Disconnecting'); + if (flush) { + _logger.info('Flushing'); + await db.batch((batch) { + db.allTables.forEach((table) { + db.delete(table).go(); + }); }); - }); + } + await db.disconnect(); + db = null; } - await db.disconnect(); - db = null; - } + }); } } diff --git a/packages/stream_chat_persistence/pubspec.yaml b/packages/stream_chat_persistence/pubspec.yaml index ea7f4e01..06f991e1 100644 --- a/packages/stream_chat_persistence/pubspec.yaml +++ b/packages/stream_chat_persistence/pubspec.yaml @@ -11,6 +11,7 @@ environment: dependencies: flutter: sdk: flutter + mutex: ^2.0.0 moor: ^3.4.0 path: ^1.7.0 path_provider: ^1.6.27