diff --git a/Core/Cleipnir.ResilientFunctions/Queuing/QueueManager.cs b/Core/Cleipnir.ResilientFunctions/Queuing/QueueManager.cs index eaf8b7bc..cc7efc57 100644 --- a/Core/Cleipnir.ResilientFunctions/Queuing/QueueManager.cs +++ b/Core/Cleipnir.ResilientFunctions/Queuing/QueueManager.cs @@ -136,7 +136,7 @@ public async Task Push(IReadOnlyList messages) if (_thrownException != null) return; - await ProcessMessages(messages); + ProcessMessages(messages); } finally { @@ -200,7 +200,7 @@ private async Task FetchAndNotify() skipPositions = _fetchedPositions.ToList(); var messages = await _messageStore.GetMessages(_storedId, skipPositions); - await ProcessMessages(messages); + ProcessMessages(messages); } finally { @@ -211,7 +211,7 @@ private async Task FetchAndNotify() // Caller must hold _fetchSemaphore. Deserializes, dedups by idempotency-key and by already-fetched // position (so pushes are idempotent), and stages messages for delivery. - private async Task ProcessMessages(IReadOnlyList messages) + private void ProcessMessages(IReadOnlyList messages) { foreach (var (messageContent, messageType, position, _, idempotencyKey, sender, receiver) in messages) {