From 09e352ec65d4ff547252b597e8fa27fee254bf8e Mon Sep 17 00:00:00 2001 From: yaroslavtykhonchuk Date: Tue, 18 Nov 2025 16:48:11 +0200 Subject: [PATCH] Copy content to new blob on quarantine MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit We had a race condition during quarantine, that could create duplicate messages in quarantine with a reference to one blob. Problem Analysis Failure Scenario: 1. Message with BlobRef quarantined → quarantine queue has copy 2. DeleteMessage fails (service crash/network issue) → message stays in main queue 3. Message becomes visible again → dequeued with DequeueCount++ 4. Same message quarantined again → 2 identical messages in quarantine with same blob ref 5. First quarantine message processed → blob deleted 6. Second quarantine message processed → blob not found → poison message Solution: create new blob on quarantine and copy its content there --- AzureBatchQueue/MessageQueue.cs | 38 +++++++++++++++++++++++++++++++-- 1 file changed, 36 insertions(+), 2 deletions(-) diff --git a/AzureBatchQueue/MessageQueue.cs b/AzureBatchQueue/MessageQueue.cs index 8845828..b7ed8f4 100644 --- a/AzureBatchQueue/MessageQueue.cs +++ b/AzureBatchQueue/MessageQueue.cs @@ -112,8 +112,22 @@ async Task QuarantineMessage(QueueMessage queueMessage, CancellationToken ct = d { try { - await quarantineQueue.SendMessageAsync(queueMessage.Body, cancellationToken: ct); - await queue.DeleteMessageAsync(queueMessage.MessageId, queueMessage.PopReceipt, ct); + BinaryData messageBody; + + if (IsBlobRef(queueMessage.Body, out var blobRef) && blobRef != null) + { + var newBlobRef = await CopyBlobForQuarantine(blobRef, queueMessage.MessageId, ct); + messageBody = Payload.OffloadedToBlob(newBlobRef).Data; + } + else + { + messageBody = queueMessage.Body; + } + + await quarantineQueue.SendMessageAsync(messageBody, cancellationToken: ct); + + var messageId = new MessageId(queueMessage.MessageId, queueMessage.PopReceipt, blobRef?.BlobName); + await DeleteMessage(messageId, ct); logger.LogInformation("QueueMessage {msgId} with {popReceipt} was quarantined after {dequeueCount} unsuccessful attempts.", queueMessage.MessageId, queueMessage.PopReceipt, queueMessage.DequeueCount); @@ -126,6 +140,26 @@ async Task QuarantineMessage(QueueMessage queueMessage, CancellationToken ct = d } } + // Why not just move the message from main queue to quarantine queue? + // Because on message quarantine, we can put the message into quarantine queue, and then the service stops and we don't delete it from the main queue. + // In this case, later we will insert a duplicate message to quarantine, with the same ref to blob. + // When one message is processed, we remove the blob. The other message will have a BlobNotFound error. + // So it's better to do some more heavy-lifting and copy blob content on quarantine. + async Task CopyBlobForQuarantine(BlobRef sourceBlobRef, string messageId, CancellationToken ct) + { + var sourceBlobClient = container.GetBlobClient(sourceBlobRef.BlobName); + var newBlobRef = BlobRef.Create(); + var destinationBlobClient = container.GetBlobClient(newBlobRef.BlobName); + + var content = await sourceBlobClient.DownloadContentAsync(ct); + await destinationBlobClient.UploadAsync(content.Value.Content, ct); + + logger.LogInformation("Copied blob {SourceBlob} to {DestBlob} for quarantine message {MsgId}", + sourceBlobRef.BlobName, newBlobRef.BlobName, messageId); + + return newBlobRef; + } + public async Task Dequarantine(CancellationToken ct = default) { do