src/MessageHandler/BatchSendPendingCampaignHandler.php line 63

Open in your IDE?
  1. <?php
  2. namespace App\MessageHandler;
  3. use App\Entity\Campaign;
  4. use App\Entity\CampaignContact;
  5. use App\Entity\Contact;
  6. use App\Entity\Task;
  7. use App\Message\BatchSendPendingCampaignMessage;
  8. use App\Service\Api\Api;
  9. use App\Service\Api\Tasks;
  10. use App\Service\Campaign\CampaignEmailContact;
  11. use App\Service\Campaign\Sender\SenderInterface;
  12. use App\Service\Campaign\Sender\SendgridSender;
  13. use Doctrine\ORM\EntityManagerInterface;
  14. use Psr\Container\ContainerInterface;
  15. use Psr\Log\LoggerInterface;
  16. use Symfony\Component\EventDispatcher\EventSubscriberInterface;
  17. use Symfony\Component\Messenger\Event\WorkerMessageFailedEvent;
  18. use Symfony\Component\Messenger\Handler\MessageHandlerInterface;
  19. use Symfony\Contracts\Service\ServiceSubscriberInterface;
  20. class BatchSendPendingCampaignHandler implements MessageHandlerInterface, ServiceSubscriberInterface, EventSubscriberInterface {
  21. /**
  22. * @var EntityManagerInterface
  23. */
  24. private $em;
  25. /**
  26. * @var ContainerInterface
  27. */
  28. private $locator;
  29. /**
  30. * @var Tasks
  31. */
  32. private $tasksApi;
  33. private $logger;
  34. public function __construct(EntityManagerInterface $em, ContainerInterface $locator, LoggerInterface $logger, Api $api)
  35. {
  36. $this->em = $em;
  37. $this->locator = $locator;
  38. $this->tasksApi = $api->get(Tasks::class);
  39. $this->logger = $logger;
  40. }
  41. static function getSubscribedServices(): array
  42. {
  43. return [
  44. SendgridSender::class
  45. ];
  46. }
  47. static function getSubscribedEvents(): array
  48. {
  49. return [
  50. WorkerMessageFailedEvent::class => [
  51. ['handleFailed', 10]
  52. ]
  53. ];
  54. }
  55. public function handleFailed(WorkerMessageFailedEvent $event)
  56. {
  57. print_r($event->getReceiverName());
  58. }
  59. public function __invoke(BatchSendPendingCampaignMessage $message)
  60. {
  61. $campaignEmail = $message->getCampaignEmail();
  62. /** @var SenderInterface $sender */
  63. $sender = $this->locator->get($message->getSenderClass());
  64. $campaign = $this->em->find(Campaign::class, $campaignEmail->getCampaignId());
  65. $iterationNum = $message->getIterationNum();
  66. $batchSize = $message->getBatchSize();
  67. // Get contacts batch
  68. $cCount = 0;
  69. $campaignEmail->clearContacts();
  70. $cContactsRepo = $this->em->getRepository(CampaignContact::class);
  71. /** @var CampaignContact[] $result */
  72. foreach ($cContactsRepo->iteratePendingByCampaign($campaign, $batchSize, $iterationNum * $batchSize) as $result) {
  73. $cContact = $result[0];
  74. $contact = $cContact->getContact();
  75. $campaignEmailContact = new CampaignEmailContact($contact->getEmail(), $contact->getId());
  76. $campaignEmail->addContact($campaignEmailContact);
  77. $cCount++;
  78. }
  79. $this->logger->debug('Pending Campaign '.$campaignEmail->getCampaignId().' batch '.$iterationNum.': ' . json_encode([
  80. 'from' => $campaignEmail->getFrom(),
  81. 'from_address' => $campaignEmail->getFromAddress(),
  82. 'subject' => $campaignEmail->getSubject(),
  83. 'count' => $cCount,
  84. 'iteration' => $iterationNum,
  85. 'offset' => $iterationNum * $batchSize,
  86. 'size' => count($campaignEmail->getContacts()),
  87. 'batch_size' => $batchSize,
  88. ]));
  89. // Send batch
  90. $sender->sendCampaign($campaignEmail);
  91. // Update task if there is one
  92. if ($task = $this->em->find(Task::class, $message->getTaskId())) {
  93. $this->tasksApi->increaseProgress($task);
  94. }
  95. }
  96. }