RoundRobinTransport.php 3.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121
  1. <?php
  2. /*
  3. * This file is part of the Symfony package.
  4. *
  5. * (c) Fabien Potencier <fabien@symfony.com>
  6. *
  7. * For the full copyright and license information, please view the LICENSE
  8. * file that was distributed with this source code.
  9. */
  10. namespace Symfony\Component\Mailer\Transport;
  11. use Symfony\Component\Mailer\Envelope;
  12. use Symfony\Component\Mailer\Exception\TransportException;
  13. use Symfony\Component\Mailer\Exception\TransportExceptionInterface;
  14. use Symfony\Component\Mailer\SentMessage;
  15. use Symfony\Component\Mime\RawMessage;
  16. /**
  17. * Uses several Transports using a round robin algorithm.
  18. *
  19. * @author Fabien Potencier <fabien@symfony.com>
  20. */
  21. class RoundRobinTransport implements TransportInterface
  22. {
  23. /**
  24. * @var \SplObjectStorage<TransportInterface, float>
  25. */
  26. private \SplObjectStorage $deadTransports;
  27. private array $transports = [];
  28. private int $retryPeriod;
  29. private int $cursor = -1;
  30. /**
  31. * @param TransportInterface[] $transports
  32. */
  33. public function __construct(array $transports, int $retryPeriod = 60)
  34. {
  35. if (!$transports) {
  36. throw new TransportException(sprintf('"%s" must have at least one transport configured.', static::class));
  37. }
  38. $this->transports = $transports;
  39. $this->deadTransports = new \SplObjectStorage();
  40. $this->retryPeriod = $retryPeriod;
  41. }
  42. public function send(RawMessage $message, Envelope $envelope = null): ?SentMessage
  43. {
  44. while ($transport = $this->getNextTransport()) {
  45. try {
  46. return $transport->send($message, $envelope);
  47. } catch (TransportExceptionInterface $e) {
  48. $this->deadTransports[$transport] = microtime(true);
  49. }
  50. }
  51. throw new TransportException('All transports failed.');
  52. }
  53. public function __toString(): string
  54. {
  55. return $this->getNameSymbol().'('.implode(' ', array_map('strval', $this->transports)).')';
  56. }
  57. /**
  58. * Rotates the transport list around and returns the first instance.
  59. */
  60. protected function getNextTransport(): ?TransportInterface
  61. {
  62. if (-1 === $this->cursor) {
  63. $this->cursor = $this->getInitialCursor();
  64. }
  65. $cursor = $this->cursor;
  66. while (true) {
  67. $transport = $this->transports[$cursor];
  68. if (!$this->isTransportDead($transport)) {
  69. break;
  70. }
  71. if ((microtime(true) - $this->deadTransports[$transport]) > $this->retryPeriod) {
  72. $this->deadTransports->detach($transport);
  73. break;
  74. }
  75. if ($this->cursor === $cursor = $this->moveCursor($cursor)) {
  76. return null;
  77. }
  78. }
  79. $this->cursor = $this->moveCursor($cursor);
  80. return $transport;
  81. }
  82. protected function isTransportDead(TransportInterface $transport): bool
  83. {
  84. return $this->deadTransports->contains($transport);
  85. }
  86. protected function getInitialCursor(): int
  87. {
  88. // the cursor initial value is randomized so that
  89. // when are not in a daemon, we are still rotating the transports
  90. return mt_rand(0, \count($this->transports) - 1);
  91. }
  92. protected function getNameSymbol(): string
  93. {
  94. return 'roundrobin';
  95. }
  96. private function moveCursor(int $cursor): int
  97. {
  98. return ++$cursor >= \count($this->transports) ? 0 : $cursor;
  99. }
  100. }