EventStreamProcessor.php 3.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106
  1. <?php
  2. /**
  3. * Refer to LICENSE.txt distributed with the Temando Shipping module for notice of license
  4. */
  5. namespace Temando\Shipping\Sync;
  6. use Magento\Framework\Exception\CouldNotDeleteException;
  7. use Magento\Framework\Exception\LocalizedException;
  8. use Psr\Log\LoggerInterface;
  9. use Temando\Shipping\Model\Config\ModuleConfigInterface;
  10. use Temando\Shipping\Model\ResourceModel\EventStream\EventRepositoryInterface;
  11. use Temando\Shipping\Model\StreamEventInterface;
  12. /**
  13. * Temando Event Stream Processor
  14. *
  15. * @package Temando\Shipping\Sync
  16. * @author Benjamin Heuer <benjamin.heuer@netresearch.de>
  17. * @license http://opensource.org/licenses/osl-3.0.php Open Software License (OSL 3.0)
  18. * @link http://www.temando.com/
  19. */
  20. class EventStreamProcessor
  21. {
  22. /**
  23. * @var ModuleConfigInterface
  24. */
  25. private $config;
  26. /**
  27. * @var EventListProcessor
  28. */
  29. private $eventListProcessor;
  30. /**
  31. * @var EventRepositoryInterface
  32. */
  33. private $streamEventRepository;
  34. /**
  35. * @var EventFilter
  36. */
  37. private $streamEventFilter;
  38. /**
  39. * @var LoggerInterface
  40. */
  41. private $logger;
  42. /**
  43. * @param ModuleConfigInterface $config
  44. * @param EventListProcessor $eventListProcessor
  45. * @param EventRepositoryInterface $streamEventRepository
  46. * @param EventFilter $streamEventFilter
  47. * @param LoggerInterface $logger
  48. */
  49. public function __construct(
  50. ModuleConfigInterface $config,
  51. EventListProcessor $eventListProcessor,
  52. EventRepositoryInterface $streamEventRepository,
  53. EventFilter $streamEventFilter,
  54. LoggerInterface $logger
  55. ) {
  56. $this->config = $config;
  57. $this->eventListProcessor = $eventListProcessor;
  58. $this->streamEventRepository = $streamEventRepository;
  59. $this->streamEventFilter = $streamEventFilter;
  60. $this->logger = $logger;
  61. }
  62. /**
  63. * @param int $iterations
  64. * @return void
  65. */
  66. public function processEvents($iterations = 10)
  67. {
  68. $streamId = $this->config->getStreamId();
  69. $eventList = EventList::fromArray([]);
  70. do {
  71. $iterations--;
  72. try {
  73. // obtain next events, pass on to list processor
  74. $streamEvents = $this->streamEventRepository->getEventList($streamId);
  75. $processableEvents = $this->streamEventFilter->filter($streamEvents);
  76. $this->eventListProcessor->processEventList($processableEvents, $eventList);
  77. } catch (LocalizedException $e) {
  78. $this->logger->error($e->getMessage(), ['exception' => $e]);
  79. continue;
  80. }
  81. $deleteStreamEvent = function (StreamEventInterface $streamEvent) use ($streamId, $eventList) {
  82. if (!$eventList->hasEvent($streamEvent->getEventId())) {
  83. $this->streamEventRepository->delete($streamId, $streamEvent->getEventId());
  84. }
  85. };
  86. try {
  87. // delete events that were processed successfully (removed from event list)
  88. array_walk($streamEvents, $deleteStreamEvent);
  89. } catch (CouldNotDeleteException $e) {
  90. $this->logger->error($e->getMessage(), ['exception' => $e]);
  91. }
  92. } while ($iterations > 0 && !empty($streamEvents));
  93. }
  94. }