<?php
namespace App\MessageHandler;
use App\Entity\Campaign;
use App\Entity\CampaignContact;
use App\Entity\Contact;
use App\Entity\Task;
use App\Message\BatchSendPendingCampaignMessage;
use App\Service\Api\Api;
use App\Service\Api\Tasks;
use App\Service\Campaign\CampaignEmailContact;
use App\Service\Campaign\Sender\SenderInterface;
use App\Service\Campaign\Sender\SendgridSender;
use Doctrine\ORM\EntityManagerInterface;
use Psr\Container\ContainerInterface;
use Psr\Log\LoggerInterface;
use Symfony\Component\EventDispatcher\EventSubscriberInterface;
use Symfony\Component\Messenger\Event\WorkerMessageFailedEvent;
use Symfony\Component\Messenger\Handler\MessageHandlerInterface;
use Symfony\Contracts\Service\ServiceSubscriberInterface;
class BatchSendPendingCampaignHandler implements MessageHandlerInterface, ServiceSubscriberInterface, EventSubscriberInterface {
/**
* @var EntityManagerInterface
*/
private $em;
/**
* @var ContainerInterface
*/
private $locator;
/**
* @var Tasks
*/
private $tasksApi;
private $logger;
public function __construct(EntityManagerInterface $em, ContainerInterface $locator, LoggerInterface $logger, Api $api)
{
$this->em = $em;
$this->locator = $locator;
$this->tasksApi = $api->get(Tasks::class);
$this->logger = $logger;
}
static function getSubscribedServices(): array
{
return [
SendgridSender::class
];
}
static function getSubscribedEvents(): array
{
return [
WorkerMessageFailedEvent::class => [
['handleFailed', 10]
]
];
}
public function handleFailed(WorkerMessageFailedEvent $event)
{
print_r($event->getReceiverName());
}
public function __invoke(BatchSendPendingCampaignMessage $message)
{
$campaignEmail = $message->getCampaignEmail();
/** @var SenderInterface $sender */
$sender = $this->locator->get($message->getSenderClass());
$campaign = $this->em->find(Campaign::class, $campaignEmail->getCampaignId());
$iterationNum = $message->getIterationNum();
$batchSize = $message->getBatchSize();
// Get contacts batch
$cCount = 0;
$campaignEmail->clearContacts();
$cContactsRepo = $this->em->getRepository(CampaignContact::class);
/** @var CampaignContact[] $result */
foreach ($cContactsRepo->iteratePendingByCampaign($campaign, $batchSize, $iterationNum * $batchSize) as $result) {
$cContact = $result[0];
$contact = $cContact->getContact();
$campaignEmailContact = new CampaignEmailContact($contact->getEmail(), $contact->getId());
$campaignEmail->addContact($campaignEmailContact);
$cCount++;
}
$this->logger->debug('Pending Campaign '.$campaignEmail->getCampaignId().' batch '.$iterationNum.': ' . json_encode([
'from' => $campaignEmail->getFrom(),
'from_address' => $campaignEmail->getFromAddress(),
'subject' => $campaignEmail->getSubject(),
'count' => $cCount,
'iteration' => $iterationNum,
'offset' => $iterationNum * $batchSize,
'size' => count($campaignEmail->getContacts()),
'batch_size' => $batchSize,
]));
// Send batch
$sender->sendCampaign($campaignEmail);
// Update task if there is one
if ($task = $this->em->find(Task::class, $message->getTaskId())) {
$this->tasksApi->increaseProgress($task);
}
}
}