From 6b29630a644db6e7cecec771850901dc63320b62 Mon Sep 17 00:00:00 2001 From: Sahil Kumar Date: Fri, 2 Jun 2023 12:24:04 +0530 Subject: [PATCH] feat(llc, persistence): add support for open and close persistence via client. Signed-off-by: xsahil03x --- .../stream_chat/lib/src/client/channel.dart | 17 ++- .../stream_chat/lib/src/client/client.dart | 117 ++++++++++++------ .../lib/src/db/chat_persistence_client.dart | 5 + .../src/db/chat_persistence_client_test.dart | 3 + packages/stream_chat/test/src/mocks.dart | 11 +- .../src/stream_chat_persistence_client.dart | 3 + 6 files changed, 103 insertions(+), 53 deletions(-) diff --git a/packages/stream_chat/lib/src/client/channel.dart b/packages/stream_chat/lib/src/client/channel.dart index 3b249fca..54feb2ab 100644 --- a/packages/stream_chat/lib/src/client/channel.dart +++ b/packages/stream_chat/lib/src/client/channel.dart @@ -1379,13 +1379,14 @@ class Channel { this.state?.updateChannelState(updatedState); return updatedState; } catch (e) { - if (!_client.persistenceEnabled) { - rethrow; + if (_client.persistenceEnabled) { + return _client.chatPersistenceClient!.getChannelStateByCid( + cid!, + messagePagination: messagesPagination, + ); } - return _client.chatPersistenceClient!.getChannelStateByCid( - cid!, - messagePagination: messagesPagination, - ); + + rethrow; } } @@ -1841,9 +1842,7 @@ class ChannelClientState { /// [isUpToDate] flag count as a stream. Stream get isUpToDateStream => _isUpToDateController.stream; - - final BehaviorSubject _isUpToDateController = - BehaviorSubject.seeded(true); + final _isUpToDateController = BehaviorSubject.seeded(true); /// The retry queue associated to this channel. late final RetryQueue _retryQueue; diff --git a/packages/stream_chat/lib/src/client/client.dart b/packages/stream_chat/lib/src/client/client.dart index 7627cf1d..54fefd57 100644 --- a/packages/stream_chat/lib/src/client/client.dart +++ b/packages/stream_chat/lib/src/client/client.dart @@ -125,10 +125,6 @@ class StreamChatClient { final _tokenManager = TokenManager(); final _connectionIdManager = ConnectionIdManager(); - set chatPersistenceClient(ChatPersistenceClient? value) { - _originalChatPersistenceClient = value; - } - /// Default user agent for all requests static String defaultUserAgent = 'stream-chat-dart-client-${CurrentPlatform.name}'; @@ -139,15 +135,15 @@ class StreamChatClient { /// The current package version static const packageVersion = PACKAGE_VERSION; - ChatPersistenceClient? _originalChatPersistenceClient; - /// Chat persistence client - ChatPersistenceClient? get chatPersistenceClient => _chatPersistenceClient; + ChatPersistenceClient? chatPersistenceClient; - ChatPersistenceClient? _chatPersistenceClient; - - /// Whether the chat persistence is available or not - bool get persistenceEnabled => _chatPersistenceClient != null; + /// Returns `True` if the [chatPersistenceClient] is available and connected. + /// Otherwise, returns `False`. + bool get persistenceEnabled { + final client = chatPersistenceClient; + return client != null && client.isConnected; + } late final RetryPolicy _retryPolicy; @@ -324,20 +320,27 @@ class StreamChatClient { final ownUser = OwnUser.fromUser(user); state.currentUser = ownUser; - if (!connectWebSocket) return ownUser; - try { - if (_originalChatPersistenceClient != null) { - _chatPersistenceClient = _originalChatPersistenceClient; - await _chatPersistenceClient!.connect(ownUser.id); + // Connect to persistence client if its set. + if (chatPersistenceClient != null) { + await openPersistenceConnection(ownUser); } - final connectedUser = await openConnection( - includeUserDetailsInConnectCall: true, - ); - return state.currentUser = connectedUser; + + // Connect to websocket if [connectWebSocket] is true. + // + // This is useful when you want to connect to websocket + // at a later stage or use the client in connection-less mode. + if (connectWebSocket) { + final connectedUser = await openConnection( + includeUserDetailsInConnectCall: true, + ); + state.currentUser = connectedUser; + } + + return state.currentUser!; } catch (e, stk) { if (e is StreamWebSocketError && e.isRetriable) { - final event = await _chatPersistenceClient?.getConnectionInfo(); + final event = await chatPersistenceClient?.getConnectionInfo(); if (event != null) return ownUser.merge(event.me); } logger.severe('error connecting user : ${ownUser.id}', e, stk); @@ -345,6 +348,40 @@ class StreamChatClient { } } + /// Connects the [chatPersistenceClient] to the given [user]. + Future openPersistenceConnection(User user) async { + final client = chatPersistenceClient; + if (client == null) { + throw const StreamChatError('Chat persistence client is not set'); + } + + if (client.isConnected) { + // If the persistence client is already connected to the userId, + // we don't need to connect again. + if (client.userId == user.id) return; + + throw const StreamChatError(''' + Chat persistence client is already connected to a different user, + please close the connection before connecting a new one.'''); + } + + // Connect the persistence client to the userId. + return client.connect(user.id); + } + + /// Disconnects the [chatPersistenceClient] from the current user. + Future closePersistenceConnection({bool flush = false}) async { + final client = chatPersistenceClient; + // If the persistence client is never connected, we don't need to close it. + if (client == null || !client.isConnected) { + logger.info('Chat persistence client is not connected'); + return; + } + + // Disconnect the persistence client. + return client.disconnect(flush: flush); + } + /// Creates a new WebSocket connection with the current user. /// If [includeUserDetailsInConnectCall] is true it will include the current /// user details in the connect call. @@ -422,7 +459,7 @@ class StreamChatClient { final connectionId = event.connectionId; if (connectionId != null) { _connectionIdManager.setConnectionId(connectionId); - _chatPersistenceClient?.updateConnectionInfo(event); + chatPersistenceClient?.updateConnectionInfo(event); } } @@ -460,9 +497,9 @@ class StreamChatClient { // channels are empty, assuming it's a fresh start // and making sure `lastSyncAt` is initialized if (persistenceEnabled) { - final lastSyncAt = await _chatPersistenceClient?.getLastSyncAt(); + final lastSyncAt = await chatPersistenceClient?.getLastSyncAt(); if (lastSyncAt == null) { - await _chatPersistenceClient?.updateLastSyncAt(DateTime.now()); + await chatPersistenceClient?.updateLastSyncAt(DateTime.now()); } } } @@ -493,13 +530,12 @@ class StreamChatClient { /// Will automatically fetch [cids] and [lastSyncedAt] if [persistenceEnabled] Future sync({List? cids, DateTime? lastSyncAt}) { return synchronized(() async { - final channels = cids ?? await _chatPersistenceClient?.getChannelCids(); + final channels = cids ?? await chatPersistenceClient?.getChannelCids(); if (channels == null || channels.isEmpty) { return; } - final syncAt = - lastSyncAt ?? await _chatPersistenceClient?.getLastSyncAt(); + final syncAt = lastSyncAt ?? await chatPersistenceClient?.getLastSyncAt(); if (syncAt == null) { return; } @@ -520,7 +556,7 @@ class StreamChatClient { final now = DateTime.now(); _lastSyncedAt = now; - _chatPersistenceClient?.updateLastSyncAt(now); + chatPersistenceClient?.updateLastSyncAt(now); } catch (e, stk) { logger.severe('Error during sync', e, stk); } @@ -679,7 +715,7 @@ class StreamChatClient { final updateData = _mapChannelStateToChannel(channels); - await _chatPersistenceClient?.updateChannelQueries( + await chatPersistenceClient?.updateChannelQueries( filter, channels.map((c) => c.channel!.cid).toList(), clearQueryCache: paginationParams.offset == 0, @@ -698,7 +734,7 @@ class StreamChatClient { List>? channelStateSort, PaginationParams paginationParams = const PaginationParams(), }) async { - final offlineChannels = (await _chatPersistenceClient?.getChannelStates( + final offlineChannels = (await chatPersistenceClient?.getChannelStates( filter: filter, // ignore: deprecated_member_use_from_same_package sort: sort, @@ -1362,7 +1398,7 @@ class StreamChatClient { final response = await _chatApi.message.deleteMessage(messageId, hard: hard); if (hard == true) { - await _chatPersistenceClient?.deleteMessageById(messageId); + await chatPersistenceClient?.deleteMessageById(messageId); } return response; } @@ -1468,34 +1504,33 @@ class StreamChatClient { Future disconnectUser({bool flushChatPersistence = false}) async { logger.info('Disconnecting user : ${state.currentUser?.id}'); - // resetting state + // resetting state. state.dispose(); state = ClientState(this); _lastSyncedAt = null; - // resetting credentials + // resetting credentials. _tokenManager.reset(); _connectionIdManager.reset(); - // disconnecting persistence client - await _chatPersistenceClient?.disconnect(flush: flushChatPersistence); - _chatPersistenceClient = null; + // closing persistence connection. + await closePersistenceConnection(flush: flushChatPersistence); // closing web-socket connection - closeConnection(); + return closeConnection(); } /// Call this function to dispose the client Future dispose() async { logger.info('Disposing new StreamChatClient'); - // disposing state + // disposing state. state.dispose(); - // disconnecting persistence client - await _chatPersistenceClient?.disconnect(); + // closing persistence connection. + await closePersistenceConnection(); - // closing web-socket connection + // closing web-socket connection. closeConnection(); await _eventController.close(); diff --git a/packages/stream_chat/lib/src/db/chat_persistence_client.dart b/packages/stream_chat/lib/src/db/chat_persistence_client.dart index 4710190e..bf0b1a76 100644 --- a/packages/stream_chat/lib/src/db/chat_persistence_client.dart +++ b/packages/stream_chat/lib/src/db/chat_persistence_client.dart @@ -17,6 +17,11 @@ abstract class ChatPersistenceClient { /// Whether the connection is established. bool get isConnected; + /// The current user id to which the client is connected. + /// + /// Returns `null` if the client is not connected. + String? get userId; + /// Creates a new connection to the client Future connect(String userId); diff --git a/packages/stream_chat/test/src/db/chat_persistence_client_test.dart b/packages/stream_chat/test/src/db/chat_persistence_client_test.dart index 62a7cf25..a0a12935 100644 --- a/packages/stream_chat/test/src/db/chat_persistence_client_test.dart +++ b/packages/stream_chat/test/src/db/chat_persistence_client_test.dart @@ -15,6 +15,9 @@ class TestPersistenceClient extends ChatPersistenceClient { @override bool get isConnected => throw UnimplementedError(); + @override + String? get userId => throw UnimplementedError(); + @override Future connect(String userId) => throw UnimplementedError(); diff --git a/packages/stream_chat/test/src/mocks.dart b/packages/stream_chat/test/src/mocks.dart index 44078519..86a77a66 100644 --- a/packages/stream_chat/test/src/mocks.dart +++ b/packages/stream_chat/test/src/mocks.dart @@ -64,11 +64,16 @@ class MockAttachmentFileUploader extends Mock implements AttachmentFileUploader {} class MockPersistenceClient extends Mock implements ChatPersistenceClient { - @override - Future connect(String userId) => Future.value(); + bool _isConnected = false; @override - Future disconnect({bool flush = false}) => Future.value(); + bool get isConnected => _isConnected; + + @override + Future connect(String userId) async => _isConnected = true; + + @override + Future disconnect({bool flush = false}) async => _isConnected = false; } class MockStreamChatClient extends Mock implements StreamChatClient { 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 374bfa59..e8737baa 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 @@ -82,6 +82,9 @@ class StreamChatPersistenceClient extends ChatPersistenceClient { @override bool get isConnected => db != null; + @override + String? get userId => db?.userId; + @override Future connect( String userId, {