Producer.php 4.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179
  1. <?php declare(strict_types=1);
  2. namespace mikemadisonweb\rabbitmq\components;
  3. use mikemadisonweb\rabbitmq\events\RabbitMQPublisherEvent;
  4. use mikemadisonweb\rabbitmq\exceptions\RuntimeException;
  5. use PhpAmqpLib\Message\AMQPMessage;
  6. use PhpAmqpLib\Wire\AMQPTable;
  7. /**
  8. * Service that sends AMQP Messages
  9. *
  10. * @package mikemadisonweb\rabbitmq\components
  11. */
  12. class Producer extends BaseRabbitMQ
  13. {
  14. protected $contentType;
  15. protected $deliveryMode;
  16. protected $serializer;
  17. protected $safe;
  18. protected $name = 'unnamed';
  19. /**
  20. * @param $contentType
  21. */
  22. public function setContentType($contentType)
  23. {
  24. $this->contentType = $contentType;
  25. }
  26. /**
  27. * @param $deliveryMode
  28. */
  29. public function setDeliveryMode($deliveryMode)
  30. {
  31. $this->deliveryMode = $deliveryMode;
  32. }
  33. /**
  34. * @param callable $serializer
  35. */
  36. public function setSerializer(callable $serializer)
  37. {
  38. $this->serializer = $serializer;
  39. }
  40. /**
  41. * @return callable
  42. */
  43. public function getSerializer(): callable
  44. {
  45. return $this->serializer;
  46. }
  47. /**
  48. * @return array
  49. */
  50. public function getBasicProperties(): array
  51. {
  52. return [
  53. 'content_type' => $this->contentType,
  54. 'delivery_mode' => $this->deliveryMode,
  55. ];
  56. }
  57. /**
  58. * @return mixed
  59. */
  60. public function getSafe(): bool
  61. {
  62. return $this->safe;
  63. }
  64. /**
  65. * @param mixed $safe
  66. */
  67. public function setSafe(bool $safe)
  68. {
  69. $this->safe = $safe;
  70. }
  71. /**
  72. * @return string
  73. */
  74. public function getName(): string
  75. {
  76. return $this->name;
  77. }
  78. /**
  79. * @param string $name
  80. */
  81. public function setName(string $name)
  82. {
  83. $this->name = $name;
  84. }
  85. /**
  86. * Publishes the message and merges additional properties with basic properties
  87. *
  88. * @param mixed $msgBody
  89. * @param string $exchangeName
  90. * @param string $routingKey
  91. * @param array $additionalProperties
  92. * @param array $headers
  93. *
  94. * @throws RuntimeException
  95. */
  96. public function publish(
  97. $msgBody,
  98. string $exchangeName,
  99. string $routingKey = '',
  100. array $additionalProperties = [],
  101. array $headers = null
  102. ) {
  103. if ($this->autoDeclare)
  104. {
  105. $this->routing->declareAll();
  106. }
  107. if ($this->safe && !$this->routing->isExchangeExists($exchangeName))
  108. {
  109. throw new RuntimeException(
  110. "Exchange `{$exchangeName}` does not declared in broker (You see this message because safe mode is ON)."
  111. );
  112. }
  113. $serialized = false;
  114. if (!is_string($msgBody))
  115. {
  116. $msgBody = call_user_func($this->serializer, $msgBody);
  117. $serialized = true;
  118. }
  119. $msg = new AMQPMessage($msgBody, array_merge($this->getBasicProperties(), $additionalProperties));
  120. if (!empty($headers) || $serialized)
  121. {
  122. if ($serialized)
  123. {
  124. $headers['rabbitmq.serialized'] = 1;
  125. }
  126. $headersTable = new AMQPTable($headers);
  127. $msg->set('application_headers', $headersTable);
  128. }
  129. \Yii::$app->rabbitmq->trigger(
  130. RabbitMQPublisherEvent::BEFORE_PUBLISH,
  131. new RabbitMQPublisherEvent(
  132. [
  133. 'message' => $msg,
  134. 'producer' => $this,
  135. ]
  136. )
  137. );
  138. $this->getChannel()->basic_publish($msg, $exchangeName, $routingKey);
  139. \Yii::$app->rabbitmq->trigger(
  140. RabbitMQPublisherEvent::AFTER_PUBLISH,
  141. new RabbitMQPublisherEvent(
  142. [
  143. 'message' => $msg,
  144. 'producer' => $this,
  145. ]
  146. )
  147. );
  148. $this->logger->log(
  149. 'AMQP message published',
  150. $msg,
  151. [
  152. 'exchange' => $exchangeName,
  153. 'routing_key' => $routingKey,
  154. ]
  155. );
  156. }
  157. }