From 3274555f5119e78b0ed78e1bb08a2bc90e80edb9 Mon Sep 17 00:00:00 2001 From: Yevgeny Tomenko Date: Tue, 22 Sep 2026 21:13:53 +0300 Subject: [PATCH] fix(queue): preserve dtoClass through failed jobs and requeue Failed jobs dropped the DTO marker: FailedJobsListener persisted only data + requeueOptions, and RequeueCommand re-pushed without dtoClass, so a requeued DTO job came back as a plain array and getDto() could no longer rehydrate it. --- ...20260922120000_AddDtoClassToFailedJobs.php | 25 ++++++ src/Command/RequeueCommand.php | 1 + src/Command/WorkerCommand.php | 2 +- src/Listener/FailedJobsListener.php | 1 + src/Model/Entity/FailedJob.php | 2 + src/Model/Table/FailedJobsTable.php | 5 ++ tests/Fixture/FailedJobsFixture.php | 3 + tests/TestCase/Command/RequeueCommandTest.php | 63 +++++++++++++++ .../Listener/FailedJobsListenerTest.php | 79 +++++++++++++++++++ tests/schema.php | 1 + 10 files changed, 181 insertions(+), 1 deletion(-) create mode 100644 config/Migrations/20260922120000_AddDtoClassToFailedJobs.php diff --git a/config/Migrations/20260922120000_AddDtoClassToFailedJobs.php b/config/Migrations/20260922120000_AddDtoClassToFailedJobs.php new file mode 100644 index 0000000..76632df --- /dev/null +++ b/config/Migrations/20260922120000_AddDtoClassToFailedJobs.php @@ -0,0 +1,25 @@ +table('queue_failed_jobs'); + $table->addColumn('dto_class', 'string', [ + 'length' => 255, + 'null' => true, + 'default' => null, + ]) + ->update(); + } +} diff --git a/src/Command/RequeueCommand.php b/src/Command/RequeueCommand.php index ce00f7c..57c47d8 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, + 'dtoClass' => $failedJob->dto_class ?? null, ], ); diff --git a/src/Command/WorkerCommand.php b/src/Command/WorkerCommand.php index 44605d1..7efa802 100644 --- a/src/Command/WorkerCommand.php +++ b/src/Command/WorkerCommand.php @@ -159,7 +159,7 @@ protected function getQueueExtension(Arguments $args, LoggerInterface $logger): protected function getLogger(Arguments $args): LoggerInterface { $logger = null; - if (!empty($args->getOption('verbose'))) { + if (!(in_array($args->getOption('verbose'), ['', '0'], true) || $args->getOption('verbose') === false || $args->getOption('verbose') === null)) { $logger = Log::engine((string)$args->getOption('logger')); } diff --git a/src/Listener/FailedJobsListener.php b/src/Listener/FailedJobsListener.php index c983fb4..fa878b2 100644 --- a/src/Listener/FailedJobsListener.php +++ b/src/Listener/FailedJobsListener.php @@ -67,6 +67,7 @@ public function storeFailedJob(object $event): void 'class' => $class, 'method' => $method, 'data' => json_encode($data), + 'dto_class' => $originalMessageBody['dtoClass'] ?? 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..6762e48 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 $dto_class * @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, + 'dto_class' => true, 'config' => true, 'priority' => true, 'queue' => true, diff --git a/src/Model/Table/FailedJobsTable.php b/src/Model/Table/FailedJobsTable.php index 208b325..6e92ee4 100644 --- a/src/Model/Table/FailedJobsTable.php +++ b/src/Model/Table/FailedJobsTable.php @@ -72,6 +72,11 @@ public function validationDefault(Validator $validator): Validator ->requirePresence('data', 'create') ->notEmptyString('data'); + $validator + ->scalar('dto_class') + ->maxLength('dto_class', 255) + ->allowEmptyString('dto_class'); + $validator ->scalar('config') ->maxLength('config', 255) diff --git a/tests/Fixture/FailedJobsFixture.php b/tests/Fixture/FailedJobsFixture.php index c44743e..636a866 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}', + 'dto_class' => 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}', + 'dto_class' => 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}', + 'dto_class' => 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..04c7449 100644 --- a/tests/TestCase/Command/RequeueCommandTest.php +++ b/tests/TestCase/Command/RequeueCommandTest.php @@ -19,9 +19,14 @@ use Cake\Console\TestSuite\ConsoleIntegrationTestTrait; use Cake\Core\Configure; use Cake\Log\Log; +use Cake\Queue\Job\Message; use Cake\Queue\QueueManager; use Cake\Queue\Test\TestCase\QueueTestTrait; use Cake\TestSuite\TestCase; +use Enqueue\Null\NullConnectionFactory; +use Enqueue\Null\NullMessage; +use TestApp\Dto\OrderDto; +use TestApp\Job\DtoJob; use TestApp\Job\LogToDebugJob; /** @@ -197,4 +202,62 @@ public function testJobsAreRequeuedByConfig() $this->assertDebugLogContains('Debug job was run'); } + + public function testRequeuedDtoJobKeepsDtoClass() + { + $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' => DtoJob::class, + 'method' => 'execute', + 'data' => json_encode(['id' => 7, 'customer' => 'Acme Corp', 'items' => []]), + 'dto_class' => OrderDto::class, + '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('dtoClass', $contents); + $this->assertStringContainsString('OrderDto', $contents); + $this->assertStringContainsString('Acme Corp', $contents); + + unlink($fsQueueFile); + } + + public function testRequeuedDtoJobHydratesAfterRequeue() + { + $parsedBody = [ + 'class' => [DtoJob::class, 'execute'], + 'data' => ['id' => 7, 'customer' => 'Acme Corp', 'items' => []], + 'dtoClass' => OrderDto::class, + ]; + $connectionFactory = new NullConnectionFactory(); + $context = $connectionFactory->createContext(); + $message = new Message(new NullMessage((string)json_encode($parsedBody)), $context); + + $dto = $message->getDto(OrderDto::class); + + $this->assertInstanceOf(OrderDto::class, $dto); + $this->assertSame(7, $dto->id); + $this->assertSame('Acme Corp', $dto->customer); + } } diff --git a/tests/TestCase/Listener/FailedJobsListenerTest.php b/tests/TestCase/Listener/FailedJobsListenerTest.php index d5a56e9..6972af7 100644 --- a/tests/TestCase/Listener/FailedJobsListenerTest.php +++ b/tests/TestCase/Listener/FailedJobsListenerTest.php @@ -33,6 +33,7 @@ use PHPUnit\Framework\Attributes\DataProvider; use RuntimeException; use stdClass; +use TestApp\Dto\OrderDto; use TestApp\Job\LogToDebugJob; class FailedJobsListenerTest extends TestCase @@ -100,6 +101,84 @@ public function testFailedJobIsAddedWhenEventIsFired() $this->assertStringContainsString('some message', $failedJob->exception); } + public function testFailedJobPreservesDtoClass() + { + $parsedBody = [ + 'class' => [LogToDebugJob::class, 'execute'], + 'data' => ['id' => 7, 'customer' => 'Acme Corp'], + 'dtoClass' => OrderDto::class, + '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(json_encode(['id' => 7, 'customer' => 'Acme Corp']), $failedJob->data); + $this->assertSame(OrderDto::class, $failedJob->dto_class); + } + + public function testFailedJobWithoutDtoClassStoresNull() + { + $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->dto_class); + } + /** * Data provider for testStoreFailedJobException * diff --git a/tests/schema.php b/tests/schema.php index 3aeba16..6aebef8 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], + 'dto_class' => ['type' => 'string', 'length' => 255, '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],