From 85db6dc3453a57c5ed7c7fb6d7cc9a8bce15e26c Mon Sep 17 00:00:00 2001 From: Jack Date: Sun, 16 Aug 2026 00:17:38 +0100 Subject: [PATCH 1/3] 13.x-queue-forwarded-event --- .../Queue/Events/QueueForwarded.php | 20 +++++++++++++ src/Illuminate/Queue/Queue.php | 28 +++++++++++++++++++ .../QueueDatabaseQueueIntegrationTest.php | 26 +++++++++++++++++ 3 files changed, 74 insertions(+) create mode 100644 src/Illuminate/Queue/Events/QueueForwarded.php diff --git a/src/Illuminate/Queue/Events/QueueForwarded.php b/src/Illuminate/Queue/Events/QueueForwarded.php new file mode 100644 index 000000000000..2c343f3d1cd9 --- /dev/null +++ b/src/Illuminate/Queue/Events/QueueForwarded.php @@ -0,0 +1,20 @@ +container->make('db.transactions')->addCallback( function () use ($queue, $job, $payload, $delay, $callback) { + $this->raiseQueueForwardedEvent($queue); + $this->raiseJobQueueingEvent($queue, $job, $payload, $delay); return tap($callback($payload, $queue, $delay), function ($jobId) use ($queue, $job, $payload, $delay) { @@ -381,6 +386,8 @@ function () use ($queue, $job, $payload, $delay, $callback) { ); } + $this->raiseQueueForwardedEvent($queue); + $this->raiseJobQueueingEvent($queue, $job, $payload, $delay); return tap($callback($payload, $queue, $delay), function ($jobId) use ($queue, $job, $payload, $delay) { @@ -487,6 +494,27 @@ protected function raiseJobQueuedEvent($queue, $jobId, $job, $payload, $delay) } } + /** + * Raise the queue forwarded event. + * + * @param \UnitEnum|string|null $queue + * @return void + */ + protected function raiseQueueForwardedEvent($queue) + { + $from = enum_value($queue) ?: ($this->default ?? null); + + if (is_null($from)) { + return; + } + + $to = $this->resolveQueue($from); + + if ($to !== $from && $this->container->bound('events')) { + $this->container['events']->dispatch(new QueueForwarded($this->connectionName, $from, $to)); + } + } + /** * Get the routed queue name for the given queue. * diff --git a/tests/Queue/QueueDatabaseQueueIntegrationTest.php b/tests/Queue/QueueDatabaseQueueIntegrationTest.php index fd5b8e7d8788..325c23393711 100644 --- a/tests/Queue/QueueDatabaseQueueIntegrationTest.php +++ b/tests/Queue/QueueDatabaseQueueIntegrationTest.php @@ -10,7 +10,9 @@ use Illuminate\Queue\DatabaseQueue; use Illuminate\Queue\Events\JobQueued; use Illuminate\Queue\Events\JobQueueing; +use Illuminate\Queue\Events\QueueForwarded; use Illuminate\Queue\Queue; +use Illuminate\Queue\QueueRoutes; use Illuminate\Support\Carbon; use Illuminate\Support\Str; use PHPUnit\Framework\TestCase; @@ -288,4 +290,28 @@ public function testJobPayloadIsAvailableOnEvents() $this->assertIsArray($jobQueuedEvent->payload()); $this->assertSame('expected-job-uuid', $jobQueuedEvent->payload()['uuid']); } + + public function testQueueForwardedEventIsDispatchedWhenQueueIsForwarded() + { + $queueForwardedEvent = null; + + Container::setInstance($this->container); + + $this->container->instance('queue.routes', tap(new QueueRoutes, function ($routes) { + $routes->forward('jobs', 'processing'); + })); + + $this->container['events']->listen(function (QueueForwarded $e) use (&$queueForwardedEvent) { + $queueForwardedEvent = $e; + }); + + $this->queue->push('MyJob', [ + 'laravel' => 'Framework', + ], 'jobs'); + + $this->assertSame('jobs', $queueForwardedEvent->from); + $this->assertSame('processing', $queueForwardedEvent->to); + + Container::setInstance(null); + } } From 6fd9ac6cffb7308723442b604f8b1feb213b7348 Mon Sep 17 00:00:00 2001 From: Jack Date: Sun, 16 Aug 2026 01:38:51 +0100 Subject: [PATCH 2/3] sqs --- src/Illuminate/Queue/SqsQueue.php | 2 ++ 1 file changed, 2 insertions(+) diff --git a/src/Illuminate/Queue/SqsQueue.php b/src/Illuminate/Queue/SqsQueue.php index 48b2678772e5..e243933dffc3 100755 --- a/src/Illuminate/Queue/SqsQueue.php +++ b/src/Illuminate/Queue/SqsQueue.php @@ -382,6 +382,8 @@ protected function prepareBatchMessages(array $jobs, $data, $queue) */ protected function sendBatchedMessages(array $messages, $queue) { + $this->raiseQueueForwardedEvent($queue); + $entries = []; foreach ($messages as $id => $message) { From 3e30e2a7e0b885c54b666d8afa194e0f864ebc23 Mon Sep 17 00:00:00 2001 From: Jack Date: Sun, 16 Aug 2026 01:45:48 +0100 Subject: [PATCH 3/3] test --- tests/Queue/QueueSqsQueueTest.php | 32 +++++++++++++++++++++++++++++++ 1 file changed, 32 insertions(+) diff --git a/tests/Queue/QueueSqsQueueTest.php b/tests/Queue/QueueSqsQueueTest.php index 8b2f0d4ad99c..1f34797ee852 100755 --- a/tests/Queue/QueueSqsQueueTest.php +++ b/tests/Queue/QueueSqsQueueTest.php @@ -13,6 +13,7 @@ use Illuminate\Contracts\Queue\ShouldBeUnique; use Illuminate\Contracts\Queue\ShouldQueue; use Illuminate\Foundation\Queue\Queueable; +use Illuminate\Queue\Events\QueueForwarded; use Illuminate\Queue\Jobs\SqsJob; use Illuminate\Queue\QueueRoutes; use Illuminate\Queue\SqsQueue; @@ -1122,6 +1123,37 @@ public function testBulkRaisesQueueingAndQueuedEventsForEachJob() $this->assertSame(['mid-0', 'mid-1'], array_map(fn ($e) => $e->id, array_values($queuedEvents))); } + public function testBulkRaisesQueueForwardedEventOnceForTheBatch() + { + Container::setInstance($container = new Container); + $routes = new QueueRoutes; + $routes->forward($this->queueName, 'processing', 'sqs'); + $container->instance('queue.routes', $routes); + + $forwarded = []; + $events = Mockery::mock(\Illuminate\Contracts\Events\Dispatcher::class); + $events->allows('dispatch')->andReturnUsing(function ($event) use (&$forwarded) { + if ($event instanceof QueueForwarded) { + $forwarded[] = $event; + } + }); + $container->instance('events', $events); + + $queue = new SqsQueue($this->sqs, $this->queueName, $this->prefix); + $queue->setContainer($container); + $queue->setConnectionName('sqs'); + + $this->sqs->expects('sendMessageBatch')->andReturn(new Result(['Successful' => [], 'Failed' => []])); + + $queue->bulk(['a', 'b'], 'data', $this->queueName); + + $this->assertCount(1, $forwarded); + $this->assertSame($this->queueName, $forwarded[0]->from); + $this->assertSame('processing', $forwarded[0]->to); + + Container::setInstance(null); + } + public function testBulkHonoursPerJobDelay() { $jobA = new FakeSqsJob;