diff --git a/config/Migrations/20260922163025_AddMetadataToFailedJobs.php b/config/Migrations/20260922163025_AddMetadataToFailedJobs.php new file mode 100644 index 0000000..3638b0f --- /dev/null +++ b/config/Migrations/20260922163025_AddMetadataToFailedJobs.php @@ -0,0 +1,25 @@ +table('queue_failed_jobs'); + $table->addColumn('metadata', 'text', [ + 'null' => true, + 'default' => null, + ]) + ->update(); + } +} diff --git a/src/Command/RequeueCommand.php b/src/Command/RequeueCommand.php index ce00f7c..a4122e0 100644 --- a/src/Command/RequeueCommand.php +++ b/src/Command/RequeueCommand.php @@ -138,6 +138,7 @@ public function execute(Arguments $args, ConsoleIo $io): int 'config' => $failedJob->config, 'priority' => $failedJob->priority, 'queue' => $failedJob->queue, + 'metadata' => $failedJob->decoded_metadata ?? [], ], ); diff --git a/src/Job/Message.php b/src/Job/Message.php index 521351d..b4fd15f 100644 --- a/src/Job/Message.php +++ b/src/Job/Message.php @@ -204,6 +204,21 @@ public function getDto(string $dtoClass): object return $dto; } + /** + * Get the envelope metadata recorded on the message body at dispatch time. + * + * Metadata carries bookkeeping fields that belong next to the payload, + * not inside it, so `data` stays pure for DTO hydration. + * + * @return array + */ + public function getMetadata(): array + { + $metadata = $this->parsedBody['metadata'] ?? []; + + return is_array($metadata) ? $metadata : []; + } + /** * The maximum number of attempts allowed by the job. */ diff --git a/src/Listener/FailedJobsListener.php b/src/Listener/FailedJobsListener.php index c983fb4..e9e98d8 100644 --- a/src/Listener/FailedJobsListener.php +++ b/src/Listener/FailedJobsListener.php @@ -67,6 +67,9 @@ public function storeFailedJob(object $event): void 'class' => $class, 'method' => $method, 'data' => json_encode($data), + 'metadata' => isset($originalMessageBody['metadata']) + ? (string)json_encode($originalMessageBody['metadata']) + : null, 'config' => $requeueOptions['config'], 'priority' => $requeueOptions['priority'], 'queue' => $requeueOptions['queue'], diff --git a/src/Model/Entity/FailedJob.php b/src/Model/Entity/FailedJob.php index 62f1fc4..164342d 100644 --- a/src/Model/Entity/FailedJob.php +++ b/src/Model/Entity/FailedJob.php @@ -12,6 +12,7 @@ * @property string $class * @property string $method * @property string $data + * @property string|null $metadata * @property string|null $config * @property string|null $priority * @property string|null $queue @@ -37,6 +38,7 @@ class FailedJob extends Entity 'class' => true, 'method' => true, 'data' => true, + 'metadata' => true, 'config' => true, 'priority' => true, 'queue' => true, @@ -52,4 +54,22 @@ protected function _getDecodedData(): array { return json_decode($this->data, true); } + + /** + * Envelope metadata as an array. Empty when the job was stored without + * a metadata envelope. + * + * @see \Cake\Queue\Model\Entity\FailedJob::$decoded_metadata + * @return array + */ + protected function _getDecodedMetadata(): array + { + if (empty($this->metadata)) { + return []; + } + + $decoded = json_decode((string)$this->metadata, true); + + return is_array($decoded) ? $decoded : []; + } } diff --git a/src/Model/Table/FailedJobsTable.php b/src/Model/Table/FailedJobsTable.php index 208b325..63d6df5 100644 --- a/src/Model/Table/FailedJobsTable.php +++ b/src/Model/Table/FailedJobsTable.php @@ -72,6 +72,10 @@ public function validationDefault(Validator $validator): Validator ->requirePresence('data', 'create') ->notEmptyString('data'); + $validator + ->scalar('metadata') + ->allowEmptyString('metadata'); + $validator ->scalar('config') ->maxLength('config', 255) diff --git a/src/QueueManager.php b/src/QueueManager.php index ff8fafd..6111c6c 100644 --- a/src/QueueManager.php +++ b/src/QueueManager.php @@ -223,6 +223,10 @@ public static function engine(string $name): SimpleClient * - `expires` - Time (in integer seconds) after which the message expires. * The message will be removed from the queue if this time is exceeded * and it has not been consumed. Default `null`. + * - `metadata` - Optional envelope data recorded on the message body next + * to `data` (e.g. `tags`, `_uniqueId`, `batch_id`). Unlike `data` it is + * never passed to the job or hydrated into a DTO. Omitted from the body + * when empty or not an array. Default `[]`. * - `priority` - Valid values: * - `\Enqueue\Client\MessagePriority::VERY_LOW` * - `\Enqueue\Client\MessagePriority::LOW` @@ -297,6 +301,11 @@ public static function push(string|array $className, array|object $data = [], ar $body['dtoClass'] = $dtoClass; } + $metadata = $options['metadata'] ?? null; + if (is_array($metadata) && $metadata !== []) { + $body['metadata'] = $metadata; + } + $message = new ClientMessage($body); if (isset($options['delay'])) { diff --git a/tests/Fixture/FailedJobsFixture.php b/tests/Fixture/FailedJobsFixture.php index c44743e..09dc47d 100644 --- a/tests/Fixture/FailedJobsFixture.php +++ b/tests/Fixture/FailedJobsFixture.php @@ -27,6 +27,7 @@ public function init(): void 'class' => LogToDebugJob::class, 'method' => 'execute', 'data' => '{"sample_data_1": "sample value", "sample_data_2": 1}', + 'metadata' => null, 'config' => 'default', 'priority' => null, 'queue' => 'default', @@ -38,6 +39,7 @@ public function init(): void 'class' => MaxAttemptsIsThreeJob::class, 'method' => 'execute', 'data' => '{"sample_data_1": "sample value", "sample_data_2": 1}', + 'metadata' => null, 'config' => 'default', 'priority' => null, 'queue' => 'default', @@ -49,6 +51,7 @@ public function init(): void 'class' => LogToDebugJob::class, 'method' => 'execute', 'data' => '{"sample_data_1": "sample value", "sample_data_2": 1}', + 'metadata' => null, 'config' => 'alternate_config', 'priority' => null, 'queue' => 'alternate_queue', diff --git a/tests/TestCase/Command/RequeueCommandTest.php b/tests/TestCase/Command/RequeueCommandTest.php index 5887b51..6dc856f 100644 --- a/tests/TestCase/Command/RequeueCommandTest.php +++ b/tests/TestCase/Command/RequeueCommandTest.php @@ -197,4 +197,44 @@ public function testJobsAreRequeuedByConfig() $this->assertDebugLogContains('Debug job was run'); } + + public function testRequeuedJobKeepsMetadata() + { + $fsQueuePath = TMP . DS . uniqid('queue'); + QueueManager::setConfig('default', [ + 'url' => 'file:///' . $fsQueuePath, + 'queue' => 'default', + ]); + + /** @var \Cake\Queue\Model\Table\FailedJobsTable $failedJobsTable */ + $failedJobsTable = $this->getTableLocator()->get('Cake/Queue.FailedJobs'); + $failedJobsTable->deleteAll(['1=1']); + + $failedJob = $failedJobsTable->newEntity([ + 'class' => LogToDebugJob::class, + 'method' => 'execute', + 'data' => json_encode(['example_key' => 'example_value']), + 'metadata' => json_encode(['tags' => ['finance'], '_uniqueId' => 'abc123']), + 'config' => 'default', + 'priority' => null, + 'queue' => 'default', + 'exception' => 'boom', + ]); + $failedJobsTable->saveOrFail($failedJob); + + $this->exec('queue requeue -f'); + + $this->assertOutputContains('Requeueing 1 jobs.'); + $this->assertOutputContains('1 jobs requeued.'); + + $fsQueueFile = $fsQueuePath . DS . 'enqueue.app.default'; + $this->assertFileExists($fsQueueFile); + + $contents = (string)file_get_contents($fsQueueFile); + $this->assertStringContainsString('metadata', $contents); + $this->assertStringContainsString('finance', $contents); + $this->assertStringContainsString('abc123', $contents); + + unlink($fsQueueFile); + } } diff --git a/tests/TestCase/Job/MessageTest.php b/tests/TestCase/Job/MessageTest.php index d8c384e..3d8a6fc 100644 --- a/tests/TestCase/Job/MessageTest.php +++ b/tests/TestCase/Job/MessageTest.php @@ -248,6 +248,48 @@ public function testGetDtoThrowsForMissingExpectedClass() $message->getDto('TestApp\Dto\DoesNotExist'); } + /** + * Test that envelope metadata is exposed separately from the payload. + * + * @return void + */ + public function testGetMetadata() + { + $parsedBody = [ + 'class' => [WelcomeMailer::class, 'welcome'], + 'data' => ['id' => 7], + 'metadata' => [ + 'tags' => ['finance', 'orders'], + '_uniqueId' => 'abc123', + ], + ]; + $connectionFactory = new NullConnectionFactory(); + $context = $connectionFactory->createContext(); + $originalMessage = new NullMessage((string)json_encode($parsedBody)); + $message = new Message($originalMessage, $context); + + $this->assertSame($parsedBody['metadata'], $message->getMetadata()); + // The payload stays pure: no envelope keys leak into the job data. + $this->assertSame(['id' => 7], $message->getArgument()); + } + + /** + * Test that missing metadata defaults to an empty array. + * + * @return void + */ + public function testGetMetadataDefaultsToEmpty() + { + $connectionFactory = new NullConnectionFactory(); + $context = $connectionFactory->createContext(); + + $plain = new Message(new NullMessage((string)json_encode([ + 'class' => [WelcomeMailer::class, 'welcome'], + 'data' => ['id' => 7], + ])), $context); + $this->assertSame([], $plain->getMetadata()); + } + /** * Test that invalid classes cannot be made into callables. * diff --git a/tests/TestCase/Listener/FailedJobsListenerTest.php b/tests/TestCase/Listener/FailedJobsListenerTest.php index d5a56e9..a1c9bc2 100644 --- a/tests/TestCase/Listener/FailedJobsListenerTest.php +++ b/tests/TestCase/Listener/FailedJobsListenerTest.php @@ -100,6 +100,84 @@ public function testFailedJobIsAddedWhenEventIsFired() $this->assertStringContainsString('some message', $failedJob->exception); } + public function testFailedJobPreservesMetadata() + { + $parsedBody = [ + 'class' => [LogToDebugJob::class, 'execute'], + 'data' => ['example_key' => 'example_value'], + 'metadata' => ['tags' => ['finance'], '_uniqueId' => 'abc123'], + 'requeueOptions' => [ + 'config' => 'example_config', + 'priority' => 'example_priority', + 'queue' => 'example_queue', + ], + ]; + $messageBody = json_encode($parsedBody); + $connectionFactory = new NullConnectionFactory(); + + $context = $connectionFactory->createContext(); + $originalMessage = new NullMessage($messageBody); + $message = new Message($originalMessage, $context); + + $event = new Event( + 'Consumption.LimitAttemptsExtension.failed', + $message, + ['exception' => 'some message'], + ); + + /** @var \Cake\Queue\Model\Table\FailedJobsTable $failedJobsTable */ + $failedJobsTable = $this->getTableLocator()->get('Cake/Queue.FailedJobs'); + $failedJobsTable->deleteAll(['1=1']); + + EventManager::instance()->on(new FailedJobsListener()); + EventManager::instance()->dispatch($event); + + $this->assertSame(1, $failedJobsTable->find()->count()); + + $failedJob = $failedJobsTable->find()->first(); + + $this->assertSame($parsedBody['metadata'], $failedJob->decoded_metadata); + } + + public function testFailedJobWithoutMetadataStoresNull() + { + $parsedBody = [ + 'class' => [LogToDebugJob::class, 'execute'], + 'data' => ['example_key' => 'example_value'], + 'requeueOptions' => [ + 'config' => 'example_config', + 'priority' => 'example_priority', + 'queue' => 'example_queue', + ], + ]; + $messageBody = json_encode($parsedBody); + $connectionFactory = new NullConnectionFactory(); + + $context = $connectionFactory->createContext(); + $originalMessage = new NullMessage($messageBody); + $message = new Message($originalMessage, $context); + + $event = new Event( + 'Consumption.LimitAttemptsExtension.failed', + $message, + ['exception' => 'some message'], + ); + + /** @var \Cake\Queue\Model\Table\FailedJobsTable $failedJobsTable */ + $failedJobsTable = $this->getTableLocator()->get('Cake/Queue.FailedJobs'); + $failedJobsTable->deleteAll(['1=1']); + + EventManager::instance()->on(new FailedJobsListener()); + EventManager::instance()->dispatch($event); + + $this->assertSame(1, $failedJobsTable->find()->count()); + + $failedJob = $failedJobsTable->find()->first(); + + $this->assertNull($failedJob->metadata); + $this->assertSame([], $failedJob->decoded_metadata); + } + /** * Data provider for testStoreFailedJobException * diff --git a/tests/TestCase/QueueManagerTest.php b/tests/TestCase/QueueManagerTest.php index 8999eb3..5ff9867 100644 --- a/tests/TestCase/QueueManagerTest.php +++ b/tests/TestCase/QueueManagerTest.php @@ -304,6 +304,41 @@ public function testPushWithoutDtoDoesNotAddDtoClass() $this->assertStringNotContainsString('dtoClass', $contents); } + public function testPushWithMetadata() + { + QueueManager::setConfig('test', [ + 'url' => $this->getFsQueueUrl(), + 'queue' => 'test', + ]); + + QueueManager::push(LogToDebugJob::class, ['id' => 7], [ + 'config' => 'test', + 'metadata' => ['tags' => ['finance'], '_uniqueId' => 'abc123'], + ]); + + $fsQueueFile = $this->getFsQueueUrl() . DS . 'enqueue.app.test'; + $this->assertFileExists($fsQueueFile); + $contents = file_get_contents($fsQueueFile); + $this->assertStringContainsString('metadata', $contents); + $this->assertStringContainsString('finance', $contents); + $this->assertStringContainsString('abc123', $contents); + } + + public function testPushWithoutMetadataDoesNotAddMetadata() + { + QueueManager::setConfig('test', [ + 'url' => $this->getFsQueueUrl(), + 'queue' => 'test', + ]); + + QueueManager::push(LogToDebugJob::class, ['id' => 7], ['config' => 'test']); + + $fsQueueFile = $this->getFsQueueUrl() . DS . 'enqueue.app.test'; + $this->assertFileExists($fsQueueFile); + $contents = file_get_contents($fsQueueFile); + $this->assertStringNotContainsString('metadata', $contents); + } + public function testUniqueMessageIsQueuedOnlyOnce() { QueueManager::setConfig('test', [ diff --git a/tests/schema.php b/tests/schema.php index 3aeba16..1dfba5f 100644 --- a/tests/schema.php +++ b/tests/schema.php @@ -9,6 +9,7 @@ 'class' => ['type' => 'string', 'length' => 255, 'null' => false, 'default' => null, 'comment' => '', 'precision' => null], 'method' => ['type' => 'string', 'length' => 255, 'null' => false, 'default' => null, 'comment' => '', 'precision' => null], 'data' => ['type' => 'text', 'length' => null, 'null' => false, 'default' => null, 'comment' => '', 'precision' => null], + 'metadata' => ['type' => 'text', 'length' => null, 'null' => true, 'default' => null, 'comment' => '', 'precision' => null], 'config' => ['type' => 'string', 'length' => 255, 'null' => true, 'default' => null, 'comment' => '', 'precision' => null], 'priority' => ['type' => 'string', 'length' => 255, 'null' => true, 'default' => null, 'comment' => '', 'precision' => null], 'queue' => ['type' => 'string', 'length' => 255, 'null' => true, 'default' => null, 'comment' => '', 'precision' => null],