TF-3334 Add more concurrence test cases for WebSocketQueueHandler

This commit is contained in:
Dat PHAM HOANG
2024-12-23 22:41:36 +07:00
committed by Dat H. Pham
parent 01b320ce89
commit e386848af8
4 changed files with 311 additions and 22 deletions
@@ -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
@@ -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
@@ -567,8 +567,6 @@ class ThreadController extends BaseController with EmailActionController {
),
beforeOption: const None(),
);
searchController.activateSimpleSearch();
final searchViewState = await _searchEmailInteractor.execute(
_session!,
_accountId!,
@@ -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<String> 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<dynamic> 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 = <String>[];
late List<dynamic> 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<dynamic> 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<void>();
final processedIds = <String>[];
errors = [];
var processingStarted = Completer<void>();
// 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 = <String>[];
final processingStarted = Completer<void>();
final batchProcessing = Completer<void>();
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 = <String>[];
final processingDelay = Completer<void>();
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 = <Future>[];
// 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 {