From e386848af8712d58b5bbdbb343749c839d7f9c7c Mon Sep 17 00:00:00 2001 From: Dat PHAM HOANG Date: Mon, 23 Dec 2024 22:41:36 +0700 Subject: [PATCH] TF-3334 Add more concurrence test cases for WebSocketQueueHandler --- .../mailbox_dashboard_controller.dart | 4 - .../websocket/web_socket_queue_handler.dart | 51 +++- .../presentation/thread_controller.dart | 2 - .../web_socket_queue_handler_test.dart | 276 ++++++++++++++++++ 4 files changed, 311 insertions(+), 22 deletions(-) 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 604724c94..a288c35ad 100644 --- a/lib/features/mailbox_dashboard/presentation/controller/mailbox_dashboard_controller.dart +++ b/lib/features/mailbox_dashboard/presentation/controller/mailbox_dashboard_controller.dart @@ -1465,10 +1465,6 @@ class MailboxDashBoardController extends ReloadableController with UserSettingPo void dispatchRoute(DashboardRoutes route) { log('MailboxDashBoardController::dispatchRoute(): $route'); dashboardRoute.value = route; - - if (dashboardRoute.value == DashboardRoutes.searchEmail) { - searchController.activateSimpleSearch(); - } } @override 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 index 455473983..34855be4a 100644 --- a/lib/features/push_notification/presentation/websocket/web_socket_queue_handler.dart +++ b/lib/features/push_notification/presentation/websocket/web_socket_queue_handler.dart @@ -38,11 +38,16 @@ class WebSocketQueueHandler { return; } - if (queueSize >= _maxQueueSize) { - log('WebSocketQueueHandler::enqueue:Queue full, removing oldest message'); - _messageQueue.removeFirst(); + try { + if (queueSize >= _maxQueueSize) { + log('WebSocketQueueHandler::enqueue:Queue full, removing oldest message'); + _messageQueue.removeFirst(); + } + } catch (e) { + logError('WebSocketQueueHandler::enqueue:Exception = $e'); } + log('WebSocketQueueHandler::enqueue(): ${message.id}'); _messageQueue.add(message); _queueController.add(message); } @@ -57,6 +62,7 @@ class WebSocketQueueHandler { try { while (queueSize > 0) { final message = _messageQueue.removeFirst(); + log('WebSocketQueueHandler::_processQueue(): processing message ${message.id}'); try { await processMessageCallback(message); @@ -78,27 +84,40 @@ class WebSocketQueueHandler { } void _addToProcessedMessages(String messageId) { - if (_processedMessageIds.length >= _maxProcessedIdsSize) { - _processedMessageIds.removeFirst(); + log('WebSocketQueueHandler::_addToProcessedMessages(): adding message $messageId to processed messages'); + try { + if (_processedMessageIds.length >= _maxProcessedIdsSize) { + _processedMessageIds.removeFirst(); + } + } catch (e) { + logError('WebSocketQueueHandler::_addToProcessedMessages:Exception = $e'); } + _processedMessageIds.add(messageId); } void removeMessagesUpToCurrent(String messageId) { - final isCurrentStateExist = _messageQueue - .any((message) => message.id == messageId); + try { + log('WebSocketQueueHandler::removeMessagesUpToCurrent(): removing messages up to $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; + if (!isCurrentStateExist) { + log('WebSocketQueueHandler::removeMessagesUpToCurrent:Current state $messageId not found in the queue.'); + return; } + + while (queueSize > 0) { + final removedMessage = _messageQueue.removeFirst(); + log('WebSocketQueueHandler::removeMessagesUpToCurrent(): removing message ${removedMessage.id} up to $messageId'); + if (removedMessage.id == messageId) { + break; + } + } + log('WebSocketQueueHandler::removeMessagesUpToCurrent:Updated Queue: $queueSize'); + } catch (e) { + logError('WebSocketQueueHandler::removeMessagesUpToCurrent:Exception = $e'); } - log('WebSocketQueueHandler::removeMessagesUpToCurrent:Updated Queue: $queueSize'); } @visibleForTesting diff --git a/lib/features/thread/presentation/thread_controller.dart b/lib/features/thread/presentation/thread_controller.dart index 278a37dca..fd0e31026 100644 --- a/lib/features/thread/presentation/thread_controller.dart +++ b/lib/features/thread/presentation/thread_controller.dart @@ -567,8 +567,6 @@ class ThreadController extends BaseController with EmailActionController { ), beforeOption: const None(), ); - searchController.activateSimpleSearch(); - final searchViewState = await _searchEmailInteractor.execute( _session!, _accountId!, 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 index 3f95a6ee8..4a72deb95 100644 --- 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 @@ -1,3 +1,4 @@ +import 'dart:async'; import 'dart:collection'; import 'package:flutter_test/flutter_test.dart'; @@ -32,6 +33,15 @@ void main() { group('Basic Operations', () { late WebSocketQueueHandler handler; + late List processedMessages; + + setUp(() { + processedMessages = []; + }); + + tearDown(() { + handler.dispose(); + }); setUp(() { handler = createHandler( @@ -55,6 +65,30 @@ void main() { expect(processedMessages, containsAllInOrder(['0', '1', '2', '3', '4'])); }); + test('Duplicate messages should be skipped', () async { + final message = MockWebSocketMessage('duplicate_msg'); + + handler.enqueue(message); + await handler.waitForEmpty(); + + handler.enqueue(message); + await handler.waitForEmpty(); + + expect(processedMessages.length, equals(1)); + expect(processedMessages, equals(['duplicate_msg'])); + }); + + test('Queue size should not exceed maximum size', () async { + // Enqueue more messages than the max queue size + final messages = List.generate(130, (i) => MockWebSocketMessage('msg_$i')); + + for (var message in messages) { + handler.enqueue(message); + } + + expect(handler.queueSize, lessThanOrEqualTo(128)); + }); + test('Should correctly remove messages up to specified ID', () async { final messages = List.generate(5, (index) => MockWebSocketMessage('$index')); @@ -73,8 +107,86 @@ void main() { }); }); + group('Queue Size Management Tests', () { + test('Queue should drop oldest message when full', () async { + late List errors = []; + + final handler = WebSocketQueueHandler( + processMessageCallback: (message) async { + await Future.delayed(const Duration(milliseconds: 10)); // Simulate processing time + processedMessages.add(message.id); + }, + onErrorCallback: (error, stackTrace) { + errors.add(error); + }, + ); + + // Fill the queue to maximum capacity (128) + for (var i = 0; i < 128; i++) { + handler.enqueue(MockWebSocketMessage('msg_$i')); + } + + expect(handler.queueSize, equals(128)); + + // Add one more message + handler.enqueue(MockWebSocketMessage('msg_128')); + + // Queue size should still be 128 + expect(handler.queueSize, equals(128)); + + // Process all messages + await handler.waitForEmpty(); + + // Verify that msg_0 (the oldest) was dropped and msg_128 (newest) was processed + expect(processedMessages.contains('msg_0'), isFalse); + expect(processedMessages.contains('msg_1'), isTrue); + expect(processedMessages.contains('msg_127'), isTrue); + expect(processedMessages.contains('msg_128'), isTrue); + }); + + test('Queue should maintain size limit during message removal', () async { + final processedIds = []; + late List errors = []; + + final handler = WebSocketQueueHandler( + processMessageCallback: (message) async { + await Future.delayed(const Duration(milliseconds: 10)); + processedIds.add(message.id); + }, + onErrorCallback: (error, stackTrace) { + errors.add(error); + }, + ); + + // Fill queue to capacity + for (var i = 0; i < 128; i++) { + handler.enqueue(MockWebSocketMessage('msg_$i')); + } + + // Remove messages up to msg_64 + handler.removeMessagesUpToCurrent('msg_64'); + + // Add new messages to fill the queue again + for (var i = 128; i < 192; i++) { + handler.enqueue(MockWebSocketMessage('msg_$i')); + } + + expect(handler.queueSize, lessThanOrEqualTo(128), + reason: 'Queue size should not exceed maximum after removal and refill'); + + await handler.waitForEmpty(); + + // Verify that earlier messages were properly removed + for (var i = 0; i <= 64; i++) { + expect(processedMessages.contains('msg_$i'), isFalse, + reason: 'Message msg_$i should have been removed'); + } + }); + }); + group('Concurrent Operations', () { late WebSocketQueueHandler handler; + late List errors; setUp(() { handler = createHandler( @@ -107,6 +219,169 @@ void main() { expect(processedMessages, containsAllInOrder(['0', '1', '2', '3', '4', '5', '6', '7', '8', '9'])); }); + + test('Should handle concurrent enqueueing while processing is blocked', () async { + final processingCompleter = Completer(); + final processedIds = []; + errors = []; + var processingStarted = Completer(); + + // Create handler with a processing delay to simulate long-running task + handler = WebSocketQueueHandler( + processMessageCallback: (message) async { + if (!processingStarted.isCompleted) { + processingStarted.complete(); + } + await processingCompleter.future; // Block processing + processedIds.add(message.id); + }, + onErrorCallback: (error, stackTrace) { + errors.add(error); + }, + ); + + // Enqueue first message to start processing + handler.enqueue(MockWebSocketMessage('initial_msg')); + + // Wait for processing to start + await processingStarted.future; + + // Concurrently enqueue messages while first message is still processing + await Future.wait( + List.generate(150, (i) => Future(() { + handler.enqueue(MockWebSocketMessage('concurrent_$i')); + })) + ); + + // Verify queue size is capped at max while processing is blocked + expect(handler.queueSize, lessThanOrEqualTo(128)); + + // Allow processing to continue + processingCompleter.complete(); + + await handler.waitForEmpty(); + // Verify process order and dropped messages + expect(processedIds[0], equals('initial_msg')); + expect(processedIds.length, lessThanOrEqualTo(129)); + expect(handler.isMessageProcessed('concurrent_149'), isTrue); + }); + + test('Should handle rapid enqueueing during active processing', () async { + final processedIds = []; + final processingStarted = Completer(); + final batchProcessing = Completer(); + errors = []; + + // Create handler with controlled processing delays + handler = WebSocketQueueHandler( + processMessageCallback: (message) async { + if (!processingStarted.isCompleted) { + processingStarted.complete(); + await batchProcessing.future; + } + processedIds.add(message.id); + }, + onErrorCallback: (error, stackTrace) { + errors.add(error); + }, + ); + + // Start with initial batch + for (var i = 0; i < 50; i++) { + handler.enqueue(MockWebSocketMessage('batch1_$i')); + } + + // Wait for the first message to start processing + await processingStarted.future; + + // Add second batch while first batch is blocked + for (var i = 0; i < 50; i++) { + handler.enqueue(MockWebSocketMessage('batch2_$i')); + } + + // Add third batch immediately + for (var i = 0; i < 50; i++) { + handler.enqueue(MockWebSocketMessage('batch3_$i')); + } + + // Allow processing to continue + batchProcessing.complete(); + + // Wait for queue to be empty + await handler.waitForEmpty(); + + // Verify results + expect(processedIds.length, lessThanOrEqualTo(129), + reason: 'Total processed messages should not exceed queue capacity'); + + // Check if we have messages from the latest batch + final lastBatchCount = processedIds + .where((id) => id.startsWith('batch3_')) + .length; + expect(lastBatchCount, greaterThan(0), + reason: 'Should have processed some messages from the latest batch'); + + // Verify that some early messages were dropped + final firstBatchCount = processedIds + .where((id) => id.startsWith('batch1_')) + .length; + expect(firstBatchCount, lessThan(50), + reason: 'Some messages from first batch should have been dropped'); + }); + + test('Should handle concurrent removeMessagesUpToCurrent during processing', () async { + final processedIds = []; + final processingDelay = Completer(); + errors = []; + + handler = WebSocketQueueHandler( + processMessageCallback: (message) async { + await processingDelay.future; + processedIds.add(message.id); + }, + onErrorCallback: (error, stackTrace) { + errors.add(error); + }, + ); + + // Fill queue + for (var i = 0; i < 128; i++) { + handler.enqueue(MockWebSocketMessage('msg_$i')); + } + + // Start concurrent operations + final futures = []; + + // Add new messages + futures.add(Future(() async { + for (var i = 128; i < 256; i++) { + handler.enqueue(MockWebSocketMessage('msg_$i')); + await Future.delayed(const Duration(microseconds: 100)); + } + })); + + // Concurrently remove messages + futures.add(Future(() async { + await Future.delayed(const Duration(milliseconds: 10)); + handler.removeMessagesUpToCurrent('msg_64'); + })); + + // Allow processing to continue after concurrent operations + await Future.delayed(const Duration(milliseconds: 50)); + processingDelay.complete(); + + await Future.wait(futures); + await handler.waitForEmpty(); + + expect(processedIds.length, lessThanOrEqualTo(192)); + + // Verify that messages after removal point were processed + for (var id in processedIds.skip(1)) { + final messageNumber = int.parse(id.split('_')[1]); + expect(messageNumber, greaterThan(64), + reason: 'Only messages after msg_64 should be processed'); + } + }); }); group('Error Handling', () { @@ -179,6 +454,7 @@ void main() { expect(processedMessages.length, burstSize); expect(processedMessages, List.generate(burstSize, (i) => '$i')); + expect(handler.queueSize, equals(0)); }); test('Should handle interleaved slow and fast messages', () async {