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 / ScheduledTaskSubscribersRepository.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
ScheduledTaskSubscribersRepository.php
343 lines
1 <?php declare(strict_types = 1);
2
3 namespace MailPoet\Newsletter\Sending;
4
5 if (!defined('ABSPATH')) exit;
6
7
8 use MailPoet\Doctrine\Repository;
9 use MailPoet\Entities\ScheduledTaskEntity;
10 use MailPoet\Entities\ScheduledTaskSubscriberEntity;
11 use MailPoet\Entities\SubscriberEntity;
12 use MailPoet\InvalidStateException;
13 use MailPoetVendor\Carbon\Carbon;
14 use MailPoetVendor\Doctrine\DBAL\ArrayParameterType;
15 use MailPoetVendor\Doctrine\DBAL\ParameterType;
16 use MailPoetVendor\Doctrine\ORM\QueryBuilder;
17
18 /**
19 * @extends Repository<ScheduledTaskSubscriberEntity>
20 */
21 class ScheduledTaskSubscribersRepository extends Repository {
22 protected function getEntityClassName() {
23 return ScheduledTaskSubscriberEntity::class;
24 }
25
26 public function isSubscriberProcessed(ScheduledTaskEntity $task, SubscriberEntity $subscriber): bool {
27 $scheduledTaskSubscriber = $this
28 ->doctrineRepository
29 ->createQueryBuilder('sts')
30 ->andWhere('sts.processed = 1')
31 ->andWhere('sts.task = :task')
32 ->andWhere('sts.subscriber = :subscriber')
33 ->setParameter('subscriber', $subscriber)
34 ->setParameter('task', $task)
35 ->getQuery()
36 ->getOneOrNullResult();
37 return !empty($scheduledTaskSubscriber);
38 }
39
40 public function createOrUpdate(array $data): ?ScheduledTaskSubscriberEntity {
41 if (!isset($data['task_id'], $data['subscriber_id'])) {
42 return null;
43 }
44
45 $taskSubscriber = $this->findOneBy(['task' => $data['task_id'], 'subscriber' => $data['subscriber_id']]);
46 if (!$taskSubscriber) {
47 $task = $this->entityManager->getReference(ScheduledTaskEntity::class, (int)$data['task_id']);
48 $subscriber = $this->entityManager->getReference(SubscriberEntity::class, (int)$data['subscriber_id']);
49 if (!$task || !$subscriber) throw new InvalidStateException('Task or subscriber not found');
50
51 $taskSubscriber = new ScheduledTaskSubscriberEntity($task, $subscriber);
52 $this->persist($taskSubscriber);
53 }
54
55 $processed = $data['processed'] ?? ScheduledTaskSubscriberEntity::STATUS_UNPROCESSED;
56 $failed = $data['failed'] ?? ScheduledTaskSubscriberEntity::FAIL_STATUS_OK;
57
58 $taskSubscriber->setProcessed($processed);
59 $taskSubscriber->setFailed($failed);
60 $this->flush();
61 return $taskSubscriber;
62 }
63
64 public function countSubscriberIdsBatchForTask(int $taskId, int $lastProcessedSubscriberId): int {
65 $queryBuilder = $this->getBaseSubscribersIdsBatchForTaskQuery($taskId, $lastProcessedSubscriberId);
66 $countSubscribers = $queryBuilder
67 ->select('count(sts.subscriber)')
68 ->getQuery()
69 ->getSingleScalarResult();
70
71 return intval($countSubscribers);
72 }
73
74 public function getSubscriberIdsBatchForTask(int $taskId, int $lastProcessedSubscriberId, int $limit): array {
75 $queryBuilder = $this->getBaseSubscribersIdsBatchForTaskQuery($taskId, $lastProcessedSubscriberId);
76 $subscribersIds = $queryBuilder
77 ->select('IDENTITY(sts.subscriber) AS subscriber_id')
78 ->orderBy('sts.subscriber', 'asc')
79 ->setMaxResults($limit)
80 ->getQuery()
81 ->getSingleColumnResult();
82
83 return $subscribersIds;
84 }
85
86 /**
87 * @param int[] $subscriberIds
88 */
89 public function updateProcessedSubscribers(ScheduledTaskEntity $task, array $subscriberIds): void {
90 if ($subscriberIds) {
91 $this->entityManager->createQueryBuilder()
92 ->update(ScheduledTaskSubscriberEntity::class, 'sts')
93 ->set('sts.processed', ScheduledTaskSubscriberEntity::STATUS_PROCESSED)
94 ->where('sts.subscriber IN (:subscriberIds)')
95 ->andWhere('sts.task = :task')
96 ->setParameter('subscriberIds', $subscriberIds, ArrayParameterType::INTEGER)
97 ->setParameter('task', $task)
98 ->getQuery()
99 ->execute();
100
101 // update was done via DQL, make sure the entities are also refreshed in the entity manager
102 $this->refreshAll(function (ScheduledTaskSubscriberEntity $entity) use ($task, $subscriberIds) {
103 return $entity->getTask() === $task && in_array($entity->getSubscriberId(), $subscriberIds, true);
104 });
105 }
106
107 $this->checkCompleted($task);
108 }
109
110 /** @param int[] $subscriberIds */
111 public function addSubscribersByIds(ScheduledTaskEntity $task, array $subscriberIds): int {
112 $subscriberIds = array_values(array_unique(array_filter(array_map('intval', $subscriberIds))));
113 if ($subscriberIds === []) {
114 return 0;
115 }
116
117 $scheduledTaskSubscribersTable = $this->entityManager->getClassMetadata(ScheduledTaskSubscriberEntity::class)->getTableName();
118 $subscribersTable = $this->entityManager->getClassMetadata(SubscriberEntity::class)->getTableName();
119
120 $result = $this->entityManager->getConnection()->executeQuery(
121 "INSERT IGNORE INTO $scheduledTaskSubscribersTable
122 (task_id, subscriber_id, processed)
123 SELECT DISTINCT ? as task_id, subscribers.`id` as subscriber_id, ? as processed
124 FROM $subscribersTable subscribers
125 WHERE subscribers.`deleted_at` IS NULL
126 AND subscribers.`status` = ?
127 AND subscribers.`id` IN (?)",
128 [
129 $task->getId(),
130 ScheduledTaskSubscriberEntity::STATUS_UNPROCESSED,
131 SubscriberEntity::STATUS_SUBSCRIBED,
132 $subscriberIds,
133 ],
134 [
135 ParameterType::INTEGER,
136 ParameterType::INTEGER,
137 ParameterType::STRING,
138 ArrayParameterType::INTEGER,
139 ]
140 );
141
142 return (int)$result->rowCount();
143 }
144
145 /** @param int[] $ids */
146 public function deleteByTaskIds(array $ids): void {
147 $this->entityManager->createQueryBuilder()
148 ->delete(ScheduledTaskSubscriberEntity::class, 'sts')
149 ->where('sts.task IN (:taskIds)')
150 ->setParameter('taskIds', $ids)
151 ->getQuery()
152 ->execute();
153
154 // delete was done via DQL, make sure the entities are also detached from the entity manager
155 $this->detachAll(function (ScheduledTaskSubscriberEntity $entity) use ($ids) {
156 $task = $entity->getTask();
157 return $task && in_array($task->getId(), $ids, true);
158 });
159 }
160
161 public function deleteByScheduledTask(ScheduledTaskEntity $scheduledTask): void {
162 $this->entityManager->createQueryBuilder()
163 ->delete(ScheduledTaskSubscriberEntity::class, 'sts')
164 ->where('sts.task = :task')
165 ->setParameter('task', $scheduledTask)
166 ->getQuery()
167 ->execute();
168
169 // delete was done via DQL, make sure the entities are also detached from the entity manager
170 $this->detachAll(function (ScheduledTaskSubscriberEntity $entity) use ($scheduledTask) {
171 return $entity->getTask() === $scheduledTask;
172 });
173 }
174
175 public function deleteByScheduledTaskAndSubscriberIds(ScheduledTaskEntity $scheduledTask, array $subscriberIds): void {
176 $this->entityManager->createQueryBuilder()
177 ->delete(ScheduledTaskSubscriberEntity::class, 'sts')
178 ->where('sts.task = :task')
179 ->andWhere('sts.subscriber IN (:subscriberIds)')
180 ->setParameter('task', $scheduledTask)
181 ->setParameter('subscriberIds', $subscriberIds, ArrayParameterType::INTEGER)
182 ->getQuery()
183 ->execute();
184
185 // delete was done via DQL, make sure the entities are also detached from the entity manager
186 $this->detachAll(function (ScheduledTaskSubscriberEntity $entity) use ($scheduledTask, $subscriberIds) {
187 return $entity->getTask() === $scheduledTask && in_array($entity->getSubscriberId(), $subscriberIds, true);
188 });
189
190 $this->checkCompleted($scheduledTask);
191 }
192
193 public function setSubscribers(ScheduledTaskEntity $task, array $subscriberIds): void {
194 $this->deleteByScheduledTask($task);
195
196 foreach ($subscriberIds as $subscriberId) {
197 $this->createOrUpdate([
198 'task_id' => $task->getId(),
199 'subscriber_id' => $subscriberId,
200 ]);
201 }
202 }
203
204 public function saveError(ScheduledTaskEntity $scheduledTask, int $subscriberId, string $errorMessage): void {
205 $scheduledTaskSubscriber = $this->findOneBy(['task' => $scheduledTask, 'subscriber' => $subscriberId]);
206
207 if ($scheduledTaskSubscriber instanceof ScheduledTaskSubscriberEntity) {
208 $scheduledTaskSubscriber->setFailed(ScheduledTaskSubscriberEntity::FAIL_STATUS_FAILED);
209 $scheduledTaskSubscriber->setProcessed(ScheduledTaskSubscriberEntity::STATUS_PROCESSED);
210 $scheduledTaskSubscriber->setError($errorMessage);
211 $this->persist($scheduledTaskSubscriber);
212 $this->flush();
213
214 $this->checkCompleted($scheduledTask);
215 }
216 }
217
218 public function countProcessed(ScheduledTaskEntity $scheduledTaskEntity): int {
219 return $this->countBy(['task' => $scheduledTaskEntity, 'processed' => ScheduledTaskSubscriberEntity::STATUS_PROCESSED]);
220 }
221
222 public function countUnprocessed(ScheduledTaskEntity $scheduledTaskEntity): int {
223 return $this->countBy(['task' => $scheduledTaskEntity, 'processed' => ScheduledTaskSubscriberEntity::STATUS_UNPROCESSED]);
224 }
225
226 public function purgeOldTaskSubscribers(int $daysToKeep, int $taskBatchSize, int $rowLimit): int {
227 $stTable = $this->entityManager->getClassMetadata(ScheduledTaskEntity::class)->getTableName();
228 $stsTable = $this->entityManager->getClassMetadata(ScheduledTaskSubscriberEntity::class)->getTableName();
229 $cutoff = Carbon::now()->subDays($daysToKeep)->toDateTimeString();
230
231 $taskIds = $this->entityManager->getConnection()->executeQuery(
232 "SELECT DISTINCT st.`id`
233 FROM `{$stTable}` st
234 INNER JOIN `{$stsTable}` sts ON sts.`task_id` = st.`id`
235 WHERE st.`type` = :type
236 AND st.`status` = :status
237 AND st.`processed_at` < :cutoff
238 AND st.`deleted_at` IS NULL
239 LIMIT :taskBatchSize",
240 [
241 'type' => 'sending',
242 'status' => ScheduledTaskEntity::STATUS_COMPLETED,
243 'cutoff' => $cutoff,
244 'taskBatchSize' => $taskBatchSize,
245 ],
246 [
247 'type' => ParameterType::STRING,
248 'status' => ParameterType::STRING,
249 'cutoff' => ParameterType::STRING,
250 'taskBatchSize' => ParameterType::INTEGER,
251 ]
252 )->fetchFirstColumn();
253
254 if (!$taskIds) {
255 return 0;
256 }
257
258 /** @var int[] $taskIds */
259 $taskIdsList = implode(',', array_map('intval', $taskIds));
260
261 $deleted = $this->entityManager->getConnection()->executeStatement(
262 "DELETE FROM `{$stsTable}`
263 WHERE `task_id` IN ({$taskIdsList})
264 LIMIT :rowLimit",
265 [
266 'rowLimit' => $rowLimit,
267 ],
268 [
269 'rowLimit' => ParameterType::INTEGER,
270 ]
271 );
272
273 return (int)$deleted;
274 }
275
276 public function purgeCompletedBounceTaskSubscribers(int $taskBatchSize, int $rowLimit): int {
277 $stTable = $this->entityManager->getClassMetadata(ScheduledTaskEntity::class)->getTableName();
278 $stsTable = $this->entityManager->getClassMetadata(ScheduledTaskSubscriberEntity::class)->getTableName();
279
280 $taskIds = $this->entityManager->getConnection()->executeQuery(
281 "SELECT DISTINCT st.`id`
282 FROM `{$stTable}` st
283 INNER JOIN `{$stsTable}` sts ON sts.`task_id` = st.`id`
284 WHERE st.`type` = :type
285 AND st.`status` = :status
286 AND st.`deleted_at` IS NULL
287 LIMIT :taskBatchSize",
288 [
289 'type' => 'bounce',
290 'status' => ScheduledTaskEntity::STATUS_COMPLETED,
291 'taskBatchSize' => $taskBatchSize,
292 ],
293 [
294 'type' => ParameterType::STRING,
295 'status' => ParameterType::STRING,
296 'taskBatchSize' => ParameterType::INTEGER,
297 ]
298 )->fetchFirstColumn();
299
300 if (!$taskIds) {
301 return 0;
302 }
303
304 /** @var int[] $taskIds */
305 $taskIdsList = implode(',', array_map('intval', $taskIds));
306
307 $deleted = $this->entityManager->getConnection()->executeStatement(
308 "DELETE FROM `{$stsTable}`
309 WHERE `task_id` IN ({$taskIdsList})
310 LIMIT :rowLimit",
311 [
312 'rowLimit' => $rowLimit,
313 ],
314 [
315 'rowLimit' => ParameterType::INTEGER,
316 ]
317 );
318
319 return (int)$deleted;
320 }
321
322 private function checkCompleted(ScheduledTaskEntity $task): void {
323 $count = $this->countUnprocessed($task);
324 if ($count === 0) {
325 $task->setStatus(ScheduledTaskEntity::STATUS_COMPLETED);
326 $task->setProcessedAt(Carbon::now()->millisecond(0));
327 $this->entityManager->flush();
328 }
329 }
330
331 private function getBaseSubscribersIdsBatchForTaskQuery(int $taskId, int $lastProcessedSubscriberId): QueryBuilder {
332 return $this->entityManager
333 ->createQueryBuilder()
334 ->from(ScheduledTaskSubscriberEntity::class, 'sts')
335 ->andWhere('sts.task = :taskId')
336 ->andWhere('sts.subscriber > :lastProcessedSubscriberId')
337 ->andWhere('sts.processed = :status')
338 ->setParameter('taskId', $taskId)
339 ->setParameter('lastProcessedSubscriberId', $lastProcessedSubscriberId)
340 ->setParameter('status', ScheduledTaskSubscriberEntity::STATUS_UNPROCESSED);
341 }
342 }
343