<?php
namespace App\MessageHandler;
use App\Entity\Contact;
use App\Entity\Task;
use App\Message\BatchSendCampaignMessage;
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 BatchSendCampaignHandler 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(BatchSendCampaignMessage $message)
{
$campaignEmail = $message->getCampaignEmail();
/** @var SenderInterface $sender */
$sender = $this->locator->get($message->getSenderClass());
$groups = $message->getGroups();
$jsonRules = $message->getJsonRules();
$iterationNum = $message->getIterationNum();
$batchSize = $message->getBatchSize();
// Get contacts batch
$cCount = 0;
$contactsRepo = $this->em->getRepository(Contact::class);
/** @var Contact[] $result */
foreach ($contactsRepo->iterateByQueryBuilder($groups, $jsonRules, $batchSize, $iterationNum * $batchSize) as $result) {
$contact = $result[0];
$campaignEmailContact = new CampaignEmailContact($contact->getEmail(), $contact->getId());
$campaignEmail->addContact($campaignEmailContact);
$cCount++;
}
$this->logger->debug('Campaign '.$campaignEmail->getCampaignId().' batch '.$iterationNum.': ' . json_encode([
'from' => $campaignEmail->getFrom(),
'from_address' => $campaignEmail->getFromAddress(),
'subject' => $campaignEmail->getSubject(),
'count' => $cCount
]));
// Send batch
$sender->sendCampaign($campaignEmail);
// Update task if there is one
if ($task = $this->em->find(Task::class, $message->getTaskId())) {
$this->tasksApi->increaseProgress($task);
}
}
}