vendor/gos/web-socket-bundle/src/Pusher/Amqp/AmqpConnectionFactory.php line 8

Open in your IDE?
  1. <?php declare(strict_types=1);
  2. namespace Gos\Bundle\WebSocketBundle\Pusher\Amqp;
  3. use Gos\Bundle\WebSocketBundle\Pusher\Exception\PusherUnsupportedException;
  4. use Symfony\Component\OptionsResolver\OptionsResolver;
  5. trigger_deprecation('gos/web-socket-bundle''3.1''The "%s" class is deprecated and will be removed in 4.0, use the symfony/messenger component instead.'AmqpConnectionFactory::class);
  6. /**
  7.  * @deprecated to be removed in 4.0, use the symfony/messenger component instead
  8.  */
  9. final class AmqpConnectionFactory implements AmqpConnectionFactoryInterface
  10. {
  11.     private array $config;
  12.     public function __construct(array $config)
  13.     {
  14.         $this->config $this->resolveConfig($config);
  15.     }
  16.     /**
  17.      * @throws PusherUnsupportedException if the AMQP pusher is not supported in the current environment
  18.      */
  19.     public function createConnection(): \AMQPConnection
  20.     {
  21.         if (!$this->isSupported()) {
  22.             throw new PusherUnsupportedException('The AMQP pusher requires the PHP amqp extension.');
  23.         }
  24.         return new \AMQPConnection($this->config);
  25.     }
  26.     public function createExchange(\AMQPConnection $connection): \AMQPExchange
  27.     {
  28.         $exchange = new \AMQPExchange(new \AMQPChannel($connection));
  29.         $exchange->setName($this->config['exchange_name']);
  30.         $exchange->setType(AMQP_EX_TYPE_DIRECT);
  31.         $exchange->setFlags(AMQP_DURABLE);
  32.         $exchange->declareExchange();
  33.         return $exchange;
  34.     }
  35.     public function createQueue(\AMQPConnection $connection): \AMQPQueue
  36.     {
  37.         $queue = new \AMQPQueue(new \AMQPChannel($connection));
  38.         $queue->setName($this->config['queue_name']);
  39.         $queue->setFlags(AMQP_DURABLE);
  40.         $queue->declareQueue();
  41.         return $queue;
  42.     }
  43.     public function isSupported(): bool
  44.     {
  45.         return \extension_loaded('amqp');
  46.     }
  47.     private function resolveConfig(array $config): array
  48.     {
  49.         $resolver = new OptionsResolver();
  50.         $resolver->setRequired(
  51.             [
  52.                 'host',
  53.                 'port',
  54.                 'login',
  55.                 'password',
  56.             ]
  57.         );
  58.         $resolver->setDefaults(
  59.             [
  60.                 'vhost' => '/',
  61.                 'read_timeout' => 0,
  62.                 'write_timeout' => 0,
  63.                 'connect_timeout' => 0,
  64.                 'queue_name' => 'gos_websocket',
  65.                 'exchange_name' => 'gos_websocket_exchange',
  66.             ]
  67.         );
  68.         $resolver->setAllowedTypes('host''string');
  69.         $resolver->setAllowedTypes('port', ['string''integer']);
  70.         $resolver->setAllowedTypes('login''string');
  71.         $resolver->setAllowedTypes('password''string');
  72.         $resolver->setAllowedTypes('vhost''string');
  73.         $resolver->setAllowedTypes('read_timeout''integer');
  74.         $resolver->setAllowedTypes('write_timeout''integer');
  75.         $resolver->setAllowedTypes('connect_timeout''integer');
  76.         $resolver->setAllowedTypes('queue_name''string');
  77.         $resolver->setAllowedTypes('exchange_name''string');
  78.         return $resolver->resolve($config);
  79.     }
  80. }