diff --git a/packages/stream_chat/CHANGELOG.md b/packages/stream_chat/CHANGELOG.md index a063de46f1..182f38bdf8 100644 --- a/packages/stream_chat/CHANGELOG.md +++ b/packages/stream_chat/CHANGELOG.md @@ -5,6 +5,7 @@ - Added `StreamChatClient.isLocalUnreadCountEnabled` (default `false`). When enabled, channels that have read events disabled (e.g. livestream channel types) track their unread count locally, on-device: incoming messages increment it, hard-deleted messages decrement it, and `Channel.markRead` / `markUnread` / `markUnreadByTimestamp` update it locally without a network request — including `Read.lastReadMessageId`, so the unread divider and jump-to-unread button anchor to the right message. Channels that support read receipts are unaffected and keep relying on server-driven unread counts. - Added `Event.watcherCount`, exposing the server-provided `watcher_count` field on events (e.g. `user.watching.start`, `user.watching.stop`, `message.new`). - Added `StreamChatNetworkError.type` (a `StreamChatNetworkErrorType` capturing the transport failure kind — connection error, timeout, cancellation, etc.). +- Added support for sending and deleting reactions while offline. ⚠️ Deprecated @@ -13,6 +14,7 @@ 🔄 Changed - Raised the minimum `dio` version to `^5.11.0`. +- `Channel.sendReaction` and `Channel.deleteReaction` now keep the optimistic change on a transient/offline error and replay it when the connection recovers, instead of reverting it. 🐞 Fixed diff --git a/packages/stream_chat/lib/src/client/channel.dart b/packages/stream_chat/lib/src/client/channel.dart index 6bacb60a94..60ca5589a8 100644 --- a/packages/stream_chat/lib/src/client/channel.dart +++ b/packages/stream_chat/lib/src/client/channel.dart @@ -5,6 +5,7 @@ import 'dart:math' as math; import 'package:collection/collection.dart'; import 'package:rxdart/rxdart.dart'; +import 'package:stream_chat/src/client/reaction_pending_operation.dart'; import 'package:stream_chat/src/client/retry_queue.dart'; import 'package:stream_chat/src/core/util/utils.dart'; import 'package:stream_chat/stream_chat.dart'; @@ -1616,19 +1617,30 @@ class Channel { state?.updateMessage(updatedMessage); try { - final reactionResp = await _client.sendReaction( + return await _client.sendReaction( messageId, reaction, skipPush: skipPush, enforceUnique: enforceUnique, ); - return reactionResp; - } catch (_) { - // Reset the message if the update fails. Use replace (not merge) - // so the rollback wins over the optimistic local state — otherwise - // `Message.updateWith`'s enrichment preservation would keep the - // optimistic `ownReactions` for messages that previously had none. - state?.replaceMessage(message); + } catch (e) { + final retriable = e is StreamChatNetworkError && e.isRetriable; + if (retriable) { + // Keep the optimistic reaction and queue it for replay on reconnect. + await _client.pendingOperationsManager.enqueue( + ReactionPendingOperation.add( + reaction, + skipPush: skipPush, + enforceUnique: enforceUnique, + ), + ); + } else { + // Reset the message on terminal failure. Use replace (not merge) + // so the rollback wins over the optimistic local state — otherwise + // `Message.updateWith`'s enrichment preservation would keep the + // optimistic `ownReactions` for messages that previously had none. + state?.replaceMessage(message); + } rethrow; } } @@ -1647,15 +1659,25 @@ class Channel { state?.updateMessage(updatedMessage); try { - final deleteResponse = await _client.deleteReaction( + return await _client.deleteReaction( message.id, reaction.type, ); - return deleteResponse; - } catch (_) { - // Reset the message if the update fails. Use replace (not merge) - // for symmetry with `sendReaction` — see that method for context. - state?.replaceMessage(message); + } catch (e) { + final retriable = e is StreamChatNetworkError && e.isRetriable; + if (retriable) { + // Keep the optimistic removal and queue it for replay on reconnect. + await _client.pendingOperationsManager.enqueue( + ReactionPendingOperation.delete( + messageId: message.id, + reactionType: reaction.type, + ), + ); + } else { + // Reset the message on terminal failure. Use replace (not merge) + // for symmetry with `sendReaction` — see that method for context. + state?.replaceMessage(message); + } rethrow; } } diff --git a/packages/stream_chat/lib/src/client/client.dart b/packages/stream_chat/lib/src/client/client.dart index d56a37c210..dacc9f2bd0 100644 --- a/packages/stream_chat/lib/src/client/client.dart +++ b/packages/stream_chat/lib/src/client/client.dart @@ -8,6 +8,7 @@ import 'package:rxdart/rxdart.dart'; import 'package:stream_chat/src/client/channel.dart'; import 'package:stream_chat/src/client/channel_delivery_reporter.dart'; import 'package:stream_chat/src/client/event_resolvers.dart' as event_resolvers; +import 'package:stream_chat/src/client/pending_operations_manager.dart'; import 'package:stream_chat/src/client/query_channels_result.dart'; import 'package:stream_chat/src/client/retry_policy.dart'; import 'package:stream_chat/src/core/api/attachment_file_uploader.dart'; @@ -165,6 +166,11 @@ class StreamChatClient { late final _appSettingsManager = AppSettingsManager(_chatApi.general); static final _systemEnvironmentManager = SystemEnvironmentManager(); + /// Owns the queue of pending operations (e.g. reactions added or removed + /// while offline) and replays them when the connection recovers. + @internal + late final pendingOperationsManager = PendingOperationsManager(this); + /// Updates the system environment information used by the client. /// /// The passed [environment] is sanitized before being applied: @@ -435,6 +441,8 @@ class StreamChatClient { // Connect to persistence client if its set. if (chatPersistenceClient != null) { await openPersistenceConnection(ownUser); + // Restore any operations that were queued before process death. + await pendingOperationsManager.hydrate(); } // Connect to websocket if [connectWebSocket] is true. @@ -607,6 +615,11 @@ class StreamChatClient { final connectionRecovered = !wasConnected && isConnected; if (connectionRecovered) { + // Replay pending offline operations (e.g. reactions) BEFORE any + // server-state refresh, so the server has each mutation before a re-query + // returns state that would otherwise clobber the optimistic change. + await pendingOperationsManager.replay(); + // connection recovered final cids = [...state.channels.keys.toSet()]; if (cids.isNotEmpty) { @@ -2538,6 +2551,10 @@ class StreamChatClient { state.dispose(); state = ClientState(this); + // clearing the in-memory pending-operation queue so a user's queued + // operations never replay under the next connected user. + pendingOperationsManager.clear(); + // clearing app settings cache. _appSettingsManager.clear(); diff --git a/packages/stream_chat/lib/src/client/pending_operations_manager.dart b/packages/stream_chat/lib/src/client/pending_operations_manager.dart new file mode 100644 index 0000000000..bfedaa485f --- /dev/null +++ b/packages/stream_chat/lib/src/client/pending_operations_manager.dart @@ -0,0 +1,205 @@ +import 'package:meta/meta.dart'; +import 'package:stream_chat/src/client/client.dart'; +import 'package:stream_chat/src/client/reaction_pending_operation.dart'; +import 'package:stream_chat/src/core/error/error.dart'; +import 'package:stream_chat/src/core/models/pending_operation.dart'; +import 'package:stream_chat/src/core/models/reaction.dart'; + +/// Owns the queue of [PendingOperation]s and replays them against the server +/// when the connection is recovered. +/// +/// The queue lives in memory for the current session and is the single source +/// replayed from. When persistence is enabled the queue is additionally +/// mirrored to [StreamChatClient.chatPersistenceClient], so operations survive +/// process death: [hydrate] loads them back into memory on the next connect. +/// Without persistence the queue is session-only, giving reactions same-session +/// replay across transient outages. +/// +/// Replay is at-least-once: an operation is removed only after the server +/// accepts or terminally rejects it, so a crash between acceptance and removal +/// re-sends it on the next recovery. Every operation type handled by +/// [_replayCallFor] must therefore be idempotent on the server — e.g. reactions +/// dedupe by (message, type, user). +@internal +class PendingOperationsManager { + /// Creates a manager for [client]'s pending-operation queue. + PendingOperationsManager(this._client); + + final StreamChatClient _client; + + /// The in-memory queue, replayed in insertion order. + /// + /// Every entry carries a non-null [PendingOperation.id]: a positive DB + /// autoincrement id when the operation is mirrored to persistence, or a + /// negative session id otherwise. The two ranges never collide. + final _operations = []; + + // Source of negative, session-only ids for operations that are not persisted. + int _memorySeq = 0; + int _nextMemoryId() => --_memorySeq; + + // Prevents overlapping replays. + bool _isReplaying = false; + + /// Appends [operation] to the queue, mirroring it to persistence when + /// enabled so it survives process death. + Future enqueue(PendingOperation operation) async { + int? id; + if (_client.persistenceEnabled) { + try { + id = await _client.chatPersistenceClient!.insertPendingOperation( + operation, + ); + } catch (error, stk) { + // Keep the operation in memory so it still replays this session. + _client.logger.warning( + 'Failed to persist pending operation', + error, + stk, + ); + } + } + _operations.add(operation.copyWith(id: id ?? _nextMemoryId())); + } + + /// Loads any persisted operations into the in-memory queue. + /// + /// Called once per connect to restore operations that survived process + /// death. A no-op when persistence is disabled. + Future hydrate() async { + if (!_client.persistenceEnabled) return; + try { + final stored = await _client.chatPersistenceClient!.getPendingOperations(); + _operations + ..clear() + ..addAll(stored); + } catch (error, stk) { + _client.logger.warning( + 'Failed to hydrate pending operations', + error, + stk, + ); + } + } + + /// Empties the in-memory queue. + /// + /// Must be called on disconnect so a user's queued operations never replay + /// under a different user. The persisted mirror is user-scoped and closed + /// separately with the persistence connection. + void clear() { + _operations.clear(); + _memorySeq = 0; + } + + /// Removes the operation with the given [id] from memory and, when persisted, + /// from the mirror. A delete of a memory-only (negative) id is a no-op. + Future _remove(int id) async { + _operations.removeWhere((it) => it.id == id); + if (!_client.persistenceEnabled || id < 0) return; + try { + await _client.chatPersistenceClient!.deletePendingOperation(id); + } catch (error, stk) { + _client.logger.warning( + 'Failed to delete pending operation $id', + error, + stk, + ); + } + } + + /// Replays each queued operation against the server in insertion order. + Future replay() async { + if (_isReplaying) return; + _isReplaying = true; + + try { + // Copy so removals during replay don't mutate the list being iterated. + final operations = List.of(_operations); + for (final operation in operations) { + try { + final Future Function()? call; + try { + call = _replayCallFor(operation); + } catch (error, stk) { + // Malformed payload for a known type — can never be replayed. + _client.logger.warning( + 'Dropping unreplayable pending operation ${operation.id}', + error, + stk, + ); + await _remove(operation.id!); + continue; + } + + if (call == null) { + // Unknown type (e.g. persisted by a newer app version) — drop it. + _client.logger.warning( + 'Dropping unknown pending operation type "${operation.type}" ' + '(${operation.id})', + ); + await _remove(operation.id!); + continue; + } + + try { + await call(); + } on StreamChatNetworkError catch (error) { + // Keep transient failures for the next recovery. + if (error.isRetriable) continue; + } + + // Accepted or terminally rejected by the server — drop it. + await _remove(operation.id!); + } catch (error, stk) { + _client.logger.warning( + 'Error replaying pending operation ${operation.id}', + error, + stk, + ); + } + } + } catch (error, stk) { + _client.logger.severe( + 'Error replaying pending operations', + error, + stk, + ); + } finally { + _isReplaying = false; + } + } + + /// Returns the server call that replays [operation], or `null` if its type + /// is unknown to this version. + Future Function()? _replayCallFor(PendingOperation operation) { + switch (operation.type) { + case ReactionPendingOperation.addType: + final targetMessageId = operation.targetMessageId; + if (targetMessageId == null) { + throw StateError('Missing targetMessageId for ${operation.type}'); + } + final reaction = Reaction.fromJson( + operation.payload[ReactionPendingOperation.reactionKey] as Map, + ); + final skipPush = operation.payload[ReactionPendingOperation.skipPushKey] as bool? ?? false; + final enforceUnique = operation.payload[ReactionPendingOperation.enforceUniqueKey] as bool? ?? false; + return () => _client.sendReaction( + targetMessageId, + reaction, + skipPush: skipPush, + enforceUnique: enforceUnique, + ); + case ReactionPendingOperation.deleteType: + final targetMessageId = operation.targetMessageId; + if (targetMessageId == null) { + throw StateError('Missing targetMessageId for ${operation.type}'); + } + final reactionType = operation.payload[ReactionPendingOperation.reactionTypeKey] as String; + return () => _client.deleteReaction(targetMessageId, reactionType); + default: + // Unknown operation type — cannot be replayed by this version. + return null; + } + } +} diff --git a/packages/stream_chat/lib/src/client/reaction_pending_operation.dart b/packages/stream_chat/lib/src/client/reaction_pending_operation.dart new file mode 100644 index 0000000000..73194c0368 --- /dev/null +++ b/packages/stream_chat/lib/src/client/reaction_pending_operation.dart @@ -0,0 +1,50 @@ +import 'package:stream_chat/src/core/models/pending_operation.dart'; +import 'package:stream_chat/src/core/models/reaction.dart'; + +/// Builds and identifies the reaction-specific forms of [PendingOperation]. +abstract class ReactionPendingOperation { + /// The [PendingOperation.type] discriminator for a reaction add. + static const addType = 'reaction.add'; + + /// The [PendingOperation.type] discriminator for a reaction delete. + static const deleteType = 'reaction.delete'; + + /// The [PendingOperation.payload] key holding the serialized reaction of an + /// add. + static const reactionKey = 'reaction'; + + /// The [PendingOperation.payload] key holding the `enforce_unique` flag of an + /// add. + static const enforceUniqueKey = 'enforce_unique'; + + /// The [PendingOperation.payload] key holding the `skip_push` flag of an add. + static const skipPushKey = 'skip_push'; + + /// The [PendingOperation.payload] key holding the reaction type of a delete. + static const reactionTypeKey = 'reaction_type'; + + /// Builds the pending operation recording an optimistic reaction add. + static PendingOperation add( + Reaction reaction, { + required bool skipPush, + required bool enforceUnique, + }) => PendingOperation( + type: addType, + targetMessageId: reaction.messageId, + payload: { + reactionKey: reaction.toJson(), + enforceUniqueKey: enforceUnique, + skipPushKey: skipPush, + }, + ); + + /// Builds the pending operation recording an optimistic reaction delete. + static PendingOperation delete({ + required String messageId, + required String reactionType, + }) => PendingOperation( + type: deleteType, + targetMessageId: messageId, + payload: {reactionTypeKey: reactionType}, + ); +} diff --git a/packages/stream_chat/lib/src/core/models/pending_operation.dart b/packages/stream_chat/lib/src/core/models/pending_operation.dart new file mode 100644 index 0000000000..b9d9e464ed --- /dev/null +++ b/packages/stream_chat/lib/src/core/models/pending_operation.dart @@ -0,0 +1,51 @@ +import 'package:equatable/equatable.dart'; + +/// {@template pendingOperation} +/// A durable record of an optimistic mutation awaiting server confirmation +/// (e.g. a reaction added or removed while offline). +/// +/// Operations are replayed at-least-once on reconnect, so the server-side +/// effect of every operation type must be idempotent. +/// {@endtemplate} +class PendingOperation extends Equatable { + /// {@macro pendingOperation} + const PendingOperation({ + required this.type, + required this.payload, + this.id, + this.targetMessageId, + }); + + /// The database autoincrement id, assigned when the operation is stored; + /// `null` until then. + final int? id; + + /// The discriminator persisted in the `type` column, e.g. `reaction.add`. + final String type; + + /// The id of the message the operation targets, if any. + final String? targetMessageId; + + /// The operation-specific value fields, stored as JSON. + final Map payload; + + /// Returns a copy of this operation with the given fields replaced. + PendingOperation copyWith({ + int? id, + String? type, + String? targetMessageId, + Map? payload, + }) => PendingOperation( + id: id ?? this.id, + type: type ?? this.type, + targetMessageId: targetMessageId ?? this.targetMessageId, + payload: payload ?? this.payload, + ); + + @override + List get props => [ + type, + targetMessageId, + payload, + ]; +} 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 98d004ea10..794f9c1359 100644 --- a/packages/stream_chat/lib/src/db/chat_persistence_client.dart +++ b/packages/stream_chat/lib/src/db/chat_persistence_client.dart @@ -11,6 +11,7 @@ import 'package:stream_chat/src/core/models/filter.dart'; import 'package:stream_chat/src/core/models/location.dart'; import 'package:stream_chat/src/core/models/member.dart'; import 'package:stream_chat/src/core/models/message.dart'; +import 'package:stream_chat/src/core/models/pending_operation.dart'; import 'package:stream_chat/src/core/models/poll.dart'; import 'package:stream_chat/src/core/models/poll_vote.dart'; import 'package:stream_chat/src/core/models/reaction.dart'; @@ -536,6 +537,24 @@ abstract class ChatPersistenceClient { ]); } + /// Inserts [operation] into the pending-operation queue, returning the id + /// assigned to the stored row (or `null` when nothing was stored). + /// + /// Pending operations are optimistic mutations queued for replay once + /// connectivity is restored. The default no-op drops them unless overridden + /// alongside [getPendingOperations] and [deletePendingOperation]. + Future insertPendingOperation(PendingOperation operation) async => null; + + /// Returns all stored pending operations ordered by insertion. + /// + /// Defaults to an empty list; see [insertPendingOperation]. + Future> getPendingOperations() async => []; + + /// Deletes the pending operation with the given [id]. + /// + /// No-op by default; see [insertPendingOperation]. + Future deletePendingOperation(int id) async {} + List _expandReactions(Message message) { final own = message.ownReactions; final latest = message.latestReactions; diff --git a/packages/stream_chat/lib/stream_chat.dart b/packages/stream_chat/lib/stream_chat.dart index 748d181189..a4812fd156 100644 --- a/packages/stream_chat/lib/stream_chat.dart +++ b/packages/stream_chat/lib/stream_chat.dart @@ -58,6 +58,7 @@ export 'src/core/models/message_state.dart'; export 'src/core/models/moderation.dart'; export 'src/core/models/mute.dart'; export 'src/core/models/own_user.dart'; +export 'src/core/models/pending_operation.dart'; export 'src/core/models/poll.dart'; export 'src/core/models/poll_option.dart'; export 'src/core/models/poll_vote.dart'; diff --git a/packages/stream_chat/test/src/client/channel_test.dart b/packages/stream_chat/test/src/client/channel_test.dart index d3330511d8..ba4c038d74 100644 --- a/packages/stream_chat/test/src/client/channel_test.dart +++ b/packages/stream_chat/test/src/client/channel_test.dart @@ -238,6 +238,10 @@ void main() { when( () => client.channelDeliveryReporter.submitForDelivery(any()), ).thenAnswer((_) async {}); + + // No persistence in this group by default; pending-operation tests + // install a MockPersistenceClient locally. + when(() => client.chatPersistenceClient).thenReturn(null); }); // Setting up a initialized channel @@ -253,6 +257,8 @@ void main() { tearDown(() { channel.dispose(); + // Restore the no-persistence default for the next test. + when(() => client.chatPersistenceClient).thenReturn(null); clearInteractions(client); }); @@ -2981,7 +2987,12 @@ void main() { when( () => client.sendReaction(message.id, reaction), - ).thenThrow(StreamChatNetworkError(ChatErrorCode.inputError)); + ).thenThrow( + StreamChatNetworkError( + ChatErrorCode.inputError, + data: ErrorResponse()..statusCode = 400, + ), + ); expectLater( // skipping first seed message list -> [] messages @@ -3105,6 +3116,35 @@ void main() { ).called(1); }, ); + + test( + 'a retriable failure keeps the optimistic reaction for replay', + () async { + const type = 'like'; + final message = Message(id: 'offline-msg', state: MessageState.sent); + final reaction = Reaction( + type: type, + messageId: message.id, + user: client.state.currentUser, + ); + + // data == null → retriable/offline error. + when( + () => client.sendReaction(message.id, reaction), + ).thenThrow(StreamChatNetworkError(ChatErrorCode.inputError)); + + await expectLater( + channel.sendReaction(message, reaction), + throwsA(isA()), + ); + + // The optimistic reaction is kept (queued for replay on reconnect) + final current = channel.state!.messages.firstWhere( + (m) => m.id == message.id, + ); + expect(current.ownReactions?.map((r) => r.type), contains(type)); + }, + ); }); group('`.sendReaction in thread`', () { @@ -3185,7 +3225,12 @@ void main() { when( () => client.sendReaction(message.id, reaction), - ).thenThrow(StreamChatNetworkError(ChatErrorCode.inputError)); + ).thenThrow( + StreamChatNetworkError( + ChatErrorCode.inputError, + data: ErrorResponse()..statusCode = 400, + ), + ); expectLater( // skipping first seed message list -> [] messages @@ -3367,7 +3412,7 @@ void main() { }); test( - 'should restore prev message state if `client.deleteReaction` throws', + 'restores the reaction if `client.deleteReaction` throws terminally', () async { const userId = 'test-user-id'; const messageId = 'test-message-id'; @@ -3388,12 +3433,23 @@ void main() { ), }, state: MessageState.sent, + // `Message.createdAt` falls back to `DateTime.now()` per call + // when not provided, which breaks merge/sort keyed on createdAt. + createdAt: DateTime.now(), ); when( () => client.deleteReaction(messageId, type), - ).thenThrow(StreamChatNetworkError(ChatErrorCode.inputError)); + ).thenThrow( + StreamChatNetworkError( + ChatErrorCode.inputError, + data: ErrorResponse()..statusCode = 400, + ), + ); + // A terminal delete rolls back the optimistic removal (restores the + // reaction), symmetric with `sendReaction`: the stream emits the + // removal, then the restore. expectLater( // skipping first seed message list -> [] messages channel.state?.messagesStream.skip(1), @@ -3419,15 +3475,52 @@ void main() { ]), ); - try { - await channel.deleteReaction(message, reaction); - } catch (e) { - expect(e, isA()); - } + await expectLater( + channel.deleteReaction(message, reaction), + throwsA(isA()), + ); verify(() => client.deleteReaction(messageId, type)).called(1); }, ); + + test( + 'a retriable failure keeps the optimistic removal for replay', + () async { + const type = 'like'; + const messageId = 'offline-del'; + final reaction = Reaction( + type: type, + messageId: messageId, + userId: 'test-user-id', + ); + final message = Message( + id: messageId, + ownReactions: [reaction], + latestReactions: [reaction], + reactionGroups: {type: ReactionGroup(count: 1, sumScores: 1)}, + state: MessageState.sent, + ); + + when( + () => client.deleteReaction(messageId, type), + ).thenThrow(StreamChatNetworkError(ChatErrorCode.inputError)); + + await expectLater( + channel.deleteReaction(message, reaction), + throwsA(isA()), + ); + + // The optimistic removal is kept (queued for replay on reconnect) + final current = channel.state!.messages.firstWhere( + (m) => m.id == messageId, + ); + expect( + current.ownReactions?.map((r) => r.type) ?? [], + isNot(contains(type)), + ); + }, + ); }); group('`.deleteReaction in thread`', () { @@ -3488,7 +3581,7 @@ void main() { }); test( - 'should restore prev message state if `client.deleteReaction` throws', + 'restores the reaction if `client.deleteReaction` throws terminally', () async { const userId = 'test-user-id'; const messageId = 'test-message-id'; @@ -3518,8 +3611,16 @@ void main() { when( () => client.deleteReaction(messageId, type), - ).thenThrow(StreamChatNetworkError(ChatErrorCode.inputError)); + ).thenThrow( + StreamChatNetworkError( + ChatErrorCode.inputError, + data: ErrorResponse()..statusCode = 400, + ), + ); + // A terminal delete rolls back the optimistic removal (restores the + // reaction), symmetric with `sendReaction`: the stream emits the + // removal, then the restore. expectLater( // skipping first seed message list -> [] messages channel.state?.threadsStream.skip(1).map((event) => event['test-parent-id']), @@ -3547,11 +3648,10 @@ void main() { ]), ); - try { - await channel.deleteReaction(message, reaction); - } catch (e) { - expect(e, isA()); - } + await expectLater( + channel.deleteReaction(message, reaction), + throwsA(isA()), + ); verify(() => client.deleteReaction(messageId, type)).called(1); }, @@ -4947,6 +5047,9 @@ void main() { when( () => client.channelDeliveryReporter.submitForDelivery(any()), ).thenAnswer((_) async {}); + + // No persistence in this group. + when(() => client.chatPersistenceClient).thenReturn(null); }); group( diff --git a/packages/stream_chat/test/src/client/client_test.dart b/packages/stream_chat/test/src/client/client_test.dart index c341e23c80..17856f410e 100644 --- a/packages/stream_chat/test/src/client/client_test.dart +++ b/packages/stream_chat/test/src/client/client_test.dart @@ -3,6 +3,7 @@ import 'dart:async'; import 'package:mocktail/mocktail.dart'; +import 'package:stream_chat/src/client/reaction_pending_operation.dart'; import 'package:stream_chat/src/core/http/token.dart'; import 'package:stream_chat/stream_chat.dart'; import 'package:test/test.dart'; @@ -1172,10 +1173,19 @@ void main() { ); }); - test('`.disconnectUser` should reset state and user', () async { + test('`.disconnectUser` should reset state, user, and pending operations', () async { + when( + () => api.message.deleteReaction(any(), any()), + ).thenAnswer((_) async => EmptyResponse()); + expect(client.state.currentUser, isNotNull); expect(client.wsConnectionStatus, ConnectionStatus.connected); + // Queue an operation as the current user. + await client.pendingOperationsManager.enqueue( + ReactionPendingOperation.delete(messageId: 'm1', reactionType: 'like'), + ); + expectLater( // skipping initial connected value client.wsConnectionStatusStream.skip(1), @@ -1186,6 +1196,11 @@ void main() { expect(client.state.currentUser, isNull); expect(client.wsConnectionStatus, ConnectionStatus.disconnected); + + // The in-memory queue was cleared, so a replay makes no calls — a user's + // queued operations can never leak into the next session. + await client.pendingOperationsManager.replay(); + verifyNever(() => api.message.deleteReaction(any(), any())); }); }); @@ -5425,6 +5440,90 @@ void main() { }); }); + group('replay pending operations on reconnect', () { + const apiKey = 'test-api-key'; + final user = User(id: 'test-user-id'); + final token = Token.development(user.id).rawValue; + + late FakeChatApi api; + late FakeWebSocket ws; + late MockPersistenceClient persistence; + late StreamChatClient client; + + setUpAll(() { + registerFallbackValue(Reaction(type: 'fallback')); + }); + + setUp(() async { + api = FakeChatApi(); + ws = FakeWebSocket(); + persistence = MockPersistenceClient(); + when(() => persistence.updateLastSyncAt(any())).thenAnswer((_) => Future.value()); + when(persistence.getLastSyncAt).thenAnswer((_) async => null); + // recoverStateOnReconnect off + no loaded channels → the reconnect does + // nothing but replay the queue, keeping these tests isolated to it. + client = StreamChatClient(apiKey, chatApi: api, ws: ws, recoverStateOnReconnect: false) + ..chatPersistenceClient = persistence; + // Initial connect replays an empty queue (a no-op); tests populate the + // queue afterwards and reconnect to exercise replay. + await client.connectUser(user, token); + await delay(300); + }); + + tearDown(() async { + await client.dispose(); + }); + + // Drives the FakeWebSocket through a connected → disconnected → connected + // transition so the client's recovery path (and replay) fires. + Future simulateReconnect() async { + ws.connectionStatus = ConnectionStatus.disconnected; + await delay(100); + ws.connectionStatus = ConnectionStatus.connected; + await delay(300); + } + + PendingOperation addOp(String messageId) => ReactionPendingOperation.add( + Reaction(type: 'like', messageId: messageId), + skipPush: false, + enforceUnique: false, + ); + + void stubSendReactionOk() { + when( + () => api.message.sendReaction( + any(), + any(), + skipPush: any(named: 'skipPush'), + enforceUnique: any(named: 'enforceUnique'), + ), + ).thenAnswer( + (invocation) async => SendReactionResponse() + ..message = Message(id: invocation.positionalArguments[0] as String) + ..reaction = Reaction(type: 'like'), + ); + } + + test('connection recovery replays the pending-operation queue', () async { + // Focused replay semantics live in pending_operations_manager_test.dart; + // this asserts only the wiring — that recovery triggers a replay. + await client.pendingOperationsManager.enqueue(addOp('m1')); + stubSendReactionOk(); + + await simulateReconnect(); + + verify( + () => api.message.sendReaction( + 'm1', + any(), + skipPush: any(named: 'skipPush'), + enforceUnique: any(named: 'enforceUnique'), + ), + ).called(1); + expect(persistence.storedPendingOperations, isEmpty); + }); + }); + group('dispose during reconnect recovery', () { const apiKey = 'test-api-key'; final user = User(id: 'test-user-id'); diff --git a/packages/stream_chat/test/src/client/pending_operations_manager_test.dart b/packages/stream_chat/test/src/client/pending_operations_manager_test.dart new file mode 100644 index 0000000000..ec967706ad --- /dev/null +++ b/packages/stream_chat/test/src/client/pending_operations_manager_test.dart @@ -0,0 +1,431 @@ +// ignore_for_file: avoid_redundant_argument_values + +import 'package:mocktail/mocktail.dart'; +import 'package:stream_chat/src/client/pending_operations_manager.dart'; +import 'package:stream_chat/src/client/reaction_pending_operation.dart'; +import 'package:stream_chat/stream_chat.dart'; +import 'package:test/test.dart'; + +import '../fakes.dart'; +import '../mocks.dart'; + +void main() { + group('PendingOperationsManager', () { + const apiKey = 'test-api-key'; + + late FakeChatApi api; + late FakeWebSocket ws; + late MockPersistenceClient persistence; + late StreamChatClient client; + late PendingOperationsManager manager; + + setUpAll(() { + registerFallbackValue(Reaction(type: 'fallback')); + }); + + setUp(() { + api = FakeChatApi(); + ws = FakeWebSocket(); + persistence = MockPersistenceClient(); + client = StreamChatClient(apiKey, chatApi: api, ws: ws)..chatPersistenceClient = persistence; + manager = PendingOperationsManager(client); + }); + + tearDown(() async { + await client.dispose(); + }); + + PendingOperation addOp(String messageId) => ReactionPendingOperation.add( + Reaction(type: 'like', messageId: messageId), + skipPush: false, + enforceUnique: false, + ); + + PendingOperation deleteOp(String messageId) => ReactionPendingOperation.delete( + messageId: messageId, + reactionType: 'like', + ); + + void stubSendReactionOk() { + when( + () => api.message.sendReaction( + any(), + any(), + skipPush: any(named: 'skipPush'), + enforceUnique: any(named: 'enforceUnique'), + ), + ).thenAnswer( + (invocation) async => SendReactionResponse() + ..message = Message(id: invocation.positionalArguments[0] as String) + ..reaction = Reaction(type: 'like'), + ); + } + + void stubSendReactionThrows(Object error) { + when( + () => api.message.sendReaction( + any(), + any(), + skipPush: any(named: 'skipPush'), + enforceUnique: any(named: 'enforceUnique'), + ), + ).thenThrow(error); + } + + // Retriable/offline: no server response, so `data == null`. + final transientError = StreamChatNetworkError(ChatErrorCode.inputError); + // Terminal: the server responded with an error. + final terminalError = StreamChatNetworkError( + ChatErrorCode.inputError, + data: ErrorResponse()..statusCode = 403, + ); + + group('with persistence enabled', () { + setUp(() async { + await persistence.connect('user-id'); + }); + + test('replays queued operations FIFO and drops them on success', () async { + await manager.enqueue(addOp('m1')); + await manager.enqueue(addOp('m2')); + stubSendReactionOk(); + + await manager.replay(); + + verifyInOrder([ + () => api.message.sendReaction( + 'm1', + any(), + skipPush: any(named: 'skipPush'), + enforceUnique: any(named: 'enforceUnique'), + ), + () => api.message.sendReaction( + 'm2', + any(), + skipPush: any(named: 'skipPush'), + enforceUnique: any(named: 'enforceUnique'), + ), + ]); + // Dropped from the persisted mirror on success. + expect(persistence.storedPendingOperations, isEmpty); + }); + + test('mirrors the enqueued op and its payload to persistence', () async { + await manager.enqueue( + ReactionPendingOperation.add( + Reaction(type: 'like', messageId: 'm1'), + skipPush: true, + enforceUnique: true, + ), + ); + + final op = persistence.storedPendingOperations.single; + expect(op.type, ReactionPendingOperation.addType); + expect(op.targetMessageId, 'm1'); + expect(op.payload[ReactionPendingOperation.skipPushKey], isTrue); + expect(op.payload[ReactionPendingOperation.enforceUniqueKey], isTrue); + }); + + test('keeps a queued operation on a transient (offline) failure', () async { + await manager.enqueue(addOp('m1')); + stubSendReactionThrows(transientError); + + await manager.replay(); + + // Kept for the next recovery — no revert, no drop. + expect(persistence.storedPendingOperations, hasLength(1)); + }); + + test('drops a queued operation on a terminal (server-rejected) failure', () async { + await manager.enqueue(addOp('m1')); + stubSendReactionThrows(terminalError); + + await manager.replay(); + + // Dropped without a revert; the reconnect refresh reconciles. + expect(persistence.storedPendingOperations, isEmpty); + }); + + test('drops an operation whose stored payload cannot be parsed', () async { + await manager.enqueue( + const PendingOperation( + type: ReactionPendingOperation.addType, + targetMessageId: 'm1', + payload: {}, // missing the reaction payload + ), + ); + + await manager.replay(); + + // Dropped without ever hitting the API — it can never be replayed. + expect(persistence.storedPendingOperations, isEmpty); + verifyNever( + () => api.message.sendReaction( + any(), + any(), + skipPush: any(named: 'skipPush'), + enforceUnique: any(named: 'enforceUnique'), + ), + ); + }); + + test('drops an operation with a missing targetMessageId', () async { + await manager.enqueue( + PendingOperation( + type: ReactionPendingOperation.addType, + // No targetMessageId — the reaction can never be routed to a message. + payload: { + ReactionPendingOperation.reactionKey: Reaction(type: 'like').toJson(), + }, + ), + ); + + await manager.replay(); + + // Dropped without ever hitting the API — it can never be replayed. + expect(persistence.storedPendingOperations, isEmpty); + verifyNever( + () => api.message.sendReaction( + any(), + any(), + skipPush: any(named: 'skipPush'), + enforceUnique: any(named: 'enforceUnique'), + ), + ); + }); + + test('drops an operation of an unknown type', () async { + await manager.enqueue( + const PendingOperation( + type: 'unknown.op', + targetMessageId: 'm1', + payload: {}, + ), + ); + + await manager.replay(); + + // Dropped without ever hitting the API — this version can't replay it. + expect(persistence.storedPendingOperations, isEmpty); + verifyNever( + () => api.message.sendReaction( + any(), + any(), + skipPush: any(named: 'skipPush'), + enforceUnique: any(named: 'enforceUnique'), + ), + ); + }); + + test('replays an operation for a channel not loaded this session', () async { + // No channel is loaded for `test:cid`; the queue still replays it. + await manager.enqueue(deleteOp('m9')); + when(() => api.message.deleteReaction('m9', 'like')).thenAnswer((_) async => EmptyResponse()); + + await manager.replay(); + + verify(() => api.message.deleteReaction('m9', 'like')).called(1); + expect(persistence.storedPendingOperations, isEmpty); + }); + + test('a persistence delete failure does not abort the rest of the batch', () async { + await manager.enqueue(addOp('m1')); // id 1 + await manager.enqueue(addOp('m2')); // id 2 + + // m1 replays terminally (server-rejected) → dropped; m2 replays OK. + when( + () => api.message.sendReaction( + 'm1', + any(), + skipPush: any(named: 'skipPush'), + enforceUnique: any(named: 'enforceUnique'), + ), + ).thenThrow(terminalError); + when( + () => api.message.sendReaction( + 'm2', + any(), + skipPush: any(named: 'skipPush'), + enforceUnique: any(named: 'enforceUnique'), + ), + ).thenAnswer( + (_) async => SendReactionResponse() + ..message = Message(id: 'm2') + ..reaction = Reaction(type: 'like'), + ); + + // Dropping the first operation (id 1) throws mid-loop. + persistence.failDeleteForIds.add(1); + + await manager.replay(); + + // The loop continued past the failed delete: m2 still replayed. + verify( + () => api.message.sendReaction( + 'm2', + any(), + skipPush: any(named: 'skipPush'), + enforceUnique: any(named: 'enforceUnique'), + ), + ).called(1); + // m1's mirror row survived (its delete failed); m2 was dropped on success. + expect(persistence.storedPendingOperations, hasLength(1)); + expect(persistence.storedPendingOperations.single.targetMessageId, 'm1'); + }); + + test('keeps the op in memory when the durable insert fails', () async { + persistence.failInsert = true; + + await manager.enqueue(addOp('m1')); + + // Nothing was mirrored, but the op still replays this session. + expect(persistence.storedPendingOperations, isEmpty); + stubSendReactionOk(); + await manager.replay(); + verify( + () => api.message.sendReaction( + 'm1', + any(), + skipPush: any(named: 'skipPush'), + enforceUnique: any(named: 'enforceUnique'), + ), + ).called(1); + }); + + group('hydrate', () { + test('loads persisted operations into memory and replays them', () async { + // Operations that survived process death, seeded straight into the DB. + await persistence.insertPendingOperation(addOp('m1')); + await persistence.insertPendingOperation(addOp('m2')); + stubSendReactionOk(); + + // A fresh manager (as after a restart) starts with an empty queue. + await manager.hydrate(); + await manager.replay(); + + verifyInOrder([ + () => api.message.sendReaction( + 'm1', + any(), + skipPush: any(named: 'skipPush'), + enforceUnique: any(named: 'enforceUnique'), + ), + () => api.message.sendReaction( + 'm2', + any(), + skipPush: any(named: 'skipPush'), + enforceUnique: any(named: 'enforceUnique'), + ), + ]); + expect(persistence.storedPendingOperations, isEmpty); + }); + }); + }); + + group('with persistence disabled', () { + // `persistence` is set on the client but never connected, so + // `client.persistenceEnabled` is false and the queue is memory-only. + + test('replays memory-only operations without touching persistence', () async { + await manager.enqueue(addOp('m1')); + await manager.enqueue(addOp('m2')); + stubSendReactionOk(); + + await manager.replay(); + + verifyInOrder([ + () => api.message.sendReaction( + 'm1', + any(), + skipPush: any(named: 'skipPush'), + enforceUnique: any(named: 'enforceUnique'), + ), + () => api.message.sendReaction( + 'm2', + any(), + skipPush: any(named: 'skipPush'), + enforceUnique: any(named: 'enforceUnique'), + ), + ]); + // Never mirrored to persistence. + expect(persistence.storedPendingOperations, isEmpty); + }); + + test('keeps a memory-only op on a transient failure for the next replay', () async { + await manager.enqueue(addOp('m1')); + stubSendReactionThrows(transientError); + await manager.replay(); + + // Still queued in memory — a later replay retries it (the first, + // failing attempt plus this successful one make two calls in total). + stubSendReactionOk(); + await manager.replay(); + + verify( + () => api.message.sendReaction( + 'm1', + any(), + skipPush: any(named: 'skipPush'), + enforceUnique: any(named: 'enforceUnique'), + ), + ).called(2); + }); + + test('drops a memory-only op on a terminal failure', () async { + await manager.enqueue(addOp('m1')); + stubSendReactionThrows(terminalError); + await manager.replay(); + + // Dropped from memory — a later replay does nothing. + stubSendReactionOk(); + await manager.replay(); + + verify( + () => api.message.sendReaction( + 'm1', + any(), + skipPush: any(named: 'skipPush'), + enforceUnique: any(named: 'enforceUnique'), + ), + ).called(1); + }); + + test('hydrate is a no-op', () async { + // A row exists in the DB, but with persistence disabled it is ignored. + await persistence.insertPendingOperation(addOp('m1')); + stubSendReactionOk(); + + await manager.hydrate(); + await manager.replay(); + + verifyNever( + () => api.message.sendReaction( + any(), + any(), + skipPush: any(named: 'skipPush'), + enforceUnique: any(named: 'enforceUnique'), + ), + ); + }); + }); + + group('clear', () { + test('empties the queue so a later replay does nothing', () async { + await manager.enqueue(addOp('m1')); + manager.clear(); + stubSendReactionOk(); + + await manager.replay(); + + verifyNever( + () => api.message.sendReaction( + any(), + any(), + skipPush: any(named: 'skipPush'), + enforceUnique: any(named: 'enforceUnique'), + ), + ); + }); + }); + }); +} diff --git a/packages/stream_chat/test/src/core/models/pending_operation_test.dart b/packages/stream_chat/test/src/core/models/pending_operation_test.dart new file mode 100644 index 0000000000..4969a1cba3 --- /dev/null +++ b/packages/stream_chat/test/src/core/models/pending_operation_test.dart @@ -0,0 +1,31 @@ +import 'package:stream_chat/src/core/models/pending_operation.dart'; +import 'package:test/test.dart'; + +void main() { + group('src/models/pending_operation', () { + PendingOperation build({ + int? id, + String type = 'reaction.add', + Map? payload, + }) => PendingOperation( + id: id, + type: type, + targetMessageId: 'm1', + payload: payload ?? const {'reaction': 'like'}, + ); + + test('equality ignores the autoincrement id', () { + // A stored operation (with an id) equals its pre-store form so equality + // compares by content, not database identity. + expect(build(id: 42), build()); + }); + + test('operations with different types are not equal', () { + expect(build(type: 'reaction.add'), isNot(build(type: 'reaction.delete'))); + }); + + test('operations with different payloads are not equal', () { + expect(build(payload: const {'a': 1}), isNot(build(payload: const {'a': 2}))); + }); + }); +} 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 555d13ecd0..bfca697cd8 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 @@ -9,6 +9,7 @@ import 'package:stream_chat/src/core/models/filter.dart'; import 'package:stream_chat/src/core/models/location.dart'; import 'package:stream_chat/src/core/models/member.dart'; import 'package:stream_chat/src/core/models/message.dart'; +import 'package:stream_chat/src/core/models/pending_operation.dart'; import 'package:stream_chat/src/core/models/poll.dart'; import 'package:stream_chat/src/core/models/poll_vote.dart'; import 'package:stream_chat/src/core/models/reaction.dart'; @@ -292,5 +293,23 @@ void main() { ); persistenceClient.updateChannelStates([channelState]); }); + + test('insertPendingOperation defaults to returning null', () async { + const operation = PendingOperation( + type: 'reaction.add', + targetMessageId: 'message-id', + payload: {'reaction': 'like'}, + ); + expect(await persistenceClient.insertPendingOperation(operation), isNull); + }); + + test('getPendingOperations defaults to an empty list', () async { + final operations = await persistenceClient.getPendingOperations(); + expect(operations, isEmpty); + }); + + test('deletePendingOperation defaults to a no-op', () async { + await persistenceClient.deletePendingOperation(1); + }); }); } diff --git a/packages/stream_chat/test/src/mocks.dart b/packages/stream_chat/test/src/mocks.dart index 848f26a48d..63944c08d7 100644 --- a/packages/stream_chat/test/src/mocks.dart +++ b/packages/stream_chat/test/src/mocks.dart @@ -4,6 +4,7 @@ import 'package:mocktail/mocktail.dart'; import 'package:stream_chat/src/client/channel.dart'; import 'package:stream_chat/src/client/channel_delivery_reporter.dart'; import 'package:stream_chat/src/client/client.dart'; +import 'package:stream_chat/src/client/pending_operations_manager.dart'; import 'package:stream_chat/src/core/api/attachment_file_uploader.dart'; import 'package:stream_chat/src/core/api/channel_api.dart'; import 'package:stream_chat/src/core/api/device_api.dart'; @@ -20,6 +21,7 @@ import 'package:stream_chat/src/core/http/stream_http_client.dart'; import 'package:stream_chat/src/core/http/token_manager.dart'; import 'package:stream_chat/src/core/models/channel_config.dart'; import 'package:stream_chat/src/core/models/event.dart'; +import 'package:stream_chat/src/core/models/pending_operation.dart'; import 'package:stream_chat/src/core/util/event_controller.dart'; import 'package:stream_chat/src/db/chat_persistence_client.dart'; import 'package:stream_chat/src/event_type.dart'; @@ -96,11 +98,66 @@ class MockPersistenceClient extends Mock implements ChatPersistenceClient { _userId = null; _isConnected = false; } + + /// In-memory pending-operation queue. Real overrides (not mocktail stubs) + /// so they behave like a working queue and don't register as interactions + /// for `verifyNoMoreInteractions`. Populate it with [insertPendingOperation]. + final List storedPendingOperations = []; + int _nextPendingOperationId = 1; + + /// Ids for which [deletePendingOperation] throws, to exercise the replay + /// loop's per-operation error isolation. + final Set failDeleteForIds = {}; + + /// When `true`, [insertPendingOperation] throws, to exercise the enqueue + /// failure fallback in `Channel.sendReaction`/`deleteReaction`. + bool failInsert = false; + + @override + Future insertPendingOperation(PendingOperation operation) async { + if (failInsert) { + throw Exception('simulated insert failure'); + } + // Assign an autoincrement id like the real DAO so replay can delete by id. + final id = _nextPendingOperationId++; + storedPendingOperations.add( + PendingOperation( + id: id, + type: operation.type, + targetMessageId: operation.targetMessageId, + payload: operation.payload, + ), + ); + return id; + } + + @override + Future> getPendingOperations() async { + // Return a snapshot (like the DAO) so a delete during replay iteration + // cannot concurrently modify the list being iterated. + return [...storedPendingOperations]; + } + + @override + Future deletePendingOperation(int id) async { + if (failDeleteForIds.contains(id)) { + throw Exception('simulated delete failure for operation $id'); + } + storedPendingOperations.removeWhere((it) => it.id == id); + } } class MockStreamChatClient extends Mock implements StreamChatClient { @override - bool get persistenceEnabled => false; + bool get persistenceEnabled => chatPersistenceClient != null; + + // A real manager backed by this mock, so `Channel.sendReaction` / + // `deleteReaction` can enqueue through it. It reads `persistenceEnabled`, + // `chatPersistenceClient` and `logger` off this mock. + late final PendingOperationsManager _pendingOperationsManager = PendingOperationsManager(this); + + @override + PendingOperationsManager get pendingOperationsManager => _pendingOperationsManager; // A plain settable field (not a `when(...)` stub) so tests can flip it // with a direct assignment, e.g. `client.isLocalUnreadCountEnabled = true`. diff --git a/packages/stream_chat_persistence/CHANGELOG.md b/packages/stream_chat_persistence/CHANGELOG.md index 25cdd2a6c2..586d4f59ce 100644 --- a/packages/stream_chat_persistence/CHANGELOG.md +++ b/packages/stream_chat_persistence/CHANGELOG.md @@ -1,3 +1,9 @@ +## Upcoming + +✅ Added + +- Added support for sending and deleting reactions while offline. + ## 10.2.0 🚀 Performance diff --git a/packages/stream_chat_persistence/lib/src/dao/dao.dart b/packages/stream_chat_persistence/lib/src/dao/dao.dart index d2f088c399..9ea362e601 100644 --- a/packages/stream_chat_persistence/lib/src/dao/dao.dart +++ b/packages/stream_chat_persistence/lib/src/dao/dao.dart @@ -5,6 +5,7 @@ export 'draft_message_dao.dart'; export 'location_dao.dart'; export 'member_dao.dart'; export 'message_dao.dart'; +export 'pending_operation_dao.dart'; export 'pinned_message_dao.dart'; export 'pinned_message_reaction_dao.dart'; export 'poll_dao.dart'; diff --git a/packages/stream_chat_persistence/lib/src/dao/pending_operation_dao.dart b/packages/stream_chat_persistence/lib/src/dao/pending_operation_dao.dart new file mode 100644 index 0000000000..022ed3cfae --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/dao/pending_operation_dao.dart @@ -0,0 +1,27 @@ +import 'package:drift/drift.dart'; +import 'package:stream_chat/stream_chat.dart'; +import 'package:stream_chat_persistence/src/db/drift_chat_database.dart'; +import 'package:stream_chat_persistence/src/entity/entity.dart'; +import 'package:stream_chat_persistence/src/mapper/mapper.dart'; + +part 'pending_operation_dao.g.dart'; + +/// The Data Access Object for operations in [PendingOperations] table. +@DriftAccessor(tables: [PendingOperations]) +class PendingOperationDao extends DatabaseAccessor with _$PendingOperationDaoMixin { + /// Creates a new pending operation dao instance + PendingOperationDao(super.db); + + /// Appends [operation] to the queue, returning its autoincrement `id`. + Future insertPendingOperation(PendingOperation operation) => + into(pendingOperations).insert(operation.toCompanion()); + + /// Returns all pending operations ordered by insertion (`id` ascending). + Future> getPendingOperations() { + final query = select(pendingOperations)..orderBy([(tbl) => OrderingTerm.asc(tbl.id)]); + return query.map((row) => row.toPendingOperation()).get(); + } + + /// Deletes the pending operation with the given autoincrement [id]. + Future deletePendingOperation(int id) => (delete(pendingOperations)..where((tbl) => tbl.id.equals(id))).go(); +} diff --git a/packages/stream_chat_persistence/lib/src/dao/pending_operation_dao.g.dart b/packages/stream_chat_persistence/lib/src/dao/pending_operation_dao.g.dart new file mode 100644 index 0000000000..22b93a62ae --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/dao/pending_operation_dao.g.dart @@ -0,0 +1,18 @@ +// GENERATED CODE - DO NOT MODIFY BY HAND + +part of 'pending_operation_dao.dart'; + +// ignore_for_file: type=lint +mixin _$PendingOperationDaoMixin on DatabaseAccessor { + $PendingOperationsTable get pendingOperations => attachedDatabase.pendingOperations; + PendingOperationDaoManager get managers => PendingOperationDaoManager(this); +} + +class PendingOperationDaoManager { + final _$PendingOperationDaoMixin _db; + PendingOperationDaoManager(this._db); + $$PendingOperationsTableTableManager get pendingOperations => $$PendingOperationsTableTableManager( + _db.attachedDatabase, + _db.pendingOperations, + ); +} diff --git a/packages/stream_chat_persistence/lib/src/db/drift_chat_database.dart b/packages/stream_chat_persistence/lib/src/db/drift_chat_database.dart index ff9e83d1b0..996a5c5ad6 100644 --- a/packages/stream_chat_persistence/lib/src/db/drift_chat_database.dart +++ b/packages/stream_chat_persistence/lib/src/db/drift_chat_database.dart @@ -15,6 +15,7 @@ part 'drift_chat_database.g.dart'; DraftMessages, Locations, Messages, + PendingOperations, PinnedMessages, Polls, PollVotes, @@ -33,6 +34,7 @@ part 'drift_chat_database.g.dart'; MessageDao, DraftMessageDao, LocationDao, + PendingOperationDao, PinnedMessageDao, PinnedMessageReactionDao, MemberDao, @@ -58,7 +60,7 @@ class DriftChatDatabase extends _$DriftChatDatabase { // you should bump this number whenever you change or add a table definition. @override - int get schemaVersion => 1000 + 35; + int get schemaVersion => 1000 + 36; // Store DateTime as ISO-8601 text to preserve sub-second precision. @override diff --git a/packages/stream_chat_persistence/lib/src/db/drift_chat_database.g.dart b/packages/stream_chat_persistence/lib/src/db/drift_chat_database.g.dart index 079b6c345a..822f108798 100644 --- a/packages/stream_chat_persistence/lib/src/db/drift_chat_database.g.dart +++ b/packages/stream_chat_persistence/lib/src/db/drift_chat_database.g.dart @@ -4272,6 +4272,306 @@ class LocationsCompanion extends UpdateCompanion { } } +class $PendingOperationsTable extends PendingOperations + with TableInfo<$PendingOperationsTable, PendingOperationEntity> { + @override + final GeneratedDatabase attachedDatabase; + final String? _alias; + $PendingOperationsTable(this.attachedDatabase, [this._alias]); + static const VerificationMeta _idMeta = const VerificationMeta('id'); + @override + late final GeneratedColumn id = GeneratedColumn( + 'id', + aliasedName, + false, + hasAutoIncrement: true, + type: DriftSqlType.int, + requiredDuringInsert: false, + defaultConstraints: GeneratedColumn.constraintIsAlways( + 'PRIMARY KEY AUTOINCREMENT', + ), + ); + static const VerificationMeta _typeMeta = const VerificationMeta('type'); + @override + late final GeneratedColumn type = GeneratedColumn( + 'type', + aliasedName, + false, + type: DriftSqlType.string, + requiredDuringInsert: true, + ); + static const VerificationMeta _targetMessageIdMeta = const VerificationMeta( + 'targetMessageId', + ); + @override + late final GeneratedColumn targetMessageId = GeneratedColumn( + 'target_message_id', + aliasedName, + true, + type: DriftSqlType.string, + requiredDuringInsert: false, + ); + @override + late final GeneratedColumnWithTypeConverter, String> payload = + GeneratedColumn( + 'payload', + aliasedName, + false, + type: DriftSqlType.string, + requiredDuringInsert: true, + ).withConverter>( + $PendingOperationsTable.$converterpayload, + ); + @override + List get $columns => [id, type, targetMessageId, payload]; + @override + String get aliasedName => _alias ?? actualTableName; + @override + String get actualTableName => $name; + static const String $name = 'pending_operations'; + @override + VerificationContext validateIntegrity( + Insertable instance, { + bool isInserting = false, + }) { + final context = VerificationContext(); + final data = instance.toColumns(true); + if (data.containsKey('id')) { + context.handle(_idMeta, id.isAcceptableOrUnknown(data['id']!, _idMeta)); + } + if (data.containsKey('type')) { + context.handle( + _typeMeta, + type.isAcceptableOrUnknown(data['type']!, _typeMeta), + ); + } else if (isInserting) { + context.missing(_typeMeta); + } + if (data.containsKey('target_message_id')) { + context.handle( + _targetMessageIdMeta, + targetMessageId.isAcceptableOrUnknown( + data['target_message_id']!, + _targetMessageIdMeta, + ), + ); + } + return context; + } + + @override + Set get $primaryKey => {id}; + @override + PendingOperationEntity map(Map data, {String? tablePrefix}) { + final effectivePrefix = tablePrefix != null ? '$tablePrefix.' : ''; + return PendingOperationEntity( + id: attachedDatabase.typeMapping.read( + DriftSqlType.int, + data['${effectivePrefix}id'], + )!, + type: attachedDatabase.typeMapping.read( + DriftSqlType.string, + data['${effectivePrefix}type'], + )!, + targetMessageId: attachedDatabase.typeMapping.read( + DriftSqlType.string, + data['${effectivePrefix}target_message_id'], + ), + payload: $PendingOperationsTable.$converterpayload.fromSql( + attachedDatabase.typeMapping.read( + DriftSqlType.string, + data['${effectivePrefix}payload'], + )!, + ), + ); + } + + @override + $PendingOperationsTable createAlias(String alias) { + return $PendingOperationsTable(attachedDatabase, alias); + } + + static TypeConverter, String> $converterpayload = MapConverter(); +} + +class PendingOperationEntity extends DataClass implements Insertable { + /// Autoincrement id. + final int id; + + /// The operation-type discriminator (e.g. `reaction.add`). + final String type; + + /// The id of the message the operation targets, if any. + final String? targetMessageId; + + /// The operation-specific value fields, stored as JSON. + final Map payload; + const PendingOperationEntity({ + required this.id, + required this.type, + this.targetMessageId, + required this.payload, + }); + @override + Map toColumns(bool nullToAbsent) { + final map = {}; + map['id'] = Variable(id); + map['type'] = Variable(type); + if (!nullToAbsent || targetMessageId != null) { + map['target_message_id'] = Variable(targetMessageId); + } + { + map['payload'] = Variable( + $PendingOperationsTable.$converterpayload.toSql(payload), + ); + } + return map; + } + + factory PendingOperationEntity.fromJson( + Map json, { + ValueSerializer? serializer, + }) { + serializer ??= driftRuntimeOptions.defaultSerializer; + return PendingOperationEntity( + id: serializer.fromJson(json['id']), + type: serializer.fromJson(json['type']), + targetMessageId: serializer.fromJson(json['targetMessageId']), + payload: serializer.fromJson>(json['payload']), + ); + } + @override + Map toJson({ValueSerializer? serializer}) { + serializer ??= driftRuntimeOptions.defaultSerializer; + return { + 'id': serializer.toJson(id), + 'type': serializer.toJson(type), + 'targetMessageId': serializer.toJson(targetMessageId), + 'payload': serializer.toJson>(payload), + }; + } + + PendingOperationEntity copyWith({ + int? id, + String? type, + Value targetMessageId = const Value.absent(), + Map? payload, + }) => PendingOperationEntity( + id: id ?? this.id, + type: type ?? this.type, + targetMessageId: targetMessageId.present ? targetMessageId.value : this.targetMessageId, + payload: payload ?? this.payload, + ); + PendingOperationEntity copyWithCompanion(PendingOperationsCompanion data) { + return PendingOperationEntity( + id: data.id.present ? data.id.value : this.id, + type: data.type.present ? data.type.value : this.type, + targetMessageId: data.targetMessageId.present ? data.targetMessageId.value : this.targetMessageId, + payload: data.payload.present ? data.payload.value : this.payload, + ); + } + + @override + String toString() { + return (StringBuffer('PendingOperationEntity(') + ..write('id: $id, ') + ..write('type: $type, ') + ..write('targetMessageId: $targetMessageId, ') + ..write('payload: $payload') + ..write(')')) + .toString(); + } + + @override + int get hashCode => Object.hash(id, type, targetMessageId, payload); + @override + bool operator ==(Object other) => + identical(this, other) || + (other is PendingOperationEntity && + other.id == this.id && + other.type == this.type && + other.targetMessageId == this.targetMessageId && + other.payload == this.payload); +} + +class PendingOperationsCompanion extends UpdateCompanion { + final Value id; + final Value type; + final Value targetMessageId; + final Value> payload; + const PendingOperationsCompanion({ + this.id = const Value.absent(), + this.type = const Value.absent(), + this.targetMessageId = const Value.absent(), + this.payload = const Value.absent(), + }); + PendingOperationsCompanion.insert({ + this.id = const Value.absent(), + required String type, + this.targetMessageId = const Value.absent(), + required Map payload, + }) : type = Value(type), + payload = Value(payload); + static Insertable custom({ + Expression? id, + Expression? type, + Expression? targetMessageId, + Expression? payload, + }) { + return RawValuesInsertable({ + if (id != null) 'id': id, + if (type != null) 'type': type, + if (targetMessageId != null) 'target_message_id': targetMessageId, + if (payload != null) 'payload': payload, + }); + } + + PendingOperationsCompanion copyWith({ + Value? id, + Value? type, + Value? targetMessageId, + Value>? payload, + }) { + return PendingOperationsCompanion( + id: id ?? this.id, + type: type ?? this.type, + targetMessageId: targetMessageId ?? this.targetMessageId, + payload: payload ?? this.payload, + ); + } + + @override + Map toColumns(bool nullToAbsent) { + final map = {}; + if (id.present) { + map['id'] = Variable(id.value); + } + if (type.present) { + map['type'] = Variable(type.value); + } + if (targetMessageId.present) { + map['target_message_id'] = Variable(targetMessageId.value); + } + if (payload.present) { + map['payload'] = Variable( + $PendingOperationsTable.$converterpayload.toSql(payload.value), + ); + } + return map; + } + + @override + String toString() { + return (StringBuffer('PendingOperationsCompanion(') + ..write('id: $id, ') + ..write('type: $type, ') + ..write('targetMessageId: $targetMessageId, ') + ..write('payload: $payload') + ..write(')')) + .toString(); + } +} + class $PinnedMessagesTable extends PinnedMessages with TableInfo<$PinnedMessagesTable, PinnedMessageEntity> { @override final GeneratedDatabase attachedDatabase; @@ -11724,6 +12024,7 @@ abstract class _$DriftChatDatabase extends GeneratedDatabase { late final $MessagesTable messages = $MessagesTable(this); late final $DraftMessagesTable draftMessages = $DraftMessagesTable(this); late final $LocationsTable locations = $LocationsTable(this); + late final $PendingOperationsTable pendingOperations = $PendingOperationsTable(this); late final $PinnedMessagesTable pinnedMessages = $PinnedMessagesTable(this); late final $PollsTable polls = $PollsTable(this); late final $PollVotesTable pollVotes = $PollVotesTable(this); @@ -11760,6 +12061,9 @@ abstract class _$DriftChatDatabase extends GeneratedDatabase { this as DriftChatDatabase, ); late final LocationDao locationDao = LocationDao(this as DriftChatDatabase); + late final PendingOperationDao pendingOperationDao = PendingOperationDao( + this as DriftChatDatabase, + ); late final PinnedMessageDao pinnedMessageDao = PinnedMessageDao( this as DriftChatDatabase, ); @@ -11783,6 +12087,7 @@ abstract class _$DriftChatDatabase extends GeneratedDatabase { messages, draftMessages, locations, + pendingOperations, pinnedMessages, polls, pollVotes, @@ -14972,6 +15277,178 @@ typedef $$LocationsTableProcessedTableManager = LocationEntity, PrefetchHooks Function({bool channelCid, bool messageId}) >; +typedef $$PendingOperationsTableCreateCompanionBuilder = + PendingOperationsCompanion Function({ + Value id, + required String type, + Value targetMessageId, + required Map payload, + }); +typedef $$PendingOperationsTableUpdateCompanionBuilder = + PendingOperationsCompanion Function({ + Value id, + Value type, + Value targetMessageId, + Value> payload, + }); + +class $$PendingOperationsTableFilterComposer extends Composer<_$DriftChatDatabase, $PendingOperationsTable> { + $$PendingOperationsTableFilterComposer({ + required super.$db, + required super.$table, + super.joinBuilder, + super.$addJoinBuilderToRootComposer, + super.$removeJoinBuilderFromRootComposer, + }); + ColumnFilters get id => $composableBuilder( + column: $table.id, + builder: (column) => ColumnFilters(column), + ); + + ColumnFilters get type => $composableBuilder( + column: $table.type, + builder: (column) => ColumnFilters(column), + ); + + ColumnFilters get targetMessageId => $composableBuilder( + column: $table.targetMessageId, + builder: (column) => ColumnFilters(column), + ); + + ColumnWithTypeConverterFilters, Map, String> get payload => $composableBuilder( + column: $table.payload, + builder: (column) => ColumnWithTypeConverterFilters(column), + ); +} + +class $$PendingOperationsTableOrderingComposer extends Composer<_$DriftChatDatabase, $PendingOperationsTable> { + $$PendingOperationsTableOrderingComposer({ + required super.$db, + required super.$table, + super.joinBuilder, + super.$addJoinBuilderToRootComposer, + super.$removeJoinBuilderFromRootComposer, + }); + ColumnOrderings get id => $composableBuilder( + column: $table.id, + builder: (column) => ColumnOrderings(column), + ); + + ColumnOrderings get type => $composableBuilder( + column: $table.type, + builder: (column) => ColumnOrderings(column), + ); + + ColumnOrderings get targetMessageId => $composableBuilder( + column: $table.targetMessageId, + builder: (column) => ColumnOrderings(column), + ); + + ColumnOrderings get payload => $composableBuilder( + column: $table.payload, + builder: (column) => ColumnOrderings(column), + ); +} + +class $$PendingOperationsTableAnnotationComposer extends Composer<_$DriftChatDatabase, $PendingOperationsTable> { + $$PendingOperationsTableAnnotationComposer({ + required super.$db, + required super.$table, + super.joinBuilder, + super.$addJoinBuilderToRootComposer, + super.$removeJoinBuilderFromRootComposer, + }); + GeneratedColumn get id => $composableBuilder(column: $table.id, builder: (column) => column); + + GeneratedColumn get type => $composableBuilder(column: $table.type, builder: (column) => column); + + GeneratedColumn get targetMessageId => $composableBuilder( + column: $table.targetMessageId, + builder: (column) => column, + ); + + GeneratedColumnWithTypeConverter, String> get payload => + $composableBuilder(column: $table.payload, builder: (column) => column); +} + +class $$PendingOperationsTableTableManager + extends + RootTableManager< + _$DriftChatDatabase, + $PendingOperationsTable, + PendingOperationEntity, + $$PendingOperationsTableFilterComposer, + $$PendingOperationsTableOrderingComposer, + $$PendingOperationsTableAnnotationComposer, + $$PendingOperationsTableCreateCompanionBuilder, + $$PendingOperationsTableUpdateCompanionBuilder, + ( + PendingOperationEntity, + BaseReferences<_$DriftChatDatabase, $PendingOperationsTable, PendingOperationEntity>, + ), + PendingOperationEntity, + PrefetchHooks Function() + > { + $$PendingOperationsTableTableManager( + _$DriftChatDatabase db, + $PendingOperationsTable table, + ) : super( + TableManagerState( + db: db, + table: table, + createFilteringComposer: () => $$PendingOperationsTableFilterComposer($db: db, $table: table), + createOrderingComposer: () => $$PendingOperationsTableOrderingComposer($db: db, $table: table), + createComputedFieldComposer: () => $$PendingOperationsTableAnnotationComposer( + $db: db, + $table: table, + ), + updateCompanionCallback: + ({ + Value id = const Value.absent(), + Value type = const Value.absent(), + Value targetMessageId = const Value.absent(), + Value> payload = const Value.absent(), + }) => PendingOperationsCompanion( + id: id, + type: type, + targetMessageId: targetMessageId, + payload: payload, + ), + createCompanionCallback: + ({ + Value id = const Value.absent(), + required String type, + Value targetMessageId = const Value.absent(), + required Map payload, + }) => PendingOperationsCompanion.insert( + id: id, + type: type, + targetMessageId: targetMessageId, + payload: payload, + ), + withReferenceMapper: (p0) => p0.map((e) => (e.readTable(table), BaseReferences(db, table, e))).toList(), + prefetchHooksCallback: null, + ), + ); +} + +typedef $$PendingOperationsTableProcessedTableManager = + ProcessedTableManager< + _$DriftChatDatabase, + $PendingOperationsTable, + PendingOperationEntity, + $$PendingOperationsTableFilterComposer, + $$PendingOperationsTableOrderingComposer, + $$PendingOperationsTableAnnotationComposer, + $$PendingOperationsTableCreateCompanionBuilder, + $$PendingOperationsTableUpdateCompanionBuilder, + ( + PendingOperationEntity, + BaseReferences<_$DriftChatDatabase, $PendingOperationsTable, PendingOperationEntity>, + ), + PendingOperationEntity, + PrefetchHooks Function() + >; typedef $$PinnedMessagesTableCreateCompanionBuilder = PinnedMessagesCompanion Function({ required String id, @@ -19218,6 +19695,8 @@ class $DriftChatDatabaseManager { $$MessagesTableTableManager get messages => $$MessagesTableTableManager(_db, _db.messages); $$DraftMessagesTableTableManager get draftMessages => $$DraftMessagesTableTableManager(_db, _db.draftMessages); $$LocationsTableTableManager get locations => $$LocationsTableTableManager(_db, _db.locations); + $$PendingOperationsTableTableManager get pendingOperations => + $$PendingOperationsTableTableManager(_db, _db.pendingOperations); $$PinnedMessagesTableTableManager get pinnedMessages => $$PinnedMessagesTableTableManager(_db, _db.pinnedMessages); $$PollsTableTableManager get polls => $$PollsTableTableManager(_db, _db.polls); $$PollVotesTableTableManager get pollVotes => $$PollVotesTableTableManager(_db, _db.pollVotes); diff --git a/packages/stream_chat_persistence/lib/src/entity/entity.dart b/packages/stream_chat_persistence/lib/src/entity/entity.dart index 8aecaf5af4..41a68d49d9 100644 --- a/packages/stream_chat_persistence/lib/src/entity/entity.dart +++ b/packages/stream_chat_persistence/lib/src/entity/entity.dart @@ -6,6 +6,7 @@ export 'draft_messages.dart'; export 'locations.dart'; export 'members.dart'; export 'messages.dart'; +export 'pending_operations.dart'; export 'pinned_message_reactions.dart'; export 'pinned_messages.dart'; export 'poll_votes.dart'; diff --git a/packages/stream_chat_persistence/lib/src/entity/pending_operations.dart b/packages/stream_chat_persistence/lib/src/entity/pending_operations.dart new file mode 100644 index 0000000000..47a39d92d1 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/entity/pending_operations.dart @@ -0,0 +1,22 @@ +// coverage:ignore-file +import 'package:drift/drift.dart'; +import 'package:stream_chat_persistence/src/converter/map_converter.dart'; + +/// Represents a [PendingOperations] table in [DriftChatDatabase]. +/// +/// Stores optimistic mutations (e.g. reactions added or removed while offline) +/// awaiting replay. +@DataClassName('PendingOperationEntity') +class PendingOperations extends Table { + /// Autoincrement id. + IntColumn get id => integer().autoIncrement()(); + + /// The operation-type discriminator (e.g. `reaction.add`). + TextColumn get type => text()(); + + /// The id of the message the operation targets, if any. + TextColumn get targetMessageId => text().nullable()(); + + /// The operation-specific value fields, stored as JSON. + TextColumn get payload => text().map(MapConverter())(); +} diff --git a/packages/stream_chat_persistence/lib/src/mapper/mapper.dart b/packages/stream_chat_persistence/lib/src/mapper/mapper.dart index 45b35dba81..05f28da331 100644 --- a/packages/stream_chat_persistence/lib/src/mapper/mapper.dart +++ b/packages/stream_chat_persistence/lib/src/mapper/mapper.dart @@ -4,6 +4,7 @@ export 'event_mapper.dart'; export 'location_mapper.dart'; export 'member_mapper.dart'; export 'message_mapper.dart'; +export 'pending_operation_mapper.dart'; export 'pinned_message_mapper.dart'; export 'pinned_message_reaction_mapper.dart'; export 'poll_mapper.dart'; diff --git a/packages/stream_chat_persistence/lib/src/mapper/pending_operation_mapper.dart b/packages/stream_chat_persistence/lib/src/mapper/pending_operation_mapper.dart new file mode 100644 index 0000000000..d2ad312104 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/mapper/pending_operation_mapper.dart @@ -0,0 +1,25 @@ +import 'package:drift/drift.dart'; +import 'package:stream_chat/stream_chat.dart'; +import 'package:stream_chat_persistence/src/db/drift_chat_database.dart'; + +/// Useful mapping functions for [PendingOperationEntity] +extension PendingOperationEntityX on PendingOperationEntity { + /// Maps a [PendingOperationEntity] into a [PendingOperation] + PendingOperation toPendingOperation() => PendingOperation( + id: id, + type: type, + targetMessageId: targetMessageId, + payload: payload, + ); +} + +/// Useful mapping functions for [PendingOperation] +extension PendingOperationX on PendingOperation { + /// Maps a [PendingOperation] into a [PendingOperationsCompanion]. + PendingOperationsCompanion toCompanion() => PendingOperationsCompanion( + id: id == null ? const Value.absent() : Value(id!), + type: Value(type), + targetMessageId: Value(targetMessageId), + payload: Value(payload), + ); +} 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 0a14996df3..6db0a9baad 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 @@ -653,6 +653,27 @@ class StreamChatPersistenceClient extends ChatPersistenceClient { } } + @override + Future insertPendingOperation(PendingOperation operation) { + assert(_debugIsConnected, ''); + _logger.info('insertPendingOperation'); + return db!.pendingOperationDao.insertPendingOperation(operation); + } + + @override + Future> getPendingOperations() { + assert(_debugIsConnected, ''); + _logger.info('getPendingOperations'); + return db!.pendingOperationDao.getPendingOperations(); + } + + @override + Future deletePendingOperation(int id) { + assert(_debugIsConnected, ''); + _logger.info('deletePendingOperation'); + return db!.pendingOperationDao.deletePendingOperation(id); + } + bool _sortRequiresMembership(SortOrder? sort) => sort?.any((opt) => opt.field == ChannelSortKey.pinnedAt) ?? false; diff --git a/packages/stream_chat_persistence/test/mock_chat_database.dart b/packages/stream_chat_persistence/test/mock_chat_database.dart index fe4facfd75..a818ebf997 100644 --- a/packages/stream_chat_persistence/test/mock_chat_database.dart +++ b/packages/stream_chat_persistence/test/mock_chat_database.dart @@ -62,6 +62,10 @@ class MockChatDatabase extends Mock implements DriftChatDatabase { LocationDao get locationDao => _locationDao ??= MockLocationDao(); LocationDao? _locationDao; + @override + PendingOperationDao get pendingOperationDao => _pendingOperationDao ??= MockPendingOperationDao(); + PendingOperationDao? _pendingOperationDao; + @override Future flush() => Future.value(); @@ -96,3 +100,5 @@ class MockPollVoteDao extends Mock implements PollVoteDao {} class MockDraftMessageDao extends Mock implements DraftMessageDao {} class MockLocationDao extends Mock implements LocationDao {} + +class MockPendingOperationDao extends Mock implements PendingOperationDao {} diff --git a/packages/stream_chat_persistence/test/src/dao/pending_operation_dao_test.dart b/packages/stream_chat_persistence/test/src/dao/pending_operation_dao_test.dart new file mode 100644 index 0000000000..4a72c41e42 --- /dev/null +++ b/packages/stream_chat_persistence/test/src/dao/pending_operation_dao_test.dart @@ -0,0 +1,74 @@ +import 'package:flutter_test/flutter_test.dart'; +import 'package:stream_chat/stream_chat.dart'; +import 'package:stream_chat_persistence/src/dao/pending_operation_dao.dart'; +import 'package:stream_chat_persistence/src/db/drift_chat_database.dart'; + +import '../../stream_chat_persistence_client_test.dart'; + +void main() { + late PendingOperationDao dao; + late DriftChatDatabase database; + + setUp(() { + database = testDatabaseProvider('testUserId'); + dao = database.pendingOperationDao; + }); + + tearDown(() async { + await database.close(); + }); + + PendingOperation operation({ + String messageId = 'm1', + String type = 'reaction.add', + Map? payload, + }) => PendingOperation( + type: type, + targetMessageId: messageId, + payload: payload ?? const {'reaction': 'like', 'enforce_unique': true}, + ); + + test('round-trips an operation, preserving identity + payload', () async { + final id = await dao.insertPendingOperation(operation()); + + final operations = await dao.getPendingOperations(); + expect(operations, hasLength(1)); + + final op = operations.single; + expect(op.type, 'reaction.add'); + expect(op.targetMessageId, 'm1'); + expect(op.payload, {'reaction': 'like', 'enforce_unique': true}); + // insert returns the assigned autoincrement id. + expect(op.id, id); + }); + + test('insert appends — the same target queues as two distinct rows', () async { + await dao.insertPendingOperation(operation()); + await dao.insertPendingOperation(operation()); + + final operations = await dao.getPendingOperations(); + expect(operations, hasLength(2)); + expect(operations[0].id, isNot(operations[1].id)); + }); + + test('getPendingOperations is ordered by insertion id', () async { + await dao.insertPendingOperation(operation(messageId: 'm1')); + await dao.insertPendingOperation(operation(messageId: 'm2')); + await dao.insertPendingOperation(operation(messageId: 'm3')); + + final operations = await dao.getPendingOperations(); + expect( + operations.map((o) => o.targetMessageId).toList(), + ['m1', 'm2', 'm3'], + ); + }); + + test('deletePendingOperation removes by id', () async { + await dao.insertPendingOperation(operation()); + final stored = await dao.getPendingOperations(); + expect(stored, hasLength(1)); + + await dao.deletePendingOperation(stored.single.id!); + expect(await dao.getPendingOperations(), isEmpty); + }); +} 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 5cc2f907f3..4404c82260 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 @@ -1373,6 +1373,37 @@ void main() { }); }); + test('insertPendingOperation', () async { + const operation = PendingOperation( + type: 'reaction.add', + targetMessageId: 'testMessageId', + payload: {'reaction': 'like'}, + ); + when(() => mockDatabase.pendingOperationDao.insertPendingOperation(operation)).thenAnswer((_) => Future.value(1)); + + expect(await client.insertPendingOperation(operation), 1); + verify(() => mockDatabase.pendingOperationDao.insertPendingOperation(operation)).called(1); + }); + + test('getPendingOperations', () async { + const operations = [ + PendingOperation(id: 1, type: 'reaction.add', payload: {}), + ]; + when(() => mockDatabase.pendingOperationDao.getPendingOperations()).thenAnswer((_) => Future.value(operations)); + + final result = await client.getPendingOperations(); + expect(result, operations); + verify(() => mockDatabase.pendingOperationDao.getPendingOperations()).called(1); + }); + + test('deletePendingOperation', () async { + const id = 42; + when(() => mockDatabase.pendingOperationDao.deletePendingOperation(id)).thenAnswer((_) => Future.value()); + + await client.deletePendingOperation(id); + verify(() => mockDatabase.pendingOperationDao.deletePendingOperation(id)).called(1); + }); + tearDown(() async { await client.disconnect(flush: true); }); diff --git a/sample_app/integration_test/reactions_test.dart b/sample_app/integration_test/reactions_test.dart index 0eee3c19b1..c08bf4a65d 100644 --- a/sample_app/integration_test/reactions_test.dart +++ b/sample_app/integration_test/reactions_test.dart @@ -165,7 +165,7 @@ void main() { streamTestWithEnv( allureId: '11287', description: 'user adds a reaction while offline', - skip: 'https://linear.app/stream/issue/FLU-506', + persistence: true, body: (env) async { step('GIVEN user opens the channel'); await env.userRobot.login().openChannel();