AutomaticEmailScheduler.php
2 years ago
AutomationEmailScheduler.php
1 year ago
LatestNewsletterScheduler.php
2 months ago
PostNotificationScheduler.php
2 months ago
ReEngagementScheduler.php
5 days ago
Scheduler.php
1 year ago
WelcomeScheduler.php
2 years ago
index.php
3 years ago
LatestNewsletterScheduler.php
309 lines
| 1 | <?php declare(strict_types = 1); |
| 2 | |
| 3 | namespace MailPoet\Newsletter\Scheduler; |
| 4 | |
| 5 | if (!defined('ABSPATH')) exit; |
| 6 | |
| 7 | |
| 8 | use MailPoet\Automation\Engine\Data\AutomationRun; |
| 9 | use MailPoet\Cron\Workers\SendingQueue\SendingQueue; |
| 10 | use MailPoet\Entities\NewsletterEntity; |
| 11 | use MailPoet\Entities\ScheduledTaskEntity; |
| 12 | use MailPoet\Entities\ScheduledTaskSubscriberEntity; |
| 13 | use MailPoet\Entities\SendingQueueEntity; |
| 14 | use MailPoet\Entities\StatisticsNewsletterEntity; |
| 15 | use MailPoet\Entities\SubscriberEntity; |
| 16 | use MailPoet\InvalidStateException; |
| 17 | use MailPoet\Newsletter\NewslettersRepository; |
| 18 | use MailPoet\Newsletter\Sending\NewsletterReplayMetadata; |
| 19 | use MailPoet\Newsletter\Sending\ScheduledTaskSubscribersRepository; |
| 20 | use MailPoetVendor\Carbon\Carbon; |
| 21 | use MailPoetVendor\Doctrine\ORM\EntityManager; |
| 22 | use MailPoetVendor\Doctrine\ORM\Query\Expr\Join; |
| 23 | |
| 24 | class LatestNewsletterScheduler { |
| 25 | public const OUTCOME_SCHEDULED = 'scheduled'; |
| 26 | public const OUTCOME_DUPLICATE = 'duplicate'; |
| 27 | public const OUTCOME_SKIPPED_NO_NEWSLETTER = 'skipped-no-newsletter'; |
| 28 | |
| 29 | private EntityManager $entityManager; |
| 30 | |
| 31 | private NewslettersRepository $newslettersRepository; |
| 32 | |
| 33 | private ScheduledTaskSubscribersRepository $scheduledTaskSubscribersRepository; |
| 34 | |
| 35 | public function __construct( |
| 36 | EntityManager $entityManager, |
| 37 | NewslettersRepository $newslettersRepository, |
| 38 | ScheduledTaskSubscribersRepository $scheduledTaskSubscribersRepository |
| 39 | ) { |
| 40 | $this->entityManager = $entityManager; |
| 41 | $this->newslettersRepository = $newslettersRepository; |
| 42 | $this->scheduledTaskSubscribersRepository = $scheduledTaskSubscribersRepository; |
| 43 | } |
| 44 | |
| 45 | /** |
| 46 | * @param array{id:mixed,run_id:mixed,step_id:mixed,run_number:mixed} $automationMeta |
| 47 | * @return array{outcome: string, newsletter: NewsletterEntity|null, task_subscriber: ScheduledTaskSubscriberEntity|null} |
| 48 | */ |
| 49 | public function schedule(SubscriberEntity $subscriber, int $segmentId, array $automationMeta): array { |
| 50 | $source = $this->newslettersRepository->findLatestSentStandardForSegment($segmentId); |
| 51 | if (!$source) { |
| 52 | return [ |
| 53 | 'outcome' => self::OUTCOME_SKIPPED_NO_NEWSLETTER, |
| 54 | 'newsletter' => null, |
| 55 | 'task_subscriber' => null, |
| 56 | ]; |
| 57 | } |
| 58 | |
| 59 | $newsletter = $source['newsletter']; |
| 60 | $subscriberId = $subscriber->getId(); |
| 61 | $newsletterId = $newsletter->getId(); |
| 62 | if (!$subscriberId || !$newsletterId) { |
| 63 | throw InvalidStateException::create(); |
| 64 | } |
| 65 | |
| 66 | $lockName = sprintf('mailpoet_latest_replay_%d_%d', $subscriberId, $newsletterId); |
| 67 | $this->acquireLock($lockName); |
| 68 | try { |
| 69 | return $this->entityManager->wrapInTransaction(function() use ($subscriber, $automationMeta, $source) { |
| 70 | $newsletter = $source['newsletter']; |
| 71 | $subscriberId = $subscriber->getId(); |
| 72 | $newsletterId = $newsletter->getId(); |
| 73 | if (!$subscriberId || !$newsletterId) { |
| 74 | throw InvalidStateException::create(); |
| 75 | } |
| 76 | |
| 77 | if ( |
| 78 | $this->hasSuccessfulProcessedSend($newsletter, $subscriber) |
| 79 | || $this->hasStatisticsNewsletter($newsletter, $subscriber) |
| 80 | || $this->hasPendingNonReplayTaskSubscriber($newsletter, $subscriber) |
| 81 | ) { |
| 82 | return [ |
| 83 | 'outcome' => self::OUTCOME_DUPLICATE, |
| 84 | 'newsletter' => $newsletter, |
| 85 | 'task_subscriber' => null, |
| 86 | ]; |
| 87 | } |
| 88 | |
| 89 | $existingReplay = $this->findExistingReplayTaskSubscriber($newsletter, $subscriber); |
| 90 | if ($existingReplay instanceof ScheduledTaskSubscriberEntity) { |
| 91 | $task = $existingReplay->getTask(); |
| 92 | $meta = $task ? $task->getMeta() : []; |
| 93 | $isSameRun = ($meta[NewsletterReplayMetadata::AUTOMATION]['run_id'] ?? null) === ($automationMeta['run_id'] ?? null); |
| 94 | return [ |
| 95 | 'outcome' => $isSameRun ? self::OUTCOME_SCHEDULED : self::OUTCOME_DUPLICATE, |
| 96 | 'newsletter' => $newsletter, |
| 97 | 'task_subscriber' => $isSameRun ? $existingReplay : null, |
| 98 | ]; |
| 99 | } |
| 100 | |
| 101 | $taskSubscriber = $this->createReplaySendingTask($source, $subscriber, $automationMeta); |
| 102 | return [ |
| 103 | 'outcome' => self::OUTCOME_SCHEDULED, |
| 104 | 'newsletter' => $newsletter, |
| 105 | 'task_subscriber' => $taskSubscriber, |
| 106 | ]; |
| 107 | }); |
| 108 | } finally { |
| 109 | $this->releaseLock($lockName); |
| 110 | } |
| 111 | } |
| 112 | |
| 113 | public function getScheduledTaskSubscriber(NewsletterEntity $newsletter, SubscriberEntity $subscriber, AutomationRun $run): ?ScheduledTaskSubscriberEntity { |
| 114 | $results = $this->entityManager->createQueryBuilder() |
| 115 | ->select('sts') |
| 116 | ->from(ScheduledTaskSubscriberEntity::class, 'sts') |
| 117 | ->join('sts.task', 'st') |
| 118 | ->join('st.sendingQueue', 'sq') |
| 119 | ->where('sq.newsletter = :newsletter') |
| 120 | ->andWhere('sts.subscriber = :subscriber') |
| 121 | ->andWhere('st.createdAt >= :runCreatedAt') |
| 122 | ->setParameter('newsletter', $newsletter) |
| 123 | ->setParameter('subscriber', $subscriber) |
| 124 | ->setParameter('runCreatedAt', $run->getCreatedAt()) |
| 125 | ->getQuery() |
| 126 | ->getResult(); |
| 127 | |
| 128 | foreach ($results as $scheduledTaskSubscriber) { |
| 129 | if (!$scheduledTaskSubscriber instanceof ScheduledTaskSubscriberEntity) { |
| 130 | continue; |
| 131 | } |
| 132 | $task = $scheduledTaskSubscriber->getTask(); |
| 133 | if (!$task instanceof ScheduledTaskEntity || !NewsletterReplayMetadata::isLatestNewsletterReplayMeta($task->getMeta())) { |
| 134 | continue; |
| 135 | } |
| 136 | $meta = $task->getMeta(); |
| 137 | if (($meta[NewsletterReplayMetadata::AUTOMATION]['run_id'] ?? null) === $run->getId()) { |
| 138 | return $scheduledTaskSubscriber; |
| 139 | } |
| 140 | } |
| 141 | return null; |
| 142 | } |
| 143 | |
| 144 | public function saveErrorAndPause(ScheduledTaskSubscriberEntity $scheduledTaskSubscriber, string $error): void { |
| 145 | $task = $scheduledTaskSubscriber->getTask(); |
| 146 | $subscriber = $scheduledTaskSubscriber->getSubscriber(); |
| 147 | if (!$task || !$subscriber || !$subscriber->getId()) { |
| 148 | return; |
| 149 | } |
| 150 | $this->scheduledTaskSubscribersRepository->saveError($task, $subscriber->getId(), $error); |
| 151 | $task->setStatus(ScheduledTaskEntity::STATUS_PAUSED); |
| 152 | $this->entityManager->flush(); |
| 153 | } |
| 154 | |
| 155 | private function hasSuccessfulProcessedSend(NewsletterEntity $newsletter, SubscriberEntity $subscriber): bool { |
| 156 | $result = $this->entityManager->createQueryBuilder() |
| 157 | ->select('COUNT(st)') |
| 158 | ->from(ScheduledTaskSubscriberEntity::class, 'sts') |
| 159 | ->join('sts.task', 'st') |
| 160 | ->join(SendingQueueEntity::class, 'sq', Join::WITH, 'sq.task = st') |
| 161 | ->where('sq.newsletter = :newsletter') |
| 162 | ->andWhere('sts.subscriber = :subscriber') |
| 163 | ->andWhere('sts.processed = :processed') |
| 164 | ->andWhere('sts.failed = :notFailed') |
| 165 | ->andWhere('st.status = :completed') |
| 166 | ->setParameter('newsletter', $newsletter) |
| 167 | ->setParameter('subscriber', $subscriber) |
| 168 | ->setParameter('processed', ScheduledTaskSubscriberEntity::STATUS_PROCESSED) |
| 169 | ->setParameter('notFailed', ScheduledTaskSubscriberEntity::FAIL_STATUS_OK) |
| 170 | ->setParameter('completed', ScheduledTaskEntity::STATUS_COMPLETED) |
| 171 | ->getQuery() |
| 172 | ->getSingleScalarResult(); |
| 173 | |
| 174 | return (int)$result > 0; |
| 175 | } |
| 176 | |
| 177 | private function hasPendingNonReplayTaskSubscriber(NewsletterEntity $newsletter, SubscriberEntity $subscriber): bool { |
| 178 | $result = $this->entityManager->createQueryBuilder() |
| 179 | ->select('COUNT(st)') |
| 180 | ->from(ScheduledTaskSubscriberEntity::class, 'sts') |
| 181 | ->join('sts.task', 'st') |
| 182 | ->join(SendingQueueEntity::class, 'sq', Join::WITH, 'sq.task = st') |
| 183 | ->where('sq.newsletter = :newsletter') |
| 184 | ->andWhere('sts.subscriber = :subscriber') |
| 185 | ->andWhere('sts.failed = :notFailed') |
| 186 | ->andWhere('(st.status = :scheduled OR st.status IS NULL)') |
| 187 | ->andWhere('st.meta IS NULL OR st.meta NOT LIKE :latestNewsletterReplayMeta') |
| 188 | ->andWhere('sq.meta IS NULL OR sq.meta NOT LIKE :latestNewsletterReplayMeta') |
| 189 | ->setParameter('newsletter', $newsletter) |
| 190 | ->setParameter('subscriber', $subscriber) |
| 191 | ->setParameter('notFailed', ScheduledTaskSubscriberEntity::FAIL_STATUS_OK) |
| 192 | ->setParameter('scheduled', ScheduledTaskEntity::STATUS_SCHEDULED) |
| 193 | ->setParameter('latestNewsletterReplayMeta', NewsletterReplayMetadata::getMetaLikePattern()) |
| 194 | ->getQuery() |
| 195 | ->getSingleScalarResult(); |
| 196 | |
| 197 | return (int)$result > 0; |
| 198 | } |
| 199 | |
| 200 | private function hasStatisticsNewsletter(NewsletterEntity $newsletter, SubscriberEntity $subscriber): bool { |
| 201 | $result = $this->entityManager->createQueryBuilder() |
| 202 | ->select('COUNT(statistics)') |
| 203 | ->from(StatisticsNewsletterEntity::class, 'statistics') |
| 204 | ->where('statistics.newsletter = :newsletter') |
| 205 | ->andWhere('statistics.subscriber = :subscriber') |
| 206 | ->setParameter('newsletter', $newsletter) |
| 207 | ->setParameter('subscriber', $subscriber) |
| 208 | ->getQuery() |
| 209 | ->getSingleScalarResult(); |
| 210 | |
| 211 | return (int)$result > 0; |
| 212 | } |
| 213 | |
| 214 | private function findExistingReplayTaskSubscriber(NewsletterEntity $newsletter, SubscriberEntity $subscriber): ?ScheduledTaskSubscriberEntity { |
| 215 | $results = $this->entityManager->createQueryBuilder() |
| 216 | ->select('sts') |
| 217 | ->from(ScheduledTaskSubscriberEntity::class, 'sts') |
| 218 | ->join('sts.task', 'st') |
| 219 | ->join(SendingQueueEntity::class, 'sq', Join::WITH, 'sq.task = st') |
| 220 | ->where('sq.newsletter = :newsletter') |
| 221 | ->andWhere('sts.subscriber = :subscriber') |
| 222 | ->andWhere('sts.failed = :notFailed') |
| 223 | ->setParameter('newsletter', $newsletter) |
| 224 | ->setParameter('subscriber', $subscriber) |
| 225 | ->setParameter('notFailed', ScheduledTaskSubscriberEntity::FAIL_STATUS_OK) |
| 226 | ->getQuery() |
| 227 | ->getResult(); |
| 228 | |
| 229 | foreach ($results as $scheduledTaskSubscriber) { |
| 230 | if (!$scheduledTaskSubscriber instanceof ScheduledTaskSubscriberEntity) { |
| 231 | continue; |
| 232 | } |
| 233 | $task = $scheduledTaskSubscriber->getTask(); |
| 234 | if (!$task instanceof ScheduledTaskEntity || !NewsletterReplayMetadata::isLatestNewsletterReplayMeta($task->getMeta())) { |
| 235 | continue; |
| 236 | } |
| 237 | $status = $task->getStatus(); |
| 238 | if (in_array($status, [ScheduledTaskEntity::STATUS_SCHEDULED, null], true)) { |
| 239 | return $scheduledTaskSubscriber; |
| 240 | } |
| 241 | if ( |
| 242 | $status === ScheduledTaskEntity::STATUS_COMPLETED |
| 243 | && $scheduledTaskSubscriber->getProcessed() === ScheduledTaskSubscriberEntity::STATUS_PROCESSED |
| 244 | ) { |
| 245 | return $scheduledTaskSubscriber; |
| 246 | } |
| 247 | } |
| 248 | return null; |
| 249 | } |
| 250 | |
| 251 | /** |
| 252 | * @param array{newsletter: NewsletterEntity, queue: SendingQueueEntity, task: ScheduledTaskEntity} $source |
| 253 | */ |
| 254 | private function createReplaySendingTask(array $source, SubscriberEntity $subscriber, array $automationMeta): ScheduledTaskSubscriberEntity { |
| 255 | $sourceTask = $source['task']; |
| 256 | $sourceQueue = $source['queue']; |
| 257 | $newsletter = $source['newsletter']; |
| 258 | |
| 259 | $meta = [ |
| 260 | NewsletterReplayMetadata::LATEST_NEWSLETTER_REPLAY => true, |
| 261 | NewsletterReplayMetadata::REPLAY_SOURCE_NEWSLETTER_ID => $newsletter->getId(), |
| 262 | NewsletterReplayMetadata::REPLAY_SOURCE_QUEUE_ID => $sourceQueue->getId(), |
| 263 | NewsletterReplayMetadata::REPLAY_SOURCE_TASK_ID => $sourceTask->getId(), |
| 264 | NewsletterReplayMetadata::REPLAY_SUBSCRIBER_ID => $subscriber->getId(), |
| 265 | NewsletterReplayMetadata::AUTOMATION => $automationMeta, |
| 266 | ]; |
| 267 | |
| 268 | $task = new ScheduledTaskEntity(); |
| 269 | $task->setType(SendingQueue::TASK_TYPE); |
| 270 | $task->setStatus(ScheduledTaskEntity::STATUS_SCHEDULED); |
| 271 | $task->setScheduledAt(Carbon::now()->millisecond(0)); |
| 272 | $task->setPriority(ScheduledTaskEntity::PRIORITY_MEDIUM); |
| 273 | $task->setMeta($meta); |
| 274 | $this->entityManager->persist($task); |
| 275 | |
| 276 | $taskSubscriber = new ScheduledTaskSubscriberEntity($task, $subscriber); |
| 277 | $this->entityManager->persist($taskSubscriber); |
| 278 | $task->getSubscribers()->add($taskSubscriber); |
| 279 | |
| 280 | $queue = new SendingQueueEntity(); |
| 281 | $queue->setTask($task); |
| 282 | $task->setSendingQueue($queue); |
| 283 | $queue->setMeta($meta); |
| 284 | $queue->setNewsletter($newsletter); |
| 285 | $queue->setCountToProcess(1); |
| 286 | $queue->setCountTotal(1); |
| 287 | $this->entityManager->persist($queue); |
| 288 | |
| 289 | return $taskSubscriber; |
| 290 | } |
| 291 | |
| 292 | private function acquireLock(string $lockName): void { |
| 293 | $result = $this->entityManager->getConnection()->executeQuery( |
| 294 | 'SELECT GET_LOCK(:lockName, 10)', |
| 295 | ['lockName' => $lockName] |
| 296 | )->fetchOne(); |
| 297 | if (!is_numeric($result) || (int)$result !== 1) { |
| 298 | throw InvalidStateException::create()->withMessage(__('Could not create sending task.', 'mailpoet')); |
| 299 | } |
| 300 | } |
| 301 | |
| 302 | private function releaseLock(string $lockName): void { |
| 303 | $this->entityManager->getConnection()->executeQuery( |
| 304 | 'SELECT RELEASE_LOCK(:lockName)', |
| 305 | ['lockName' => $lockName] |
| 306 | ); |
| 307 | } |
| 308 | } |
| 309 |