diff --git a/Classes/Queue/AzureQueueStorage.php b/Classes/Queue/AzureQueueStorage.php index ff21986..7b2ffce 100644 --- a/Classes/Queue/AzureQueueStorage.php +++ b/Classes/Queue/AzureQueueStorage.php @@ -100,6 +100,11 @@ class AzureQueueStorage implements QueueInterface */ protected string $poisonSuffix = '-poison'; + /** + * Whether to preserve original payload in poison queue (if false, only metadata will be stored) + */ + protected bool $preservePoisonPayload = false; + /** * @Flow\Inject * @var AzureStorageClientFactory @@ -150,6 +155,9 @@ public function __construct(string $name, array $options = []) $this->validateQueueSuffix($options['poisonSuffix']); $this->poisonSuffix = $options['poisonSuffix']; } + if (isset($options['preservePoisonPayload'])) { + $this->preservePoisonPayload = (bool)$options['preservePoisonPayload']; + } $this->normalPriorityQueueName = $name; $this->priorityQueueName = $name . $this->prioritySuffix; @@ -327,6 +335,7 @@ public function waitAndReserve(?int $timeout = null): ?Message 'popReceipt' => $message->getPopReceipt(), 'blobName' => $message->getBlobName(), 'queueName' => $message->getQueueName(), + 'payload' => $message->getPayload(), ]; return $message; @@ -356,6 +365,13 @@ public function release(string $messageId, array $options = []): void unset($this->reservedMessages[$messageId]); } catch (Exception $e) { + if ($e->getCode() === 404) { + $this->systemLogger->warning('Message no longer available during release(), likely taken by another worker', [ + 'messageId' => $messageId, + ]); + unset($this->reservedMessages[$messageId]); + return; + } throw new JobQueueException('Failed to release message: ' . $e->getMessage(), 1234567897); } } @@ -373,35 +389,63 @@ public function abort(string $messageId): void try { // Push a failed record to the poison queue - if ($this->usePoisonQueue) { - $failedPayload = json_encode([ + $isProcessingPoisonQueue = $messageInfo['queueName'] === $this->poisonQueueName; + if ($this->usePoisonQueue && !$isProcessingPoisonQueue) { + $poisonPayload = [ 'messageId' => $messageId, 'queueMessageId' => $messageInfo['queueMessageId'], 'originalQueue' => $messageInfo['queueName'], - 'blobName' => $messageInfo['blobName'] ?? null, - 'timestamp' => time(), - ]); - $this->getQueueService()->createMessage($this->poisonQueueName, $failedPayload); + 'timestamp' => time(), + ]; + if ($this->preservePoisonPayload) { + $poisonPayload['payload'] = $messageInfo['payload']; + $serialized = json_encode($poisonPayload); + + if (strlen($serialized) > $this->claimCheckThreshold) { + unset($poisonPayload['payload']); + $poisonBlobName = $this->generateBlobName('poison-' . $messageId); + $this->getBlobService()->createBlockBlob( + $this->containerName, + $poisonBlobName, + json_encode($messageInfo['payload']) + ); + $poisonPayload['blobName'] = $poisonBlobName; + $poisonPayload['isClaimCheck'] = true; + } + } + $this->getQueueService()->createMessage($this->poisonQueueName, json_encode($poisonPayload)); } - // Delete the message from original queue - $this->getQueueService()->deleteMessage( - $messageInfo['queueName'], - $messageInfo['queueMessageId'], - $messageInfo['popReceipt'] - ); + if (!$isProcessingPoisonQueue) { + // Delete the message from original queue + $this->getQueueService()->deleteMessage( + $messageInfo['queueName'], + $messageInfo['queueMessageId'], + $messageInfo['popReceipt'] + ); - // Clean up blob if it exists - if (!empty($messageInfo['blobName'])) { - try { - $this->getBlobService()->deleteBlob($this->containerName, $messageInfo['blobName']); - } catch (Exception $e) { - // Log but don't throw - message is already aborted + // Clean up blob if it exists + if (!empty($messageInfo['blobName'])) { + try { + $this->getBlobService()->deleteBlob($this->containerName, $messageInfo['blobName']); + } catch (Exception $e) { + $this->systemLogger->warning('Failed to delete blob after message abort', [ + 'blobName' => $messageInfo['blobName'], + 'error' => $e->getMessage(), + ]); + } } } unset($this->reservedMessages[$messageId]); } catch (Exception $e) { + if ($e->getCode() === 404) { + $this->systemLogger->warning('Message no longer available during abort(), likely taken by another worker', [ + 'messageId' => $messageId, + ]); + unset($this->reservedMessages[$messageId]); + return; + } throw new JobQueueException('Failed to abort message: ' . $e->getMessage(), 1234567899); } } @@ -431,13 +475,23 @@ public function finish(string $messageId): bool try { $this->getBlobService()->deleteBlob($this->containerName, $messageInfo['blobName']); } catch (Exception $e) { - // Log but don't throw - message is already finished + $this->systemLogger->warning('Failed to delete blob after message finish', [ + 'blobName' => $messageInfo['blobName'], + 'error' => $e->getMessage(), + ]); } } unset($this->reservedMessages[$messageId]); return true; } catch (Exception $e) { + if ($e->getCode() === 404) { + $this->systemLogger->info('Message already deleted by another worker during finish', [ + 'messageId' => $messageId, + ]); + unset($this->reservedMessages[$messageId]); + return true; + } throw new JobQueueException('Failed to finish message: ' . $e->getMessage(), 1234567900); } } @@ -569,7 +623,7 @@ public function flush(): void $this->getQueueService()->clearMessages($this->priorityQueueName); } - // Clear poison queue only if priority queue feature is enabled + // Clear poison queue only if poison queue feature is enabled if ($this->usePoisonQueue) { $this->getQueueService()->clearMessages($this->poisonQueueName); } @@ -781,6 +835,8 @@ protected function extractPayload(string $messageText) throw new JobQueueException('Invalid JSON in blob: ' . json_last_error_msg(), 1234567912); } return $payload; + } catch (JobQueueException $e) { + throw $e; } catch (Exception $e) { throw new JobQueueException('Failed to retrieve claim check blob: ' . $e->getMessage(), 1234567904); } diff --git a/README.md b/README.md index 1d57057..f95d23e 100644 --- a/README.md +++ b/README.md @@ -24,7 +24,7 @@ Flowpack: JobQueue: Common: queues: - 'my-azure-storage-queuee': + 'my-azure-storage-queue': className: 'Oniva\JobQueue\AzureQueueStorage\Queue\AzureQueueStorage' options: connectionString: DefaultEndpointsProtocol=https;AccountName=myaccountname;AccountKey=myaccountkey;EndpointSuffix=core.windows.net @@ -50,6 +50,7 @@ Flowpack: blobContainer: jobqueue-blobs # Blob container name for claim check messages usePriorityQueue: true # Enable priority queueing usePoisonQueue: true # Enable poison queue for failed jobs + preservePoisonPayload: true # Preserve the original payload for failed jobs in the poison queue prioritySuffix: '-priority' # Suffix for priority queue poisonSuffix: '-poison' # Suffix for poison queue ``` @@ -62,9 +63,26 @@ This allows you to submit high-priority jobs that will be processed before regul $queue->submit($payload, ['priority' => true]); ``` +## Poison Queue + +To enable the poison queue, set `usePoisonQueue` to `true`. Failed jobs (those that exceed +`maximumNumberOfReleases`) are automatically moved to a dead-letter queue for inspection. + +By default, only metadata is stored in the poison queue. To preserve the original payload for +retry or debugging, also set `preservePoisonPayload: true`. + +To retry failed jobs, run a worker directly against the poison queue: +```bash +./flow flowpack.jobqueue.common:job:work my-azure-storage-queue-poison +``` + +Note: jobs being processed from the poison queue will not be re-poisoned on failure — they +are left visible in the queue for manual inspection instead. + ## Caveats * 7-day message limit - Azure Queue Storage automatically deletes messages after 7 days maximum, even if unprocessed * 64KB queue message size - While the claim check pattern handles larger payloads, it adds latency and blob storage costs * No native queue priorities - Priority queues are simulated by polling multiple queues, which increases API calls * Approximate counts only - Queue metrics like countReady() are estimates, not exact counts, due to Azure's distributed nature * No message ordering guarantee - Azure Queue Storage doesn't guarantee FIFO ordering, messages may be processed out of sequence +* Poison queue retry is manual - Failed jobs are moved to a dead-letter queue but not automatically retried. \ No newline at end of file diff --git a/Tests/Unit/Queue/AzureStorageQueueTest.php b/Tests/Unit/Queue/AzureStorageQueueTest.php index 1618ae7..fda070b 100644 --- a/Tests/Unit/Queue/AzureStorageQueueTest.php +++ b/Tests/Unit/Queue/AzureStorageQueueTest.php @@ -575,6 +575,31 @@ public function releaseThrowsExceptionForUnknownMessage(): void $queue->release('unknown-message'); } + /** + * @test + */ + public function releaseHandles404Gracefully(): void + { + $queue = $this->createQueue(); + $this->setupReservedMessage($queue, 'msg_123'); + + $this->queueService->expects($this->once()) + ->method('updateMessage') + ->willThrowException(new Exception('Not found', 404)); + + $this->logger->expects($this->once()) + ->method('warning') + ->with( + 'Message no longer available during release(), likely taken by another worker', + ['messageId' => 'msg_123'] + ); + + // Should not throw + $queue->release('msg_123'); + + $this->assertEquals(0, $queue->countReserved()); + } + /** * @test */ @@ -628,6 +653,136 @@ public function abortPushesToPoisonQueueIfEnabled(): void $queue->abort('msg_123'); } + + /** + * @test + */ + public function abortOnPoisonQueueSkipsDeleteAndRepoison(): void + { + $queue = $this->createQueue('test-queue', ['usePoisonQueue' => true]); + + // Manually set up a reserved message that came from the poison queue + $reflection = new ReflectionClass($queue); + $property = $reflection->getProperty('reservedMessages'); + $property->setAccessible(true); + $property->setValue($queue, [ + 'msg_123' => [ + 'queueMessageId' => 'azure-msg-id', + 'popReceipt' => 'pop-receipt', + 'blobName' => null, + 'queueName' => 'test-queue-poison', + 'payload' => ['test' => 'data'], + ], + ]); + + // Neither createMessage nor deleteMessage should be called + $this->queueService->expects($this->never())->method('createMessage'); + $this->queueService->expects($this->never())->method('deleteMessage'); + + $queue->abort('msg_123'); + + // Message should still be untracked locally after abort + $this->assertEquals(0, $queue->countReserved()); + } + + /** + * @test + */ + public function abortHandles404Gracefully(): void + { + $queue = $this->createQueue(); + $this->setupReservedMessage($queue, 'msg_123'); + + $this->queueService->expects($this->once()) + ->method('deleteMessage') + ->willThrowException(new Exception('Not found', 404)); + + $this->logger->expects($this->once()) + ->method('warning') + ->with( + 'Message no longer available during abort(), likely taken by another worker', + ['messageId' => 'msg_123'] + ); + + $queue->abort('msg_123'); + + $this->assertEquals(0, $queue->countReserved()); + } + + /** + * @test + */ + public function abortPreservesPoisonPayloadInline(): void + { + $queue = $this->createQueue('test-queue', [ + 'usePoisonQueue' => true, + 'preservePoisonPayload' => true, + ]); + + $this->setupReservedMessage($queue, 'msg_123'); + + // Manually inject payload into reserved message + $reflection = new ReflectionClass($queue); + $property = $reflection->getProperty('reservedMessages'); + $property->setAccessible(true); + $messages = $property->getValue($queue); + $messages['msg_123']['payload'] = ['important' => 'data']; + $property->setValue($queue, $messages); + + $this->queueService->expects($this->once()) + ->method('createMessage') + ->with( + 'test-queue-poison', + $this->callback(function ($envelope) { + $data = json_decode($envelope, true); + return isset($data['payload']) && $data['payload'] === ['important' => 'data']; + }) + ); + + $this->queueService->method('deleteMessage'); + + $queue->abort('msg_123'); + } + + /** + * @test + */ + public function abortPreservesPoisonPayloadViaClaimCheckForLargePayload(): void + { + $queue = $this->createQueue('test-queue', [ + 'usePoisonQueue' => true, + 'preservePoisonPayload' => true, + 'claimCheckThreshold' => 100, + ]); + + $this->setupReservedMessage($queue, 'msg_123'); + + $reflection = new ReflectionClass($queue); + $property = $reflection->getProperty('reservedMessages'); + $property->setAccessible(true); + $messages = $property->getValue($queue); + $messages['msg_123']['payload'] = ['data' => str_repeat('x', 200)]; + $property->setValue($queue, $messages); + + $this->blobService->expects($this->once()) + ->method('createBlockBlob') + ->with('jobqueue-blobs', $this->stringContains('poison-msg_123')); + + $this->queueService->expects($this->once()) + ->method('createMessage') + ->with( + 'test-queue-poison', + $this->callback(function ($envelope) { + $data = json_decode($envelope, true); + return ($data['isClaimCheck'] ?? false) === true && isset($data['blobName']); + }) + ); + + $this->queueService->method('deleteMessage'); + + $queue->abort('msg_123'); + } + /** * @test */ @@ -662,6 +817,28 @@ public function finishReturnsFalseForUnknownMessage(): void $this->assertFalse($result); } + /** + * @test + */ + public function finishReturnsTrueWhenMessageAlreadyDeletedByAnotherWorker(): void + { + $queue = $this->createQueue(); + $this->setupReservedMessage($queue, 'msg_123'); + + $this->queueService->expects($this->once()) + ->method('deleteMessage') + ->willThrowException(new Exception('Not found', 404)); + + $this->logger->expects($this->once()) + ->method('info') + ->with('Message already deleted by another worker during finish', ['messageId' => 'msg_123']); + + $result = $queue->finish('msg_123'); + + $this->assertTrue($result); + $this->assertEquals(0, $queue->countReserved()); + } + /** * @test */ @@ -932,19 +1109,6 @@ public function receiveMessageChecksHighPriorityFirst(): void $this->assertEquals(['priority' => 'high'], $message->getPayload()); } - /** - * Helper method to create a mock queue message - */ - protected function createMockQueueMessage(string $messageText, string $messageId, string $popReceipt): QueueMessage&MockObject - { - $queueMessage = $this->createMock(QueueMessage::class); - $queueMessage->method('getMessageText')->willReturn($messageText); - $queueMessage->method('getMessageId')->willReturn($messageId); - $queueMessage->method('getPopReceipt')->willReturn($popReceipt); - $queueMessage->method('getDequeueCount')->willReturn(1); - - return $queueMessage; - } /** * @test @@ -969,6 +1133,20 @@ public function dequeueSetsVisibilityTimeout(): void $queue->waitAndReserve(0); } + /** + * Helper method to create a mock queue message + */ + protected function createMockQueueMessage(string $messageText, string $messageId, string $popReceipt): QueueMessage&MockObject + { + $queueMessage = $this->createMock(QueueMessage::class); + $queueMessage->method('getMessageText')->willReturn($messageText); + $queueMessage->method('getMessageId')->willReturn($messageId); + $queueMessage->method('getPopReceipt')->willReturn($popReceipt); + $queueMessage->method('getDequeueCount')->willReturn(1); + + return $queueMessage; + } + /** * Helper method to setup a reserved message in the queue's internal tracking */ @@ -980,6 +1158,7 @@ protected function setupReservedMessage(AzureQueueStorage $queue, string $messag 'popReceipt' => 'pop-receipt', 'blobName' => $blobName, 'queueName' => 'test-queue', + 'payload' => ['test' => 'data'], ], ];