diff --git a/core/lib/core.dart b/core/lib/core.dart index bfe0664e6..1d92e3eef 100644 --- a/core/lib/core.dart +++ b/core/lib/core.dart @@ -16,6 +16,7 @@ export 'presentation/extensions/string_extension.dart'; export 'presentation/extensions/tap_down_details_extension.dart'; export 'domain/extensions/media_type_extension.dart'; export 'presentation/extensions/map_extensions.dart'; +export 'presentation/extensions/either_view_state_extension.dart'; // Exceptions export 'domain/exceptions/download_file_exception.dart'; diff --git a/core/lib/presentation/extensions/either_view_state_extension.dart b/core/lib/presentation/extensions/either_view_state_extension.dart new file mode 100644 index 000000000..76faf2f41 --- /dev/null +++ b/core/lib/presentation/extensions/either_view_state_extension.dart @@ -0,0 +1,12 @@ +import 'package:core/presentation/state/failure.dart'; +import 'package:core/presentation/state/success.dart'; +import 'package:dartz/dartz.dart'; + +extension EitherViewStateExtension on Either { + dynamic foldSuccessWithResult() { + return fold( + (failure) => failure, + (success) => success is T ? success as T : null, + ); + } +} \ No newline at end of file diff --git a/lib/features/base/base_controller.dart b/lib/features/base/base_controller.dart index 1c8f904e8..9c8ca1668 100644 --- a/lib/features/base/base_controller.dart +++ b/lib/features/base/base_controller.dart @@ -143,20 +143,7 @@ abstract class BaseController extends GetxController void onData(Either newState) { viewState.value = newState; - viewState.value.fold( - (failure) { - if (failure is FeatureFailure) { - final isUrgentException = validateUrgentException(failure.exception); - if (isUrgentException) { - handleUrgentException(failure: failure, exception: failure.exception); - } else { - handleFailureViewState(failure); - } - } else { - handleFailureViewState(failure); - } - }, - handleSuccessViewState); + viewState.value.fold(onDataFailureViewState, handleSuccessViewState); } void onError(dynamic error, StackTrace stackTrace) { @@ -272,6 +259,19 @@ abstract class BaseController extends GetxController } } + void onDataFailureViewState(Failure failure) { + if (failure is FeatureFailure) { + final isUrgentException = validateUrgentException(failure.exception); + if (isUrgentException) { + handleUrgentException(failure: failure, exception: failure.exception); + } else { + handleFailureViewState(failure); + } + } else { + handleFailureViewState(failure); + } + } + void handleFailureViewState(Failure failure) async { logError('$runtimeType::handleFailureViewState():Failure = $failure'); if (failure is LogoutOidcFailure) { diff --git a/lib/features/base/base_mailbox_controller.dart b/lib/features/base/base_mailbox_controller.dart index 854eca61a..aa62a5145 100644 --- a/lib/features/base/base_mailbox_controller.dart +++ b/lib/features/base/base_mailbox_controller.dart @@ -104,8 +104,7 @@ abstract class BaseMailboxController extends BaseController { teamMailboxesTree.value = tupleTree.value3; } - Future syncAllMailboxWithDisplayName(BuildContext context) async { - log("BaseMailboxController::syncAllMailboxWithDisplayName"); + void syncAllMailboxWithDisplayName(BuildContext context) { final syncedMailbox = allMailboxes .map((mailbox) => mailbox.withDisplayName(mailbox.getDisplayName(context))) .toList(); diff --git a/lib/features/destination_picker/presentation/destination_picker_controller.dart b/lib/features/destination_picker/presentation/destination_picker_controller.dart index eb77e50a2..28d6ede52 100644 --- a/lib/features/destination_picker/presentation/destination_picker_controller.dart +++ b/lib/features/destination_picker/presentation/destination_picker_controller.dart @@ -112,12 +112,12 @@ class DestinationPickerController extends BaseMailboxController { await buildTree(success.mailboxList.listSubscribedMailboxesAndDefaultMailboxes); } if (currentContext != null) { - await syncAllMailboxWithDisplayName(currentContext!); + syncAllMailboxWithDisplayName(currentContext!); } } else if (success is RefreshChangesAllMailboxSuccess) { await refreshTree(success.mailboxList.listSubscribedMailboxesAndDefaultMailboxes); if (currentContext != null) { - await syncAllMailboxWithDisplayName(currentContext!); + syncAllMailboxWithDisplayName(currentContext!); } } else if (success is SearchMailboxSuccess) { _searchMailboxSuccess(success); diff --git a/lib/features/email/presentation/controller/single_email_controller.dart b/lib/features/email/presentation/controller/single_email_controller.dart index 8988f2df7..6918ee525 100644 --- a/lib/features/email/presentation/controller/single_email_controller.dart +++ b/lib/features/email/presentation/controller/single_email_controller.dart @@ -1711,57 +1711,54 @@ class SingleEmailController extends BaseController with AppLoaderMixin { void _rejectCalendarEventAction(EmailId emailId) { if (_rejectCalendarEventInteractor == null || _displayingEventBlobId == null - || mailboxDashBoardController.accountId.value == null - || mailboxDashBoardController.sessionCurrent == null - || mailboxDashBoardController.sessionCurrent - !.validateCalendarEventCapability(mailboxDashBoardController.accountId.value!) - .isAvailable == false + || accountId == null + || session == null + || session!.validateCalendarEventCapability(accountId!).isAvailable == false ) { consumeState(Stream.value(Left(CalendarEventRejectFailure()))); } else { consumeState(_rejectCalendarEventInteractor!.execute( - mailboxDashBoardController.accountId.value!, + accountId!, {_displayingEventBlobId!}, emailId, - mailboxDashBoardController.sessionCurrent!.getLanguageForCalendarEvent( + session!.getLanguageForCalendarEvent( LocalizationService.getLocaleFromLanguage(), - mailboxDashBoardController.accountId.value!))); + accountId!, + ), + )); } } void _maybeCalendarEventAction(EmailId emailId) { if (_maybeCalendarEventInteractor == null || _displayingEventBlobId == null - || mailboxDashBoardController.accountId.value == null - || mailboxDashBoardController.sessionCurrent == null - || mailboxDashBoardController.sessionCurrent - !.validateCalendarEventCapability(mailboxDashBoardController.accountId.value!) - .isAvailable == false + || accountId == null + || session == null + || session!.validateCalendarEventCapability(accountId!).isAvailable == false ) { consumeState(Stream.value(Left(CalendarEventMaybeFailure()))); } else { consumeState(_maybeCalendarEventInteractor!.execute( - mailboxDashBoardController.accountId.value!, + accountId!, {_displayingEventBlobId!}, emailId, - mailboxDashBoardController.sessionCurrent!.getLanguageForCalendarEvent( + session!.getLanguageForCalendarEvent( LocalizationService.getLocaleFromLanguage(), - mailboxDashBoardController.accountId.value!))); + accountId!, + ), + )); } } void calendarEventSuccess(CalendarEventReplySuccess success) { - final session = mailboxDashBoardController.sessionCurrent; - final accountId = mailboxDashBoardController.accountId.value; - if (session == null || accountId == null) { consumeState(Stream.value(Left(StoreEventAttendanceStatusFailure(exception: NotFoundSessionException())))); return; } consumeState(_storeEventAttendanceStatusInteractor.execute( - session, - accountId, + session!, + accountId!, success.emailId, success.getEventActionType() )); @@ -1839,16 +1836,15 @@ class SingleEmailController extends BaseController with AppLoaderMixin { } Future previewPDFFileAction(BuildContext context, Attachment attachment) async { - final accountId = mailboxDashBoardController.accountId.value; - final downloadUrl = mailboxDashBoardController.sessionCurrent - ?.getDownloadUrl(jmapUrl: dynamicUrlInterceptors.jmapUrl); - - if (accountId == null || downloadUrl == null) { + if (accountId == null || session == null) { appToast.showToastErrorMessage( context, AppLocalizations.of(context).noPreviewAvailable); return; } + final downloadUrl = session!.getDownloadUrl( + jmapUrl: dynamicUrlInterceptors.jmapUrl, + ); await Get.generalDialog( barrierColor: Colors.black.withOpacity(0.8), @@ -1856,7 +1852,7 @@ class SingleEmailController extends BaseController with AppLoaderMixin { return PointerInterceptor( child: PDFViewer( attachment: attachment, - accountId: accountId, + accountId: accountId!, downloadUrl: downloadUrl, downloadAction: _downloadPDFFile, printAction: _printPDFFile, @@ -1876,7 +1872,10 @@ class SingleEmailController extends BaseController with AppLoaderMixin { } final listEmailAddressMailTo = listEmailAddressAttendees - .where((emailAddress) => emailAddress.emailAddress.isNotEmpty && emailAddress.emailAddress != mailboxDashBoardController.sessionCurrent?.username.value) + .where((emailAddress) { + return emailAddress.emailAddress.isNotEmpty && + emailAddress.emailAddress != session?.username.value; + }) .toSet() .toList(); diff --git a/lib/features/mailbox/presentation/mailbox_controller.dart b/lib/features/mailbox/presentation/mailbox_controller.dart index d6ac44031..f117eae64 100644 --- a/lib/features/mailbox/presentation/mailbox_controller.dart +++ b/lib/features/mailbox/presentation/mailbox_controller.dart @@ -64,6 +64,8 @@ import 'package:tmail_ui_user/features/mailbox_creator/presentation/model/new_ma import 'package:tmail_ui_user/features/mailbox_dashboard/presentation/action/dashboard_action.dart'; import 'package:tmail_ui_user/features/mailbox_dashboard/presentation/controller/mailbox_dashboard_controller.dart'; import 'package:tmail_ui_user/features/mailbox_dashboard/presentation/model/dashboard_routes.dart'; +import 'package:tmail_ui_user/features/push_notification/presentation/websocket/web_socket_message.dart'; +import 'package:tmail_ui_user/features/push_notification/presentation/websocket/web_socket_queue_handler.dart'; import 'package:tmail_ui_user/features/search/mailbox/presentation/search_mailbox_bindings.dart'; import 'package:tmail_ui_user/features/thread/domain/model/search_query.dart'; import 'package:tmail_ui_user/main/localizations/app_localizations.dart'; @@ -94,6 +96,7 @@ class MailboxController extends BaseMailboxController with MailboxActionHandlerM MailboxId? _newFolderId; NavigationRouter? _navigationRouter; + WebSocketQueueHandler? _webSocketQueueHandler; final _openMailboxEventController = StreamController(); final mailboxListScrollController = ScrollController(); @@ -128,6 +131,7 @@ class MailboxController extends BaseMailboxController with MailboxActionHandlerM @override void onInit() { _registerObxStreamListener(); + _initWebSocketQueueHandler(); super.onInit(); } @@ -145,6 +149,7 @@ class MailboxController extends BaseMailboxController with MailboxActionHandlerM void onClose() { _openMailboxEventController.close(); mailboxListScrollController.dispose(); + _webSocketQueueHandler?.dispose(); super.onClose(); } @@ -153,8 +158,6 @@ class MailboxController extends BaseMailboxController with MailboxActionHandlerM super.handleSuccessViewState(success); if (success is GetAllMailboxSuccess) { _handleGetAllMailboxSuccess(success); - } else if (success is RefreshChangesAllMailboxSuccess) { - _handleRefreshChangesAllMailboxSuccess(success); } else if (success is CreateNewMailboxSuccess) { _createNewMailboxSuccess(success); } else if (success is DeleteMultipleMailboxAllSuccess) { @@ -181,8 +184,6 @@ class MailboxController extends BaseMailboxController with MailboxActionHandlerM _renameMailboxFailure(failure); } else if (failure is DeleteMultipleMailboxFailure) { _deleteMailboxFailure(failure); - } else if (failure is RefreshChangesAllMailboxFailure) { - _clearNewFolderId(); } } @@ -204,13 +205,6 @@ class MailboxController extends BaseMailboxController with MailboxActionHandlerM if (PlatformInfo.isIOS) { _updateMailboxIdsBlockNotificationToKeychain(success.mailboxList); } - } else if (success is RefreshChangesAllMailboxSuccess) { - _selectSelectedMailboxDefault(); - mailboxDashBoardController.refreshSpamReportBanner(); - - if (_newFolderId != null) { - _redirectToNewFolder(); - } } }); } @@ -258,6 +252,13 @@ class MailboxController extends BaseMailboxController with MailboxActionHandlerM }); } + void _initWebSocketQueueHandler() { + _webSocketQueueHandler = WebSocketQueueHandler( + processMessageCallback: _handleWebSocketMessage, + onErrorCallback: onError, + ); + } + void _initCollapseMailboxCategories() { if (kIsWeb && currentContext != null && (responsiveUtils.isMobile(currentContext!) || responsiveUtils.isTablet(currentContext!))) { @@ -283,17 +284,66 @@ class MailboxController extends BaseMailboxController with MailboxActionHandlerM if (accountId == null || session == null || currentMailboxState == null || - newState == currentMailboxState) { + newState == null) { _newFolderId = null; return; } - refreshMailboxChanges( - session!, - accountId!, - currentMailboxState!, - properties: MailboxConstants.propertiesDefault, - ); + _webSocketQueueHandler?.enqueue(WebSocketMessage(newState: newState)); + } + + Future _handleWebSocketMessage(WebSocketMessage message) async { + try { + if (currentMailboxState == message.newState) { + log('MailboxController::_handleWebSocketMessage:Skipping redundant state: ${message.newState}'); + return Future.value(); + } + + final refreshViewState = await refreshAllMailboxInteractor!.execute( + session!, + accountId!, + currentMailboxState!, + properties: MailboxConstants.propertiesDefault, + ).last; + + final refreshState = refreshViewState + .foldSuccessWithResult(); + + if (refreshState is RefreshChangesAllMailboxSuccess) { + await _handleRefreshChangeMailboxSuccess(refreshState); + } else { + _clearNewFolderId(); + onDataFailureViewState(refreshState); + } + } catch (e, stackTrace) { + logError('MailboxController::_processMailboxStateQueue:Error processing state: $e'); + onError(e, stackTrace); + } + if (currentMailboxState != null) { + _webSocketQueueHandler?.removeMessagesUpToCurrent(currentMailboxState!.value); + } + } + + Future _handleRefreshChangeMailboxSuccess(RefreshChangesAllMailboxSuccess success) async { + currentMailboxState = success.currentMailboxState; + log('MailboxController::_handleRefreshChangeMailboxSuccess:currentMailboxState: $currentMailboxState'); + final listMailboxDisplayed = success + .mailboxList + .listSubscribedMailboxesAndDefaultMailboxes; + + await refreshTree(listMailboxDisplayed); + + if (currentContext != null) { + syncAllMailboxWithDisplayName(currentContext!); + } + _setMapMailbox(); + _setOutboxMailbox(); + _selectSelectedMailboxDefault(); + mailboxDashBoardController.refreshSpamReportBanner(); + + if (_newFolderId != null) { + _redirectToNewFolder(); + } } void _setMapMailbox() { @@ -1078,7 +1128,7 @@ class MailboxController extends BaseMailboxController with MailboxActionHandlerM final listMailboxDisplayed = success.mailboxList.listSubscribedMailboxesAndDefaultMailboxes; await buildTree(listMailboxDisplayed); if (currentContext != null) { - await syncAllMailboxWithDisplayName(currentContext!); + syncAllMailboxWithDisplayName(currentContext!); } _setMapMailbox(); _setOutboxMailbox(); @@ -1105,18 +1155,6 @@ class MailboxController extends BaseMailboxController with MailboxActionHandlerM mailboxIds: mailboxIdsBlockNotification); } - void _handleRefreshChangesAllMailboxSuccess(RefreshChangesAllMailboxSuccess success) async { - currentMailboxState = success.currentMailboxState; - log('MailboxController::_handleRefreshChangesAllMailboxSuccess:currentMailboxState: $currentMailboxState'); - final listMailboxDisplayed = success.mailboxList.listSubscribedMailboxesAndDefaultMailboxes; - await refreshTree(listMailboxDisplayed); - if (currentContext != null) { - await syncAllMailboxWithDisplayName(currentContext!); - } - _setMapMailbox(); - _setOutboxMailbox(); - } - void _unsubscribeMailboxAction(MailboxId mailboxId) { if (session != null && accountId != null) { final subscribeRequest = generateSubscribeRequest( diff --git a/lib/features/mailbox_dashboard/presentation/controller/mailbox_dashboard_controller.dart b/lib/features/mailbox_dashboard/presentation/controller/mailbox_dashboard_controller.dart index 7a8161165..d97e43b1a 100644 --- a/lib/features/mailbox_dashboard/presentation/controller/mailbox_dashboard_controller.dart +++ b/lib/features/mailbox_dashboard/presentation/controller/mailbox_dashboard_controller.dart @@ -243,6 +243,7 @@ class MailboxDashBoardController extends ReloadableController with UserSettingPo PresentationMailbox? outboxMailbox; ComposerArguments? composerArguments; List? _identities; + jmap.State? _currentEmailState; ScrollController? listSearchFilterScrollController; StreamSubscription? _pendingSharedFileInfoSubscription; StreamSubscription? _receivingFileSharingStreamSubscription; @@ -2930,6 +2931,12 @@ class MailboxDashBoardController extends ReloadableController with UserSettingPo accountId: accountId.value!); } + void setCurrentEmailState(jmap.State? newState) { + _currentEmailState = newState; + } + + jmap.State? get currentEmailState => _currentEmailState; + @override void onClose() { if (PlatformInfo.isWeb) { @@ -2957,6 +2964,7 @@ class MailboxDashBoardController extends ReloadableController with UserSettingPo mapMailboxById = {}; mapDefaultMailboxIdByRole = {}; WebSocketController.instance.onClose(); + _currentEmailState = null; super.onClose(); } } \ No newline at end of file diff --git a/lib/features/manage_account/presentation/mailbox_visibility/mailbox_visibility_controller.dart b/lib/features/manage_account/presentation/mailbox_visibility/mailbox_visibility_controller.dart index ceb7b0669..d405b92fa 100644 --- a/lib/features/manage_account/presentation/mailbox_visibility/mailbox_visibility_controller.dart +++ b/lib/features/manage_account/presentation/mailbox_visibility/mailbox_visibility_controller.dart @@ -73,7 +73,7 @@ class MailboxVisibilityController extends BaseMailboxController { currentMailboxState = success.currentMailboxState; await refreshTree(success.mailboxList); if (currentContext != null) { - await syncAllMailboxWithDisplayName(currentContext!); + syncAllMailboxWithDisplayName(currentContext!); } } else if (success is SubscribeMailboxSuccess) { _subscribeMailboxSuccess(success); @@ -99,7 +99,7 @@ class MailboxVisibilityController extends BaseMailboxController { await buildTree(mailboxList); dispatchState(Right(BuildTreeMailboxVisibilitySuccess())); if (currentContext != null) { - await syncAllMailboxWithDisplayName(currentContext!); + syncAllMailboxWithDisplayName(currentContext!); } } diff --git a/lib/features/push_notification/presentation/controller/web_socket_controller.dart b/lib/features/push_notification/presentation/controller/web_socket_controller.dart index 821a873cb..afdada5c6 100644 --- a/lib/features/push_notification/presentation/controller/web_socket_controller.dart +++ b/lib/features/push_notification/presentation/controller/web_socket_controller.dart @@ -5,6 +5,7 @@ import 'package:core/presentation/state/failure.dart'; import 'package:core/presentation/state/success.dart'; import 'package:core/utils/app_logger.dart'; import 'package:core/utils/platform_info.dart'; +import 'package:debounce_throttle/debounce_throttle.dart'; import 'package:fcm/model/type_name.dart'; import 'package:flutter/material.dart'; import 'package:jmap_dart_client/jmap/account_id.dart'; @@ -18,6 +19,7 @@ import 'package:tmail_ui_user/features/push_notification/presentation/controller import 'package:tmail_ui_user/features/push_notification/presentation/extensions/state_change_extension.dart'; import 'package:tmail_ui_user/features/push_notification/presentation/listener/email_change_listener.dart'; import 'package:tmail_ui_user/features/push_notification/presentation/listener/mailbox_change_listener.dart'; +import 'package:tmail_ui_user/features/push_notification/presentation/utils/fcm_utils.dart'; import 'package:tmail_ui_user/main/routes/route_navigation.dart'; import 'package:tmail_ui_user/main/utils/app_config.dart'; import 'package:web_socket_channel/web_socket_channel.dart'; @@ -36,6 +38,7 @@ class WebSocketController extends PushBaseController { Timer? _webSocketPingTimer; StreamSubscription? _webSocketSubscription; AppLifecycleListener? _appLifecycleListener; + Debouncer? _stateChangeDebouncer; @override void handleFailureViewState(Failure failure) { @@ -105,8 +108,9 @@ class WebSocketController extends PushBaseController { _webSocketSubscription?.cancel(); _webSocketChannel = null; _webSocketPingTimer?.cancel(); + _stateChangeDebouncer?.cancel(); } - + void _handleWebSocketConnectionSuccess(WebSocketConnectionSuccess success) { log('WebSocketController::_handleWebSocketConnectionSuccess(): $success'); _cleanUpWebSocketResources(); @@ -117,6 +121,7 @@ class WebSocketController extends PushBaseController { _pingWebSocket(); } _listenToWebSocket(); + _initStateChangeDeouncerTimer(); } void _handleWebSocketConnectionRetry() { @@ -150,14 +155,7 @@ class WebSocketController extends PushBaseController { try { final stateChange = StateChange.fromJson(data); - final mapTypeState = stateChange.getMapTypeState(accountId!); - mappingTypeStateToAction( - mapTypeState, - accountId!, - emailChangeListener: EmailChangeListener.instance, - mailboxChangeListener: MailboxChangeListener.instance, - session!.username, - session: session); + _stateChangeDebouncer?.value = stateChange; } catch (e) { logError('WebSocketController::_listenToWebSocket(): Data is not StateChange'); } @@ -173,4 +171,31 @@ class WebSocketController extends PushBaseController { }, ); } + + void _initStateChangeDeouncerTimer() { + _stateChangeDebouncer = Debouncer( + const Duration(milliseconds: FcmUtils.durationMessageComing), + initialValue: null, + ); + + _stateChangeDebouncer?.values.listen(_handleStateChange); + } + + void _handleStateChange(StateChange? stateChange) { + try { + if (stateChange == null || accountId == null || session == null) return; + + final mapTypeState = stateChange.getMapTypeState(accountId!); + mappingTypeStateToAction( + mapTypeState, + accountId!, + emailChangeListener: EmailChangeListener.instance, + mailboxChangeListener: MailboxChangeListener.instance, + session!.username, + session: session, + ); + } catch (e) { + logError('WebSocketController::_handleStateChange:Exception = $e'); + } + } } \ No newline at end of file diff --git a/lib/features/push_notification/presentation/websocket/web_socket_message.dart b/lib/features/push_notification/presentation/websocket/web_socket_message.dart new file mode 100644 index 000000000..58532f855 --- /dev/null +++ b/lib/features/push_notification/presentation/websocket/web_socket_message.dart @@ -0,0 +1,13 @@ +import 'package:equatable/equatable.dart'; +import 'package:jmap_dart_client/jmap/core/state.dart' as jmap; + +class WebSocketMessage with EquatableMixin { + final jmap.State newState; + + WebSocketMessage({required this.newState}); + + String get id => newState.value; + + @override + List get props => [newState]; +} \ No newline at end of file diff --git a/lib/features/push_notification/presentation/websocket/web_socket_queue_handler.dart b/lib/features/push_notification/presentation/websocket/web_socket_queue_handler.dart new file mode 100644 index 000000000..155b0deb6 --- /dev/null +++ b/lib/features/push_notification/presentation/websocket/web_socket_queue_handler.dart @@ -0,0 +1,121 @@ + +import 'dart:async'; +import 'dart:collection'; + +import 'package:core/utils/app_logger.dart'; +import 'package:flutter/cupertino.dart'; +import 'package:tmail_ui_user/features/push_notification/presentation/websocket/web_socket_message.dart'; + +typedef ProcessMessageCallback = Future Function(WebSocketMessage message); +typedef OnErrorCallback = void Function(dynamic error, StackTrace stackTrace); + +class WebSocketQueueHandler { + static const int _maxQueueSize = 128; + static const int _maxProcessedIdsSize = 128; + + final Queue _messageQueue = Queue(); + final Queue _processedMessageIds = Queue(); + + Completer? _processingLock; + + final _queueController = StreamController.broadcast(); + + final ProcessMessageCallback processMessageCallback; + final OnErrorCallback? onErrorCallback; + + WebSocketQueueHandler({ + required this.processMessageCallback, + this.onErrorCallback, + }) { + _queueController.stream.listen((_) { + _processQueue(); + }); + } + + void enqueue(WebSocketMessage message) { + if (isMessageProcessed(message.id)) { + log('WebSocketQueueHandler::enqueue:Message ${message.id} already processed, skipping'); + return; + } + + if (queueSize >= _maxQueueSize) { + log('WebSocketQueueHandler::enqueue:Queue full, removing oldest message'); + _messageQueue.removeFirst(); + } + + _messageQueue.add(message); + _queueController.add(message); + } + + Future _processQueue() async { + if (_processingLock != null) { + return; + } + + _processingLock = Completer(); + + try { + while (queueSize > 0) { + final message = _messageQueue.removeFirst(); + + try { + await processMessageCallback(message); + } catch (e, stackTrace) { + logError('WebSocketQueueHandler::_processQueue:Error processing message ${message.id}: $e'); + onErrorCallback?.call(e, stackTrace); + } finally { + _addToProcessedMessages(message.id); + } + } + } finally { + _processingLock?.complete(); + _processingLock = null; + + if (queueSize > 0) { + scheduleMicrotask(() => _queueController.add(_messageQueue.first)); + } + } + } + + void _addToProcessedMessages(String messageId) { + if (_processedMessageIds.length >= _maxProcessedIdsSize) { + _processedMessageIds.removeFirst(); + } + _processedMessageIds.add(messageId); + } + + void removeMessagesUpToCurrent(String messageId) { + final isCurrentStateExist = _messageQueue + .any((message) => message.id == messageId); + + if (!isCurrentStateExist) { + log('WebSocketQueueHandler::removeMessagesUpToCurrent:Current state $messageId not found in the queue.'); + return; + } + while (queueSize > 0) { + final removedMessage = _messageQueue.removeFirst(); + if (removedMessage.id == messageId) { + break; + } + } + log('WebSocketQueueHandler::removeMessagesUpToCurrent:Updated Queue: $queueSize'); + } + + @visibleForTesting + Future waitForEmpty() async { + while (_messageQueue.isNotEmpty || _processingLock != null) { + if (_processingLock != null) { + await _processingLock!.future; + } + await Future.delayed(const Duration(milliseconds: 100)); + } + } + + int get queueSize => _messageQueue.length; + + bool isMessageProcessed(String messageId) => _processedMessageIds.contains(messageId); + + void dispose() { + _queueController.close(); + } +} diff --git a/lib/features/rules_filter_creator/presentation/rules_filter_creator_controller.dart b/lib/features/rules_filter_creator/presentation/rules_filter_creator_controller.dart index e9f535dea..bea072347 100644 --- a/lib/features/rules_filter_creator/presentation/rules_filter_creator_controller.dart +++ b/lib/features/rules_filter_creator/presentation/rules_filter_creator_controller.dart @@ -127,7 +127,7 @@ class RulesFilterCreatorController extends BaseMailboxController { if (success is GetAllMailboxSuccess) { await buildTree(success.mailboxList); if (currentContext != null) { - await syncAllMailboxWithDisplayName(currentContext!); + syncAllMailboxWithDisplayName(currentContext!); } } else if (success is GetAllRulesSuccess) { log('RulesFilterCreatorController::handleSuccessViewState():GetAllRulesSuccess: ${success.rules}'); diff --git a/lib/features/search/email/presentation/search_email_controller.dart b/lib/features/search/email/presentation/search_email_controller.dart index f96f957ea..37c9081eb 100644 --- a/lib/features/search/email/presentation/search_email_controller.dart +++ b/lib/features/search/email/presentation/search_email_controller.dart @@ -1,4 +1,5 @@ +import 'package:core/presentation/extensions/either_view_state_extension.dart'; import 'package:core/presentation/state/failure.dart'; import 'package:core/presentation/state/success.dart'; import 'package:core/presentation/utils/keyboard_utils.dart'; @@ -10,6 +11,7 @@ import 'package:flutter/material.dart'; import 'package:get/get.dart'; import 'package:jmap_dart_client/jmap/account_id.dart'; import 'package:jmap_dart_client/jmap/core/session/session.dart'; +import 'package:jmap_dart_client/jmap/core/state.dart' as jmap; import 'package:jmap_dart_client/jmap/core/unsigned_int.dart'; import 'package:jmap_dart_client/jmap/core/utc_date.dart'; import 'package:jmap_dart_client/jmap/mail/email/email_address.dart'; @@ -50,6 +52,8 @@ import 'package:tmail_ui_user/features/mailbox_dashboard/presentation/model/sear import 'package:tmail_ui_user/features/manage_account/presentation/extensions/datetime_extension.dart'; import 'package:tmail_ui_user/features/network_connection/presentation/network_connection_controller.dart' if (dart.library.html) 'package:tmail_ui_user/features/network_connection/presentation/web_network_connection_controller.dart'; +import 'package:tmail_ui_user/features/push_notification/presentation/websocket/web_socket_message.dart'; +import 'package:tmail_ui_user/features/push_notification/presentation/websocket/web_socket_queue_handler.dart'; import 'package:tmail_ui_user/features/search/email/domain/state/refresh_changes_search_email_state.dart'; import 'package:tmail_ui_user/features/search/email/domain/usecases/refresh_changes_search_email_interactor.dart'; import 'package:tmail_ui_user/features/search/email/presentation/model/search_more_state.dart'; @@ -105,6 +109,7 @@ class SearchEmailController extends BaseController late Worker dashBoardActionWorker; late SearchMoreState searchMoreState; late bool canSearchMore; + WebSocketQueueHandler? _webSocketQueueHandler; PresentationMailbox? get currentMailbox => mailboxDashBoardController.selectedMailbox.value; @@ -145,6 +150,7 @@ class SearchEmailController extends BaseController _initializeDebounceTimeTextSearchChange(); _initializeTextInputFocus(); _initWorkerListener(); + _initWebSocketQueueHandler(); } @override @@ -165,8 +171,6 @@ class SearchEmailController extends BaseController searchMoreState = SearchMoreState.waiting; } else if (success is SearchMoreEmailSuccess) { _searchMoreEmailsSuccess(success); - } else if (success is RefreshChangesSearchEmailSuccess) { - _refreshChangesSearchEmailsSuccess(success); } else if (success is SearchingState) { resultSearchViewState.value = Right(success); } @@ -252,40 +256,88 @@ class SearchEmailController extends BaseController mailboxDashBoardController.emailUIAction, (action) { if (action is RefreshChangeEmailAction) { - _refreshEmailChanges(); + _refreshEmailChanges(newState: action.newState); } }, ); } + void _refreshEmailChanges({jmap.State? newState}) { + log('SearchEmailController::_refreshEmailChanges(): newState: $newState'); + if (accountId == null || + session == null || + mailboxDashBoardController.currentEmailState == null || + newState == null || + searchIsRunning.isFalse) { + return; + } + + _webSocketQueueHandler?.enqueue(WebSocketMessage(newState: newState)); + } + + void _initWebSocketQueueHandler() { + _webSocketQueueHandler = WebSocketQueueHandler( + processMessageCallback: _handleWebSocketMessage, + onErrorCallback: onError, + ); + } + + Future _handleWebSocketMessage(WebSocketMessage message) async { + try { + if (mailboxDashBoardController.currentEmailState == null || + mailboxDashBoardController.currentEmailState == message.newState) { + log('SearchEmailController::_handleWebSocketMessage:Skipping redundant state: ${message.newState}'); + return Future.value(); + } + + final limit = listResultSearch.isNotEmpty + ? UnsignedInt(listResultSearch.length) + : ThreadConstants.defaultLimit; + + if (limit.value > ThreadConstants.maximumEmailQueryLimit && + resultSearchScrollController.hasClients) { + resultSearchScrollController.jumpTo(0); + } + + _updateSimpleSearchFilter( + beforeOption: const None(), + positionOption: option(searchEmailFilter.value.sortOrderType.isScrollByPosition(), 0), + ); + + final searchViewState = await _refreshChangesSearchEmailInteractor.execute( + session!, + accountId!, + limit: limit, + position: searchEmailFilter.value.position, + sort: searchEmailFilter.value.sortOrderType.getSortOrder().toNullable(), + filter: searchEmailFilter.value.mappingToEmailFilterCondition(), + properties: EmailUtils.getPropertiesForEmailGetMethod(session!, accountId!), + ).last; + + final searchState = searchViewState + .foldSuccessWithResult(); + + if (searchState is RefreshChangesSearchEmailSuccess) { + _handleRefreshChangesSearchEmailsSuccess(searchState); + } + } catch (e, stackTrace) { + logError('SearchEmailController::_handleWebSocketMessage:Error processing state: $e'); + onError(e, stackTrace); + } finally { + if (mailboxDashBoardController.currentEmailState != null) { + _webSocketQueueHandler?.removeMessagesUpToCurrent( + mailboxDashBoardController.currentEmailState!.value); + } + } + } + void _onSearchTextInputListener() { if (textInputSearchFocus.hasFocus) { searchIsRunning.value = false; } } - void _refreshEmailChanges() { - if (searchIsRunning.isTrue && session != null && accountId != null) { - final limit = listResultSearch.isNotEmpty - ? UnsignedInt(listResultSearch.length) - : ThreadConstants.defaultLimit; - _updateSimpleSearchFilter( - beforeOption: const None(), - positionOption: option(searchEmailFilter.value.sortOrderType.isScrollByPosition(), 0) - ); - consumeState(_refreshChangesSearchEmailInteractor.execute( - session!, - accountId!, - limit: limit, - position: searchEmailFilter.value.position, - sort: searchEmailFilter.value.sortOrderType.getSortOrder().toNullable(), - filter: searchEmailFilter.value.mappingToEmailFilterCondition(), - properties: EmailUtils.getPropertiesForEmailGetMethod(session!, accountId!), - )); - } - } - - void _refreshChangesSearchEmailsSuccess(RefreshChangesSearchEmailSuccess success) { + void _handleRefreshChangesSearchEmailsSuccess(RefreshChangesSearchEmailSuccess success) { final resultEmailSearchList = success.emailList .map((email) => email.toSearchPresentationEmail(mailboxDashBoardController.mapMailboxById)) .toList(); @@ -306,8 +358,8 @@ class SearchEmailController extends BaseController return _getAllRecentSearchLatestInteractor .execute(pattern: pattern) .then((result) => result.fold( - (failure) => [], - (success) => success is GetAllRecentSearchLatestSuccess + (failure) => [], + (success) => success is GetAllRecentSearchLatestSuccess ? success.listRecentSearch : [])); } diff --git a/lib/features/search/mailbox/presentation/search_mailbox_controller.dart b/lib/features/search/mailbox/presentation/search_mailbox_controller.dart index 3447b9ff5..2ae4bc445 100644 --- a/lib/features/search/mailbox/presentation/search_mailbox_controller.dart +++ b/lib/features/search/mailbox/presentation/search_mailbox_controller.dart @@ -136,13 +136,13 @@ class SearchMailboxController extends BaseMailboxController with MailboxActionHa currentMailboxState = success.currentMailboxState; await buildTree(success.mailboxList); if (currentContext != null) { - await syncAllMailboxWithDisplayName(currentContext!); + syncAllMailboxWithDisplayName(currentContext!); } } else if (success is RefreshChangesAllMailboxSuccess) { currentMailboxState = success.currentMailboxState; await refreshTree(success.mailboxList); if (currentContext != null) { - await syncAllMailboxWithDisplayName(currentContext!); + syncAllMailboxWithDisplayName(currentContext!); } searchMailboxAction(); } else if (success is SearchMailboxSuccess) { diff --git a/lib/features/thread/domain/constants/thread_constants.dart b/lib/features/thread/domain/constants/thread_constants.dart index 2e5791697..9a5618b92 100644 --- a/lib/features/thread/domain/constants/thread_constants.dart +++ b/lib/features/thread/domain/constants/thread_constants.dart @@ -5,6 +5,7 @@ import 'package:model/email/email_property.dart'; class ThreadConstants { static const maxCountEmails = 20; + static const maximumEmailQueryLimit = 256; static final defaultLimit = UnsignedInt(maxCountEmails); static final propertiesDefault = Properties({ EmailProperty.id, diff --git a/lib/features/thread/presentation/thread_controller.dart b/lib/features/thread/presentation/thread_controller.dart index 8aac45a65..c23994fcd 100644 --- a/lib/features/thread/presentation/thread_controller.dart +++ b/lib/features/thread/presentation/thread_controller.dart @@ -1,5 +1,6 @@ import 'dart:async'; +import 'package:core/presentation/extensions/either_view_state_extension.dart'; import 'package:core/presentation/state/failure.dart'; import 'package:core/presentation/state/success.dart'; import 'package:core/utils/app_logger.dart'; @@ -32,6 +33,8 @@ import 'package:tmail_ui_user/features/manage_account/domain/state/create_new_ru import 'package:tmail_ui_user/features/manage_account/domain/usecases/create_new_email_rule_filter_interactor.dart'; import 'package:tmail_ui_user/features/network_connection/presentation/network_connection_controller.dart' if (dart.library.html) 'package:tmail_ui_user/features/network_connection/presentation/web_network_connection_controller.dart'; +import 'package:tmail_ui_user/features/push_notification/presentation/websocket/web_socket_message.dart'; +import 'package:tmail_ui_user/features/push_notification/presentation/websocket/web_socket_queue_handler.dart'; import 'package:tmail_ui_user/features/rules_filter_creator/presentation/model/rules_filter_creator_arguments.dart'; import 'package:tmail_ui_user/features/search/email/presentation/search_email_bindings.dart'; import 'package:tmail_ui_user/features/thread/domain/constants/thread_constants.dart'; @@ -90,12 +93,12 @@ class ThreadController extends BaseController with EmailActionController { bool canLoadMore = false; bool canSearchMore = false; MailboxId? _currentMemoryMailboxId; - jmap.State? _currentEmailState; final ScrollController listEmailController = ScrollController(); final FocusNode focusNodeKeyBoard = FocusNode(); final latestEmailSelectedOrUnselected = Rxn(); @visibleForTesting bool isListEmailScrollViewJumping = false; + WebSocketQueueHandler? _webSocketQueueHandler; StreamSubscription? _resizeBrowserStreamSubscription; @@ -128,6 +131,7 @@ class ThreadController extends BaseController with EmailActionController { if (PlatformInfo.isWeb) { _registerBrowserResizeListener(); } + _initWebSocketQueueHandler(); super.onInit(); } @@ -140,7 +144,6 @@ class ThreadController extends BaseController with EmailActionController { @override void onClose() { _currentMemoryMailboxId = null; - _currentEmailState = null; listEmailController.dispose(); focusNodeKeyBoard.dispose(); if (PlatformInfo.isWeb) { @@ -154,8 +157,6 @@ class ThreadController extends BaseController with EmailActionController { super.handleSuccessViewState(success); if (success is GetAllEmailSuccess) { _getAllEmailSuccess(success); - } else if (success is RefreshChangesAllEmailSuccess) { - _refreshChangesAllEmailSuccess(success); } else if (success is LoadMoreEmailsSuccess) { _loadMoreEmailsSuccess(success); } else if (success is SearchEmailSuccess) { @@ -231,6 +232,13 @@ class ThreadController extends BaseController with EmailActionController { }); } + void _initWebSocketQueueHandler() { + _webSocketQueueHandler = WebSocketQueueHandler( + processMessageCallback: _handleWebSocketMessage, + onErrorCallback: onError, + ); + } + void _resetLoadingMore() { if (loadingMoreStatus.value == LoadingMoreStatus.running) { loadingMoreStatus.value = LoadingMoreStatus.idle; @@ -383,8 +391,8 @@ class ThreadController extends BaseController with EmailActionController { log('ThreadController::_getAllEmailSuccess: GetAllForMailboxId = ${success.currentMailboxId?.asString} | SELECTED_MAILBOX_ID = ${selectedMailboxId?.asString} | SELECTED_MAILBOX_NAME = ${selectedMailbox?.name?.name}'); return; } - _currentEmailState = success.currentEmailState; - log('ThreadController::_getAllEmailSuccess():COUNT = ${success.emailList.length} | EMAIL_STATE = $_currentEmailState'); + mailboxDashBoardController.setCurrentEmailState(success.currentEmailState); + log('ThreadController::_getAllEmailSuccess():COUNT = ${success.emailList.length} | EMAIL_STATE = ${mailboxDashBoardController.currentEmailState}'); final newListEmail = success.emailList.syncPresentationEmail( mapMailboxById: mailboxDashBoardController.mapMailboxById, selectedMailbox: selectedMailbox, @@ -413,8 +421,7 @@ class ThreadController extends BaseController with EmailActionController { log('ThreadController::_refreshChangesAllEmailSuccess: RefreshedMailboxId = ${success.currentMailboxId?.asString} | SELECTED_MAILBOX_ID = ${selectedMailboxId?.asString} | SELECTED_MAILBOX_NAME = ${selectedMailbox?.name?.name}'); return; } - - _currentEmailState = success.currentEmailState; + mailboxDashBoardController.setCurrentEmailState(success.currentEmailState); log('ThreadController::_refreshChangesAllEmailSuccess: COUNT = ${success.emailList.length}'); final emailsBeforeChanges = mailboxDashBoardController.emailsInCurrentMailbox; final emailsAfterChanges = success.emailList; @@ -515,31 +522,112 @@ class ThreadController extends BaseController with EmailActionController { void _refreshEmailChanges({jmap.State? newState}) { log('ThreadController::_refreshEmailChanges(): newState: $newState'); - if (searchController.isSearchEmailRunning) { - _searchEmail(limit: limitEmailFetched, needRefreshSearchState: true); - } else { - if (_currentEmailState == null || - _currentEmailState == newState || - _session == null || - _accountId == null) { - return; + if (_accountId == null || + _session == null || + mailboxDashBoardController.currentEmailState == null || + newState == null) { + return; + } + + _webSocketQueueHandler?.enqueue(WebSocketMessage(newState: newState)); + } + + Future _handleWebSocketMessage(WebSocketMessage message) async { + try { + if (mailboxDashBoardController.currentEmailState == null || + mailboxDashBoardController.currentEmailState == message.newState) { + log('ThreadController::_handleWebSocketMessage:Skipping redundant state: ${message.newState}'); + return Future.value(); } - consumeState(_refreshChangesEmailsInMailboxInteractor.execute( + + if (searchController.isSearchEmailRunning) { + await _refreshChangeSearchEmail(); + } else { + await _refreshChangeListEmail(); + } + } catch (e, stackTrace) { + logError('ThreadController::_handleWebSocketMessage:Error processing state: $e'); + onError(e, stackTrace); + } finally { + if (mailboxDashBoardController.currentEmailState != null) { + _webSocketQueueHandler?.removeMessagesUpToCurrent( + mailboxDashBoardController.currentEmailState!.value); + } + } + } + + Future _refreshChangeSearchEmail() async { + log('ThreadController::_refreshChangeSearchEmail:'); + if (limitEmailFetched.value > ThreadConstants.maximumEmailQueryLimit && + listEmailController.hasClients) { + listEmailController.jumpTo(0); + } + canSearchMore = false; + searchController.updateFilterEmail( + positionOption: option( + _searchEmailFilter.sortOrderType.isScrollByPosition(), + 0, + ), + beforeOption: const None(), + ); + searchController.activateSimpleSearch(); + + final searchViewState = await _searchEmailInteractor.execute( + _session!, + _accountId!, + limit: limitEmailFetched, + position: _searchEmailFilter.position, + sort: _searchEmailFilter.sortOrderType.getSortOrder().toNullable(), + filter: _searchEmailFilter.mappingToEmailFilterCondition( + moreFilterCondition: _getFilterCondition(), + ), + properties: EmailUtils.getPropertiesForEmailGetMethod( _session!, _accountId!, - _currentEmailState!, - sort: EmailSortOrderType.mostRecent.getSortOrder().toNullable(), - propertiesCreated: EmailUtils.getPropertiesForEmailGetMethod( - _session!, - _accountId!, - ), - propertiesUpdated: ThreadConstants.propertiesUpdatedDefault, - emailFilter: EmailFilter( - filter: _getFilterCondition(mailboxIdSelected: selectedMailboxId), - filterOption: mailboxDashBoardController.filterMessageOption.value, - mailboxId: selectedMailboxId, - ), - )); + ), + needRefreshSearchState: true, + ).last; + + final searchState = searchViewState + .foldSuccessWithResult(); + + if (searchState is SearchEmailSuccess) { + _searchEmailsSuccess(searchState); + } else { + mailboxDashBoardController.updateRefreshAllEmailState( + Left(RefreshAllEmailFailure())); + canSearchMore = false; + mailboxDashBoardController.emailsInCurrentMailbox.clear(); + onDataFailureViewState(searchState); + } + } + + Future _refreshChangeListEmail() async { + log('ThreadController::_refreshChangeListEmail:'); + final refreshViewState = await _refreshChangesEmailsInMailboxInteractor.execute( + _session!, + _accountId!, + mailboxDashBoardController.currentEmailState!, + sort: EmailSortOrderType.mostRecent.getSortOrder().toNullable(), + propertiesCreated: EmailUtils.getPropertiesForEmailGetMethod( + _session!, + _accountId!, + ), + propertiesUpdated: ThreadConstants.propertiesUpdatedDefault, + emailFilter: EmailFilter( + filter: _getFilterCondition(mailboxIdSelected: selectedMailboxId), + filterOption: mailboxDashBoardController.filterMessageOption.value, + mailboxId: selectedMailboxId, + ), + ).last; + + final refreshState = refreshViewState + .foldSuccessWithResult(); + + if (refreshState is RefreshChangesAllEmailSuccess) { + _refreshChangesAllEmailSuccess(refreshState); + } else { + onDataFailureViewState(refreshState); } } @@ -1168,12 +1256,10 @@ class ThreadController extends BaseController with EmailActionController { } void goToCreateEmailRuleView() async { - final accountId = mailboxDashBoardController.accountId.value; - final session = mailboxDashBoardController.sessionCurrent; - if (accountId != null && session != null) { + if (_accountId != null && _session != null) { final arguments = RulesFilterCreatorArguments( - accountId, - session, + _accountId!, + _session!, mailboxDestination: selectedMailbox ); @@ -1182,7 +1268,7 @@ class ThreadController extends BaseController with EmailActionController { : await push(AppRoutes.rulesFilterCreator, arguments: arguments); if (newRuleFilterRequest is CreateNewEmailRuleFilterRequest) { - _createNewRuleFilterAction(accountId, newRuleFilterRequest); + _createNewRuleFilterAction(_accountId!, newRuleFilterRequest); } } else { logError('ThreadController::goToCreateEmailRuleView: Account or Session is NULL'); diff --git a/test/features/push_notification/presentation/websocket/web_socket_queue_handler_test.dart b/test/features/push_notification/presentation/websocket/web_socket_queue_handler_test.dart new file mode 100644 index 000000000..3f95a6ee8 --- /dev/null +++ b/test/features/push_notification/presentation/websocket/web_socket_queue_handler_test.dart @@ -0,0 +1,197 @@ +import 'dart:collection'; + +import 'package:flutter_test/flutter_test.dart'; +import 'package:jmap_dart_client/jmap/core/state.dart'; +import 'package:tmail_ui_user/features/push_notification/presentation/websocket/web_socket_message.dart'; +import 'package:tmail_ui_user/features/push_notification/presentation/websocket/web_socket_queue_handler.dart'; + +class MockWebSocketMessage extends WebSocketMessage { + MockWebSocketMessage(String message) + : super( + newState: State(message), + ); +} + +void main() { + group('WebSocketQueueHandler::test', () { + late Queue processedMessages; + + setUp(() { + processedMessages = Queue(); + }); + + WebSocketQueueHandler createHandler({ + required ProcessMessageCallback processMessageCallback, + OnErrorCallback? onErrorCallback, + }) { + return WebSocketQueueHandler( + processMessageCallback: processMessageCallback, + onErrorCallback: onErrorCallback, + ); + } + + group('Basic Operations', () { + late WebSocketQueueHandler handler; + + setUp(() { + handler = createHandler( + processMessageCallback: (message) async { + processedMessages.add(message.id); + }, + ); + }); + + tearDown(() => handler.dispose()); + + test('Should process messages in correct order', () async { + final messages = List.generate(5, (index) => MockWebSocketMessage('$index')); + + for (var message in messages) { + handler.enqueue(message); + } + + await handler.waitForEmpty(); + + expect(processedMessages, containsAllInOrder(['0', '1', '2', '3', '4'])); + }); + + test('Should correctly remove messages up to specified ID', () async { + final messages = List.generate(5, (index) => MockWebSocketMessage('$index')); + + for (var message in messages) { + handler.enqueue(message); + } + + handler.removeMessagesUpToCurrent('2'); + + expect(handler.queueSize, 2); + + await handler.waitForEmpty(); + + expect(processedMessages.length, 2); + expect(processedMessages.first, '3'); + }); + }); + + group('Concurrent Operations', () { + late WebSocketQueueHandler handler; + + setUp(() { + handler = createHandler( + processMessageCallback: (message) async { + processedMessages.add(message.id); + }, + ); + }); + + tearDown(() => handler.dispose()); + + test('Should handle concurrent message enqueueing', () async { + final messages = List.generate(5, (index) => MockWebSocketMessage('$index')); + + await Future.wait(messages.map((message) => Future(() => handler.enqueue(message)))); + + await handler.waitForEmpty(); + + expect(processedMessages, containsAllInOrder(['0', '1', '2', '3', '4'])); + }); + + test('Should maintain order under high concurrency', () async { + final messages = List.generate(10, (index) => MockWebSocketMessage('$index')); + + for (var message in messages) { + handler.enqueue(message); + } + + await handler.waitForEmpty(); + + expect(processedMessages, containsAllInOrder(['0', '1', '2', '3', '4', '5', '6', '7', '8', '9'])); + }); + }); + + group('Error Handling', () { + test('Should continue processing after message failure', () async { + final handler = createHandler( + processMessageCallback: (message) async { + if (message.id == '2') throw Exception('Simulated Failure'); + processedMessages.add(message.id); + }, + ); + + handler.enqueue(MockWebSocketMessage('1')); + handler.enqueue(MockWebSocketMessage('2')); + handler.enqueue(MockWebSocketMessage('3')); + + await handler.waitForEmpty(); + + expect(processedMessages, ['1', '3']); + handler.dispose(); + }); + + test('Should handle exception in process callback', () async { + final handler = createHandler( + processMessageCallback: (message) async { + if (message.id == '1') throw Exception('Simulated Failure'); + processedMessages.add(message.id); + }, + onErrorCallback: (error, stackTrace) { + expect(error, isA()); + }, + ); + + handler.enqueue(MockWebSocketMessage('1')); + handler.enqueue(MockWebSocketMessage('2')); + + await handler.waitForEmpty(); + + expect(processedMessages, ['2']); + handler.dispose(); + }); + }); + + group('Stress Testing', () { + late WebSocketQueueHandler handler; + + setUp(() { + handler = createHandler( + processMessageCallback: (message) async { + processedMessages.add(message.id); + }, + ); + }); + + tearDown(() => handler.dispose()); + + test('Should handle large bursts of messages', () async { + handler = createHandler( + processMessageCallback: (message) async { + processedMessages.add(message.id); + }, + ); + + const int burstSize = 1000; + for (int i = 0; i < burstSize; i++) { + await Future.delayed(const Duration(milliseconds: 10)); + handler.enqueue(MockWebSocketMessage('$i')); + } + + await handler.waitForEmpty(); + + expect(processedMessages.length, burstSize); + expect(processedMessages, List.generate(burstSize, (i) => '$i')); + }); + + test('Should handle interleaved slow and fast messages', () async { + handler.enqueue(MockWebSocketMessage('1')); + + Future.delayed(const Duration(milliseconds: 10), () { + handler.enqueue(MockWebSocketMessage('2')); + }); + + await handler.waitForEmpty(); + + expect(processedMessages, ['1', '2']); + }); + }); + }); +} diff --git a/test/features/thread/presentation/controller/thread_controller_test.dart b/test/features/thread/presentation/controller/thread_controller_test.dart index 7f8ad860d..e582c7db1 100644 --- a/test/features/thread/presentation/controller/thread_controller_test.dart +++ b/test/features/thread/presentation/controller/thread_controller_test.dart @@ -286,6 +286,7 @@ void main() { when(mockMailboxDashBoardController.emailUIAction).thenReturn(Rxn(null)); when(mockMailboxDashBoardController.viewState).thenReturn(Rx(Right(UIState.idle))); when(mockMailboxDashBoardController.filterMessageOption).thenReturn(Rx(FilterMessageOption.all)); + when(mockMailboxDashBoardController.currentEmailState).thenReturn(State('old-state')); when(mockSearchController.searchState).thenReturn(Rx(SearchState.initial())); when(mockSearchController.isAdvancedSearchViewOpen).thenReturn(RxBool(false)); when(mockSearchController.isSearchEmailRunning).thenReturn(true);