| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179 |
- <?php declare(strict_types=1);
- namespace mikemadisonweb\rabbitmq\components;
- use mikemadisonweb\rabbitmq\events\RabbitMQPublisherEvent;
- use mikemadisonweb\rabbitmq\exceptions\RuntimeException;
- use PhpAmqpLib\Message\AMQPMessage;
- use PhpAmqpLib\Wire\AMQPTable;
- /**
- * Service that sends AMQP Messages
- *
- * @package mikemadisonweb\rabbitmq\components
- */
- class Producer extends BaseRabbitMQ
- {
- protected $contentType;
- protected $deliveryMode;
- protected $serializer;
- protected $safe;
- protected $name = 'unnamed';
- /**
- * @param $contentType
- */
- public function setContentType($contentType)
- {
- $this->contentType = $contentType;
- }
- /**
- * @param $deliveryMode
- */
- public function setDeliveryMode($deliveryMode)
- {
- $this->deliveryMode = $deliveryMode;
- }
- /**
- * @param callable $serializer
- */
- public function setSerializer(callable $serializer)
- {
- $this->serializer = $serializer;
- }
- /**
- * @return callable
- */
- public function getSerializer(): callable
- {
- return $this->serializer;
- }
- /**
- * @return array
- */
- public function getBasicProperties(): array
- {
- return [
- 'content_type' => $this->contentType,
- 'delivery_mode' => $this->deliveryMode,
- ];
- }
- /**
- * @return mixed
- */
- public function getSafe(): bool
- {
- return $this->safe;
- }
- /**
- * @param mixed $safe
- */
- public function setSafe(bool $safe)
- {
- $this->safe = $safe;
- }
- /**
- * @return string
- */
- public function getName(): string
- {
- return $this->name;
- }
- /**
- * @param string $name
- */
- public function setName(string $name)
- {
- $this->name = $name;
- }
- /**
- * Publishes the message and merges additional properties with basic properties
- *
- * @param mixed $msgBody
- * @param string $exchangeName
- * @param string $routingKey
- * @param array $additionalProperties
- * @param array $headers
- *
- * @throws RuntimeException
- */
- public function publish(
- $msgBody,
- string $exchangeName,
- string $routingKey = '',
- array $additionalProperties = [],
- array $headers = null
- ) {
- if ($this->autoDeclare)
- {
- $this->routing->declareAll();
- }
- if ($this->safe && !$this->routing->isExchangeExists($exchangeName))
- {
- throw new RuntimeException(
- "Exchange `{$exchangeName}` does not declared in broker (You see this message because safe mode is ON)."
- );
- }
- $serialized = false;
- if (!is_string($msgBody))
- {
- $msgBody = call_user_func($this->serializer, $msgBody);
- $serialized = true;
- }
- $msg = new AMQPMessage($msgBody, array_merge($this->getBasicProperties(), $additionalProperties));
- if (!empty($headers) || $serialized)
- {
- if ($serialized)
- {
- $headers['rabbitmq.serialized'] = 1;
- }
- $headersTable = new AMQPTable($headers);
- $msg->set('application_headers', $headersTable);
- }
- \Yii::$app->rabbitmq->trigger(
- RabbitMQPublisherEvent::BEFORE_PUBLISH,
- new RabbitMQPublisherEvent(
- [
- 'message' => $msg,
- 'producer' => $this,
- ]
- )
- );
- $this->getChannel()->basic_publish($msg, $exchangeName, $routingKey);
- \Yii::$app->rabbitmq->trigger(
- RabbitMQPublisherEvent::AFTER_PUBLISH,
- new RabbitMQPublisherEvent(
- [
- 'message' => $msg,
- 'producer' => $this,
- ]
- )
- );
- $this->logger->log(
- 'AMQP message published',
- $msg,
- [
- 'exchange' => $exchangeName,
- 'routing_key' => $routingKey,
- ]
- );
- }
- }
|