fix: persistence (#329)

* fix(llc): Save original passed persistence for client reconnection.

Signed-off-by: Sahil Kumar <[email protected]>

* fix(persistence): Check for empty channels before applying offset

Signed-off-by: Sahil Kumar <[email protected]>

* fix(llc): Remove debug logs

Signed-off-by: Sahil Kumar <[email protected]>

* fix(persistence): Reset isolate on successful disconnect

Signed-off-by: Sahil Kumar <[email protected]>

* use mutexes in persistence client

* fix(persistence): use `super.updateChannelStates`

Signed-off-by: Sahil Kumar <[email protected]>

* fix(persistence): fix database multiple times creation.

Signed-off-by: Sahil Kumar <[email protected]>

Co-authored-by: Sahil Kumar <[email protected]>
This commit is contained in:
Salvatore Giordano
2021-03-12 14:16:41 +01:00
committed by GitHub
co-authored by Sahil Kumar
parent ca1f70948e
commit 9952265109
11 changed files with 332 additions and 154 deletions
@@ -7,11 +7,11 @@ import 'package:logging/logging.dart';
import 'package:rxdart/rxdart.dart'; import 'package:rxdart/rxdart.dart';
import 'package:stream_chat/src/api/retry_queue.dart'; import 'package:stream_chat/src/api/retry_queue.dart';
import 'package:stream_chat/src/event_type.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/attachment_file.dart';
import 'package:stream_chat/src/models/channel_state.dart'; import 'package:stream_chat/src/models/channel_state.dart';
import 'package:stream_chat/src/models/user.dart'; import 'package:stream_chat/src/models/user.dart';
import 'package:stream_chat/stream_chat.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. /// This a the class that manages a specific channel.
class Channel { class Channel {
@@ -1209,9 +1209,9 @@ class ChannelClientState {
ChannelClientState( ChannelClientState(
this._channel, this._channel,
ChannelState channelState, ChannelState channelState,
) : _debouncedUpdatePersistenceChannelState = _channel ) : _debouncedUpdatePersistenceChannelState = ((ChannelState state) {
?._client?.chatPersistenceClient?.updateChannelState _channel?._client?.chatPersistenceClient?.updateChannelState(state);
?.debounced(const Duration(seconds: 1)) { }).debounced(const Duration(seconds: 1)) {
retryQueue = RetryQueue( retryQueue = RetryQueue(
channel: _channel, channel: _channel,
logger: Logger('RETRY QUEUE ${_channel.cid}'), logger: Logger('RETRY QUEUE ${_channel.cid}'),
+47 -35
View File
@@ -1,6 +1,7 @@
// ignore_for_file: unnecessary_getters_setters
import 'dart:async'; import 'dart:async';
import 'dart:convert'; import 'dart:convert';
import 'package:stream_chat/src/extensions/map_extension.dart';
import 'package:dio/dio.dart'; import 'package:dio/dio.dart';
import 'package:logging/logging.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/db/chat_persistence_client.dart';
import 'package:stream_chat/src/event_type.dart'; import 'package:stream_chat/src/event_type.dart';
import 'package:stream_chat/src/exceptions.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/attachment_file.dart';
import 'package:stream_chat/src/models/channel_model.dart'; import 'package:stream_chat/src/models/channel_model.dart';
import 'package:stream_chat/src/models/channel_state.dart'; import 'package:stream_chat/src/models/channel_state.dart';
@@ -106,14 +108,22 @@ class StreamChatClient {
logger.info('instantiating new client'); logger.info('instantiating new client');
} }
set chatPersistenceClient(ChatPersistenceClient value) {
_originalChatPersistenceClient = value;
}
ChatPersistenceClient _originalChatPersistenceClient;
/// Chat persistence client /// Chat persistence client
ChatPersistenceClient chatPersistenceClient; ChatPersistenceClient get chatPersistenceClient => _chatPersistenceClient;
ChatPersistenceClient _chatPersistenceClient;
/// Attachment uploader /// Attachment uploader
AttachmentFileUploader attachmentFileUploader; AttachmentFileUploader attachmentFileUploader;
/// Whether the chat persistence is available or not /// Whether the chat persistence is available or not
bool get persistenceEnabled => chatPersistenceClient != null; bool get persistenceEnabled => _chatPersistenceClient != null;
RetryPolicy _retryPolicy; RetryPolicy _retryPolicy;
@@ -357,7 +367,7 @@ class StreamChatClient {
/// Call this function to dispose the client /// Call this function to dispose the client
void dispose() async { void dispose() async {
await chatPersistenceClient?.disconnect(); await _chatPersistenceClient?.disconnect();
await _disconnect(); await _disconnect();
httpClient.close(); httpClient.close();
await _controller.close(); await _controller.close();
@@ -446,8 +456,8 @@ class StreamChatClient {
if (!event.isLocal) { if (!event.isLocal) {
if (_synced && event.createdAt != null) { if (_synced && event.createdAt != null) {
await chatPersistenceClient?.updateConnectionInfo(event); await _chatPersistenceClient?.updateConnectionInfo(event);
await chatPersistenceClient?.updateLastSyncAt(event.createdAt); await _chatPersistenceClient?.updateLastSyncAt(event.createdAt);
} }
} }
@@ -478,8 +488,9 @@ class StreamChatClient {
_wsConnectionStatus = ConnectionStatus.connecting; _wsConnectionStatus = ConnectionStatus.connecting;
if (persistenceEnabled) { if (_originalChatPersistenceClient != null) {
await chatPersistenceClient.connect(state.user.id); _chatPersistenceClient = _originalChatPersistenceClient;
await _chatPersistenceClient.connect(state.user.id);
} }
_ws = WebSocket( _ws = WebSocket(
@@ -508,34 +519,35 @@ class StreamChatClient {
), ),
); );
if (status == ConnectionStatus.connected && if (status == ConnectionStatus.connected) {
state.channels?.isNotEmpty == true) { handleEvent(Event(
// ignore: unawaited_futures type: EventType.connectionRecovered,
queryChannelsOnline(filter: { online: true,
'cid': { ));
'\$in': state.channels.keys.toList(), if (state.channels?.isNotEmpty == true) {
}, // ignore: unawaited_futures
}).then( queryChannelsOnline(filter: {
(_) async { 'cid': {
await resync(); '\$in': state.channels.keys.toList(),
handleEvent(Event( },
type: EventType.connectionRecovered, }).then(
online: true, (_) async {
)); await resync();
}, },
); );
} else { } else {
_synced = false; _synced = false;
}
} }
}; };
_connectionStatusSubscription = _connectionStatusSubscription =
_ws.connectionStatusStream.listen(_connectionStatusHandler); _ws.connectionStatusStream.listen(_connectionStatusHandler);
var event = await chatPersistenceClient?.getConnectionInfo(); var event = await _chatPersistenceClient?.getConnectionInfo();
await _ws.connect().then((e) async { await _ws.connect().then((e) async {
await chatPersistenceClient?.updateConnectionInfo(e); await _chatPersistenceClient?.updateConnectionInfo(e);
event = e; event = e;
await resync(); await resync();
}).catchError((err, stacktrace) { }).catchError((err, stacktrace) {
@@ -551,14 +563,14 @@ class StreamChatClient {
/// Get the events missed while offline to sync the offline storage /// Get the events missed while offline to sync the offline storage
Future<void> resync([List<String> cids]) async { Future<void> resync([List<String> cids]) async {
final lastSyncAt = await chatPersistenceClient?.getLastSyncAt(); final lastSyncAt = await _chatPersistenceClient?.getLastSyncAt();
if (lastSyncAt == null) { if (lastSyncAt == null) {
_synced = true; _synced = true;
return; return;
} }
cids ??= await chatPersistenceClient?.getChannelCids(); cids ??= await _chatPersistenceClient?.getChannelCids();
if (cids?.isEmpty == true) { if (cids?.isEmpty == true) {
return; return;
@@ -586,7 +598,7 @@ class StreamChatClient {
res.events.forEach(handleEvent); res.events.forEach(handleEvent);
await chatPersistenceClient?.updateLastSyncAt(DateTime.now()); await _chatPersistenceClient?.updateLastSyncAt(DateTime.now());
_synced = true; _synced = true;
} catch (error) { } catch (error) {
logger.severe('Error during resync $error'); logger.severe('Error during resync $error');
@@ -723,7 +735,7 @@ class StreamChatClient {
final updateData = _mapChannelStateToChannel(channels); final updateData = _mapChannelStateToChannel(channels);
await chatPersistenceClient?.updateChannelQueries( await _chatPersistenceClient?.updateChannelQueries(
filter, filter,
channels.map((c) => c.channel.cid).toList(), channels.map((c) => c.channel.cid).toList(),
paginationParams?.offset == null || paginationParams.offset == 0, paginationParams?.offset == null || paginationParams.offset == 0,
@@ -739,7 +751,7 @@ class StreamChatClient {
@required List<SortOption<ChannelModel>> sort, @required List<SortOption<ChannelModel>> sort,
PaginationParams paginationParams = const PaginationParams(), PaginationParams paginationParams = const PaginationParams(),
}) async { }) async {
final offlineChannels = await chatPersistenceClient?.getChannelStates( final offlineChannels = await _chatPersistenceClient?.getChannelStates(
filter: filter, filter: filter,
sort: sort, sort: sort,
paginationParams: paginationParams, paginationParams: paginationParams,
@@ -960,8 +972,8 @@ class StreamChatClient {
logger.info('Disconnecting flushOfflineStorage: $flushChatPersistence; ' logger.info('Disconnecting flushOfflineStorage: $flushChatPersistence; '
'clearUser: $clearUser'); 'clearUser: $clearUser');
await chatPersistenceClient?.disconnect(flush: flushChatPersistence); await _chatPersistenceClient?.disconnect(flush: flushChatPersistence);
chatPersistenceClient = null; _chatPersistenceClient = null;
_connectCompleter = null; _connectCompleter = null;
@@ -35,6 +35,14 @@ dependencies:
# Use with the CupertinoIcons class for iOS style icons. # Use with the CupertinoIcons class for iOS style icons.
cupertino_icons: ^1.0.0 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: dev_dependencies:
flutter_test: flutter_test:
sdk: flutter sdk: flutter
@@ -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/channel_queries.dart';
import 'package:stream_chat_persistence/src/entity/channels.dart'; import 'package:stream_chat_persistence/src/entity/channels.dart';
import 'package:stream_chat_persistence/src/entity/users.dart'; import 'package:stream_chat_persistence/src/entity/users.dart';
import '../mapper/mapper.dart'; import '../mapper/mapper.dart';
part 'channel_query_dao.g.dart'; part 'channel_query_dao.g.dart';
@@ -33,24 +34,26 @@ class ChannelQueryDao extends DatabaseAccessor<MoorChatDatabase>
List<String> cids, List<String> cids,
bool clearQueryCache, bool clearQueryCache,
) async { ) async {
final hash = _computeHash(filter); return transaction(() async {
if (clearQueryCache) { final hash = _computeHash(filter);
if (clearQueryCache) {
await batch((it) {
it.deleteWhere<ChannelQueries, ChannelQueryEntity>(
channelQueries,
(c) => c.queryHash.equals(hash),
);
});
}
await batch((it) { await batch((it) {
it.deleteWhere<ChannelQueries, ChannelQueryEntity>( it.insertAll(
channelQueries, 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<MoorChatDatabase>
cachedChannels.sort(chainedComparator); cachedChannels.sort(chainedComparator);
if (paginationParams?.offset != null) { if (paginationParams?.offset != null && cachedChannels.isNotEmpty) {
cachedChannels.removeRange(0, paginationParams.offset); cachedChannels.removeRange(0, paginationParams.offset);
} }
@@ -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/db/moor_chat_database.dart';
import 'package:stream_chat_persistence/src/entity/connection_events.dart'; import 'package:stream_chat_persistence/src/entity/connection_events.dart';
import 'package:stream_chat_persistence/src/entity/users.dart'; import 'package:stream_chat_persistence/src/entity/users.dart';
import '../mapper/mapper.dart'; import '../mapper/mapper.dart';
part 'connection_event_dao.g.dart'; part 'connection_event_dao.g.dart';
@@ -27,20 +28,23 @@ class ConnectionEventDao extends DatabaseAccessor<MoorChatDatabase>
} }
/// Update stored connection event with latest data /// Update stored connection event with latest data
Future<int> updateConnectionEvent(Event event) async { Future<void> updateConnectionEvent(Event event) async {
final connectionInfo = await select(connectionEvents).getSingle(); return transaction(() async {
return into(connectionEvents).insert( final connectionInfo = await select(connectionEvents).getSingle();
ConnectionEventEntity( await into(connectionEvents).insert(
id: 1, ConnectionEventEntity(
lastSyncAt: connectionInfo?.lastSyncAt, id: 1,
lastEventAt: event.createdAt ?? connectionInfo?.lastEventAt, lastSyncAt: connectionInfo?.lastSyncAt,
totalUnreadCount: lastEventAt: event.createdAt ?? connectionInfo?.lastEventAt,
event.totalUnreadCount ?? connectionInfo?.totalUnreadCount, totalUnreadCount:
ownUser: event.me?.toJson() ?? connectionInfo?.ownUser, event.totalUnreadCount ?? connectionInfo?.totalUnreadCount,
unreadChannels: event.unreadChannels ?? connectionInfo?.unreadChannels, ownUser: event.me?.toJson() ?? connectionInfo?.ownUser,
), unreadChannels:
mode: InsertMode.insertOrReplace, event.unreadChannels ?? connectionInfo?.unreadChannels,
); ),
mode: InsertMode.insertOrReplace,
);
});
} }
/// Update stored lastSyncAt with latest data /// Update stored lastSyncAt with latest data
@@ -60,7 +60,6 @@ class MoorChatDatabase extends _$MoorChatDatabase {
/// Instantiate a new database instance /// Instantiate a new database instance
MoorChatDatabase.connect( MoorChatDatabase.connect(
this._userId, this._userId,
this._isolate,
DatabaseConnection connection, DatabaseConnection connection,
) : super.connect(connection); ) : super.connect(connection);
@@ -69,8 +68,6 @@ class MoorChatDatabase extends _$MoorChatDatabase {
/// User id to which the database is connected /// User id to which the database is connected
String get userId => _userId; String get userId => _userId;
MoorIsolate _isolate;
// you should bump this number whenever you change or add a table definition. // you should bump this number whenever you change or add a table definition.
@override @override
int get schemaVersion => 2; int get schemaVersion => 2;
@@ -89,8 +86,5 @@ class MoorChatDatabase extends _$MoorChatDatabase {
); );
/// Closes the database instance /// Closes the database instance
Future<void> disconnect() async { Future<void> disconnect() => close();
await _isolate?.shutdownAll();
await close();
}
} }
@@ -76,17 +76,21 @@ class SharedDB {
/// [MoorChatDatabase.connect] created on a background isolate. /// [MoorChatDatabase.connect] created on a background isolate.
/// ///
/// Generally used with [ConnectionMode.background]. /// Generally used with [ConnectionMode.background].
static Future<MoorChatDatabase> constructMoorChatDatabase( static MoorChatDatabase constructMoorChatDatabase(
String userId, { String userId, {
bool logStatements = false, bool logStatements = false,
}) async { }) {
final dbName = 'db_$userId'; final dbName = 'db_$userId';
final isolate = await _createMoorIsolate( return MoorChatDatabase.connect(
dbName, userId,
logStatements: logStatements, 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);
} }
} }
@@ -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'; import 'package:stream_chat_persistence/stream_chat_persistence.dart';
/// A Helper class to construct new instances of [MoorChatDatabase] /// A Helper class to construct new instances of [MoorChatDatabase]
@@ -5,7 +7,7 @@ class SharedDB {
/// Returns a new instance of database. /// Returns a new instance of database.
/// ///
/// Generally used with [ConnectionMode.regular]. /// Generally used with [ConnectionMode.regular].
static dynamic constructDatabase( static Future<DelegatedDatabase> constructDatabase(
String userId, { String userId, {
bool logStatements = false, bool logStatements = false,
bool persistOnDisk = true, bool persistOnDisk = true,
@@ -16,7 +18,7 @@ class SharedDB {
/// Return a new instance of moor chat database. /// Return a new instance of moor chat database.
/// ///
/// Generally used with [ConnectionMode.background]. /// Generally used with [ConnectionMode.background].
static dynamic constructMoorChatDatabase( static MoorChatDatabase constructMoorChatDatabase(
String userId, { String userId, {
bool logStatements = false, bool logStatements = false,
}) { }) {
@@ -22,10 +22,10 @@ class SharedDB {
/// default constructor. /// default constructor.
/// ///
/// Generally used with [ConnectionMode.background]. /// Generally used with [ConnectionMode.background].
static Future<MoorChatDatabase> constructMoorChatDatabase( static MoorChatDatabase constructMoorChatDatabase(
String userId, { String userId, {
bool logStatements = false, bool logStatements = false,
}) async { }) {
final dbName = 'db_$userId'; final dbName = 'db_$userId';
return MoorChatDatabase(dbName, logStatements: logStatements); return MoorChatDatabase(dbName, logStatements: logStatements);
} }
@@ -1,4 +1,6 @@
import 'package:logging/logging.dart' show LogRecord;
import 'package:meta/meta.dart'; import 'package:meta/meta.dart';
import 'package:mutex/mutex.dart';
import 'package:stream_chat/stream_chat.dart'; import 'package:stream_chat/stream_chat.dart';
import 'db/moor_chat_database.dart'; import 'db/moor_chat_database.dart';
@@ -13,6 +15,12 @@ enum ConnectionMode {
background, background,
} }
final levelEmojiMapper = {
Level.INFO: '',
Level.WARNING: '⚠️',
Level.SEVERE: '🚨',
};
/// A [MoorChatDatabase] based implementation of the [ChatPersistenceClient] /// A [MoorChatDatabase] based implementation of the [ChatPersistenceClient]
class StreamChatPersistenceClient extends ChatPersistenceClient { class StreamChatPersistenceClient extends ChatPersistenceClient {
/// Creates a new instance of the stream chat persistence client /// 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 /// Connection mode on which the client will work
ConnectionMode connectionMode = ConnectionMode.regular, ConnectionMode connectionMode = ConnectionMode.regular,
Level logLevel = Level.WARNING, Level logLevel = Level.WARNING,
LogHandlerFunction logHandlerFunction,
}) : assert(connectionMode != null), }) : assert(connectionMode != null),
assert(logLevel != null), assert(logLevel != null),
_connectionMode = connectionMode, _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 @visibleForTesting
MoorChatDatabase db; MoorChatDatabase db;
final Logger _logger; final Logger _logger;
final ConnectionMode _connectionMode; 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<T> readProtected<T>(Future<T> Function() f) async {
T ret;
await _mutex.protectRead(() async {
if (db == null) {
return;
}
ret = await f();
});
return ret;
}
@override @override
Future<void> connect(String userId) async { Future<void> connect(String userId) async {
@@ -45,67 +97,105 @@ class StreamChatPersistenceClient extends ChatPersistenceClient {
return; return;
case ConnectionMode.background: case ConnectionMode.background:
_logger.info('Connecting on background isolate'); _logger.info('Connecting on background isolate');
db = await SharedDB.constructMoorChatDatabase(userId); db = SharedDB.constructMoorChatDatabase(userId);
return; return;
} }
} }
@override @override
Future<Event> getConnectionInfo() { Future<Event> getConnectionInfo() {
return db.connectionEventDao.connectionEvent; return readProtected(() {
_logger.info('getConnectionInfo');
return db.connectionEventDao.connectionEvent;
});
} }
@override @override
Future<void> updateConnectionInfo(Event event) { Future<void> updateConnectionInfo(Event event) {
return db.connectionEventDao.updateConnectionEvent(event); return readProtected(() {
_logger.info('updateConnectionInfo');
return db.connectionEventDao.updateConnectionEvent(event);
});
} }
@override @override
Future<void> updateLastSyncAt(DateTime lastSyncAt) { Future<void> updateLastSyncAt(DateTime lastSyncAt) {
return db.connectionEventDao.updateLastSyncAt(lastSyncAt); return readProtected(() {
_logger.info('updateLastSyncAt');
return db.connectionEventDao.updateLastSyncAt(lastSyncAt);
});
} }
@override @override
Future<DateTime> getLastSyncAt() { Future<DateTime> getLastSyncAt() {
return db.connectionEventDao.lastSyncAt; return readProtected(() {
_logger.info('getLastSyncAt');
return db.connectionEventDao.lastSyncAt;
});
} }
@override @override
Future<void> deleteChannels(List<String> cids) { Future<void> deleteChannels(List<String> cids) {
return db.channelDao.deleteChannelByCids(cids); return readProtected(() {
_logger.info('deleteChannels');
return db.channelDao.deleteChannelByCids(cids);
});
} }
@override @override
Future<List<String>> getChannelCids() => db.channelDao.cids; Future<List<String>> getChannelCids() {
return readProtected(() {
_logger.info('getChannelCids');
return db.channelDao.cids;
});
}
@override @override
Future<void> deleteMessageByIds(List<String> messageIds) { Future<void> deleteMessageByIds(List<String> messageIds) {
return db.messageDao.deleteMessageByIds(messageIds); return readProtected(() {
_logger.info('deleteMessageByIds');
return db.messageDao.deleteMessageByIds(messageIds);
});
} }
@override @override
Future<void> deletePinnedMessageByIds(List<String> messageIds) { Future<void> deletePinnedMessageByIds(List<String> messageIds) {
return db.pinnedMessageDao.deleteMessageByIds(messageIds); return readProtected(() {
_logger.info('deletePinnedMessageByIds');
return db.pinnedMessageDao.deleteMessageByIds(messageIds);
});
} }
@override @override
Future<void> deleteMessageByCids(List<String> cids) { Future<void> deleteMessageByCids(List<String> cids) {
return db.messageDao.deleteMessageByCids(cids); return readProtected(() {
_logger.info('deleteMessageByCids');
return db.messageDao.deleteMessageByCids(cids);
});
} }
@override @override
Future<void> deletePinnedMessageByCids(List<String> cids) { Future<void> deletePinnedMessageByCids(List<String> cids) {
return db.pinnedMessageDao.deleteMessageByCids(cids); return readProtected(() {
_logger.info('deletePinnedMessageByCids');
return db.pinnedMessageDao.deleteMessageByCids(cids);
});
} }
@override @override
Future<List<Member>> getMembersByCid(String cid) { Future<List<Member>> getMembersByCid(String cid) {
return db.memberDao.getMembersByCid(cid); return readProtected(() {
_logger.info('getMembersByCid');
return db.memberDao.getMembersByCid(cid);
});
} }
@override @override
Future<ChannelModel> getChannelByCid(String cid) { Future<ChannelModel> getChannelByCid(String cid) {
return db.channelDao.getChannelByCid(cid); return readProtected(() {
_logger.info('getChannelByCid');
return db.channelDao.getChannelByCid(cid);
});
} }
@override @override
@@ -113,10 +203,13 @@ class StreamChatPersistenceClient extends ChatPersistenceClient {
String cid, { String cid, {
PaginationParams messagePagination, PaginationParams messagePagination,
}) { }) {
return db.messageDao.getMessagesByCid( return readProtected(() {
cid, _logger.info('getMessagesByCid');
messagePagination: messagePagination, return db.messageDao.getMessagesByCid(
); cid,
messagePagination: messagePagination,
);
});
} }
@override @override
@@ -124,29 +217,38 @@ class StreamChatPersistenceClient extends ChatPersistenceClient {
String cid, { String cid, {
PaginationParams messagePagination, PaginationParams messagePagination,
}) { }) {
return db.pinnedMessageDao.getMessagesByCid( return readProtected(() {
cid, _logger.info('getPinnedMessagesByCid');
messagePagination: messagePagination, return db.pinnedMessageDao.getMessagesByCid(
); cid,
messagePagination: messagePagination,
);
});
} }
@override @override
Future<List<Read>> getReadsByCid(String cid) { Future<List<Read>> getReadsByCid(String cid) {
return db.readDao.getReadsByCid(cid); return readProtected(() {
_logger.info('getReadsByCid');
return db.readDao.getReadsByCid(cid);
});
} }
@override @override
Future<Map<String, List<Message>>> getChannelThreads(String cid) async { Future<Map<String, List<Message>>> getChannelThreads(String cid) async {
final messages = await db.messageDao.getThreadMessages(cid); return readProtected(() async {
final messageByParentIdDictionary = <String, List<Message>>{}; _logger.info('getChannelThreads');
for (final message in messages) { final messages = await db.messageDao.getThreadMessages(cid);
final parentId = message.parentId; final messageByParentIdDictionary = <String, List<Message>>{};
messageByParentIdDictionary[parentId] = [ for (final message in messages) {
...messageByParentIdDictionary[parentId] ?? [], final parentId = message.parentId;
message messageByParentIdDictionary[parentId] = [
]; ...messageByParentIdDictionary[parentId] ?? [],
} message
return messageByParentIdDictionary; ];
}
return messageByParentIdDictionary;
});
} }
@override @override
@@ -154,10 +256,13 @@ class StreamChatPersistenceClient extends ChatPersistenceClient {
String parentId, { String parentId, {
PaginationParams options, PaginationParams options,
}) { }) {
return db.messageDao.getThreadMessagesByParentId( return readProtected(() async {
parentId, _logger.info('getReplies');
options: options, return db.messageDao.getThreadMessagesByParentId(
); parentId,
options: options,
);
});
} }
@override @override
@@ -166,12 +271,15 @@ class StreamChatPersistenceClient extends ChatPersistenceClient {
List<SortOption<ChannelModel>> sort = const [], List<SortOption<ChannelModel>> sort = const [],
PaginationParams paginationParams, PaginationParams paginationParams,
}) async { }) async {
final channels = await db.channelQueryDao.getChannels( return readProtected(() async {
filter: filter, _logger.info('getChannelStates');
sort: sort, final channels = await db.channelQueryDao.getChannels(
paginationParams: paginationParams, filter: filter,
); sort: sort,
return Future.wait(channels.map((e) => getChannelStateByCid(e.cid))); paginationParams: paginationParams,
);
return Future.wait(channels.map((e) => getChannelStateByCid(e.cid)));
});
} }
@override @override
@@ -180,72 +288,114 @@ class StreamChatPersistenceClient extends ChatPersistenceClient {
List<String> cids, List<String> cids,
bool clearQueryCache, bool clearQueryCache,
) { ) {
return db.channelQueryDao.updateChannelQueries( return readProtected(() async {
filter, _logger.info('updateChannelQueries');
cids, return db.channelQueryDao.updateChannelQueries(
clearQueryCache, filter,
); cids,
clearQueryCache,
);
});
} }
@override @override
Future<void> updateChannels(List<ChannelModel> channels) { Future<void> updateChannels(List<ChannelModel> channels) {
return db.channelDao.updateChannels(channels); return readProtected(() async {
_logger.info('updateChannels');
return db.channelDao.updateChannels(channels);
});
} }
@override @override
Future<void> updateMembers(String cid, List<Member> members) { Future<void> updateMembers(String cid, List<Member> members) {
return db.memberDao.updateMembers(cid, members); return readProtected(() async {
_logger.info('updateMembers');
return db.memberDao.updateMembers(cid, members);
});
} }
@override @override
Future<void> updateMessages(String cid, List<Message> messages) { Future<void> updateMessages(String cid, List<Message> messages) {
return db.messageDao.updateMessages(cid, messages); return readProtected(() async {
_logger.info('updateMessages');
return db.messageDao.updateMessages(cid, messages);
});
} }
@override @override
Future<void> updatePinnedMessages(String cid, List<Message> messages) { Future<void> updatePinnedMessages(String cid, List<Message> messages) {
return db.pinnedMessageDao.updateMessages(cid, messages); return readProtected(() async {
_logger.info('updatePinnedMessages');
return db.pinnedMessageDao.updateMessages(cid, messages);
});
} }
@override @override
Future<void> updateReactions(List<Reaction> reactions) { Future<void> updateReactions(List<Reaction> reactions) {
return db.reactionDao.updateReactions(reactions); return readProtected(() async {
_logger.info('updateReactions');
return db.reactionDao.updateReactions(reactions);
});
} }
@override @override
Future<void> updateReads(String cid, List<Read> reads) { Future<void> updateReads(String cid, List<Read> reads) {
return db.readDao.updateReads(cid, reads); return readProtected(() async {
_logger.info('updateReads');
return db.readDao.updateReads(cid, reads);
});
} }
@override @override
Future<void> updateUsers(List<User> users) { Future<void> updateUsers(List<User> users) {
return db.userDao.updateUsers(users); return readProtected(() async {
_logger.info('updateUsers');
return db.userDao.updateUsers(users);
});
} }
@override @override
Future<void> deleteReactionsByMessageId(List<String> messageIds) { Future<void> deleteReactionsByMessageId(List<String> messageIds) {
return db.reactionDao.deleteReactionsByMessageIds(messageIds); return readProtected(() async {
_logger.info('deleteReactionsByMessageId');
return db.reactionDao.deleteReactionsByMessageIds(messageIds);
});
} }
@override @override
Future<void> deleteMembersByCids(List<String> cids) { Future<void> deleteMembersByCids(List<String> cids) {
return db.memberDao.deleteMemberByCids(cids); return readProtected(() async {
_logger.info('deleteMembersByCids');
return db.memberDao.deleteMemberByCids(cids);
});
}
@override
Future<void> updateChannelStates(List<ChannelState> channelStates) {
return readProtected(() async {
return db.transaction(() async {
await super.updateChannelStates(channelStates);
});
});
} }
@override @override
Future<void> disconnect({bool flush = false}) async { Future<void> disconnect({bool flush = false}) async {
if (db != null) { return _mutex.protectWrite(() async {
_logger.info('Disconnecting'); _logger.info('disconnect');
if (flush) { if (db != null) {
_logger.info('Flushing'); _logger.info('Disconnecting');
await db.batch((batch) { if (flush) {
db.allTables.forEach((table) { _logger.info('Flushing');
db.delete(table).go(); await db.batch((batch) {
db.allTables.forEach((table) {
db.delete(table).go();
});
}); });
}); }
await db.disconnect();
db = null;
} }
await db.disconnect(); });
db = null;
}
} }
} }
@@ -11,6 +11,7 @@ environment:
dependencies: dependencies:
flutter: flutter:
sdk: flutter sdk: flutter
mutex: ^2.0.0
moor: ^3.4.0 moor: ^3.4.0
path: ^1.7.0 path: ^1.7.0
path_provider: ^1.6.27 path_provider: ^1.6.27