Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions composer.json
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,7 @@
"illuminate/workbench": "self.version"
},
"require-dev": {
"aws/aws-sdk-php": "^3.322.9",
"mockery/mockery": "~1.3",
"phpspec/prophecy-phpunit": "~2.0",
"phpunit/phpunit": "~9.6",
Expand Down
31 changes: 29 additions & 2 deletions src/Illuminate/Queue/Connectors/SqsConnector.php
Original file line number Diff line number Diff line change
Expand Up @@ -8,14 +8,41 @@ class SqsConnector implements ConnectorInterface {
/**
* Establish a queue connection.
*
* The client is built the AWS SDK v3 way (explicit `version` + nested
* `credentials`); the removed v2 `SqsClient::factory()` is no longer used.
*
* @param array $config
* @return \Illuminate\Queue\QueueInterface
*/
public function connect(array $config)
{
$sqs = SqsClient::factory($config);
$clientConfig = array(
'region' => isset($config['region']) ? $config['region'] : 'us-east-1',
'version' => isset($config['version']) ? $config['version'] : 'latest',
);

// Custom endpoint for SQS-compatible services (ElasticMQ/LocalStack) in
// local/dev; omitted in production so the SDK targets real AWS.
if ( ! empty($config['endpoint']))
{
$clientConfig['endpoint'] = $config['endpoint'];
}

// Credentials are optional: when absent the SDK falls back to its
// default provider chain (env vars, IAM instance/task role).
if ( ! empty($config['key']) && ! empty($config['secret']))
{
$clientConfig['credentials'] = array(
'key' => $config['key'],
'secret' => $config['secret'],
);
}

return new SqsQueue($sqs, $config['queue']);
return new SqsQueue(
new SqsClient($clientConfig),
$config['queue'],
isset($config['prefix']) ? $config['prefix'] : ''
);
}

}
25 changes: 23 additions & 2 deletions src/Illuminate/Queue/SqsQueue.php
Original file line number Diff line number Diff line change
Expand Up @@ -19,17 +19,26 @@ class SqsQueue extends Queue implements QueueInterface {
*/
protected $default;

/**
* The queue URL prefix (account base URL).
*
* @var string
*/
protected $prefix;

/**
* Create a new Amazon SQS queue instance.
*
* @param \Aws\Sqs\SqsClient $sqs
* @param string $default
* @param string $prefix
* @return void
*/
public function __construct(SqsClient $sqs, $default)
public function __construct(SqsClient $sqs, $default, $prefix = '')
{
$this->sqs = $sqs;
$this->default = $default;
$this->prefix = $prefix;
}

/**
Expand Down Expand Up @@ -105,12 +114,24 @@ public function pop($queue = null)
/**
* Get the queue or return the default.
*
* A bare queue name is resolved into a full URL using the configured
* prefix (the account base URL); full URLs are returned untouched.
*
* @param string|null $queue
* @return string
*/
public function getQueue($queue)
{
return $queue ?: $this->default;
$queue = $queue ?: $this->default;

if (filter_var($queue, FILTER_VALIDATE_URL) !== false)
{
return $queue;
}

return $this->prefix !== ''
? rtrim($this->prefix, '/') . '/' . $queue
: $queue;
}

/**
Expand Down
177 changes: 85 additions & 92 deletions tests/Queue/QueueSqsQueueTest.php
Original file line number Diff line number Diff line change
@@ -1,7 +1,8 @@
<?php

use Aws\Result;
use Aws\Sqs\SqsClient;
use Guzzle\Service\Resource\Model;
use Carbon\Carbon;
use Illuminate\Container\Container;
use Illuminate\Queue\Jobs\SqsJob;
use Illuminate\Queue\SqsQueue;
Expand All @@ -15,15 +16,16 @@ class QueueSqsQueueTest extends BackwardCompatibleTestCase
private string $account;
private string $queueName;
private string $baseUrl;
private string $prefix;
private string $queueUrl;
private string $mockedJob;
private array $mockedData;
private string|false $mockedPayload;
private int $mockedDelay;
private string $mockedMessageId;
private string $mockedReceiptHandle;
private Model $mockedSendMessageResponseModel;
private Model $mockedReceiveMessageResponseModel;
private Result $mockedSendMessageResponseModel;
private Result $mockedReceiveMessageResponseModel;

protected function tearDown(): void
{
Expand All @@ -32,107 +34,98 @@ protected function tearDown(): void

protected function setUp(): void
{
$this->markTestSkipped();

// Use Mockery to mock the SqsClient
$this->sqs = m::mock('Aws\Sqs\SqsClient');
$this->sqs = m::mock(SqsClient::class);

$this->account = '1234567891011';
$this->queueName = 'emails';
$this->baseUrl = 'https://sqs.someregion.amazonaws.com';

// This is how the modified getQueue builds the queueUrl
$this->queueUrl = $this->baseUrl . '/' . $this->account . '/' . $this->queueName;
// This is how the modified getQueue builds the queueUrl.
$this->prefix = $this->baseUrl . '/' . $this->account . '/';
$this->queueUrl = $this->prefix . $this->queueName;

$this->mockedJob = 'foo';
$this->mockedData = ['data'];
$this->mockedPayload = json_encode(['job' => $this->mockedJob, 'data' => $this->mockedData]);
$this->mockedDelay = 10;
$this->mockedMessageId = 'e3cd03ee-59a3-4ad8-b0aa-ee2e3808ac81';
$this->mockedReceiptHandle = '0NNAq8PwvXuWv5gMtS9DJ8qEdyiUwbAjpp45w2m6M4SJ1Y+PxCh7R930NRB8ylSacEmoSnW18bgd4nK\/O6ctE+VFVul4eD23mA07vVoSnPI4F\/voI1eNCp6Iax0ktGmhlNVzBwaZHEr91BRtqTRM3QKd2ASF8u+IQaSwyl\/DGK+P1+dqUOodvOVtExJwdyDLy1glZVgm85Yw9Jf5yZEEErqRwzYz\/qSigdvW4sm2l7e4phRol\/+IjMtovOyH\/ukueYdlVbQ4OshQLENhUKe7RNN5i6bE\/e5x9bnPhfj2gbM';
$this->mockedJob = 'foo';
$this->mockedData = ['data'];
$this->mockedPayload = json_encode(['job' => $this->mockedJob, 'data' => $this->mockedData]);
$this->mockedDelay = 10;
$this->mockedMessageId = 'e3cd03ee-59a3-4ad8-b0aa-ee2e3808ac81';
$this->mockedReceiptHandle = '0NNAq8PwvXuWv5gMtS9DJ8qEdyiUwbAjpp45w2m6M4SJ1Y+PxCh7R930NRB8ylSacEmoSnW18bgd4nK/O6ctE';

$this->mockedSendMessageResponseModel = new Model([
$this->mockedSendMessageResponseModel = new Result([
'Body' => $this->mockedPayload,
'MD5OfBody' => md5((string) $this->mockedPayload),
'ReceiptHandle' => $this->mockedReceiptHandle,
'MessageId' => $this->mockedMessageId,
'Attributes' => ['ApproximateReceiveCount' => 1]
'MD5OfBody' => md5((string) $this->mockedPayload),
'ReceiptHandle' => $this->mockedReceiptHandle,
'MessageId' => $this->mockedMessageId,
'Attributes' => ['ApproximateReceiveCount' => 1],
]);

$this->mockedReceiveMessageResponseModel = new Model([
$this->mockedReceiveMessageResponseModel = new Result([
'Messages' => [
0 => [
'Body' => $this->mockedPayload,
'MD5OfBody' => md5((string) $this->mockedPayload),
'ReceiptHandle' => $this->mockedReceiptHandle,
'MessageId' => $this->mockedMessageId
]
]
'Body' => $this->mockedPayload,
'MD5OfBody' => md5((string) $this->mockedPayload),
'ReceiptHandle' => $this->mockedReceiptHandle,
'MessageId' => $this->mockedMessageId,
],
],
]);
}


public function testPopProperlyPopsJobOffOfSqs()
{
$queue = $this->getMock(SqsQueue::class, ['getQueue'], [$this->sqs, $this->queueName, $this->account]);
$queue->setContainer(m::mock(Container::class));
$queue->expects($this->once())->method('getQueue')->with($this->queueName)->willReturn($this->queueUrl);
$this->sqs->shouldReceive('receiveMessage')->once()->with(
['QueueUrl' => $this->queueUrl, 'AttributeNames' => ['ApproximateReceiveCount']]
)->andReturn($this->mockedReceiveMessageResponseModel);
$result = $queue->pop($this->queueName);
$this->assertInstanceOf(SqsJob::class, $result);
}


public function testDelayedPushWithDateTimeProperlyPushesJobOntoSqs()
{
$now = Carbon::now();
$queue = $this->getMock(SqsQueue::class, ['createPayload', 'getSeconds', 'getQueue'], [$this->sqs, $this->queueName, $this->account]
);
$queue->expects($this->once())->method('createPayload')->with($this->mockedJob, $this->mockedData)->willReturn(
$this->mockedPayload
);
$queue->expects($this->once())->method('getSeconds')->with($now)->willReturn(5);
$queue->expects($this->once())->method('getQueue')->with($this->queueName)->willReturn($this->queueUrl);
$this->sqs->shouldReceive('sendMessage')->once()->with(
['QueueUrl' => $this->queueUrl, 'MessageBody' => $this->mockedPayload, 'DelaySeconds' => 5]
)->andReturn($this->mockedSendMessageResponseModel);
$id = $queue->later($now->addSeconds(5), $this->mockedJob, $this->mockedData, $this->queueName);
$this->assertEquals($this->mockedMessageId, $id);
}


public function testDelayedPushProperlyPushesJobOntoSqs()
{
$queue = $this->getMock(SqsQueue::class, ['createPayload', 'getSeconds', 'getQueue'], [$this->sqs, $this->queueName, $this->account]
);
$queue->expects($this->once())->method('createPayload')->with($this->mockedJob, $this->mockedData)->willReturn(
$this->mockedPayload
);
$queue->expects($this->once())->method('getSeconds')->with($this->mockedDelay)->willReturn($this->mockedDelay);
$queue->expects($this->once())->method('getQueue')->with($this->queueName)->willReturn($this->queueUrl);
$this->sqs->shouldReceive('sendMessage')->once()->with(
['QueueUrl' => $this->queueUrl, 'MessageBody' => $this->mockedPayload, 'DelaySeconds' => $this->mockedDelay]
)->andReturn($this->mockedSendMessageResponseModel);
$id = $queue->later($this->mockedDelay, $this->mockedJob, $this->mockedData, $this->queueName);
$this->assertEquals($this->mockedMessageId, $id);
}


public function testPushProperlyPushesJobOntoSqs()
{
$queue = $this->getMock(SqsQueue::class, ['createPayload', 'getQueue'], [$this->sqs, $this->queueName, $this->account]
);
$queue->expects($this->once())->method('createPayload')->with($this->mockedJob, $this->mockedData)->willReturn(
$this->mockedPayload
);
$queue->expects($this->once())->method('getQueue')->with($this->queueName)->willReturn($this->queueUrl);
$this->sqs->shouldReceive('sendMessage')->once()->with(
['QueueUrl' => $this->queueUrl, 'MessageBody' => $this->mockedPayload]
)->andReturn($this->mockedSendMessageResponseModel);
$id = $queue->push($this->mockedJob, $this->mockedData, $this->queueName);
$this->assertEquals($this->mockedMessageId, $id);
}
}

public function testPopProperlyPopsJobOffOfSqs()
{
$queue = $this->getMockBuilder(SqsQueue::class)->onlyMethods(['getQueue'])->setConstructorArgs([$this->sqs, $this->queueName, $this->account])->getMock();
$queue->setContainer(m::mock(Container::class));
$queue->expects($this->once())->method('getQueue')->with($this->queueName)->willReturn($this->queueUrl);
$this->sqs->shouldReceive('receiveMessage')->once()->with(['QueueUrl' => $this->queueUrl, 'AttributeNames' => ['ApproximateReceiveCount']])->andReturn($this->mockedReceiveMessageResponseModel);
$result = $queue->pop($this->queueName);
$this->assertInstanceOf(SqsJob::class, $result);
}

public function testDelayedPushWithDateTimeProperlyPushesJobOntoSqs()
{
$now = Carbon::now();
$queue = $this->getMockBuilder(SqsQueue::class)->onlyMethods(['createPayload', 'getSeconds', 'getQueue'])->setConstructorArgs([$this->sqs, $this->queueName, $this->account])->getMock();
$queue->expects($this->once())->method('createPayload')->with($this->mockedJob, $this->mockedData)->willReturn($this->mockedPayload);
$queue->expects($this->once())->method('getSeconds')->with($now)->willReturn(5);
$queue->expects($this->once())->method('getQueue')->with($this->queueName)->willReturn($this->queueUrl);
$this->sqs->shouldReceive('sendMessage')->once()->with(['QueueUrl' => $this->queueUrl, 'MessageBody' => $this->mockedPayload, 'DelaySeconds' => 5])->andReturn($this->mockedSendMessageResponseModel);
$id = $queue->later($now, $this->mockedJob, $this->mockedData, $this->queueName);
$this->assertEquals($this->mockedMessageId, $id);
}

public function testDelayedPushProperlyPushesJobOntoSqs()
{
$queue = $this->getMockBuilder(SqsQueue::class)->onlyMethods(['createPayload', 'getSeconds', 'getQueue'])->setConstructorArgs([$this->sqs, $this->queueName, $this->account])->getMock();
$queue->expects($this->once())->method('createPayload')->with($this->mockedJob, $this->mockedData)->willReturn($this->mockedPayload);
$queue->expects($this->once())->method('getSeconds')->with($this->mockedDelay)->willReturn($this->mockedDelay);
$queue->expects($this->once())->method('getQueue')->with($this->queueName)->willReturn($this->queueUrl);
$this->sqs->shouldReceive('sendMessage')->once()->with(['QueueUrl' => $this->queueUrl, 'MessageBody' => $this->mockedPayload, 'DelaySeconds' => $this->mockedDelay])->andReturn($this->mockedSendMessageResponseModel);
$id = $queue->later($this->mockedDelay, $this->mockedJob, $this->mockedData, $this->queueName);
$this->assertEquals($this->mockedMessageId, $id);
}

public function testPushProperlyPushesJobOntoSqs()
{
$queue = $this->getMockBuilder(SqsQueue::class)->onlyMethods(['createPayload', 'getQueue'])->setConstructorArgs([$this->sqs, $this->queueName, $this->account])->getMock();
$queue->expects($this->once())->method('createPayload')->with($this->mockedJob, $this->mockedData)->willReturn($this->mockedPayload);
$queue->expects($this->once())->method('getQueue')->with($this->queueName)->willReturn($this->queueUrl);
$this->sqs->shouldReceive('sendMessage')->once()->with(['QueueUrl' => $this->queueUrl, 'MessageBody' => $this->mockedPayload])->andReturn($this->mockedSendMessageResponseModel);
$id = $queue->push($this->mockedJob, $this->mockedData, $this->queueName);
$this->assertEquals($this->mockedMessageId, $id);
}

public function testGetQueueProperlyResolvesUrlWithPrefix()
{
$queue = new SqsQueue($this->sqs, $this->queueName, $this->prefix);
$this->assertEquals($this->queueUrl, $queue->getQueue(null));
$this->assertEquals($this->baseUrl . '/' . $this->account . '/test', $queue->getQueue('test'));
}

public function testGetQueueProperlyResolvesUrlWithoutPrefix()
{
$queue = new SqsQueue($this->sqs, $this->queueUrl);
$this->assertEquals($this->queueUrl, $queue->getQueue(null));
$this->assertEquals($this->queueUrl, $queue->getQueue($this->queueUrl));
}

}
Loading