PluginProbe ʕ •ᴥ•ʔ
MailPoet – Newsletters, Email Marketing, and Automation / 5.34.2
MailPoet – Newsletters, Email Marketing, and Automation v5.34.2
5.34.2 5.34.1 5.34.0 5.33.1 5.33.0 5.32.0 5.31.0 5.30.0 5.29.0 5.28.1 5.28.0 5.27.0 5.26.0 5.26.1 5.25.0 5.24.0 4.43.0 4.43.1 4.44.0 4.44.1 4.45.0 4.46.0 4.47.0 4.48.0 4.48.1 4.48.2 4.49.0 4.49.1 4.5.0 4.5.1 4.5.2 4.50.0 4.50.1 4.51.0 4.51.1 4.51.2 4.52.0 4.53.0 4.54.0 4.55.0 4.56.0 4.57.0 4.58.0 4.58.1 4.58.2 4.6.0 4.6.1 4.6.2 4.7.0 4.7.1 4.8.0 4.8.1 4.9.0 5.0.0 5.0.1 5.0.2 5.1.0 5.1.1 5.10.0 5.10.1 5.11.0 5.12.0 5.12.1 5.12.10 5.12.11 5.12.12 5.12.13 5.12.2 5.12.3 5.12.4 5.12.5 5.12.6 5.12.7 5.12.8 5.12.9 5.13.0 5.13.1 5.13.2 5.14.0 5.14.1 5.14.2 5.14.3 5.15.0 5.15.1 5.16.0 5.16.1 5.16.2 5.16.3 5.16.4 5.17.0 5.17.1 5.17.2 5.17.3 5.17.4 5.17.5 5.17.6 5.18.0 5.19.0 5.2.0 5.2.1 5.2.2 5.2.3 5.20.0 5.21.0 5.21.1 5.21.2 5.21.3 5.22.0 5.22.1 5.22.2 5.22.3 5.22.4 5.23.0 5.23.1 5.23.2 5.3.0 5.3.1 5.3.2 5.3.3 5.3.4 5.3.5 5.3.6 5.3.7 5.4.0 5.4.1 5.4.2 5.5.0 5.5.1 5.5.2 5.6.0 5.6.1 5.6.2 5.6.3 5.6.4 5.7.0 5.7.1 5.8.0 5.8.1 5.9.0 3.0.0-beta.15 3.7.1 3.0.0-beta.16 3.7.2 3.0.0-beta.17 3.7.3 3.0.0-beta.18 3.7.4 3.0.0-beta.19 3.7.5 3.0.0-beta.2 3.7.6 3.0.0-beta.20 3.7.8 3.0.0-beta.21 3.70.0 3.0.0-beta.22 3.71.0 3.0.0-beta.23 3.71.1 3.0.0-beta.23.1 3.71.2 3.0.0-beta.23.2 3.71.3 3.0.0-beta.24 3.72.0 3.0.0-beta.25 3.73.0 3.0.0-beta.26 3.73.1 3.0.0-beta.27 3.73.2 3.0.0-beta.28 3.74.0 3.0.0-beta.29 3.74.1 3.0.0-beta.3 3.74.2 3.0.0-beta.30 3.74.3 3.0.0-beta.31 3.75.0 3.0.0-beta.32 3.75.1 3.0.0-beta.33 3.76.0 3.0.0-beta.33.1 3.77.0 3.0.0-beta.34.0.0 3.77.1 3.0.0-beta.36.0.0 3.78.0 3.0.0-beta.36.0.1 3.79.0 3.0.0-beta.36.2.0 3.8 3.0.0-beta.36.3.0 3.8.1 3.0.0-beta.36.3.1 3.8.2 3.0.0-beta.37.0.0 3.8.3 3.0.0-beta.4 3.8.4 3.0.0-beta.5 3.8.5 3.0.0-beta.6 3.8.6 3.0.0-beta.7 3.80.0 3.0.0-beta.7.1 3.81.0 3.0.0-beta.8 3.82.0 3.0.0-beta.9 3.83.0 3.0.0-rc.1.0.0 3.84.0 3.0.0-rc.1.0.1 3.84.1 3.0.0-rc.1.0.2 3.85.0 3.0.0-rc.1.0.3 3.85.1 3.0.0-rc.1.0.4 3.86.0 3.0.0-rc.2.0.0 3.87.0 3.0.0-rc.2.0.1 3.87.1 3.0.0-rc.2.0.2 3.87.2 3.0.0-rc.2.0.3 3.88.0 3.0.1 3.88.1 3.0.2 3.88.2 3.0.3 3.89.0 3.0.4 3.89.1 3.0.5 3.89.2 3.0.6 3.89.3 3.0.7 3.89.4 3.0.8 3.9.0 3.0.9 3.9.1 3.1.0 3.90.0 3.10 3.90.1 3.10.1 3.90.2 3.100.0 3.91.0 3.100.1 3.91.1 3.100.2 3.92.0 3.101.0 3.92.1 3.101.1 3.93.0 3.102.0 3.93.1 3.102.1 3.94.0 3.103.0 3.95.0 3.103.1 3.95.1 3.11.0 3.96.0 3.11.1 3.96.1 3.11.2 3.97.0 3.11.3 3.98.0 3.11.4 3.98.1 3.11.5 3.99.0 3.12.0 3.99.1 3.12.1 4.0.0 3.13.0 4.0.1 3.14.0 4.1.0 3.14.1 4.1.1 3.15.0 4.10.0 3.16.0 4.11.0 3.16.1 4.11.1 3.16.2 4.12.0 3.16.3 4.12.1 3.17.0 4.12.2 3.17.1 4.13.0 3.17.2 4.14.0 3.18.0 4.15.0 3.18.1 4.16.0 3.18.2 4.17.0 3.19.0 4.17.1 3.19.1 4.18.0 3.19.2 4.18.1 3.19.3 4.19.0 3.2.0 4.2.0 3.2.1 4.20.0 3.2.2 4.20.1 3.2.3 4.20.2 3.2.4 4.21.0 3.2.5 4.22.0 3.20.0 4.22.1 3.21.0 4.22.2 3.21.1 4.23.0 3.22.0 4.24.0 3.23.0 4.25.0 3.23.1 4.26.0 3.23.2 4.26.1 3.24.0 4.27.0 3.25.0 4.28.0 3.25.1 4.29.0 3.26.0 4.3.0 3.26.1 4.3.1 3.27.0 4.30.0 3.28.0 4.31.0 3.29.0 4.31.1 3.3.0 4.32.0 3.3.1 4.33.0 3.3.2 4.34.0 3.3.3 4.35.0 3.3.4 4.35.1 3.3.5 4.36.0 3.3.6 4.37.0 3.30.0 4.38.0 3.31.0 4.39.0 3.31.1 4.4.0 3.32.0 4.40.0 3.32.1 4.41.0 3.32.2 4.41.1 3.33.0 4.41.2 3.34.0 4.41.3 3.34.1 4.42.0 3.34.2 4.42.1 3.34.3 3.34.4 3.35.0 3.35.1 3.35.3 3.35.4 3.36.0 3.37.0 3.37.1 3.37.2 3.37.3 3.38.0 3.38.1 3.39.0 3.39.1 3.39.2 3.4.0 3.4.1 3.4.2 3.4.3 3.4.4 3.40.0 3.40.1 3.41.0 3.41.1 3.41.2 3.42.0 3.42.1 3.42.2 3.42.3 3.43.0 3.43.1 3.44.0 3.45.0 3.45.1 3.46.0 3.46.1 3.46.10 3.46.11 3.46.12 3.46.13 3.46.14 3.46.2 3.46.3 3.46.4 3.46.5 3.46.6 3.46.7 3.46.8 3.46.9 3.47.0 3.47.1 3.47.10 3.47.11 3.47.2 3.47.3 3.47.5 3.47.6 3.47.7 3.47.9 3.48.0 3.48.1 3.49.0 3.49.1 3.5.0 3.5.1 3.50.0 3.51.0 3.51.1 3.51.2 3.52.0 3.53.0 3.54.0 3.54.1 3.54.2 3.54.3 3.55.0 3.55.1 3.56.0 3.56.1 3.56.2 3.57.0 3.57.1 3.58.0 3.59.0 3.59.1 3.59.2 3.6.0 3.6.1 3.6.2 3.6.3 3.6.4 3.6.5 3.6.6 3.6.7 3.60.0 3.60.1 3.60.10 3.60.11 3.60.12 3.60.2 3.60.3 3.60.4 3.60.6 3.60.7 3.60.8 3.60.9 3.61.0 3.62.0 3.62.1 3.63.0 3.64.0 3.64.1 3.64.2 3.64.3 3.65.0 trunk 3.65.1 3.0.0 3.66.0 3.0.0-beta.1 3.67.0 3.0.0-beta.10 3.67.1 3.0.0-beta.11 3.68.0 3.0.0-beta.12 3.69.0 3.0.0-beta.13 3.69.1 3.0.0-beta.14 3.7.0
mailpoet / lib / Newsletter / Sending / ScheduledTasksRepository.php
mailpoet / lib / Newsletter / Sending Last commit date
NewsletterReplayMetadata.php 2 months ago ScheduledTaskSubscribersListingRepository.php 2 years ago ScheduledTaskSubscribersRepository.php 1 month ago ScheduledTasksRepository.php 1 month ago SendingQueuesRepository.php 1 week ago TimeZoneCampaignScheduler.php 1 week ago index.php 3 years ago
ScheduledTasksRepository.php
474 lines
1 <?php // phpcs:ignore SlevomatCodingStandard.TypeHints.DeclareStrictTypes.DeclareStrictTypesMissing
2
3 namespace MailPoet\Newsletter\Sending;
4
5 if (!defined('ABSPATH')) exit;
6
7
8 use MailPoet\Cron\Workers\SendingQueue\SendingQueue;
9 use MailPoet\Doctrine\Repository;
10 use MailPoet\Entities\NewsletterEntity;
11 use MailPoet\Entities\ScheduledTaskEntity;
12 use MailPoet\Entities\ScheduledTaskSubscriberEntity;
13 use MailPoet\Entities\SendingQueueEntity;
14 use MailPoet\Entities\SubscriberEntity;
15 use MailPoetVendor\Carbon\Carbon;
16 use MailPoetVendor\Doctrine\DBAL\ArrayParameterType;
17 use MailPoetVendor\Doctrine\ORM\EntityManager;
18 use MailPoetVendor\Doctrine\ORM\Query\Expr\Join;
19
20 /**
21 * @extends Repository<ScheduledTaskEntity>
22 */
23 class ScheduledTasksRepository extends Repository {
24 const TASK_BATCH_SIZE = 20;
25 const CANCELLABLE_STATUSES = [
26 ScheduledTaskEntity::STATUS_SCHEDULED,
27 ScheduledTaskEntity::VIRTUAL_STATUS_RUNNING,
28 null,
29 ];
30
31 private SendingQueuesRepository $sendingQueuesRepository;
32
33 public function __construct(
34 EntityManager $entityManager,
35 SendingQueuesRepository $sendingQueuesRepository
36 ) {
37 $this->sendingQueuesRepository = $sendingQueuesRepository;
38 parent::__construct($entityManager);
39 }
40
41 /**
42 * @param NewsletterEntity $newsletter
43 * @return ScheduledTaskEntity[]
44 */
45 public function findByNewsletterAndStatus(NewsletterEntity $newsletter, string $status): array {
46 return $this->doctrineRepository->createQueryBuilder('st')
47 ->select('st')
48 ->join(SendingQueueEntity::class, 'sq', Join::WITH, 'st = sq.task')
49 ->andWhere('st.status = :status')
50 ->andWhere('sq.newsletter = :newsletter')
51 ->setParameter('status', $status)
52 ->setParameter('newsletter', $newsletter)
53 ->getQuery()
54 ->getResult();
55 }
56
57 /**
58 * @param NewsletterEntity $newsletter
59 */
60 public function findOneByNewsletter(NewsletterEntity $newsletter): ?ScheduledTaskEntity {
61 $scheduledTask = $this->doctrineRepository->createQueryBuilder('st')
62 ->join(SendingQueueEntity::class, 'sq', Join::WITH, 'st = sq.task')
63 ->andWhere('sq.newsletter = :newsletter')
64 ->orderBy('sq.updatedAt', 'desc')
65 ->setMaxResults(1)
66 ->setParameter('newsletter', $newsletter)
67 ->getQuery()
68 ->getOneOrNullResult();
69 // for phpstan because it detects mixed instead of entity
70 return ($scheduledTask instanceof ScheduledTaskEntity) ? $scheduledTask : null;
71 }
72
73 public function findOneBySendingQueue(SendingQueueEntity $sendingQueue): ?ScheduledTaskEntity {
74 $scheduledTask = $this->doctrineRepository->createQueryBuilder('st')
75 ->join(SendingQueueEntity::class, 'sq', Join::WITH, 'st = sq.task')
76 ->andWhere('sq.id = :sendingQueue')
77 ->setMaxResults(1)
78 ->setParameter('sendingQueue', $sendingQueue)
79 ->getQuery()
80 ->getOneOrNullResult();
81 // for phpstan because it detects mixed instead of entity
82 return ($scheduledTask instanceof ScheduledTaskEntity) ? $scheduledTask : null;
83 }
84
85 /**
86 * @param NewsletterEntity $newsletter
87 * @return ScheduledTaskEntity[]
88 */
89 public function findByScheduledAndRunningForNewsletter(NewsletterEntity $newsletter): array {
90 return $this->doctrineRepository->createQueryBuilder('st')
91 ->select('st')
92 ->join(SendingQueueEntity::class, 'sq', Join::WITH, 'st = sq.task')
93 ->andWhere('st.status = :status OR st.status IS NULL')
94 ->andWhere('sq.newsletter = :newsletter')
95 ->setParameter('status', NewsletterEntity::STATUS_SCHEDULED)
96 ->setParameter('newsletter', $newsletter)
97 ->getQuery()
98 ->getResult();
99 }
100
101 /**
102 * @param NewsletterEntity $newsletter
103 * @return ScheduledTaskEntity[]
104 */
105 public function findByNewsletterAndSubscriberId(NewsletterEntity $newsletter, int $subscriberId): array {
106 return $this->doctrineRepository->createQueryBuilder('st')
107 ->select('st')
108 ->join(SendingQueueEntity::class, 'sq', Join::WITH, 'st = sq.task')
109 ->join(ScheduledTaskSubscriberEntity::class, 'sts', Join::WITH, 'st = sts.task')
110 ->andWhere('sq.newsletter = :newsletter')
111 ->andWhere('sts.subscriber = :subscriber')
112 ->setParameter('newsletter', $newsletter)
113 ->setParameter('subscriber', $subscriberId)
114 ->getQuery()
115 ->getResult();
116 }
117
118 public function findOneScheduledByNewsletterAndSubscriber(NewsletterEntity $newsletter, SubscriberEntity $subscriber): ?ScheduledTaskEntity {
119 $scheduledTask = $this->doctrineRepository->createQueryBuilder('st')
120 ->join(SendingQueueEntity::class, 'sq', Join::WITH, 'st = sq.task')
121 ->join(ScheduledTaskSubscriberEntity::class, 'sts', Join::WITH, 'st = sts.task')
122 ->andWhere('st.status = :status')
123 ->andWhere('sq.newsletter = :newsletter')
124 ->andWhere('sts.subscriber = :subscriber')
125 ->setMaxResults(1)
126 ->setParameter('status', ScheduledTaskEntity::STATUS_SCHEDULED)
127 ->setParameter('newsletter', $newsletter)
128 ->setParameter('subscriber', $subscriber)
129 ->getQuery()
130 ->getOneOrNullResult();
131 // for phpstan because it detects mixed instead of entity
132 return ($scheduledTask instanceof ScheduledTaskEntity) ? $scheduledTask : null;
133 }
134
135 public function findScheduledOrRunningTask(?string $type): ?ScheduledTaskEntity {
136 $queryBuilder = $this->doctrineRepository->createQueryBuilder('st')
137 ->select('st')
138 ->where('((st.status = :scheduledStatus) OR (st.status is NULL))')
139 ->andWhere('st.deletedAt IS NULL')
140 ->setParameter('scheduledStatus', ScheduledTaskEntity::STATUS_SCHEDULED)
141 ->setMaxResults(1)
142 ->orderBy('st.scheduledAt', 'DESC');
143 if (!empty($type)) {
144 $queryBuilder
145 ->andWhere('st.type = :type')
146 ->setParameter('type', $type);
147 }
148 return $queryBuilder->getQuery()->getOneOrNullResult();
149 }
150
151 public function findScheduledTask(?string $type): ?ScheduledTaskEntity {
152 $queryBuilder = $this->doctrineRepository->createQueryBuilder('st')
153 ->select('st')
154 ->where('st.status = :scheduledStatus')
155 ->andWhere('st.deletedAt IS NULL')
156 ->setParameter('scheduledStatus', ScheduledTaskEntity::STATUS_SCHEDULED)
157 ->setMaxResults(1)
158 ->orderBy('st.scheduledAt', 'DESC');
159 if (!empty($type)) {
160 $queryBuilder
161 ->andWhere('st.type = :type')
162 ->setParameter('type', $type);
163 }
164 return $queryBuilder->getQuery()->getOneOrNullResult();
165 }
166
167 public function findSoonestScheduledTaskByType(string $type): ?ScheduledTaskEntity {
168 return $this->doctrineRepository->createQueryBuilder('st')
169 ->select('st')
170 ->where('st.status = :scheduledStatus')
171 ->andWhere('st.type = :type')
172 ->andWhere('st.deletedAt IS NULL')
173 ->setParameter('scheduledStatus', ScheduledTaskEntity::STATUS_SCHEDULED)
174 ->setParameter('type', $type)
175 ->setMaxResults(1)
176 ->orderBy('st.scheduledAt', 'ASC')
177 ->getQuery()
178 ->getOneOrNullResult();
179 }
180
181 /**
182 * Atomically claims a scheduled/paused row for a CLI run by flipping it to STATUS_CLI in a single
183 * guarded UPDATE. Returns false when no row was changed (another process already claimed or ran it),
184 * so concurrent CLI runs never both process the same task.
185 */
186 public function claimAsCli(ScheduledTaskEntity $task): bool {
187 $updated = (int)$this->entityManager->createQuery(
188 'UPDATE ' . ScheduledTaskEntity::class . ' st ' .
189 'SET st.status = :cli ' .
190 'WHERE st.id = :id AND st.status IN (:claimable) AND st.deletedAt IS NULL'
191 )
192 ->setParameter('cli', ScheduledTaskEntity::STATUS_CLI)
193 ->setParameter('id', $task->getId())
194 ->setParameter('claimable', [ScheduledTaskEntity::STATUS_SCHEDULED, ScheduledTaskEntity::STATUS_PAUSED])
195 ->execute();
196
197 if ($updated === 0) {
198 return false;
199 }
200
201 // The DQL UPDATE bypasses the entity manager, so refresh the in-memory entity to the claimed state.
202 // Otherwise its status snapshot stays 'scheduled' and a later hand-back to 'scheduled' is seen as a
203 // no-op change, leaving the row stuck as 'cli'.
204 $this->entityManager->refresh($task);
205 return true;
206 }
207
208 public function findPreviousTask(ScheduledTaskEntity $task): ?ScheduledTaskEntity {
209 return $this->doctrineRepository->createQueryBuilder('st')
210 ->select('st')
211 ->where('st.type = :type')
212 ->setParameter('type', $task->getType())
213 ->andWhere('st.createdAt < :created')
214 ->setParameter('created', $task->getCreatedAt())
215 ->orderBy('st.scheduledAt', 'DESC')
216 ->setMaxResults(1)
217 ->getQuery()
218 ->getOneOrNullResult();
219 }
220
221 public function findDueByType($type, $limit = null) {
222 return $this->findByTypeAndStatus($type, ScheduledTaskEntity::STATUS_SCHEDULED, $limit);
223 }
224
225 public function findRunningByType($type, $limit = null) {
226 return $this->findByTypeAndStatus($type, null, $limit);
227 }
228
229 public function findCompletedByType($type, $limit = null) {
230 return $this->findByTypeAndStatus($type, ScheduledTaskEntity::STATUS_COMPLETED, $limit);
231 }
232
233 public function findFutureScheduledByType($type, $limit = null) {
234 return $this->findByTypeAndStatus($type, ScheduledTaskEntity::STATUS_SCHEDULED, $limit, true);
235 }
236
237 public function getCountsPerStatus(string $type = 'sending') {
238 $stats = [
239 ScheduledTaskEntity::STATUS_COMPLETED => 0,
240 ScheduledTaskEntity::STATUS_PAUSED => 0,
241 ScheduledTaskEntity::STATUS_SCHEDULED => 0,
242 ScheduledTaskEntity::STATUS_CANCELLED => 0,
243 ScheduledTaskEntity::VIRTUAL_STATUS_RUNNING => 0,
244 ];
245
246 $counts = $this->doctrineRepository->createQueryBuilder('st')
247 ->select('COUNT(st.id) as value')
248 ->addSelect('st.status')
249 ->where('st.deletedAt IS NULL')
250 ->andWhere('st.type = :type')
251 ->setParameter('type', $type)
252 ->addGroupBy('st.status')
253 ->getQuery()
254 ->getResult();
255
256 foreach ($counts as $count) {
257 if ($count['status'] === null) {
258 $stats[ScheduledTaskEntity::VIRTUAL_STATUS_RUNNING] = (int)$count['value'];
259 continue;
260 }
261 $stats[$count['status']] = (int)$count['value'];
262 }
263 return $stats;
264 }
265
266 /**
267 * @param string|null $type
268 * @param array $statuses
269 * @param int $limit
270 * @return array<ScheduledTaskEntity>
271 */
272 public function getLatestTasks(
273 $type = null,
274 $statuses = [
275 ScheduledTaskEntity::STATUS_COMPLETED,
276 ScheduledTaskEntity::STATUS_CANCELLED,
277 ScheduledTaskEntity::STATUS_SCHEDULED,
278 ScheduledTaskEntity::STATUS_PAUSED,
279 ScheduledTaskEntity::VIRTUAL_STATUS_RUNNING,
280 ],
281 $limit = self::TASK_BATCH_SIZE
282 ) {
283 $result = [];
284 foreach ($statuses as $status) {
285 $tasksQuery = $this->doctrineRepository->createQueryBuilder('st')
286 ->select('st')
287 ->where('st.deletedAt IS NULL');
288
289 if ($status === ScheduledTaskEntity::VIRTUAL_STATUS_RUNNING) {
290 $tasksQuery = $tasksQuery->andWhere('st.status = :status OR st.status IS NULL');
291 } else {
292 $tasksQuery = $tasksQuery->andWhere('st.status = :status');
293 }
294
295 if ($type) {
296 $tasksQuery = $tasksQuery->andWhere('st.type = :type')
297 ->setParameter('type', $type);
298 }
299
300 $tasks = $tasksQuery
301 ->setParameter('status', $status)
302 ->setMaxResults($limit)
303 ->orderBy('st.id', 'desc')
304 ->getQuery()
305 ->getResult();
306 $result = array_merge($result, $tasks);
307 }
308
309 return $result;
310 }
311
312 /**
313 * @return ScheduledTaskEntity[]
314 */
315 public function findRunningSendingTasks(?int $limit = null): array {
316 return $this->doctrineRepository->createQueryBuilder('st')
317 ->select('st')
318 ->join('st.sendingQueue', 'sq')
319 ->where('st.type = :type')
320 ->andWhere('st.status IS NULL')
321 ->andWhere('st.deletedAt IS NULL')
322 ->orderBy('st.priority', 'ASC')
323 ->addOrderBy('st.updatedAt', 'ASC')
324 ->setMaxResults($limit)
325 ->setParameter('type', SendingQueue::TASK_TYPE)
326 ->getQuery()
327 ->getResult();
328 }
329
330 /**
331 * @param string $type
332 * @param SubscriberEntity $subscriber
333 * @return ScheduledTaskEntity[]
334 * @throws \MailPoetVendor\Doctrine\ORM\NonUniqueResultException
335 */
336 public function findByTypeAndSubscriber(string $type, SubscriberEntity $subscriber): array {
337 $query = $this->doctrineRepository->createQueryBuilder('st')
338 ->select('st')
339 ->join(ScheduledTaskSubscriberEntity::class, 'sts', Join::WITH, 'st = sts.task')
340 ->where('st.type = :type')
341 ->andWhere('sts.subscriber = :subscriber')
342 ->andWhere('st.deletedAt IS NULL')
343 ->andWhere('st.status = :status')
344 ->setParameter('type', $type)
345 ->setParameter('subscriber', $subscriber->getId())
346 ->setParameter('status', ScheduledTaskEntity::STATUS_SCHEDULED)
347 ->getQuery();
348 $tasks = $query->getResult();
349 return $tasks;
350 }
351
352 public function touchAllByIds(array $ids): void {
353 $now = Carbon::now()->millisecond(0);
354 $this->entityManager->createQueryBuilder()
355 ->update(ScheduledTaskEntity::class, 'st')
356 ->set('st.updatedAt', ':updatedAt')
357 ->setParameter('updatedAt', $now)
358 ->where('st.id IN (:ids)')
359 ->setParameter('ids', $ids, ArrayParameterType::INTEGER)
360 ->getQuery()
361 ->execute();
362
363 // update was done via DQL, make sure the entities are also refreshed in the entity manager
364 $this->refreshAll(function (ScheduledTaskEntity $entity) use ($ids) {
365 return in_array($entity->getId(), $ids, true);
366 });
367 }
368
369 /**
370 * @return ScheduledTaskEntity[]
371 */
372 public function findScheduledSendingTasks(?int $limit = null): array {
373 $now = Carbon::now()->millisecond(0);
374 return $this->doctrineRepository->createQueryBuilder('st')
375 ->select('st')
376 ->join('st.sendingQueue', 'sq')
377 ->where('st.deletedAt IS NULL')
378 ->andWhere('st.status = :status')
379 ->andWhere('st.scheduledAt <= :now')
380 ->andWhere('st.type = :type')
381 ->orderBy('st.updatedAt', 'ASC')
382 ->setMaxResults($limit)
383 ->setParameter('status', ScheduledTaskEntity::STATUS_SCHEDULED)
384 ->setParameter('now', $now)
385 ->setParameter('type', SendingQueue::TASK_TYPE)
386 ->getQuery()
387 ->getResult();
388 }
389
390 public function invalidateTask(ScheduledTaskEntity $task): void {
391 $task->setStatus(ScheduledTaskEntity::STATUS_INVALID);
392 $this->persist($task);
393 $this->flush();
394 }
395
396 public function cancelTask(ScheduledTaskEntity $task): void {
397 if (!in_array($task->getStatus(), self::CANCELLABLE_STATUSES)) {
398 throw new \Exception(__('Only scheduled and running tasks can be cancelled', 'mailpoet'), 400);
399 }
400 $task->setStatus(ScheduledTaskEntity::STATUS_CANCELLED);
401 $task->setCancelledAt(Carbon::now()->millisecond(0));
402 $this->persist($task);
403 $this->flush();
404 }
405
406 public function rescheduleTask(ScheduledTaskEntity $task): void {
407 if ($task->getStatus() !== ScheduledTaskEntity::STATUS_CANCELLED) {
408 throw new \Exception(__('Only cancelled tasks can be rescheduled', 'mailpoet'), 400);
409 }
410 if ($task->getScheduledAt() <= Carbon::now()->millisecond(0)) {
411 $task->setStatus(ScheduledTaskEntity::VIRTUAL_STATUS_RUNNING);
412 $queue = $task->getSendingQueue();
413 if ($queue) {
414 $this->sendingQueuesRepository->resume($queue);
415 }
416 } else {
417 $task->setStatus(ScheduledTaskEntity::STATUS_SCHEDULED);
418 }
419 $task->setCancelledAt(null);
420 $this->persist($task);
421 $this->flush();
422 }
423
424 /** @param int[] $ids */
425 public function deleteByIds(array $ids): void {
426 $this->entityManager->createQueryBuilder()
427 ->delete(ScheduledTaskEntity::class, 't')
428 ->where('t.id IN (:ids)')
429 ->setParameter('ids', $ids)
430 ->getQuery()
431 ->execute();
432
433 // delete was done via DQL, make sure the entities are also detached from the entity manager
434 $this->detachAll(function (ScheduledTaskEntity $entity) use ($ids) {
435 return in_array($entity->getId(), $ids, true);
436 });
437 }
438
439 protected function findByTypeAndStatus($type, $status, $limit = null, $future = false) {
440 $queryBuilder = $this->doctrineRepository->createQueryBuilder('st')
441 ->select('st')
442 ->where('st.type = :type')
443 ->setParameter('type', $type)
444 ->andWhere('st.deletedAt IS NULL');
445
446 if (is_null($status)) {
447 $queryBuilder->andWhere('st.status IS NULL');
448 } else {
449 $queryBuilder
450 ->andWhere('st.status = :status')
451 ->setParameter('status', $status);
452 }
453
454 if ($future) {
455 $queryBuilder->andWhere('st.scheduledAt > :now');
456 } else {
457 $queryBuilder->andWhere('st.scheduledAt <= :now');
458 }
459
460 $now = Carbon::now()->millisecond(0);
461 $queryBuilder->setParameter('now', $now);
462
463 if ($limit) {
464 $queryBuilder->setMaxResults($limit);
465 }
466
467 return $queryBuilder->getQuery()->getResult();
468 }
469
470 protected function getEntityClassName() {
471 return ScheduledTaskEntity::class;
472 }
473 }
474