diff --git a/README.md b/README.md index 7913874b0..d774c78cf 100644 --- a/README.md +++ b/README.md @@ -9,7 +9,14 @@ Features: * [JMS](https://docs.oracle.com/javaee/7/api/javax/jms/package-summary.html) like transport [abstraction](https://github.com/php-enqueue/psr-queue). * Feature rich. -* Supports [AMQP](docs/transport/amqp.md) (RabbitMQ, ActiveMQ and others), [STOMP](docs/transport/stomp.md) (RabbitMQ, ActiveMQ and others), [Redis](docs/transport/redis.md), [Doctrine DBAL](docs/transport/dbal.md), [Filesystem](docs/transport/filesystem.md), [Null](docs/transport/null.md) transports. +* Supports transports: + - [AMQP](docs/transport/amqp.md) (RabbitMQ, ActiveMQ and others), + - [STOMP](docs/transport/stomp.md) + - [Amazon SQS](docs/transport/sqs.md) + - [Redis](docs/transport/redis.md) + - [Doctrine DBAL](docs/transport/dbal.md) + - [Filesystem](docs/transport/filesystem.md) + - [Null](docs/transport/null.md). * Generic purpose abstraction level (the transport level). * "Opinionated" easy to use abstraction level (the client level). * [Message bus](http://www.enterpriseintegrationpatterns.com/patterns/messaging/MessageBus.html) support. diff --git a/bin/dev b/bin/dev index b1d9bfbb3..d57602c96 100755 --- a/bin/dev +++ b/bin/dev @@ -3,7 +3,7 @@ set -x set -e -while getopts "bustefc" OPTION; do +while getopts "bustefcd" OPTION; do case $OPTION in b) COMPOSE_PROJECT_NAME=mqdev docker-compose build @@ -27,6 +27,8 @@ while getopts "bustefc" OPTION; do COMPOSE_PROJECT_NAME=mqdev docker-compose run -e CHANGELOG_GITHUB_TOKEN=${CHANGELOG_GITHUB_TOKEN:-""} --workdir="/mqdev" --rm generate-changelog github_changelog_generator --future-release "$2" --simple-list ;; + d) COMPOSE_PROJECT_NAME=mqdev docker-compose run --workdir="/mqdev" --rm dev php pkg/enqueue-bundle/Tests/Functional/app/console.php config:dump-reference enqueue + ;; \?) echo "Invalid option: -$OPTARG" >&2 exit 1 diff --git a/bin/subtree-split b/bin/subtree-split index dd972a54e..53e46cce1 100755 --- a/bin/subtree-split +++ b/bin/subtree-split @@ -51,6 +51,7 @@ remote fs git@github.com:php-enqueue/fs.git remote redis git@github.com:php-enqueue/redis.git remote dbal git@github.com:php-enqueue/dbal.git remote null git@github.com:php-enqueue/null.git +remote sqs git@github.com:php-enqueue/sqs.git remote enqueue-bundle git@github.com:php-enqueue/enqueue-bundle.git remote job-queue git@github.com:php-enqueue/job-queue.git remote test git@github.com:php-enqueue/test.git @@ -63,6 +64,7 @@ split 'pkg/fs' fs split 'pkg/redis' redis split 'pkg/dbal' dbal split 'pkg/null' null +split 'pkg/sqs' sqs split 'pkg/enqueue-bundle' enqueue-bundle split 'pkg/job-queue' job-queue split 'pkg/test' test diff --git a/composer.json b/composer.json index 53088983e..2afbec114 100644 --- a/composer.json +++ b/composer.json @@ -12,6 +12,7 @@ "enqueue/fs": "*@dev", "enqueue/null": "*@dev", "enqueue/dbal": "*@dev", + "enqueue/sqs": "*@dev", "enqueue/enqueue-bundle": "*@dev", "enqueue/job-queue": "*@dev", "enqueue/test": "*@dev", @@ -72,6 +73,10 @@ { "type": "path", "url": "pkg/dbal" + }, + { + "type": "path", + "url": "pkg/sqs" } ] } diff --git a/docker-compose.yml b/docker-compose.yml index a892eafe8..8a55a7ea0 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -24,6 +24,9 @@ services: - SYMFONY__DB__PASSWORD=rootpass - SYMFONY__REDIS__HOST=redis - SYMFONY__REDIS__PORT=6379 + - AWS__SQS__KEY=$ENQUEUE_AWS__SQS__KEY + - AWS__SQS__SECRET=$ENQUEUE_AWS__SQS__SECRET + - AWS__SQS__REGION=$ENQUEUE_AWS__SQS__REGION rabbitmq: image: enqueue/rabbitmq:latest diff --git a/docs/bundle/config_reference.md b/docs/bundle/config_reference.md index 57f935b8e..48d01b645 100644 --- a/docs/bundle/config_reference.md +++ b/docs/bundle/config_reference.md @@ -130,6 +130,16 @@ enqueue: # How often query for new messages. polling_interval: 1000 lazy: true + sqs: + key: null + secret: null + token: null + region: ~ # Required + retries: 3 + version: '2012-11-05' + + # the connection will be performed as later as possible, if the option set to true + lazy: true client: traceable_producer: false prefix: enqueue diff --git a/docs/index.md b/docs/index.md index 6fa020ddd..c0d54db4b 100644 --- a/docs/index.md +++ b/docs/index.md @@ -3,7 +3,8 @@ * [Quick tour](quick_tour.md) * Transports - [Amqp (RabbitMQ, ActiveMQ)](transport/amqp.md) - - [Stomp (RabbitMQ, ActiveMQ)](transport/stomp.md) + - [Amazon SQS](transport/sqs.md) + - [Stomp](transport/stomp.md) - [Redis](transport/redis.md) - [Doctrine DBAL](transport/dbal.md) - [Filesystem](transport/filesystem.md) diff --git a/docs/transport/sqs.md b/docs/transport/sqs.md new file mode 100644 index 000000000..86e3971e5 --- /dev/null +++ b/docs/transport/sqs.md @@ -0,0 +1,89 @@ +# Amazon SQS transport + +A transport for [Amazon SQS](https://aws.amazon.com/sqs/) broker. +It uses internally official [aws sdk library](https://packagist.org/packages/aws/aws-sdk-php) + +* [Installation](#installation) +* [Create context](#create-context) +* [Declare queue](#decalre-queue) +* [Send message to queue](#send-message-to-queue) +* [Consume message](#consume-message) +* [Purge queue messages](#purge-queue-messages) + +## Installation + +```bash +$ composer require enqueue/sqs +``` + +## Create context + +```php + 'aKey', + 'secret' => 'aSecret', + 'region' => 'aRegion', +]); + +$psrContext = $connectionFactory->createContext(); +``` + +## Declare queue. + +Declare queue operation creates a queue on a broker side. + +```php +createQueue('foo'); +$psrContext->declareQueue($fooQueue); + +// to remove queue use deleteQueue method +//$psrContext->deleteQueue($fooQueue); +``` + +## Send message to queue + +```php +createQueue('foo'); +$message = $psrContext->createMessage('Hello world!'); + +$psrContext->createProducer()->send($fooQueue, $message); +``` + +## Consume message: + +```php +createQueue('foo'); +$consumer = $psrContext->createConsumer($fooQueue); + +$message = $consumer->receive(); + +// process a message + +$consumer->acknowledge($message); +// $consumer->reject($message); +``` + +## Purge queue messages: + +```php +createQueue('foo'); + +$psrContext->purge($fooQueue); +``` + +[back to index](../index.md) \ No newline at end of file diff --git a/phpunit.xml.dist b/phpunit.xml.dist index b1ed3898d..f2faadb8e 100644 --- a/phpunit.xml.dist +++ b/phpunit.xml.dist @@ -45,6 +45,10 @@ pkg/null/Tests + + pkg/sqs/Tests + + pkg/enqueue-bundle/Tests diff --git a/pkg/enqueue-bundle/EnqueueBundle.php b/pkg/enqueue-bundle/EnqueueBundle.php index bdf738dde..4c8e5806d 100644 --- a/pkg/enqueue-bundle/EnqueueBundle.php +++ b/pkg/enqueue-bundle/EnqueueBundle.php @@ -17,6 +17,8 @@ use Enqueue\Fs\Symfony\FsTransportFactory; use Enqueue\Redis\RedisContext; use Enqueue\Redis\Symfony\RedisTransportFactory; +use Enqueue\Sqs\SqsContext; +use Enqueue\Sqs\Symfony\SqsTransportFactory; use Enqueue\Stomp\StompContext; use Enqueue\Stomp\Symfony\RabbitMqStompTransportFactory; use Enqueue\Stomp\Symfony\StompTransportFactory; @@ -64,5 +66,9 @@ public function build(ContainerBuilder $container) if (class_exists(DbalContext::class)) { $extension->addTransportFactory(new DbalTransportFactory()); } + + if (class_exists(SqsContext::class)) { + $extension->addTransportFactory(new SqsTransportFactory()); + } } } diff --git a/pkg/enqueue-bundle/Tests/Functional/UseCasesTest.php b/pkg/enqueue-bundle/Tests/Functional/UseCasesTest.php index 063cf05a3..40d1eb778 100644 --- a/pkg/enqueue-bundle/Tests/Functional/UseCasesTest.php +++ b/pkg/enqueue-bundle/Tests/Functional/UseCasesTest.php @@ -86,6 +86,16 @@ public function provideEnqueueConfigs() ] ] ]], + ['sqs' => [ + 'transport' => [ + 'default' => 'sqs', + 'sqs' => [ + 'key' => getenv('AWS__SQS__KEY'), + 'secret' => getenv('AWS__SQS__SECRET'), + 'region' => getenv('AWS__SQS__REGION'), + ] + ] + ]], ]; } diff --git a/pkg/enqueue-bundle/Tests/Unit/EnqueueBundleTest.php b/pkg/enqueue-bundle/Tests/Unit/EnqueueBundleTest.php index b4bd8d688..ee10168f6 100644 --- a/pkg/enqueue-bundle/Tests/Unit/EnqueueBundleTest.php +++ b/pkg/enqueue-bundle/Tests/Unit/EnqueueBundleTest.php @@ -14,6 +14,7 @@ use Enqueue\Dbal\Symfony\DbalTransportFactory; use Enqueue\Fs\Symfony\FsTransportFactory; use Enqueue\Redis\Symfony\RedisTransportFactory; +use Enqueue\Sqs\Symfony\SqsTransportFactory; use Enqueue\Stomp\Symfony\RabbitMqStompTransportFactory; use Enqueue\Stomp\Symfony\StompTransportFactory; use Enqueue\Symfony\DefaultTransportFactory; @@ -194,6 +195,23 @@ public function testShouldRegisterDbalTransportFactory() $bundle->build($container); } + public function testShouldRegisterSqsTransportFactory() + { + $extensionMock = $this->createEnqueueExtensionMock(); + + $container = new ContainerBuilder(); + $container->registerExtension($extensionMock); + + $extensionMock + ->expects($this->at(9)) + ->method('addTransportFactory') + ->with($this->isInstanceOf(SqsTransportFactory::class)) + ; + + $bundle = new EnqueueBundle(); + $bundle->build($container); + } + /** * @return \PHPUnit_Framework_MockObject_MockObject|EnqueueExtension */ diff --git a/pkg/sqs/.gitignore b/pkg/sqs/.gitignore new file mode 100644 index 000000000..a770439e5 --- /dev/null +++ b/pkg/sqs/.gitignore @@ -0,0 +1,6 @@ +*~ +/composer.lock +/composer.phar +/phpunit.xml +/vendor/ +/.idea/ diff --git a/pkg/sqs/.travis.yml b/pkg/sqs/.travis.yml new file mode 100644 index 000000000..42374ddc7 --- /dev/null +++ b/pkg/sqs/.travis.yml @@ -0,0 +1,21 @@ +sudo: false + +git: + depth: 1 + +language: php + +php: + - '5.6' + - '7.0' + +cache: + directories: + - $HOME/.composer/cache + +install: + - composer self-update + - composer install --prefer-source + +script: + - vendor/bin/phpunit --exclude-group=functional diff --git a/pkg/sqs/Client/SqsDriver.php b/pkg/sqs/Client/SqsDriver.php new file mode 100644 index 000000000..44725e890 --- /dev/null +++ b/pkg/sqs/Client/SqsDriver.php @@ -0,0 +1,170 @@ +context = $context; + $this->config = $config; + $this->queueMetaRegistry = $queueMetaRegistry; + } + + /** + * {@inheritdoc} + */ + public function sendToRouter(Message $message) + { + if (false == $message->getProperty(Config::PARAMETER_TOPIC_NAME)) { + throw new \LogicException('Topic name parameter is required but is not set'); + } + + $queue = $this->createQueue($this->config->getRouterQueueName()); + $transportMessage = $this->createTransportMessage($message); + + $this->context->createProducer()->send($queue, $transportMessage); + } + + /** + * {@inheritdoc} + */ + public function sendToProcessor(Message $message) + { + if (false == $message->getProperty(Config::PARAMETER_PROCESSOR_NAME)) { + throw new \LogicException('Processor name parameter is required but is not set'); + } + + if (false == $queueName = $message->getProperty(Config::PARAMETER_PROCESSOR_QUEUE_NAME)) { + throw new \LogicException('Queue name parameter is required but is not set'); + } + + $transportMessage = $this->createTransportMessage($message); + $destination = $this->createQueue($queueName); + + $this->context->createProducer()->send($destination, $transportMessage); + } + + /** + * {@inheritdoc} + * + * @return SqsDestination + */ + public function createQueue($queueName) + { + $transportName = $this->queueMetaRegistry->getQueueMeta($queueName)->getTransportName(); + $transportName = str_replace('.', '_dot_', $transportName); + + return $this->context->createQueue($transportName); + } + + /** + * {@inheritdoc} + */ + public function setupBroker(LoggerInterface $logger = null) + { + $logger = $logger ?: new NullLogger(); + $log = function ($text, ...$args) use ($logger) { + $logger->debug(sprintf('[AmqpDriver] '.$text, ...$args)); + }; + + // setup router + $routerQueue = $this->createQueue($this->config->getRouterQueueName()); + $log('Declare router queue: %s', $routerQueue->getQueueName()); + $this->context->declareQueue($routerQueue); + + // setup queues + foreach ($this->queueMetaRegistry->getQueuesMeta() as $meta) { + $queue = $this->createQueue($meta->getClientName()); + + $log('Declare processor queue: %s', $queue->getQueueName()); + $this->context->declareQueue($queue); + } + } + + /** + * {@inheritdoc} + * + * @return SqsMessage + */ + public function createTransportMessage(Message $message) + { + $properties = $message->getProperties(); + + $headers = $message->getHeaders(); + $headers['content_type'] = $message->getContentType(); + + $transportMessage = $this->context->createMessage(); + $transportMessage->setBody($message->getBody()); + $transportMessage->setHeaders($headers); + $transportMessage->setProperties($properties); + $transportMessage->setMessageId($message->getMessageId()); + $transportMessage->setTimestamp($message->getTimestamp()); + $transportMessage->setReplyTo($message->getReplyTo()); + $transportMessage->setCorrelationId($message->getCorrelationId()); + + return $transportMessage; + } + + /** + * @param SqsMessage $message + * + * {@inheritdoc} + */ + public function createClientMessage(PsrMessage $message) + { + $clientMessage = new Message(); + + $clientMessage->setBody($message->getBody()); + $clientMessage->setHeaders($message->getHeaders()); + $clientMessage->setProperties($message->getProperties()); + + $clientMessage->setContentType($message->getHeader('content_type')); + $clientMessage->setMessageId($message->getMessageId()); + $clientMessage->setTimestamp($message->getTimestamp()); + $clientMessage->setPriority(MessagePriority::NORMAL); + $clientMessage->setReplyTo($message->getReplyTo()); + $clientMessage->setCorrelationId($message->getCorrelationId()); + + return $clientMessage; + } + + /** + * @return Config + */ + public function getConfig() + { + return $this->config; + } +} diff --git a/pkg/sqs/LICENSE b/pkg/sqs/LICENSE new file mode 100644 index 000000000..f1e6a22fe --- /dev/null +++ b/pkg/sqs/LICENSE @@ -0,0 +1,20 @@ +The MIT License (MIT) +Copyright (c) 2016 Kotliar Maksym + +Permission is hereby granted, free of charge, to any person obtaining a copy +of this software and associated documentation files (the "Software"), to deal +in the Software without restriction, including without limitation the rights +to use, copy, modify, merge, publish, distribute, sublicense, and/or sell +copies of the Software, and to permit persons to whom the Software is furnished +to do so, subject to the following conditions: + +The above copyright notice and this permission notice shall be included in all +copies or substantial portions of the Software. + +THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR +IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, +FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE +AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER +LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, +OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN +THE SOFTWARE. diff --git a/pkg/sqs/README.md b/pkg/sqs/README.md new file mode 100644 index 000000000..ab867e6e6 --- /dev/null +++ b/pkg/sqs/README.md @@ -0,0 +1,18 @@ +# Amazon SQS Transport + +[![Gitter](https://badges.gitter.im/php-enqueue/Lobby.svg)](https://gitter.im/php-enqueue/Lobby) +[![Build Status](https://travis-ci.org/php-enqueue/sqs.png?branch=master)](https://travis-ci.org/php-enqueue/sqs) +[![Total Downloads](https://poser.pugx.org/enqueue/sqs/d/total.png)](https://packagist.org/packages/enqueue/sqs) +[![Latest Stable Version](https://poser.pugx.org/enqueue/sqs/version.png)](https://packagist.org/packages/enqueue/sqs) + +This is an implementation of PSR specification. It allows you to send and consume message through Amazon SQS library. + +## Resources + +* [Documentation](https://github.com/php-enqueue/enqueue-dev/blob/master/docs/index.md) +* [Questions](https://gitter.im/php-enqueue/Lobby) +* [Issue Tracker](https://github.com/php-enqueue/enqueue-dev/issues) + +## License + +It is released under the [MIT License](LICENSE). \ No newline at end of file diff --git a/pkg/sqs/SqsConnectionFactory.php b/pkg/sqs/SqsConnectionFactory.php new file mode 100644 index 000000000..5b7180454 --- /dev/null +++ b/pkg/sqs/SqsConnectionFactory.php @@ -0,0 +1,98 @@ + null - AWS credentials. If no credentials are provided, the SDK will attempt to load them from the environment. + * 'secret' => null, - AWS credentials. If no credentials are provided, the SDK will attempt to load them from the environment. + * 'token' => null, - AWS credentials. If no credentials are provided, the SDK will attempt to load them from the environment. + * 'region' => null, - (string, required) Region to connect to. See http://docs.aws.amazon.com/general/latest/gr/rande.html for a list of available regions. + * 'retries' => 3, - (int, default=int(3)) Configures the maximum number of allowed retries for a client (pass 0 to disable retries). + * 'version' => '2012-11-05', - (string, required) The version of the webservice to utilize + * 'lazy' => true, - Enable lazy connection (boolean) + * ] + * + * @param $config + */ + public function __construct(array $config = []) + { + $this->config = array_replace([ + 'key' => null, + 'secret' => null, + 'token' => null, + 'region' => null, + 'retries' => 3, + 'version' => '2012-11-05', + 'lazy' => true, + ], $config); + } + + /** + * {@inheritdoc} + * + * @return SqsContext + */ + public function createContext() + { + if ($this->config['lazy']) { + return new SqsContext(function () { + return $this->establishConnection(); + }); + } + + return new SqsContext($this->establishConnection()); + } + + /** + * @return SqsClient + */ + private function establishConnection() + { + if ($this->client) { + return $this->client; + } + + $config = [ + 'version' => $this->config['version'], + 'retries' => $this->config['retries'], + 'region' => $this->config['region'], + ]; + + if ($this->config['key'] && $this->config['secret']) { + $config['credentials'] = [ + 'key' => $this->config['key'], + 'secret' => $this->config['secret'], + ]; + + if ($this->config['token']) { + $config['credentials']['token'] = $this->config['token']; + } + } + + $this->client = new SqsClient($config); + + return $this->client; + } + + /** + * {@inheritdoc} + */ + public function close() + { + } +} diff --git a/pkg/sqs/SqsConsumer.php b/pkg/sqs/SqsConsumer.php new file mode 100644 index 000000000..7dd3850f8 --- /dev/null +++ b/pkg/sqs/SqsConsumer.php @@ -0,0 +1,206 @@ +context = $context; + $this->queue = $queue; + $this->messages = []; + $this->maxNumberOfMessages = 1; + } + + /** + * @return int|null + */ + public function getVisibilityTimeout() + { + return $this->visibilityTimeout; + } + + /** + * The duration (in seconds) that the received messages are hidden from subsequent retrieve + * requests after being retrieved by a ReceiveMessage request. + * + * @param int|null $visibilityTimeout + */ + public function setVisibilityTimeout($visibilityTimeout) + { + $this->visibilityTimeout = is_null($visibilityTimeout) ? null : (int) $visibilityTimeout; + } + + /** + * @return int + */ + public function getMaxNumberOfMessages() + { + return $this->maxNumberOfMessages; + } + + /** + * The maximum number of messages to return. Amazon SQS never returns more messages than this value + * (however, fewer messages might be returned). Valid values are 1 to 10. Default is 1. + * + * @param int $maxNumberOfMessages + */ + public function setMaxNumberOfMessages($maxNumberOfMessages) + { + $this->maxNumberOfMessages = (int) $maxNumberOfMessages; + } + + /** + * {@inheritdoc} + * + * @return SqsDestination + */ + public function getQueue() + { + return $this->queue; + } + + /** + * {@inheritdoc} + */ + public function receive($timeout = 0) + { + $timeout /= 1000; + + return $this->receiveMessage($timeout); + } + + /** + * {@inheritdoc} + */ + public function receiveNoWait() + { + return $this->receiveMessage(0); + } + + /** + * {@inheritdoc} + * + * @param SqsMessage $message + */ + public function acknowledge(PsrMessage $message) + { + InvalidMessageException::assertMessageInstanceOf($message, SqsMessage::class); + + $this->context->getClient()->deleteMessage([ + 'QueueUrl' => $this->context->getQueueUrl($this->queue), + 'ReceiptHandle' => $message->getReceiptHandle(), + ]); + } + + /** + * {@inheritdoc} + * + * @param SqsMessage $message + */ + public function reject(PsrMessage $message, $requeue = false) + { + InvalidMessageException::assertMessageInstanceOf($message, SqsMessage::class); + + $this->context->getClient()->deleteMessage([ + 'QueueUrl' => $this->context->getQueueUrl($this->queue), + 'ReceiptHandle' => $message->getReceiptHandle(), + ]); + + if ($requeue) { + $this->context->createProducer()->send($this->queue, $message); + } + } + + /** + * @param int $timeoutSeconds + * + * @return SqsMessage|null + */ + protected function receiveMessage($timeoutSeconds) + { + if ($message = array_pop($this->messages)) { + return $this->convertMessage($message); + } + + $arguments = [ + 'AttributeNames' => ['All'], + 'MessageAttributeNames' => ['All'], + 'MaxNumberOfMessages' => $this->maxNumberOfMessages, + 'QueueUrl' => $this->context->getQueueUrl($this->queue), + 'WaitTimeSeconds' => $timeoutSeconds, + ]; + + if ($this->visibilityTimeout) { + $arguments['VisibilityTimeout'] = $this->visibilityTimeout; + } + + $result = $this->context->getClient()->receiveMessage($arguments); + + if ($result->hasKey('Messages')) { + $this->messages = $result->get('Messages'); + } + + if ($message = array_pop($this->messages)) { + return $this->convertMessage($message); + } + } + + /** + * @param array $sqsMessage + * + * @return SqsMessage + */ + protected function convertMessage(array $sqsMessage) + { + $message = $this->context->createMessage(); + + $message->setBody($sqsMessage['Body']); + $message->setReceiptHandle($sqsMessage['ReceiptHandle']); + + if (isset($sqsMessage['Attributes']['ApproximateReceiveCount'])) { + $message->setRedelivered(((int) $sqsMessage['Attributes']['ApproximateReceiveCount']) > 1); + } + + if (isset($sqsMessage['MessageAttributes']['Headers'])) { + $headers = json_decode($sqsMessage['MessageAttributes']['Headers']['StringValue'], true); + + $message->setHeaders($headers[0]); + $message->setProperties($headers[1]); + } + + return $message; + } +} diff --git a/pkg/sqs/SqsContext.php b/pkg/sqs/SqsContext.php new file mode 100644 index 000000000..c6e6ce7af --- /dev/null +++ b/pkg/sqs/SqsContext.php @@ -0,0 +1,196 @@ +client = $client; + } elseif (is_callable($client)) { + $this->clientFactory = $client; + } else { + throw new \InvalidArgumentException(sprintf( + 'The $client argument must be either %s or callable that returns %s once called.', + SqsClient::class, + SqsClient::class + )); + } + } + + /** + * {@inheritdoc} + * + * @return SqsMessage + */ + public function createMessage($body = '', array $properties = [], array $headers = []) + { + return new SqsMessage($body, $properties, $headers); + } + + /** + * {@inheritdoc} + * + * @return SqsDestination + */ + public function createTopic($topicName) + { + return new SqsDestination($topicName); + } + + /** + * {@inheritdoc} + * + * @return SqsDestination + */ + public function createQueue($queueName) + { + return new SqsDestination($queueName); + } + + /** + * {@inheritdoc} + */ + public function createTemporaryQueue() + { + throw new \BadMethodCallException('SQS transport does not support temporary queues'); + } + + /** + * {@inheritdoc} + * + * @return SqsProducer + */ + public function createProducer() + { + return new SqsProducer($this); + } + + /** + * {@inheritdoc} + * + * @param SqsDestination $destination + * + * @return SqsConsumer + */ + public function createConsumer(PsrDestination $destination) + { + InvalidDestinationException::assertDestinationInstanceOf($destination, SqsDestination::class); + + return new SqsConsumer($this, $destination); + } + + /** + * {@inheritdoc} + */ + public function close() + { + } + + /** + * @return SqsClient + */ + public function getClient() + { + if (false == $this->client) { + $client = call_user_func($this->clientFactory); + if (false == $client instanceof SqsClient) { + throw new \LogicException(sprintf( + 'The factory must return instance of "%s". But it returns %s', + SqsClient::class, + is_object($client) ? get_class($client) : gettype($client) + )); + } + + $this->client = $client; + } + + return $this->client; + } + + /** + * @param SqsDestination $destination + * + * @return string + */ + public function getQueueUrl(SqsDestination $destination) + { + if (isset($this->queueUrls[$destination->getQueueName()])) { + return $this->queueUrls[$destination->getQueueName()]; + } + + $result = $this->getClient()->getQueueUrl([ + 'QueueName' => $destination->getQueueName(), + ]); + + if (false == $result->hasKey('QueueUrl')) { + throw new \RuntimeException(sprintf('QueueUrl cannot be resolved. queueName: "%s"', $destination->getQueueName())); + } + + return $this->queueUrls[$destination->getQueueName()] = $result->get('QueueUrl'); + } + + /** + * @param SqsDestination $dest + */ + public function declareQueue(SqsDestination $dest) + { + $result = $this->getClient()->createQueue([ + 'Attributes' => $dest->getAttributes(), + 'QueueName' => $dest->getQueueName(), + ]); + + if (false == $result->hasKey('QueueUrl')) { + throw new \RuntimeException(sprintf('Cannot create queue. queueName: "%s"', $dest->getQueueName())); + } + + $this->queueUrls[$dest->getQueueName()] = $result->get('QueueUrl'); + } + + /** + * @param SqsDestination $dest + */ + public function deleteQueue(SqsDestination $dest) + { + $this->getClient()->deleteQueue([ + 'QueueUrl' => $this->getQueueUrl($dest), + ]); + + unset($this->queueUrls[$dest->getQueueName()]); + } + + /** + * @param SqsDestination $dest + */ + public function purge(SqsDestination $dest) + { + $this->getClient()->purgeQueue([ + 'QueueUrl' => $this->getQueueUrl($dest), + ]); + } +} diff --git a/pkg/sqs/SqsDestination.php b/pkg/sqs/SqsDestination.php new file mode 100644 index 000000000..d820f2e35 --- /dev/null +++ b/pkg/sqs/SqsDestination.php @@ -0,0 +1,193 @@ +name = $name; + $this->attributes = []; + } + + /** + * {@inheritdoc} + */ + public function getQueueName() + { + return $this->name; + } + + /** + * {@inheritdoc} + */ + public function getTopicName() + { + return $this->name; + } + + /** + * @return array + */ + public function getAttributes() + { + return $this->attributes; + } + + /** + * The number of seconds for which the delivery of all messages in the queue is delayed. + * Valid values: An integer from 0 to 900 seconds (15 minutes). The default is 0 (zero). + * + * @param int $seconds + */ + public function setDelaySeconds($seconds) + { + $this->attributes['DelaySeconds'] = (int) $seconds; + } + + /** + * The limit of how many bytes a message can contain before Amazon SQS rejects it. + * Valid values: An integer from 1,024 bytes (1 KiB) to 262,144 bytes (256 KiB). + * The default is 262,144 (256 KiB). + * + * @param int $bytes + */ + public function setMaximumMessageSize($bytes) + { + $this->attributes['MaximumMessageSize'] = (int) $bytes; + } + + /** + * The number of seconds for which Amazon SQS retains a message. + * Valid values: An integer from 60 seconds (1 minute) to 1,209,600 seconds (14 days). + * The default is 345,600 (4 days). + * + * @param int $seconds + */ + public function setMessageRetentionPeriod($seconds) + { + $this->attributes['MessageRetentionPeriod'] = (int) $seconds; + } + + /** + * The queue's policy. A valid AWS policy. For more information about policy structure, + * see http://docs.aws.amazon.com/IAM/latest/UserGuide/access_policies.html. + * + * @param string $policy + */ + public function setPolicy($policy) + { + $this->attributes['Policy'] = $policy; + } + + /** + * The number of seconds for which a ReceiveMessage action waits for a message to arrive. + * Valid values: An integer from 0 to 20 (seconds). The default is 0 (zero). + * + * @param int $seconds + */ + public function setReceiveMessageWaitTimeSeconds($seconds) + { + $this->attributes['ReceiveMessageWaitTimeSeconds'] = (int) $seconds; + } + + /** + * The parameters for the dead letter queue functionality of the source queue. + * For more information about the redrive policy and dead letter queues, + * see http://docs.aws.amazon.com/AWSSimpleQueueService/latest/SQSDeveloperGuide/sqs-dead-letter-queues.html. + * The dead letter queue of a FIFO queue must also be a FIFO queue. + * Similarly, the dead letter queue of a standard queue must also be a standard queue. + * + * @param int $maxReceiveCount + * @param string $deadLetterTargetArn + */ + public function setRedrivePolicy($maxReceiveCount, $deadLetterTargetArn) + { + $this->attributes['RedrivePolicy'] = json_encode([ + 'maxReceiveCount' => (string) $maxReceiveCount, + 'deadLetterTargetArn' => (string) $deadLetterTargetArn, + ]); + } + + /** + * The visibility timeout for the queue. Valid values: An integer from 0 to 43,200 (12 hours). + * The default is 30. For more information about the visibility timeout, + * see http://docs.aws.amazon.com/AWSSimpleQueueService/latest/SQSDeveloperGuide/sqs-visibility-timeout.html. + * + * @param int $seconds + */ + public function setVisibilityTimeout($seconds) + { + $this->attributes['VisibilityTimeout'] = (int) $seconds; + } + + /** + * Only FIFO + * + * Designates a queue as FIFO. You can provide this attribute only during queue creation. + * You can't change it for an existing queue. When you set this attribute, you must provide a MessageGroupId explicitly. + * For more information, see http://docs.aws.amazon.com/AWSSimpleQueueService/latest/SQSDeveloperGuide/FIFO-queues.html#FIFO-queues-understanding-logic. + * + * @param bool $enable + */ + public function setFifoQueue($enable) + { + if ($enable) { + $this->attributes['FifoQueue'] = 'true'; + } else { + unset($this->attributes['FifoQueue']); + } + } + + /** + * Only FIFO + * + * Enables content-based deduplication. + * For more information, see http://docs.aws.amazon.com/AWSSimpleQueueService/latest/SQSDeveloperGuide/FIFO-queues.html#FIFO-queues-exactly-once-processing. + * * Every message must have a unique MessageDeduplicationId, + * * You may provide a MessageDeduplicationId explicitly. + * * If you aren't able to provide a MessageDeduplicationId and you enable ContentBasedDeduplication for your queue, + * Amazon SQS uses a SHA-256 hash to generate the MessageDeduplicationId using the body of the message (but not the attributes of the message). + * * If you don't provide a MessageDeduplicationId and the queue doesn't have ContentBasedDeduplication set, + * the action fails with an error. + * * If the queue has ContentBasedDeduplication set, your MessageDeduplicationId overrides the generated one. + * * When ContentBasedDeduplication is in effect, messages with identical content sent within the deduplication + * interval are treated as duplicates and only one copy of the message is delivered. + * * You can also use ContentBasedDeduplication for messages with identical content to be treated as duplicates. + * * If you send one message with ContentBasedDeduplication enabled and then another message with a MessageDeduplicationId + * that is the same as the one generated for the first MessageDeduplicationId, the two messages are treated as + * duplicates and only one copy of the message is delivered. + * + * @param bool $enable + */ + public function setContentBasedDeduplication($enable) + { + if ($enable) { + $this->attributes['ContentBasedDeduplication'] = 'true'; + } else { + unset($this->attributes['ContentBasedDeduplication']); + } + } +} diff --git a/pkg/sqs/SqsMessage.php b/pkg/sqs/SqsMessage.php new file mode 100644 index 000000000..0e99fdb28 --- /dev/null +++ b/pkg/sqs/SqsMessage.php @@ -0,0 +1,312 @@ +body = $body; + $this->properties = $properties; + $this->headers = $headers; + $this->redelivered = false; + $this->delaySeconds = 0; + } + + /** + * @param string $body + */ + public function setBody($body) + { + $this->body = $body; + } + + /** + * {@inheritdoc} + */ + public function getBody() + { + return $this->body; + } + + /** + * {@inheritdoc} + */ + public function setProperties(array $properties) + { + $this->properties = $properties; + } + + /** + * {@inheritdoc} + */ + public function setProperty($name, $value) + { + $this->properties[$name] = $value; + } + + /** + * {@inheritdoc} + */ + public function getProperties() + { + return $this->properties; + } + + /** + * {@inheritdoc} + */ + public function getProperty($name, $default = null) + { + return array_key_exists($name, $this->properties) ? $this->properties[$name] : $default; + } + + /** + * {@inheritdoc} + */ + public function setHeader($name, $value) + { + $this->headers[$name] = $value; + } + + /** + * @param array $headers + */ + public function setHeaders(array $headers) + { + $this->headers = $headers; + } + + /** + * {@inheritdoc} + */ + public function getHeaders() + { + return $this->headers; + } + + /** + * {@inheritdoc} + */ + public function getHeader($name, $default = null) + { + return array_key_exists($name, $this->headers) ?$this->headers[$name] : $default; + } + + /** + * {@inheritdoc} + */ + public function isRedelivered() + { + return $this->redelivered; + } + + /** + * {@inheritdoc} + */ + public function setRedelivered($redelivered) + { + $this->redelivered = $redelivered; + } + + /** + * {@inheritdoc} + */ + public function setReplyTo($replyTo) + { + $this->setHeader('reply_to', $replyTo); + } + + /** + * {@inheritdoc} + */ + public function getReplyTo() + { + return $this->getHeader('reply_to'); + } + + /** + * {@inheritdoc} + */ + public function setCorrelationId($correlationId) + { + $this->setHeader('correlation_id', $correlationId); + } + + /** + * {@inheritdoc} + */ + public function getCorrelationId() + { + return $this->getHeader('correlation_id', ''); + } + + /** + * {@inheritdoc} + */ + public function setMessageId($messageId) + { + $this->setHeader('message_id', $messageId); + } + + /** + * {@inheritdoc} + */ + public function getMessageId() + { + return $this->getHeader('message_id', ''); + } + + /** + * {@inheritdoc} + */ + public function getTimestamp() + { + return $this->getHeader('timestamp'); + } + + /** + * {@inheritdoc} + */ + public function setTimestamp($timestamp) + { + $this->setHeader('timestamp', (int) $timestamp); + } + + /** + * The number of seconds to delay a specific message. Valid values: 0 to 900. Maximum: 15 minutes. + * Messages with a positive DelaySeconds value become available for processing after the delay period is finished. + * If you don't specify a value, the default value for the queue applies. + * When you set FifoQueue, you can't set DelaySeconds per message. You can set this parameter only on a queue level. + * + * Set delay in seconds + * + * @param int $seconds + */ + public function setDelaySeconds($seconds) + { + $this->delaySeconds = (int) $seconds; + } + + /** + * @return int + */ + public function getDelaySeconds() + { + return $this->delaySeconds; + } + + /** + * Only FIFO + * + * The token used for deduplication of sent messages. If a message with a particular MessageDeduplicationId is sent successfully, + * any messages sent with the same MessageDeduplicationId are accepted successfully but aren't delivered during the 5-minute + * deduplication interval. For more information, see http://docs.aws.amazon.com/AWSSimpleQueueService/latest/SQSDeveloperGuide/FIFO-queues.html#FIFO-queues-exactly-once-processing. + * + * @param string|null $id + */ + public function setMessageDeduplicationId($id) + { + $this->messageDeduplicationId = $id; + } + + /** + * @return string|null + */ + public function getMessageDeduplicationId() + { + return $this->messageDeduplicationId; + } + + /** + * Only FIFO + * + * The tag that specifies that a message belongs to a specific message group. Messages that belong to the same message group + * are processed in a FIFO manner (however, messages in different message groups might be processed out of order). + * To interleave multiple ordered streams within a single queue, use MessageGroupId values (for example, session data + * for multiple users). In this scenario, multiple readers can process the queue, but the session data + * of each user is processed in a FIFO fashion. + * + * @param string|null $id + */ + public function setMessageGroupId($id) + { + $this->messageGroupId = $id; + } + + /** + * @return string|null + */ + public function getMessageGroupId() + { + return $this->messageGroupId; + } + + /** + * This handle is associated with the action of receiving the message, not with the message itself. + * To delete the message or to change the message visibility, you must provide the receipt handle (not the message ID). + * + * If you receive a message more than once, each time you receive it, you get a different receipt handle. + * You must provide the most recently received receipt handle when you request to delete the message (otherwise, the message might not be deleted). + * + * @param string $receipt + */ + public function setReceiptHandle($receipt) + { + $this->receiptHandle = $receipt; + } + + /** + * @return string + */ + public function getReceiptHandle() + { + return $this->receiptHandle; + } +} diff --git a/pkg/sqs/SqsProducer.php b/pkg/sqs/SqsProducer.php new file mode 100644 index 000000000..8d92dbf7a --- /dev/null +++ b/pkg/sqs/SqsProducer.php @@ -0,0 +1,75 @@ +context = $context; + } + + /** + * {@inheritdoc} + * + * @param SqsDestination $destination + * @param SqsMessage $message + */ + public function send(PsrDestination $destination, PsrMessage $message) + { + InvalidDestinationException::assertDestinationInstanceOf($destination, SqsDestination::class); + InvalidMessageException::assertMessageInstanceOf($message, SqsMessage::class); + + $body = $message->getBody(); + if (is_scalar($body) || is_null($body)) { + $body = (string) $body; + } else { + throw new InvalidMessageException(sprintf( + 'The message body must be a scalar or null. Got: %s', + is_object($body) ? get_class($body) : gettype($body) + )); + } + + $arguments = [ + 'MessageAttributes' => [ + 'Headers' => [ + 'DataType' => 'String', + 'StringValue' => json_encode([$message->getHeaders(), $message->getProperties()]), + ], + ], + 'MessageBody' => $body, + 'QueueUrl' => $this->context->getQueueUrl($destination), + ]; + + if ($message->getDelaySeconds()) { + $arguments['DelaySeconds'] = $message->getDelaySeconds(); + } + + if ($message->getMessageDeduplicationId()) { + $arguments['MessageDeduplicationId'] = $message->getMessageDeduplicationId(); + } + + if ($message->getMessageGroupId()) { + $arguments['MessageGroupId'] = $message->getMessageGroupId(); + } + + $result = $this->context->getClient()->sendMessage($arguments); + + if (false == $result->hasKey('MessageId')) { + throw new \RuntimeException('Message was not sent'); + } + } +} diff --git a/pkg/sqs/Symfony/SqsTransportFactory.php b/pkg/sqs/Symfony/SqsTransportFactory.php new file mode 100644 index 000000000..134a49931 --- /dev/null +++ b/pkg/sqs/Symfony/SqsTransportFactory.php @@ -0,0 +1,103 @@ +name = $name; + } + + /** + * {@inheritdoc} + */ + public function addConfiguration(ArrayNodeDefinition $builder) + { + $builder + ->children() + ->scalarNode('key')->defaultNull()->end() + ->scalarNode('secret')->defaultNull()->end() + ->scalarNode('token')->defaultNull()->end() + ->scalarNode('region')->isRequired()->end() + ->integerNode('retries')->defaultValue(3)->end() + ->scalarNode('version')->cannotBeEmpty()->defaultValue('2012-11-05')->end() + ->booleanNode('lazy') + ->defaultTrue() + ->info('the connection will be performed as later as possible, if the option set to true') + ->end() + ; + } + + /** + * {@inheritdoc} + */ + public function createConnectionFactory(ContainerBuilder $container, array $config) + { + $factory = new Definition(SqsConnectionFactory::class); + $factory->setArguments([$config]); + + $factoryId = sprintf('enqueue.transport.%s.connection_factory', $this->getName()); + $container->setDefinition($factoryId, $factory); + + return $factoryId; + } + + /** + * {@inheritdoc} + */ + public function createContext(ContainerBuilder $container, array $config) + { + $factoryId = sprintf('enqueue.transport.%s.connection_factory', $this->getName()); + + $context = new Definition(SqsContext::class); + $context->setFactory([new Reference($factoryId), 'createContext']); + + $contextId = sprintf('enqueue.transport.%s.context', $this->getName()); + $container->setDefinition($contextId, $context); + + return $contextId; + } + + /** + * {@inheritdoc} + */ + public function createDriver(ContainerBuilder $container, array $config) + { + $driver = new Definition(SqsDriver::class); + $driver->setArguments([ + new Reference(sprintf('enqueue.transport.%s.context', $this->getName())), + new Reference('enqueue.client.config'), + new Reference('enqueue.client.meta.queue_meta_registry'), + ]); + + $driverId = sprintf('enqueue.client.%s.driver', $this->getName()); + $container->setDefinition($driverId, $driver); + + return $driverId; + } + + /** + * {@inheritdoc} + */ + public function getName() + { + return $this->name; + } +} diff --git a/pkg/sqs/Tests/Client/SqsDriverTest.php b/pkg/sqs/Tests/Client/SqsDriverTest.php new file mode 100644 index 000000000..596219582 --- /dev/null +++ b/pkg/sqs/Tests/Client/SqsDriverTest.php @@ -0,0 +1,407 @@ +assertClassImplements(DriverInterface::class, SqsDriver::class); + } + + public function testCouldBeConstructedWithRequiredArguments() + { + new SqsDriver( + $this->createPsrContextMock(), + Config::create(), + $this->createQueueMetaRegistryMock() + ); + } + + public function testShouldReturnConfigObject() + { + $config = Config::create();; + + $driver = new SqsDriver($this->createPsrContextMock(), $config, $this->createQueueMetaRegistryMock()); + + $this->assertSame($config, $driver->getConfig()); + } + + public function testShouldCreateAndReturnQueueInstance() + { + $expectedQueue = new SqsDestination('aQueueName'); + + $context = $this->createPsrContextMock(); + $context + ->expects($this->once()) + ->method('createQueue') + ->with('aprefix_dot_afooqueue') + ->willReturn($expectedQueue) + ; + + $driver = new SqsDriver($context, $this->createDummyConfig(), $this->createDummyQueueMetaRegistry()); + $queue = $driver->createQueue('aFooQueue'); + + $this->assertSame($expectedQueue, $queue); + $this->assertSame('aQueueName', $queue->getQueueName()); + } + + public function testShouldCreateAndReturnQueueInstanceWithHardcodedTransportName() + { + $expectedQueue = new SqsDestination('aQueueName'); + + $context = $this->createPsrContextMock(); + $context + ->expects($this->once()) + ->method('createQueue') + ->with('aBarQueue') + ->willReturn($expectedQueue) + ; + + $driver = new SqsDriver($context, $this->createDummyConfig(), $this->createDummyQueueMetaRegistry()); + + $queue = $driver->createQueue('aBarQueue'); + $this->assertSame($expectedQueue, $queue); + } + + public function testShouldConvertTransportMessageToClientMessage() + { + $transportMessage = new SqsMessage(); + $transportMessage->setBody('body'); + $transportMessage->setHeaders(['hkey' => 'hval']); + $transportMessage->setProperties(['key' => 'val']); + $transportMessage->setHeader('content_type', 'ContentType'); + $transportMessage->setMessageId('MessageId'); + $transportMessage->setTimestamp(1000); + $transportMessage->setReplyTo('theReplyTo'); + $transportMessage->setCorrelationId('theCorrelationId'); + $transportMessage->setReplyTo('theReplyTo'); + $transportMessage->setCorrelationId('theCorrelationId'); + + $driver = new SqsDriver( + $this->createPsrContextMock(), + Config::create(), + $this->createQueueMetaRegistryMock() + ); + + $clientMessage = $driver->createClientMessage($transportMessage); + + $this->assertInstanceOf(Message::class, $clientMessage); + $this->assertSame('body', $clientMessage->getBody()); + $this->assertSame([ + 'hkey' => 'hval', + 'content_type' => 'ContentType', + 'message_id' => 'MessageId', + 'timestamp' => 1000, + 'reply_to' => 'theReplyTo', + 'correlation_id' => 'theCorrelationId', + ], $clientMessage->getHeaders()); + $this->assertSame([ + 'key' => 'val', + ], $clientMessage->getProperties()); + $this->assertSame('MessageId', $clientMessage->getMessageId()); + $this->assertSame('ContentType', $clientMessage->getContentType()); + $this->assertSame(1000, $clientMessage->getTimestamp()); + $this->assertSame('theReplyTo', $clientMessage->getReplyTo()); + $this->assertSame('theCorrelationId', $clientMessage->getCorrelationId()); + + $this->assertNull($clientMessage->getExpire()); + $this->assertSame(MessagePriority::NORMAL, $clientMessage->getPriority()); + } + + public function testShouldConvertClientMessageToTransportMessage() + { + $clientMessage = new Message(); + $clientMessage->setBody('body'); + $clientMessage->setHeaders(['hkey' => 'hval']); + $clientMessage->setProperties(['key' => 'val']); + $clientMessage->setContentType('ContentType'); + $clientMessage->setExpire(123); + $clientMessage->setPriority(MessagePriority::VERY_HIGH); + $clientMessage->setMessageId('MessageId'); + $clientMessage->setTimestamp(1000); + $clientMessage->setReplyTo('theReplyTo'); + $clientMessage->setCorrelationId('theCorrelationId'); + + $context = $this->createPsrContextMock(); + $context + ->expects($this->once()) + ->method('createMessage') + ->willReturn(new SqsMessage()) + ; + + $driver = new SqsDriver( + $context, + Config::create(), + $this->createQueueMetaRegistryMock() + ); + + $transportMessage = $driver->createTransportMessage($clientMessage); + + $this->assertInstanceOf(SqsMessage::class, $transportMessage); + $this->assertSame('body', $transportMessage->getBody()); + $this->assertSame([ + 'hkey' => 'hval', + 'content_type' => 'ContentType', + 'message_id' => 'MessageId', + 'timestamp' => 1000, + 'reply_to' => 'theReplyTo', + 'correlation_id' => 'theCorrelationId', + ], $transportMessage->getHeaders()); + $this->assertSame([ + 'key' => 'val', + ], $transportMessage->getProperties()); + $this->assertSame('MessageId', $transportMessage->getMessageId()); + $this->assertSame(1000, $transportMessage->getTimestamp()); + $this->assertSame('theReplyTo', $transportMessage->getReplyTo()); + $this->assertSame('theCorrelationId', $transportMessage->getCorrelationId()); + } + + public function testShouldSendMessageToRouterQueue() + { + $topic = new SqsDestination('aDestinationName'); + $transportMessage = new SqsMessage(); + $config = $this->createConfigMock(); + + $producer = $this->createPsrProducerMock(); + $producer + ->expects($this->once()) + ->method('send') + ->with($this->identicalTo($topic), $this->identicalTo($transportMessage)) + ; + $context = $this->createPsrContextMock(); + $context + ->expects($this->once()) + ->method('createQueue') + ->with('theTransportName') + ->willReturn($topic) + ; + $context + ->expects($this->once()) + ->method('createProducer') + ->willReturn($producer) + ; + $context + ->expects($this->once()) + ->method('createMessage') + ->willReturn($transportMessage) + ; + + $meta = $this->createQueueMetaRegistryMock(); + $meta + ->expects($this->once()) + ->method('getQueueMeta') + ->willReturn(new QueueMeta('theClientName', 'theTransportName')) + ; + + $driver = new SqsDriver($context, $config, $meta); + + $message = new Message(); + $message->setProperty(Config::PARAMETER_TOPIC_NAME, 'topic'); + + $driver->sendToRouter($message); + } + + public function testShouldThrowExceptionIfTopicParameterIsNotSet() + { + $driver = new SqsDriver( + $this->createPsrContextMock(), + Config::create(), + $this->createQueueMetaRegistryMock() + ); + + $this->expectException(\LogicException::class); + $this->expectExceptionMessage('Topic name parameter is required but is not set'); + + $driver->sendToRouter(new Message()); + } + + public function testShouldSendMessageToProcessor() + { + $queue = new SqsDestination('aDestinationName'); + $transportMessage = new SqsMessage(); + + $producer = $this->createPsrProducerMock(); + $producer + ->expects($this->once()) + ->method('send') + ->with($this->identicalTo($queue), $this->identicalTo($transportMessage)) + ; + $context = $this->createPsrContextMock(); + $context + ->expects($this->once()) + ->method('createQueue') + ->willReturn($queue) + ; + $context + ->expects($this->once()) + ->method('createProducer') + ->willReturn($producer) + ; + $context + ->expects($this->once()) + ->method('createMessage') + ->willReturn($transportMessage) + ; + + $meta = $this->createQueueMetaRegistryMock(); + $meta + ->expects($this->once()) + ->method('getQueueMeta') + ->willReturn(new QueueMeta('theClientName', 'theTransportName')) + ; + + $driver = new SqsDriver($context, Config::create(), $meta); + + $message = new Message(); + $message->setProperty(Config::PARAMETER_PROCESSOR_NAME, 'processor'); + $message->setProperty(Config::PARAMETER_PROCESSOR_QUEUE_NAME, 'queue'); + + $driver->sendToProcessor($message); + } + + public function testShouldThrowExceptionIfProcessorNameParameterIsNotSet() + { + $driver = new SqsDriver( + $this->createPsrContextMock(), + Config::create(), + $this->createQueueMetaRegistryMock() + ); + + $this->expectException(\LogicException::class); + $this->expectExceptionMessage('Processor name parameter is required but is not set'); + + $driver->sendToProcessor(new Message()); + } + + public function testShouldThrowExceptionIfProcessorQueueNameParameterIsNotSet() + { + $driver = new SqsDriver( + $this->createPsrContextMock(), + Config::create(), + $this->createQueueMetaRegistryMock() + ); + + $this->expectException(\LogicException::class); + $this->expectExceptionMessage('Queue name parameter is required but is not set'); + + $message = new Message(); + $message->setProperty(Config::PARAMETER_PROCESSOR_NAME, 'processor'); + + $driver->sendToProcessor($message); + } + + public function testShouldSetupBroker() + { + $routerQueue = new SqsDestination(''); + $processorQueue = new SqsDestination(''); + + $context = $this->createPsrContextMock(); + // setup router + $context + ->expects($this->at(0)) + ->method('createQueue') + ->willReturn($routerQueue) + ; + $context + ->expects($this->at(1)) + ->method('declareQueue') + ->with($this->identicalTo($routerQueue)) + ; + // setup processor queue + $context + ->expects($this->at(2)) + ->method('createQueue') + ->willReturn($processorQueue) + ; + $context + ->expects($this->at(3)) + ->method('declareQueue') + ->with($this->identicalTo($processorQueue)) + ; + + $metaRegistry = $this->createQueueMetaRegistryMock(); + $metaRegistry + ->expects($this->once()) + ->method('getQueuesMeta') + ->willReturn([new QueueMeta('theClientName', 'theTransportName')]) + ; + $metaRegistry + ->expects($this->exactly(2)) + ->method('getQueueMeta') + ->willReturn(new QueueMeta('theClientName', 'theTransportName')) + ; + + $driver = new SqsDriver($context, $this->createConfigMock(), $metaRegistry); + + $driver->setupBroker(); + } + + /** + * @return \PHPUnit_Framework_MockObject_MockObject|SqsContext + */ + private function createPsrContextMock() + { + return $this->createMock(SqsContext::class); + } + + /** + * @return \PHPUnit_Framework_MockObject_MockObject|PsrProducer + */ + private function createPsrProducerMock() + { + return $this->createMock(PsrProducer::class); + } + + /** + * @return \PHPUnit_Framework_MockObject_MockObject|QueueMetaRegistry + */ + private function createQueueMetaRegistryMock() + { + return $this->createMock(QueueMetaRegistry::class); + } + + /** + * @return \PHPUnit_Framework_MockObject_MockObject|Config + */ + private function createConfigMock() + { + return $this->createMock(Config::class); + } + + /** + * @return Config + */ + private function createDummyConfig() + { + return Config::create('aPrefix'); + } + + /** + * @return QueueMetaRegistry + */ + private function createDummyQueueMetaRegistry() + { + $registry = new QueueMetaRegistry($this->createDummyConfig(), []); + $registry->add('default'); + $registry->add('aFooQueue'); + $registry->add('aBarQueue', 'aBarQueue'); + + return $registry; + } +} diff --git a/pkg/sqs/Tests/Functional/SqsCommonUseCasesTest.php b/pkg/sqs/Tests/Functional/SqsCommonUseCasesTest.php new file mode 100644 index 000000000..7ff1db9ac --- /dev/null +++ b/pkg/sqs/Tests/Functional/SqsCommonUseCasesTest.php @@ -0,0 +1,124 @@ +context = $this->buildSqsContext(); + + $this->queueName = str_replace('.', '_dot_', uniqid('enqueue_test_queue_', true));; + $this->queue = $this->context->createQueue($this->queueName); + + $this->context->declareQueue($this->queue); + } + + protected function tearDown() + { + parent::tearDown(); + + if ($this->context && $this->queue) { + $this->context->deleteQueue($this->queue); + } + } + + public function testWaitsForTwoSecondsAndReturnNullOnReceive() + { + $queue = $this->context->createQueue($this->queueName); + + $startAt = microtime(true); + + $consumer = $this->context->createConsumer($queue); + $message = $consumer->receive(2000); + + $endAt = microtime(true); + + $this->assertNull($message); + + $this->assertGreaterThan(1.5, $endAt - $startAt); + $this->assertLessThan(2.5, $endAt - $startAt); + } + + public function testReturnNullImmediatelyOnReceiveNoWait() + { + $queue = $this->context->createQueue($this->queueName); + + $startAt = microtime(true); + + $consumer = $this->context->createConsumer($queue); + $message = $consumer->receiveNoWait(); + + $endAt = microtime(true); + + $this->assertNull($message); + + $this->assertLessThan(2, $endAt - $startAt); + } + + public function testProduceAndReceiveOneMessageSentDirectlyToQueue() + { + $queue = $this->context->createQueue($this->queueName); + + $message = $this->context->createMessage( + __METHOD__, + ['FooProperty' => 'FooVal'], + ['BarHeader' => 'BarVal'] + ); + + $producer = $this->context->createProducer(); + $producer->send($queue, $message); + + $consumer = $this->context->createConsumer($queue); + $message = $consumer->receive(1000); + + $this->assertInstanceOf(SqsMessage::class, $message); + $consumer->acknowledge($message); + + $this->assertEquals(__METHOD__, $message->getBody()); + $this->assertEquals(['FooProperty' => 'FooVal'], $message->getProperties()); + $this->assertEquals(['BarHeader' => 'BarVal'], $message->getHeaders()); + } + + public function testProduceAndReceiveOneMessageSentDirectlyToTopic() + { + $topic = $this->context->createTopic($this->queueName); + + $message = $this->context->createMessage(__METHOD__); + + $producer = $this->context->createProducer(); + $producer->send($topic, $message); + + $consumer = $this->context->createConsumer($topic); + $message = $consumer->receive(1000); + + $this->assertInstanceOf(SqsMessage::class, $message); + $consumer->acknowledge($message); + + $this->assertEquals(__METHOD__, $message->getBody()); + } +} diff --git a/pkg/sqs/Tests/Functional/SqsConsumptionUseCasesTest.php b/pkg/sqs/Tests/Functional/SqsConsumptionUseCasesTest.php new file mode 100644 index 000000000..4943084da --- /dev/null +++ b/pkg/sqs/Tests/Functional/SqsConsumptionUseCasesTest.php @@ -0,0 +1,112 @@ +context = $this->buildSqsContext(); + + $queue = $this->context->createQueue('enqueue_test_queue'); + $replyQueue = $this->context->createQueue('enqueue_test_queue_reply'); + + $this->context->declareQueue($queue); + $this->context->declareQueue($replyQueue); + + try { + $this->context->purge($queue); + $this->context->purge($replyQueue); + } catch (\Exception $e) {} + } + + public function testConsumeOneMessageAndExit() + { + $queue = $this->context->createQueue('enqueue_test_queue'); + + $message = $this->context->createMessage(__METHOD__); + $this->context->createProducer()->send($queue, $message); + + $queueConsumer = new QueueConsumer($this->context, new ChainExtension([ + new LimitConsumedMessagesExtension(1), + new LimitConsumptionTimeExtension(new \DateTime('+3sec')), + ])); + + $processor = new StubProcessor(); + $queueConsumer->bind($queue, $processor); + + $queueConsumer->consume(); + + $this->assertInstanceOf(PsrMessage::class, $processor->lastProcessedMessage); + $this->assertEquals(__METHOD__, $processor->lastProcessedMessage->getBody()); + } + + public function testConsumeOneMessageAndSendReplyExit() + { + $queue = $this->context->createQueue('enqueue_test_queue'); + $replyQueue = $this->context->createQueue('enqueue_test_queue_reply'); + + $message = $this->context->createMessage(__METHOD__); + $message->setReplyTo($replyQueue->getQueueName()); + $this->context->createProducer()->send($queue, $message); + + $queueConsumer = new QueueConsumer($this->context, new ChainExtension([ + new LimitConsumedMessagesExtension(2), + new LimitConsumptionTimeExtension(new \DateTime('+3sec')), + new ReplyExtension(), + ])); + + $replyMessage = $this->context->createMessage(__METHOD__.'.reply'); + + $processor = new StubProcessor(); + $processor->result = Result::reply($replyMessage); + + $replyProcessor = new StubProcessor(); + + $queueConsumer->bind($queue, $processor); + $queueConsumer->bind($replyQueue, $replyProcessor); + $queueConsumer->consume(); + + $this->assertInstanceOf(PsrMessage::class, $processor->lastProcessedMessage); + $this->assertEquals(__METHOD__, $processor->lastProcessedMessage->getBody()); + + $this->assertInstanceOf(PsrMessage::class, $replyProcessor->lastProcessedMessage); + $this->assertEquals(__METHOD__.'.reply', $replyProcessor->lastProcessedMessage->getBody()); + } +} + +class StubProcessor implements PsrProcessor +{ + public $result = self::ACK; + + /** @var PsrMessage */ + public $lastProcessedMessage; + + public function process(PsrMessage $message, PsrContext $context) + { + $this->lastProcessedMessage = $message; + + return $this->result; + } +} diff --git a/pkg/sqs/Tests/SqsConnectionFactoryTest.php b/pkg/sqs/Tests/SqsConnectionFactoryTest.php new file mode 100644 index 000000000..05d1af16d --- /dev/null +++ b/pkg/sqs/Tests/SqsConnectionFactoryTest.php @@ -0,0 +1,59 @@ +assertClassImplements(PsrConnectionFactory::class, SqsConnectionFactory::class); + } + + public function testCouldBeConstructedWithEmptyConfiguration() + { + $factory = new SqsConnectionFactory([]); + + $this->assertAttributeEquals([ + 'lazy' => true, + 'key' => null, + 'secret' => null, + 'token' => null, + 'region' => null, + 'retries' => 3, + 'version' => '2012-11-05', + ], 'config', $factory); + } + + public function testCouldBeConstructedWithCustomConfiguration() + { + $factory = new SqsConnectionFactory(['key' => 'theKey']); + + $this->assertAttributeEquals([ + 'lazy' => true, + 'key' => 'theKey', + 'secret' => null, + 'token' => null, + 'region' => null, + 'retries' => 3, + 'version' => '2012-11-05', + ], 'config', $factory); + } + + public function testShouldCreateLazyContext() + { + $factory = new SqsConnectionFactory(['lazy' => true]); + + $context = $factory->createContext(); + + $this->assertInstanceOf(SqsContext::class, $context); + + $this->assertAttributeEquals(null, 'client', $context); + $this->assertInternalType('callable', $this->readAttribute($context, 'clientFactory')); + } +} diff --git a/pkg/sqs/Tests/SqsConsumerTest.php b/pkg/sqs/Tests/SqsConsumerTest.php new file mode 100644 index 000000000..23b0bc68f --- /dev/null +++ b/pkg/sqs/Tests/SqsConsumerTest.php @@ -0,0 +1,286 @@ +assertClassImplements(PsrConsumer::class, SqsConsumer::class); + } + + public function testCouldBeConstructedWithRequiredArguments() + { + new SqsConsumer($this->createContextMock(), new SqsDestination('queue')); + } + + public function testShouldReturnInstanceOfDestination() + { + $destination = new SqsDestination('queue'); + + $consumer = new SqsConsumer($this->createContextMock(), $destination); + + $this->assertSame($destination, $consumer->getQueue()); + } + + public function testAcknowledgeShouldThrowIfInstanceOfMessageIsInvalid() + { + $this->expectException(InvalidMessageException::class); + $this->expectExceptionMessage('The message must be an instance of Enqueue\Sqs\SqsMessage but it is Mock_PsrMessage'); + + $consumer = new SqsConsumer($this->createContextMock(), new SqsDestination('queue')); + $consumer->acknowledge($this->createMock(PsrMessage::class)); + } + + public function testCouldAcknowledgeMessage() + { + $client = $this->createSqsClientMock(); + $client + ->expects($this->once()) + ->method('deleteMessage') + ->with($this->identicalTo(['QueueUrl' => 'theQueueUrl', 'ReceiptHandle' => 'theReceipt'])) + ; + + $context = $this->createContextMock(); + $context + ->expects($this->once()) + ->method('getClient') + ->willReturn($client) + ; + $context + ->expects($this->once()) + ->method('getQueueUrl') + ->willReturn('theQueueUrl') + ; + + $message = new SqsMessage(); + $message->setReceiptHandle('theReceipt'); + + $consumer = new SqsConsumer($context, new SqsDestination('queue')); + $consumer->acknowledge($message); + } + + public function testRejectShouldThrowIfInstanceOfMessageIsInvalid() + { + $this->expectException(InvalidMessageException::class); + $this->expectExceptionMessage('The message must be an instance of Enqueue\Sqs\SqsMessage but it is Mock_PsrMessage'); + + $consumer = new SqsConsumer($this->createContextMock(), new SqsDestination('queue')); + $consumer->reject($this->createMock(PsrMessage::class)); + } + + public function testShouldRejectMessage() + { + $client = $this->createSqsClientMock(); + $client + ->expects($this->once()) + ->method('deleteMessage') + ->with($this->identicalTo(['QueueUrl' => 'theQueueUrl', 'ReceiptHandle' => 'theReceipt'])) + ; + + $context = $this->createContextMock(); + $context + ->expects($this->once()) + ->method('getClient') + ->willReturn($client) + ; + $context + ->expects($this->once()) + ->method('getQueueUrl') + ->willReturn('theQueueUrl') + ; + $context + ->expects($this->never()) + ->method('createProducer') + ; + + $message = new SqsMessage(); + $message->setReceiptHandle('theReceipt'); + + $consumer = new SqsConsumer($context, new SqsDestination('queue')); + $consumer->reject($message); + } + + public function testShouldRejectMessageAndRequeue() + { + $client = $this->createSqsClientMock(); + $client + ->expects($this->once()) + ->method('deleteMessage') + ->with($this->identicalTo(['QueueUrl' => 'theQueueUrl', 'ReceiptHandle' => 'theReceipt'])) + ; + + $message = new SqsMessage(); + $message->setReceiptHandle('theReceipt'); + + $destination = new SqsDestination('queue'); + + $producer = $this->createProducerMock(); + $producer + ->expects($this->once()) + ->method('send') + ->with($this->identicalTo($destination), $this->identicalTo($message)) + ; + + $context = $this->createContextMock(); + $context + ->expects($this->once()) + ->method('getClient') + ->willReturn($client) + ; + $context + ->expects($this->once()) + ->method('getQueueUrl') + ->willReturn('theQueueUrl') + ; + $context + ->expects($this->once()) + ->method('createProducer') + ->willReturn($producer) + ; + + $consumer = new SqsConsumer($context, $destination); + $consumer->reject($message, true); + } + + public function testShouldReceiveMessage() + { + $expectedAttributes = [ + 'AttributeNames' => ['All'], + 'MessageAttributeNames' => ['All'], + 'MaxNumberOfMessages' => 1, + 'QueueUrl' => 'theQueueUrl', + 'WaitTimeSeconds' => 0, + ]; + + $expectedSqsMessage = [ + 'Body' => 'The Body', + 'ReceiptHandle' => 'The Receipt', + 'Attributes' => [ + 'ApproximateReceiveCount' => 3, + ], + 'MessageAttributes' => [ + 'Headers' => [ + 'StringValue' => json_encode([['hkey' => 'hvalue'], ['key' => 'value']]), + 'DataType' => 'String' + ], + ] + ]; + + $client = $this->createSqsClientMock(); + $client + ->expects($this->once()) + ->method('receiveMessage') + ->with($this->identicalTo($expectedAttributes)) + ->willReturn(new Result(['Messages' => [$expectedSqsMessage]])) + ; + + $context = $this->createContextMock(); + $context + ->expects($this->once()) + ->method('getClient') + ->willReturn($client) + ; + $context + ->expects($this->once()) + ->method('getQueueUrl') + ->willReturn('theQueueUrl') + ; + $context + ->expects($this->once()) + ->method('createMessage') + ->willReturn(new SqsMessage()) + ; + + $consumer = new SqsConsumer($context, new SqsDestination('queue')); + $result = $consumer->receiveNoWait(); + + $this->assertInstanceOf(SqsMessage::class, $result); + $this->assertEquals('The Body', $result->getBody()); + $this->assertEquals(['hkey' => 'hvalue'], $result->getHeaders()); + $this->assertEquals(['key' => 'value'], $result->getProperties()); + $this->assertTrue($result->isRedelivered()); + $this->assertEquals('The Receipt', $result->getReceiptHandle()); + } + + public function testShouldReturnNullIfThereIsNoNewMessage() + { + $expectedAttributes = [ + 'AttributeNames' => ['All'], + 'MessageAttributeNames' => ['All'], + 'MaxNumberOfMessages' => 1, + 'QueueUrl' => 'theQueueUrl', + 'WaitTimeSeconds' => 10, + ]; + + $client = $this->createSqsClientMock(); + $client + ->expects($this->once()) + ->method('receiveMessage') + ->with($this->identicalTo($expectedAttributes)) + ->willReturn(new Result()) + ; + + $context = $this->createContextMock(); + $context + ->expects($this->once()) + ->method('getClient') + ->willReturn($client) + ; + $context + ->expects($this->once()) + ->method('getQueueUrl') + ->willReturn('theQueueUrl') + ; + $context + ->expects($this->never()) + ->method('createMessage') + ; + + $consumer = new SqsConsumer($context, new SqsDestination('queue')); + $result = $consumer->receive(10000); + + $this->assertNull($result); + } + + /** + * @return \PHPUnit_Framework_MockObject_MockObject|SqsProducer + */ + private function createProducerMock() + { + return $this->createMock(SqsProducer::class); + } + + /** + * @return \PHPUnit_Framework_MockObject_MockObject|SqsClient + */ + private function createSqsClientMock() + { + return $this->getMockBuilder(SqsClient::class) + ->disableOriginalConstructor() + ->setMethods(['deleteMessage', 'receiveMessage']) + ->getMock() + ; + } + + /** + * @return \PHPUnit_Framework_MockObject_MockObject|SqsContext + */ + private function createContextMock() + { + return $this->createMock(SqsContext::class); + } +} diff --git a/pkg/sqs/Tests/SqsContextTest.php b/pkg/sqs/Tests/SqsContextTest.php new file mode 100644 index 000000000..faf350cf8 --- /dev/null +++ b/pkg/sqs/Tests/SqsContextTest.php @@ -0,0 +1,236 @@ +assertClassImplements(PsrContext::class, SqsContext::class); + } + + public function testCouldBeConstructedWithSqsClientAsFirstArgument() + { + new SqsContext($this->createSqsClientMock()); + } + + public function testCouldBeConstructedWithSqsClientFactoryAsFirstArgument() + { + new SqsContext(function() { + return $this->createSqsClientMock(); + }); + } + + public function testThrowIfNeitherSqsClientNorFactoryGiven() + { + $this->expectException(\InvalidArgumentException::class); + $this->expectExceptionMessage('The $client argument must be either Aws\Sqs\SqsClient or callable that returns Aws\Sqs\SqsClient once called.'); + new SqsContext(new \stdClass()); + } + + public function testShouldAllowCreateEmptyMessage() + { + $context = new SqsContext($this->createSqsClientMock()); + + $message = $context->createMessage(); + + $this->assertInstanceOf(SqsMessage::class, $message); + + $this->assertSame('', $message->getBody()); + $this->assertSame([], $message->getProperties()); + $this->assertSame([], $message->getHeaders()); + } + + public function testShouldAllowCreateCustomMessage() + { + $context = new SqsContext($this->createSqsClientMock()); + + $message = $context->createMessage('theBody', ['aProp' => 'aPropVal'], ['aHeader' => 'aHeaderVal']); + + $this->assertInstanceOf(SqsMessage::class, $message); + + $this->assertSame('theBody', $message->getBody()); + $this->assertSame(['aProp' => 'aPropVal'], $message->getProperties()); + $this->assertSame(['aHeader' => 'aHeaderVal'], $message->getHeaders()); + } + + public function testShouldCreateQueue() + { + $context = new SqsContext($this->createSqsClientMock()); + + $queue = $context->createQueue('aQueue'); + + $this->assertInstanceOf(SqsDestination::class, $queue); + $this->assertSame('aQueue', $queue->getQueueName()); + } + + public function testShouldAllowCreateTopic() + { + $context = new SqsContext($this->createSqsClientMock()); + + $topic = $context->createTopic('aTopic'); + + $this->assertInstanceOf(SqsDestination::class, $topic); + $this->assertSame('aTopic', $topic->getTopicName()); + } + + public function testThrowNotImplementedOnCreateTmpQueueCall() + { + $context = new SqsContext($this->createSqsClientMock()); + + $this->expectException(\BadMethodCallException::class); + $this->expectExceptionMessage('SQS transport does not support temporary queues'); + $context->createTemporaryQueue(); + } + + public function testShouldCreateProducer() + { + $context = new SqsContext($this->createSqsClientMock()); + + $producer = $context->createProducer(); + + $this->assertInstanceOf(SqsProducer::class, $producer); + } + + public function testShouldThrowIfNotSqsDestinationGivenOnCreateConsumer() + { + $context = new SqsContext($this->createSqsClientMock()); + + $this->expectException(InvalidDestinationException::class); + $this->expectExceptionMessage('The destination must be an instance of Enqueue\Sqs\SqsDestination but got Mock_PsrQueue'); + + $context->createConsumer($this->createMock(PsrQueue::class)); + } + + public function testShouldCreateConsumer() + { + $context = new SqsContext($this->createSqsClientMock()); + + $queue = $context->createQueue('aQueue'); + + $consumer = $context->createConsumer($queue); + + $this->assertInstanceOf(SqsConsumer::class, $consumer); + } + + public function testShouldAllowDeclareQueue() + { + $sqsClient = $this->createSqsClientMock(); + $sqsClient + ->expects($this->once()) + ->method('createQueue') + ->with($this->identicalTo(['Attributes' => [], 'QueueName' => 'aQueueName'])) + ->willReturn(new Result(['QueueUrl' => 'theQueueUrl'])) + ; + + $context = new SqsContext($sqsClient); + + $queue = $context->createQueue('aQueueName'); + + $context->declareQueue($queue); + } + + public function testShouldAllowDeleteQueue() + { + $sqsClient = $this->createSqsClientMock(); + $sqsClient + ->expects($this->once()) + ->method('getQueueUrl') + ->with($this->identicalTo(['QueueName' => 'aQueueName'])) + ->willReturn(new Result(['QueueUrl' => 'theQueueUrl'])) + ; + $sqsClient + ->expects($this->once()) + ->method('deleteQueue') + ->with($this->identicalTo(['QueueUrl' => 'theQueueUrl'])) + ->willReturn(new Result()) + ; + + $context = new SqsContext($sqsClient); + + $queue = $context->createQueue('aQueueName'); + + $context->deleteQueue($queue); + } + + public function testShouldAllowPurgeQueue() + { + $sqsClient = $this->createSqsClientMock(); + $sqsClient + ->expects($this->once()) + ->method('getQueueUrl') + ->with($this->identicalTo(['QueueName' => 'aQueueName'])) + ->willReturn(new Result(['QueueUrl' => 'theQueueUrl'])) + ; + $sqsClient + ->expects($this->once()) + ->method('purgeQueue') + ->with($this->identicalTo(['QueueUrl' => 'theQueueUrl'])) + ->willReturn(new Result()) + ; + + $context = new SqsContext($sqsClient); + + $queue = $context->createQueue('aQueueName'); + + $context->purge($queue); + } + + public function testShouldAllowGetQueueUrl() + { + $sqsClient = $this->createSqsClientMock(); + $sqsClient + ->expects($this->once()) + ->method('getQueueUrl') + ->with($this->identicalTo(['QueueName' => 'aQueueName'])) + ->willReturn(new Result(['QueueUrl' => 'theQueueUrl'])) + ; + + $context = new SqsContext($sqsClient); + + $context->getQueueUrl(new SqsDestination('aQueueName')); + } + + public function testShouldThrowExceptionIfGetQueueUrlResultHasNoQueueUrlProperty() + { + $sqsClient = $this->createSqsClientMock(); + $sqsClient + ->expects($this->once()) + ->method('getQueueUrl') + ->with($this->identicalTo(['QueueName' => 'aQueueName'])) + ->willReturn(new Result([])) + ; + + $context = new SqsContext($sqsClient); + + $this->expectException(\RuntimeException::class); + $this->expectExceptionMessage('QueueUrl cannot be resolved. queueName: "aQueueName"'); + + $context->getQueueUrl(new SqsDestination('aQueueName')); + } + + /** + * @return \PHPUnit_Framework_MockObject_MockObject|SqsClient + */ + private function createSqsClientMock() + { + return $this->getMockBuilder(SqsClient::class) + ->disableOriginalConstructor() + ->setMethods(['deleteQueue', 'purgeQueue', 'createQueue', 'getQueueUrl']) + ->getMock() + ; + } +} diff --git a/pkg/sqs/Tests/SqsDestinationTest.php b/pkg/sqs/Tests/SqsDestinationTest.php new file mode 100644 index 000000000..7cf63bdc3 --- /dev/null +++ b/pkg/sqs/Tests/SqsDestinationTest.php @@ -0,0 +1,104 @@ +assertClassImplements(PsrTopic::class, SqsDestination::class); + $this->assertClassImplements(PsrQueue::class, SqsDestination::class); + } + + public function testShouldReturnNameSetInConstructor() + { + $destination = new SqsDestination('aDestinationName'); + + $this->assertSame('aDestinationName', $destination->getQueueName()); + $this->assertSame('aDestinationName', $destination->getTopicName()); + } + + public function testCouldSetDelaySecondsAttribute() + { + $destination = new SqsDestination('aDestinationName'); + $destination->setDelaySeconds(12345); + + $this->assertSame(['DelaySeconds' => 12345], $destination->getAttributes()); + } + + public function testCouldSetMaximumMessageSizeAttribute() + { + $destination = new SqsDestination('aDestinationName'); + $destination->setMaximumMessageSize(12345); + + $this->assertSame(['MaximumMessageSize' => 12345], $destination->getAttributes()); + } + + public function testCouldSetMessageRetentionPeriodAttribute() + { + $destination = new SqsDestination('aDestinationName'); + $destination->setMessageRetentionPeriod(12345); + + $this->assertSame(['MessageRetentionPeriod' => 12345], $destination->getAttributes()); + } + + public function testCouldSetPolicyAttribute() + { + $destination = new SqsDestination('aDestinationName'); + $destination->setPolicy('thePolicy'); + + $this->assertSame(['Policy' => 'thePolicy'], $destination->getAttributes()); + } + + public function testCouldSetReceiveMessageWaitTimeSecondsAttribute() + { + $destination = new SqsDestination('aDestinationName'); + $destination->setReceiveMessageWaitTimeSeconds(12345); + + $this->assertSame(['ReceiveMessageWaitTimeSeconds' => 12345], $destination->getAttributes()); + } + + public function testCouldSetRedrivePolicyAttribute() + { + $destination = new SqsDestination('aDestinationName'); + $destination->setRedrivePolicy(12345, 'theDeadQueueArn'); + + $this->assertSame(['RedrivePolicy' => '{"maxReceiveCount":"12345","deadLetterTargetArn":"theDeadQueueArn"}'], $destination->getAttributes()); + } + + public function testCouldSetVisibilityTimeoutAttribute() + { + $destination = new SqsDestination('aDestinationName'); + $destination->setVisibilityTimeout(12345); + + $this->assertSame(['VisibilityTimeout' => 12345], $destination->getAttributes()); + } + + public function testCouldSetFifoQueueAttributeAndUnsetIt() + { + $destination = new SqsDestination('aDestinationName'); + + $destination->setFifoQueue(true); + $this->assertSame(['FifoQueue' => 'true'], $destination->getAttributes()); + + $destination->setFifoQueue(false); + $this->assertSame([], $destination->getAttributes()); + } + + public function testCouldSetContentBasedDeduplicationAttributeAndUnsetIt() + { + $destination = new SqsDestination('aDestinationName'); + + $destination->setContentBasedDeduplication(true); + $this->assertSame(['ContentBasedDeduplication' => 'true'], $destination->getAttributes()); + + $destination->setContentBasedDeduplication(false); + $this->assertSame([], $destination->getAttributes()); + } +} diff --git a/pkg/sqs/Tests/SqsMessageTest.php b/pkg/sqs/Tests/SqsMessageTest.php new file mode 100644 index 000000000..dc3eb0472 --- /dev/null +++ b/pkg/sqs/Tests/SqsMessageTest.php @@ -0,0 +1,175 @@ +assertClassImplements(PsrMessage::class, SqsMessage::class); + } + + public function testCouldConstructMessageWithBody() + { + $message = new SqsMessage('body'); + + $this->assertSame('body', $message->getBody()); + } + + public function testCouldConstructMessageWithProperties() + { + $message = new SqsMessage('', ['key' => 'value']); + + $this->assertSame(['key' => 'value'], $message->getProperties()); + } + + public function testCouldConstructMessageWithHeaders() + { + $message = new SqsMessage('', [], ['key' => 'value']); + + $this->assertSame(['key' => 'value'], $message->getHeaders()); + } + + public function testCouldSetGetBody() + { + $message = new SqsMessage(); + $message->setBody('body'); + + $this->assertSame('body', $message->getBody()); + } + + public function testCouldSetGetProperties() + { + $message = new SqsMessage(); + $message->setProperties(['key' => 'value']); + + $this->assertSame(['key' => 'value'], $message->getProperties()); + } + + public function testCouldSetGetHeaders() + { + $message = new SqsMessage(); + $message->setHeaders(['key' => 'value']); + + $this->assertSame(['key' => 'value'], $message->getHeaders()); + } + + public function testCouldSetGetRedelivered() + { + $message = new SqsMessage(); + + $message->setRedelivered(true); + $this->assertTrue($message->isRedelivered()); + + $message->setRedelivered(false); + $this->assertFalse($message->isRedelivered()); + } + + public function testCouldSetGetCorrelationId() + { + $message = new SqsMessage(); + $message->setCorrelationId('the-correlation-id'); + + $this->assertSame('the-correlation-id', $message->getCorrelationId()); + } + + public function testShouldSetCorrelationIdAsHeader() + { + $message = new SqsMessage(); + $message->setCorrelationId('the-correlation-id'); + + $this->assertSame(['correlation_id' => 'the-correlation-id'], $message->getHeaders()); + } + + public function testCouldSetGetMessageId() + { + $message = new SqsMessage(); + $message->setMessageId('the-message-id'); + + $this->assertSame('the-message-id', $message->getMessageId()); + } + + public function testCouldSetMessageIdAsHeader() + { + $message = new SqsMessage(); + $message->setMessageId('the-message-id'); + + $this->assertSame(['message_id' => 'the-message-id'], $message->getHeaders()); + } + + public function testCouldSetGetTimestamp() + { + $message = new SqsMessage(); + $message->setTimestamp(12345); + + $this->assertSame(12345, $message->getTimestamp()); + } + + public function testCouldSetTimestampAsHeader() + { + $message = new SqsMessage(); + $message->setTimestamp(12345); + + $this->assertSame(['timestamp' => 12345], $message->getHeaders()); + } + + public function testShouldReturnNullAsDefaultReplyTo() + { + $message = new SqsMessage(); + + $this->assertSame(null, $message->getReplyTo()); + } + + public function testShouldAllowGetPreviouslySetReplyTo() + { + $message = new SqsMessage(); + $message->setReplyTo('theQueueName'); + + $this->assertSame('theQueueName', $message->getReplyTo()); + } + + public function testShouldAllowGetPreviouslySetReplyToAsHeader() + { + $message = new SqsMessage(); + $message->setReplyTo('theQueueName'); + + $this->assertSame(['reply_to' => 'theQueueName'], $message->getHeaders()); + } + + public function testShouldAllowGetDelaySeconds() + { + $message = new SqsMessage(); + $message->setDelaySeconds(12345); + + $this->assertSame(12345, $message->getDelaySeconds()); + } + + public function testShouldAllowGetMessageDeduplicationId() + { + $message = new SqsMessage(); + $message->setMessageDeduplicationId('theId'); + + $this->assertSame('theId', $message->getMessageDeduplicationId()); + } + + public function testShouldAllowGetMessageGroupId() + { + $message = new SqsMessage(); + $message->setMessageGroupId('theId'); + + $this->assertSame('theId', $message->getMessageGroupId()); + } + + public function testShouldAllowGetReceiptHandle() + { + $message = new SqsMessage(); + $message->setReceiptHandle('theId'); + + $this->assertSame('theId', $message->getReceiptHandle()); + } +} diff --git a/pkg/sqs/Tests/SqsProducerTest.php b/pkg/sqs/Tests/SqsProducerTest.php new file mode 100644 index 000000000..d982ac9ab --- /dev/null +++ b/pkg/sqs/Tests/SqsProducerTest.php @@ -0,0 +1,152 @@ +assertClassImplements(PsrProducer::class, SqsProducer::class); + } + + public function testCouldBeConstructedWithRequiredArguments() + { + new SqsProducer($this->createSqsContextMock()); + } + + public function testShouldThrowIfBodyOfInvalidType() + { + $this->expectException(InvalidMessageException::class); + $this->expectExceptionMessage('The message body must be a scalar or null. Got: stdClass'); + + $producer = new SqsProducer($this->createSqsContextMock()); + + $message = new SqsMessage(new \stdClass()); + + $producer->send(new SqsDestination(''), $message); + } + + public function testShouldThrowIfDestinationOfInvalidType() + { + $this->expectException(InvalidDestinationException::class); + $this->expectExceptionMessage('The destination must be an instance of Enqueue\Sqs\SqsDestination but got Mock_PsrDestinat'); + + $producer = new SqsProducer($this->createSqsContextMock()); + + $producer->send($this->createMock(PsrDestination::class), new SqsMessage()); + } + + public function testShouldThrowIfSendMessageFailed() + { + $client = $this->createSqsClientMock(); + $client + ->expects($this->once()) + ->method('sendMessage') + ->willReturn(new Result()) + ; + + $context = $this->createSqsContextMock(); + $context + ->expects($this->once()) + ->method('getQueueUrl') + ->willReturn('theQueueUrl') + ; + $context + ->expects($this->once()) + ->method('getClient') + ->will($this->returnValue($client)) + ; + + $destination = new SqsDestination('queue-name'); + $message = new SqsMessage(); + + $this->expectException(\RuntimeException::class); + $this->expectExceptionMessage('Message was not sent'); + + $producer = new SqsProducer($context); + $producer->send($destination, $message); + } + + public function testShouldSendMessage() + { + $expectedArguments = [ + 'MessageAttributes' => [ + 'Headers' => [ + 'DataType' => 'String', + 'StringValue' => '[{"hkey":"hvaleu"},{"key":"value"}]', + ], + ], + 'MessageBody' => 'theBody', + 'QueueUrl' => 'theQueueUrl', + 'DelaySeconds' => 12345, + 'MessageDeduplicationId' => 'theDeduplicationId', + 'MessageGroupId' => 'groupId', + ]; + + $client = $this->createSqsClientMock(); + $client + ->expects($this->once()) + ->method('sendMessage') + ->with($this->identicalTo($expectedArguments)) + ->willReturn(new Result()) + ; + + $context = $this->createSqsContextMock(); + $context + ->expects($this->once()) + ->method('getQueueUrl') + ->willReturn('theQueueUrl') + ; + $context + ->expects($this->once()) + ->method('getClient') + ->will($this->returnValue($client)) + ; + + $destination = new SqsDestination('queue-name'); + $message = new SqsMessage('theBody', ['key' => 'value'], ['hkey' => 'hvaleu']); + $message->setDelaySeconds(12345); + $message->setMessageDeduplicationId('theDeduplicationId'); + $message->setMessageGroupId('groupId'); + + $this->expectException(\RuntimeException::class); + $this->expectExceptionMessage('Message was not sent'); + + $producer = new SqsProducer($context); + $producer->send($destination, $message); + } + + /** + * @return \PHPUnit_Framework_MockObject_MockObject|SqsContext + */ + private function createSqsContextMock() + { + return $this->createMock(SqsContext::class); + } + + /** + * @return \PHPUnit_Framework_MockObject_MockObject|SqsClient + */ + private function createSqsClientMock() + { + return $this + ->getMockBuilder(SqsClient::class) + ->disableOriginalConstructor() + ->setMethods(['sendMessage']) + ->getMock() + ; + } +} diff --git a/pkg/sqs/Tests/Symfony/SqsTransportFactoryTest.php b/pkg/sqs/Tests/Symfony/SqsTransportFactoryTest.php new file mode 100644 index 000000000..08cb1755c --- /dev/null +++ b/pkg/sqs/Tests/Symfony/SqsTransportFactoryTest.php @@ -0,0 +1,130 @@ +assertClassImplements(TransportFactoryInterface::class, SqsTransportFactory::class); + } + + public function testCouldBeConstructedWithDefaultName() + { + $transport = new SqsTransportFactory(); + + $this->assertEquals('sqs', $transport->getName()); + } + + public function testCouldBeConstructedWithCustomName() + { + $transport = new SqsTransportFactory('theCustomName'); + + $this->assertEquals('theCustomName', $transport->getName()); + } + + public function testShouldAllowAddConfiguration() + { + $transport = new SqsTransportFactory(); + $tb = new TreeBuilder(); + $rootNode = $tb->root('foo'); + + $transport->addConfiguration($rootNode); + $processor = new Processor(); + $config = $processor->process($tb->buildTree(), [[ + 'key' => 'theKey', + 'secret' => 'theSecret', + 'token' => 'theToken', + 'region' => 'theRegion', + 'retries' => 5, + 'version' => 'theVersion', + 'lazy' => false, + ]]); + + $this->assertEquals([ + 'key' => 'theKey', + 'secret' => 'theSecret', + 'token' => 'theToken', + 'region' => 'theRegion', + 'retries' => 5, + 'version' => 'theVersion', + 'lazy' => false, + ], $config); + } + + public function testShouldCreateConnectionFactory() + { + $container = new ContainerBuilder(); + + $transport = new SqsTransportFactory(); + + $serviceId = $transport->createConnectionFactory($container, [ + 'key' => 'theKey', + 'secret' => 'theSecret', + ]); + + $this->assertTrue($container->hasDefinition($serviceId)); + $factory = $container->getDefinition($serviceId); + $this->assertEquals(SqsConnectionFactory::class, $factory->getClass()); + $this->assertSame([[ + 'key' => 'theKey', + 'secret' => 'theSecret', + ]], $factory->getArguments()); + } + + public function testShouldCreateContext() + { + $container = new ContainerBuilder(); + + $transport = new SqsTransportFactory(); + + $serviceId = $transport->createContext($container, [ + 'key' => 'theKey', + 'secret' => 'theSecret', + ]); + + $this->assertEquals('enqueue.transport.sqs.context', $serviceId); + $this->assertTrue($container->hasDefinition($serviceId)); + + $context = $container->getDefinition('enqueue.transport.sqs.context'); + $this->assertInstanceOf(Reference::class, $context->getFactory()[0]); + $this->assertEquals('enqueue.transport.sqs.connection_factory', (string) $context->getFactory()[0]); + $this->assertEquals('createContext', $context->getFactory()[1]); + } + + public function testShouldCreateDriver() + { + $container = new ContainerBuilder(); + + $transport = new SqsTransportFactory(); + + $serviceId = $transport->createDriver($container, []); + + $this->assertEquals('enqueue.client.sqs.driver', $serviceId); + $this->assertTrue($container->hasDefinition($serviceId)); + + $driver = $container->getDefinition($serviceId); + $this->assertSame(SqsDriver::class, $driver->getClass()); + + $this->assertInstanceOf(Reference::class, $driver->getArgument(0)); + $this->assertEquals('enqueue.transport.sqs.context', (string) $driver->getArgument(0)); + + $this->assertInstanceOf(Reference::class, $driver->getArgument(1)); + $this->assertEquals('enqueue.client.config', (string) $driver->getArgument(1)); + + $this->assertInstanceOf(Reference::class, $driver->getArgument(2)); + $this->assertEquals('enqueue.client.meta.queue_meta_registry', (string) $driver->getArgument(2)); + } +} diff --git a/pkg/sqs/composer.json b/pkg/sqs/composer.json new file mode 100644 index 000000000..5c948620d --- /dev/null +++ b/pkg/sqs/composer.json @@ -0,0 +1,41 @@ +{ + "name": "enqueue/sqs", + "type": "library", + "description": "Message Queue Amazon SQS Transport", + "keywords": ["messaging", "queue", "amazon", "aws", "sqs"], + "license": "MIT", + "repositories": [ + { + "type": "vcs", + "url": "git@github.com:php-enqueue/test.git" + } + ], + "require": { + "php": ">=5.6", + "enqueue/psr-queue": "^0.3", + "aws/aws-sdk-php": "~3.26", + "psr/log": "^1" + }, + "require-dev": { + "phpunit/phpunit": "~5.4.0", + "enqueue/test": "^0.3", + "enqueue/enqueue": "^0.3", + "symfony/dependency-injection": "^2.8|^3", + "symfony/config": "^2.8|^3" + }, + "autoload": { + "psr-4": { "Enqueue\\Sqs\\": "" }, + "exclude-from-classmap": [ + "/Tests/" + ] + }, + "suggest": { + "enqueue/enqueue": "If you'd like to use advanced features like Client abstract layer or Symfony integration features" + }, + "minimum-stability": "dev", + "extra": { + "branch-alias": { + "dev-master": "0.3.x-dev" + } + } +} diff --git a/pkg/sqs/examples/consume.php b/pkg/sqs/examples/consume.php new file mode 100644 index 000000000..a8914d165 --- /dev/null +++ b/pkg/sqs/examples/consume.php @@ -0,0 +1,39 @@ + getenv('AWS__SQS__KEY'), + 'secret' => getenv('AWS__SQS__SECRET'), + 'region' => getenv('AWS__SQS__REGION'), +]; + +$factory = new SqsConnectionFactory($config); +$context = $factory->createContext(); + +$queue = $context->createQueue('enqueue'); +$consumer = $context->createConsumer($queue); + +while (true) { + if ($m = $consumer->receive(20000)) { + $consumer->acknowledge($m); + echo 'Received message: '.$m->getBody().PHP_EOL; + } +} + +echo 'Done'."\n"; diff --git a/pkg/sqs/examples/produce.php b/pkg/sqs/examples/produce.php new file mode 100644 index 000000000..d08a40196 --- /dev/null +++ b/pkg/sqs/examples/produce.php @@ -0,0 +1,40 @@ + getenv('AWS__SQS__KEY'), + 'secret' => getenv('AWS__SQS__SECRET'), + 'region' => getenv('AWS__SQS__REGION'), +]; + +$factory = new SqsConnectionFactory($config); +$context = $factory->createContext(); + +$queue = $context->createQueue('enqueue'); +$message = $context->createMessage('Hello Bar!'); + +$context->declareQueue($queue); + +while (true) { + $context->createProducer()->send($queue, $message); + echo 'Sent message: ' . $message->getBody() . PHP_EOL; + sleep(1); +} + +echo 'Done'."\n"; diff --git a/pkg/sqs/phpunit.xml.dist b/pkg/sqs/phpunit.xml.dist new file mode 100644 index 000000000..7c026b4e6 --- /dev/null +++ b/pkg/sqs/phpunit.xml.dist @@ -0,0 +1,30 @@ + + + + + + + ./Tests + + + + + + . + + ./vendor + ./Tests + + + + diff --git a/pkg/test/SqsExtension.php b/pkg/test/SqsExtension.php new file mode 100644 index 000000000..1d58dbde8 --- /dev/null +++ b/pkg/test/SqsExtension.php @@ -0,0 +1,27 @@ + getenv('AWS__SQS__KEY'), + 'secret' => getenv('AWS__SQS__SECRET'), + 'region' => getenv('AWS__SQS__REGION'), + 'lazy' => false, + ]; + + return (new SqsConnectionFactory($config))->createContext(); + } +} \ No newline at end of file