StompConsumer.php 5.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211
  1. <?php
  2. declare(strict_types=1);
  3. namespace Enqueue\Stomp;
  4. use Interop\Queue\Consumer;
  5. use Interop\Queue\Exception\InvalidMessageException;
  6. use Interop\Queue\Message;
  7. use Interop\Queue\Queue;
  8. use Stomp\Client;
  9. use Stomp\Transport\Frame;
  10. class StompConsumer implements Consumer
  11. {
  12. const ACK_AUTO = 'auto';
  13. const ACK_CLIENT = 'client';
  14. const ACK_CLIENT_INDIVIDUAL = 'client-individual';
  15. /**
  16. * @var StompDestination
  17. */
  18. private $queue;
  19. /**
  20. * @var Client
  21. */
  22. private $stomp;
  23. /**
  24. * @var bool
  25. */
  26. private $isSubscribed;
  27. /**
  28. * @var string
  29. */
  30. private $ackMode;
  31. /**
  32. * @var int
  33. */
  34. private $prefetchCount;
  35. /**
  36. * @var string
  37. */
  38. private $subscriptionId;
  39. public function __construct(BufferedStompClient $stomp, StompDestination $queue)
  40. {
  41. $this->stomp = $stomp;
  42. $this->queue = $queue;
  43. $this->isSubscribed = false;
  44. $this->ackMode = self::ACK_CLIENT_INDIVIDUAL;
  45. $this->prefetchCount = 1;
  46. $this->subscriptionId = StompDestination::TYPE_TEMP_QUEUE == $queue->getType() ?
  47. $queue->getQueueName() :
  48. uniqid('', true)
  49. ;
  50. }
  51. public function setAckMode(string $mode): void
  52. {
  53. if (false === in_array($mode, [self::ACK_AUTO, self::ACK_CLIENT, self::ACK_CLIENT_INDIVIDUAL], true)) {
  54. throw new \LogicException(sprintf('Ack mode is not valid: "%s"', $mode));
  55. }
  56. $this->ackMode = $mode;
  57. }
  58. public function getAckMode(): string
  59. {
  60. return $this->ackMode;
  61. }
  62. public function getPrefetchCount(): int
  63. {
  64. return $this->prefetchCount;
  65. }
  66. public function setPrefetchCount(int $prefetchCount): void
  67. {
  68. $this->prefetchCount = $prefetchCount;
  69. }
  70. /**
  71. * @return StompDestination
  72. */
  73. public function getQueue(): Queue
  74. {
  75. return $this->queue;
  76. }
  77. public function receive(int $timeout = 0): ?Message
  78. {
  79. $this->subscribe();
  80. if (0 === $timeout) {
  81. while (true) {
  82. if ($message = $this->stomp->readMessageFrame($this->subscriptionId, 100)) {
  83. return $this->convertMessage($message);
  84. }
  85. }
  86. } else {
  87. if ($message = $this->stomp->readMessageFrame($this->subscriptionId, $timeout)) {
  88. return $this->convertMessage($message);
  89. }
  90. }
  91. return null;
  92. }
  93. public function receiveNoWait(): ?Message
  94. {
  95. $this->subscribe();
  96. if ($message = $this->stomp->readMessageFrame($this->subscriptionId, 0)) {
  97. return $this->convertMessage($message);
  98. }
  99. return null;
  100. }
  101. /**
  102. * @param StompMessage $message
  103. */
  104. public function acknowledge(Message $message): void
  105. {
  106. InvalidMessageException::assertMessageInstanceOf($message, StompMessage::class);
  107. $this->stomp->sendFrame(
  108. $this->stomp->getProtocol()->getAckFrame($message->getFrame())
  109. );
  110. }
  111. /**
  112. * @param StompMessage $message
  113. */
  114. public function reject(Message $message, bool $requeue = false): void
  115. {
  116. InvalidMessageException::assertMessageInstanceOf($message, StompMessage::class);
  117. $nackFrame = $this->stomp->getProtocol()->getNackFrame($message->getFrame());
  118. // rabbitmq STOMP protocol extension
  119. $nackFrame->addHeaders([
  120. 'requeue' => $requeue ? 'true' : 'false',
  121. ]);
  122. $this->stomp->sendFrame($nackFrame);
  123. }
  124. private function subscribe(): void
  125. {
  126. if (StompDestination::TYPE_TEMP_QUEUE == $this->queue->getType()) {
  127. $this->isSubscribed = true;
  128. return;
  129. }
  130. if (false == $this->isSubscribed) {
  131. $this->isSubscribed = true;
  132. $frame = $this->stomp->getProtocol()->getSubscribeFrame(
  133. $this->queue->getQueueName(),
  134. $this->subscriptionId,
  135. $this->ackMode
  136. );
  137. // rabbitmq STOMP protocol extension
  138. $headers = $this->queue->getHeaders();
  139. $headers['prefetch-count'] = $this->prefetchCount;
  140. $headers = StompHeadersEncoder::encode($headers);
  141. foreach ($headers as $key => $value) {
  142. $frame[$key] = $value;
  143. }
  144. $this->stomp->sendFrame($frame);
  145. }
  146. }
  147. private function convertMessage(Frame $frame): StompMessage
  148. {
  149. if ('MESSAGE' !== $frame->getCommand()) {
  150. throw new \LogicException(sprintf('Frame is not MESSAGE frame but: "%s"', $frame->getCommand()));
  151. }
  152. list($headers, $properties) = StompHeadersEncoder::decode($frame->getHeaders());
  153. $redelivered = isset($headers['redelivered']) && 'true' === $headers['redelivered'];
  154. unset(
  155. $headers['redelivered'],
  156. $headers['destination'],
  157. $headers['message-id'],
  158. $headers['ack'],
  159. $headers['receipt'],
  160. $headers['subscription'],
  161. $headers['content-length']
  162. );
  163. $message = new StompMessage((string) $frame->getBody(), $properties, $headers);
  164. $message->setRedelivered($redelivered);
  165. $message->setFrame($frame);
  166. return $message;
  167. }
  168. }