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 / Cron / CliCommands / TaskRunner.php
mailpoet / lib / Cron / CliCommands Last commit date
ClaimedTaskRunner.php 1 month ago Cli.php 1 month ago CronCommand.php 1 month ago DaemonRunner.php 1 month ago ExecutionLimitOverride.php 1 month ago ScheduledTaskResolver.php 1 month ago ScheduledTasksLister.php 1 month ago TaskAdder.php 1 month ago TaskCanceller.php 1 month ago TaskRunner.php 1 month ago TaskTrigger.php 1 month ago WorkerTypesCatalog.php 1 month ago index.php 1 month ago
TaskRunner.php
288 lines
1 <?php declare(strict_types = 1);
2
3 namespace MailPoet\Cron\CliCommands;
4
5 if (!defined('ABSPATH')) exit;
6
7
8 use InvalidArgumentException;
9 use MailPoet\Cron\CronHelper;
10 use MailPoet\Cron\CronWorkerInterface;
11 use MailPoet\Cron\Workers\SendingQueue\SendingQueue as SendingQueueWorker;
12 use MailPoet\Cron\Workers\StatsNotifications\Worker as StatsNotificationsWorker;
13 use MailPoet\Cron\Workers\WorkersFactory;
14 use MailPoet\Entities\ScheduledTaskEntity;
15 use MailPoet\InvalidStateException;
16 use MailPoet\Newsletter\Sending\ScheduledTasksRepository;
17 use RuntimeException;
18 use Throwable;
19
20 /**
21 * Runs a MailPoet cron worker inside the current process.
22 *
23 * Bulk `run <type>` snapshots the type's currently-due tasks, then claims each as a CLI row (status
24 * 'cli', invisible to the site daemon) and runs it once through the shared ClaimedTaskRunner — the same
25 * claim model as `run --task-id`. This closes the concurrency hole (the daemon can never double-process
26 * a row the CLI owns) and stops self-rescheduling batched workers from running away: continuations they
27 * create mid-run are not in the snapshot, so each due task processes one batch and the continuation is
28 * left for the site cron. Mailing workers run via process(), and 'sending' runs the Scheduler then the
29 * SendingQueue. `run --task-id` claims one exact row through ClaimedTaskRunner. The 20-second execution
30 * limit is lifted by default (or capped via $timeout).
31 *
32 * See doc/wp-cli-cron-commands.md for command behaviour and the cli-claim rationale.
33 */
34 class TaskRunner {
35 private WorkerTypesCatalog $workerTypesCatalog;
36
37 private WorkersFactory $workersFactory;
38
39 private ExecutionLimitOverride $executionLimitOverride;
40
41 private ScheduledTasksRepository $scheduledTasksRepository;
42
43 private ScheduledTaskResolver $taskResolver;
44
45 private ClaimedTaskRunner $claimedTaskRunner;
46
47 public function __construct(
48 WorkerTypesCatalog $workerTypesCatalog,
49 WorkersFactory $workersFactory,
50 ExecutionLimitOverride $executionLimitOverride,
51 ScheduledTasksRepository $scheduledTasksRepository,
52 ScheduledTaskResolver $taskResolver,
53 ClaimedTaskRunner $claimedTaskRunner
54 ) {
55 $this->workerTypesCatalog = $workerTypesCatalog;
56 $this->workersFactory = $workersFactory;
57 $this->executionLimitOverride = $executionLimitOverride;
58 $this->scheduledTasksRepository = $scheduledTasksRepository;
59 $this->taskResolver = $taskResolver;
60 $this->claimedTaskRunner = $claimedTaskRunner;
61 }
62
63 /**
64 * @return array{completed: int, message: string, limit_reached: bool, backlog_drained: bool}
65 */
66 public function run(string $type, ?int $taskId = null, ?int $timeout = null): array {
67 $this->workerTypesCatalog->assertValidType($type);
68
69 if ($taskId !== null) {
70 return $this->runClaimedTask($type, $taskId, $timeout);
71 }
72
73 if ($type === SendingQueueWorker::TASK_TYPE || $type === StatsNotificationsWorker::TASK_TYPE) {
74 return $this->runMailing($type, $timeout);
75 }
76
77 return $this->runDueTasks($type, $timeout);
78 }
79
80 /**
81 * --task-id: resolve the exact row (scheduled/paused only, type must match), claim it as a CLI row
82 * preserving its meta, and run it through the shared ClaimedTaskRunner. The shared CronWorkerRunner
83 * cannot see a STATUS_CLI row, so claiming is the only way the exact row is processed in-CLI.
84 *
85 * @return array{completed: int, message: string, limit_reached: bool, backlog_drained: bool}
86 */
87 private function runClaimedTask(string $type, int $taskId, ?int $timeout): array {
88 $worker = $this->workerTypesCatalog->getWorkerByType($type);
89 if ($worker === null) {
90 throw new InvalidArgumentException("Task type '{$type}' has no runnable worker, so --task-id cannot run it in-process. Use `wp mailpoet cron trigger` for mailing types.");
91 }
92
93 $task = $this->taskResolver->resolveById($taskId, $type);
94 $status = $task->getStatus();
95 if (!in_array($status, [ScheduledTaskEntity::STATUS_SCHEDULED, ScheduledTaskEntity::STATUS_PAUSED], true)) {
96 $current = $this->taskResolver->nameStatus($task);
97 throw new InvalidArgumentException("Task {$taskId} is '{$current}' and cannot be run. Only scheduled or paused tasks can be run by ID.");
98 }
99
100 if (!$this->claimedTaskRunner->claimExisting($task)) {
101 throw new RuntimeException(sprintf('Task %d was claimed by another process; nothing ran.', $taskId));
102 }
103
104 // ClaimedTaskRunner installs the execution-limit override itself, so the timeout is passed through
105 // rather than wrapped here — wrapping would be defeated by its own (inner) override.
106 $runResult = $this->claimedTaskRunner->run($worker, $task, $timeout);
107
108 return [
109 'completed' => $runResult['completed'] ? 1 : 0,
110 'limit_reached' => $runResult['limit_reached'],
111 'backlog_drained' => $runResult['completed'],
112 'message' => $runResult['message'],
113 ];
114 }
115
116 /**
117 * Bulk `run <type>` for standard workers: snapshot the currently-due tasks, then claim and run each one
118 * exactly once through the shared ClaimedTaskRunner. Tasks are claimed as 'cli' so the site daemon
119 * cannot double-process them, and continuations created during the run (e.g. self-rescheduling batched
120 * workers) are not in the snapshot, so they are left for the site cron instead of being chased.
121 *
122 * @return array{completed: int, message: string, limit_reached: bool, backlog_drained: bool}
123 */
124 private function runDueTasks(string $type, ?int $timeout): array {
125 $worker = $this->resolveWorker($type);
126 if ($worker === null) {
127 // Should be unreachable: assertValidType allows only mailing + standard types, and mailing types
128 // are handled before this method is called.
129 throw new InvalidArgumentException("Task type '{$type}' has no runnable worker.");
130 }
131
132 // Pre-check requirements once so we never claim a task only to remove it on the requirements path.
133 try {
134 $requirementsMet = $worker->checkProcessingRequirements();
135 } catch (Throwable $e) {
136 throw new RuntimeException(sprintf("Requirements check for '%s' failed: %s. Nothing ran.", $type, $e->getMessage()), 0, $e);
137 }
138 if (!$requirementsMet) {
139 // Nothing ran and the due tasks are left scheduled, so the backlog is not drained: surface a
140 // warning rather than a green success.
141 return [
142 'completed' => 0,
143 'limit_reached' => false,
144 'backlog_drained' => false,
145 'message' => sprintf("Ran '%s': requirements not met, nothing ran.", $type),
146 ];
147 }
148
149 $dueTasks = $this->scheduledTasksRepository->findDueByType($type);
150
151 $completed = 0;
152 $handedBack = 0;
153 $limitReached = false;
154 $start = microtime(true);
155
156 foreach ($dueTasks as $task) {
157 // With --timeout, each task gets the run's remaining budget (capping the worker's own execution
158 // limit too, not just the gap between tasks); once it is spent we stop starting new tasks.
159 $remaining = null;
160 if ($timeout !== null) {
161 $remaining = $timeout - (microtime(true) - $start);
162 if ($remaining <= 0) {
163 $limitReached = true;
164 break;
165 }
166 }
167
168 // Atomic claim: another CLI run (or the daemon) may have taken this row since the snapshot, so a
169 // lost claim is skipped rather than double-processed.
170 if (!$this->claimedTaskRunner->claimExisting($task)) {
171 continue;
172 }
173
174 // A worker failure aborts the bulk run: already-processed tasks stay done, the failing one is
175 // handed back by ClaimedTaskRunner, and the RuntimeException propagates to the caller.
176 $result = $this->claimedTaskRunner->run($worker, $task, $remaining === null ? null : (int)ceil($remaining));
177 if ($result['limit_reached']) {
178 // The task hit the cap and was handed back; stop the run so it is not counted as completed.
179 $limitReached = true;
180 break;
181 }
182 $result['completed'] ? $completed++ : $handedBack++;
183 }
184
185 // Handed-back tasks (partial work, not ready, or a failed worker) are not "drained" — the command
186 // surfaces this as a warning so it is visible to operators and scripts.
187 $backlogDrained = !$limitReached && $handedBack === 0;
188
189 return [
190 'completed' => $completed,
191 'limit_reached' => $limitReached,
192 'backlog_drained' => $backlogDrained,
193 'message' => $this->buildMessage($type, $completed, $handedBack, $limitReached, $timeout),
194 ];
195 }
196
197 /**
198 * Mailing types ('sending', 'stats_notification') run their own mailer-driven flow instead of
199 * CronWorkerInterface, via their own process step under the execution-limit override, surfacing a hit
200 * limit as limit_reached.
201 *
202 * @return array{completed: int, message: string, limit_reached: bool, backlog_drained: bool}
203 */
204 private function runMailing(string $type, ?int $timeout): array {
205 $limitReached = false;
206
207 try {
208 $this->executionLimitOverride->overrideDuring($timeout, function () use ($type): void {
209 if ($type === SendingQueueWorker::TASK_TYPE) {
210 $this->runSending();
211 return;
212 }
213 $this->workersFactory->createStatsNotificationsWorker()->process();
214 });
215 } catch (\Exception $e) {
216 if ($e->getCode() === CronHelper::DAEMON_EXECUTION_LIMIT_REACHED) {
217 $limitReached = true;
218 } else {
219 throw $e;
220 }
221 }
222
223 if ($limitReached) {
224 return [
225 'completed' => 0,
226 'limit_reached' => true,
227 'backlog_drained' => false,
228 'message' => sprintf("Execution limit of %d seconds reached while running '%s'.", (int)$timeout, $type),
229 ];
230 }
231
232 return [
233 'completed' => 0,
234 'limit_reached' => false,
235 'backlog_drained' => true,
236 'message' => sprintf("Ran '%s'.", $type),
237 ];
238 }
239
240 /**
241 * Seam over the catalog lookup so tests can drive the run with a stub worker.
242 */
243 protected function resolveWorker(string $type): ?CronWorkerInterface {
244 return $this->workerTypesCatalog->getWorkerByType($type);
245 }
246
247 private function runSending(): void {
248 // Scheduling first, then sending: the Scheduler turns scheduled newsletters into running sending
249 // tasks the SendingQueue then picks up. Instantiated lazily because the queue worker resolves the
250 // mailer eagerly and throws on a site with no sender configured.
251 try {
252 $scheduler = $this->workersFactory->createScheduleWorker();
253 $queue = $this->workersFactory->createQueueWorker();
254 } catch (InvalidStateException $e) {
255 throw new RuntimeException('Sending is not configured on this site: ' . $e->getMessage());
256 }
257
258 $scheduler->process();
259 $queue->process();
260 }
261
262 private function buildMessage(string $type, int $completed, int $handedBack, bool $limitReached, ?int $timeout): string {
263 if ($limitReached) {
264 return sprintf(
265 "Execution limit of %d seconds reached while running '%s'. %d task(s) completed; remaining due tasks will run on the next invocation.",
266 (int)$timeout,
267 $type,
268 $completed
269 );
270 }
271
272 if ($handedBack > 0) {
273 return sprintf(
274 "Ran '%s': %d task(s) completed; %d task(s) handed back to the site cron (not ready or failed).",
275 $type,
276 $completed,
277 $handedBack
278 );
279 }
280
281 if ($completed > 0) {
282 return sprintf("Ran '%s': %d task(s) completed.", $type, $completed);
283 }
284
285 return sprintf("Ran '%s': no tasks completed (nothing was due).", $type);
286 }
287 }
288