PluginProbe
ووسلام – همگام سازی ووکامرس و باسلام / 1.10.22
ووسلام – همگام سازی ووکامرس و باسلام v1.10.22
1.10.22 1.10.21 1.10.19 1.10.20 1.10.18 1.10.17 1.10.15 1.10.14 1.10.13 1.10.12 1.10.10 1.10.9 1.10.8 1.10.7 1.10.6 1.10.5 1.10.4 1.10.3 1.10.2 1.10.1 1.10.0 1.9.2 1.9.1 1.9.0 1.8.8 All 55 releases
sync-basalam / JobsRunner.php

JobsRunner.php in ووسلام – همگام سازی ووکامرس و باسلام 1.10.22, at JobsRunner.php

280 lines 9.0 KB
No matching file
Up and down to move Enter to open Esc to close
Raw Download Zip
1 <?php
2
3 namespace SyncBasalam;
4
5 use SyncBasalam\Admin\Settings;
6 use SyncBasalam\Services\Api\CircuitBreaker;
7
8 defined('ABSPATH') || exit;
9
10 class JobsRunner
11 {
12 private const ASYNC_ACTION = 'sync_basalam_run_jobs_async';
13 private const ASYNC_DISPATCH_LOCK_TRANSIENT = 'sync_basalam_jobs_runner_async_dispatch_lock';
14 private const ASYNC_IDLE_PROBE_TRANSIENT = 'sync_basalam_jobs_runner_idle_probe_lock';
15 // Keep the dispatch lease longer than a normal async batch. Without this,
16 // every frontend request can boot another full WordPress AJAX worker while
17 // a large product queue is active.
18 private const ASYNC_DISPATCH_LOCK_SECONDS = 25;
19 private const ASYNC_IDLE_PROBE_SECONDS = 5;
20 private const ASYNC_TIME_LIMIT_SECONDS = 20;
21 private const GLOBAL_RUNNER_LAST_RUN_OPTION = 'sync_basalam_jobs_runner_last_run';
22 private const STALE_PROCESSING_TIMEOUT_SECONDS = 120;
23
24 private $jobExecutor;
25 private $jobManager;
26 private $discountScheduler;
27 private $CheckHttpBlockService;
28
29 public function __construct(
30 $jobManager,
31 $jobExecutor,
32 $discountScheduler,
33 $CheckHttpBlockService
34 ) {
35 add_action('wp_ajax_' . self::ASYNC_ACTION, [$this, 'handleAsyncRequest']);
36 add_action('wp_ajax_nopriv_' . self::ASYNC_ACTION, [$this, 'handleAsyncRequest']);
37 // Run before later shutdown callbacks can leave an unread mysqli result.
38 add_action('shutdown', [$this, 'maybeDispatchAsyncRequest'], 2);
39 add_action('sync_basalam_job_created', [$this, 'maybeDispatchAsyncRequest'], 10, 0);
40
41 $this->jobManager = $jobManager;
42 $this->jobExecutor = $jobExecutor;
43 $this->discountScheduler = $discountScheduler;
44 $this->CheckHttpBlockService = $CheckHttpBlockService;
45 }
46
47 public function maybeDispatchAsyncRequest(): void
48 {
49 if ($this->isCurrentAsyncRequest()) return;
50
51 if ($this->isShutdownCallback()) {
52 $this->finishFastCgiResponse();
53 $this->repairDatabaseConnection();
54 }
55
56 if ($this->CheckHttpBlockService->SyncBasalamHttpBlock()) return;
57 $isNewJob = function_exists('current_filter') && current_filter() === 'sync_basalam_job_created';
58 if (!$isNewJob && get_transient(self::ASYNC_IDLE_PROBE_TRANSIENT)) return;
59 if (get_transient(self::ASYNC_DISPATCH_LOCK_TRANSIENT)) return;
60
61 if (!$this->jobManager->hasPendingOrStaleProcessingJobs(self::STALE_PROCESSING_TIMEOUT_SECONDS)) {
62 if (!$isNewJob) {
63 set_transient(self::ASYNC_IDLE_PROBE_TRANSIENT, 1, self::ASYNC_IDLE_PROBE_SECONDS);
64 }
65 return;
66 }
67
68 // An empty-queue probe must not block a job created moments later.
69 // The global database lock still prevents duplicate workers from
70 // processing the same queue when requests race here.
71 if (get_transient(self::ASYNC_DISPATCH_LOCK_TRANSIENT)) return;
72 set_transient(self::ASYNC_DISPATCH_LOCK_TRANSIENT, 1, self::ASYNC_DISPATCH_LOCK_SECONDS);
73
74 $this->dispatchAsyncRequest();
75 }
76
77 private function isShutdownCallback(): bool
78 {
79 return function_exists('current_filter') && current_filter() === 'shutdown';
80 }
81
82 private function finishFastCgiResponse(): void
83 {
84 if (function_exists('fastcgi_finish_request')) {
85 fastcgi_finish_request();
86 }
87 }
88
89 private function repairDatabaseConnection(): void
90 {
91 global $wpdb;
92
93 if (!isset($wpdb) || !is_object($wpdb)) return;
94
95 if (method_exists($wpdb, 'flush')) {
96 $wpdb->flush();
97 }
98
99 if (method_exists($wpdb, 'check_connection')) {
100 $wpdb->check_connection(false);
101 }
102 }
103
104 public function handleAsyncRequest(): void
105 {
106 if (!check_ajax_referer(self::ASYNC_ACTION, 'nonce', false)) {
107 wp_send_json_error(['message' => 'Invalid async jobs runner nonce.'], 403);
108 }
109
110 if (function_exists('ignore_user_abort')) {
111 ignore_user_abort(true);
112 }
113
114 if (function_exists('session_write_close')) {
115 session_write_close();
116 }
117
118 $processed = $this->runAsyncBatch();
119
120 wp_send_json_success(['processed' => $processed]);
121 }
122
123 public function checkAndRunJobs(): bool
124 {
125 if ($this->CheckHttpBlockService->SyncBasalamHttpBlock()) return false;
126 if (!$this->jobExecutor->acquireGlobalJobsLock(0)) return false;
127
128 try {
129 return $this->runEligibleJobs();
130 } finally {
131 $this->jobExecutor->releaseGlobalJobsLock();
132 }
133 }
134
135 private function runAsyncBatch(): int
136 {
137 if ($this->CheckHttpBlockService->SyncBasalamHttpBlock()) return 0;
138
139 // Hold the advisory lock for the whole batch, including rate-limit
140 // waits. Previously it was released after every job, so duplicate
141 // async requests could pile up and sleep in parallel until the next
142 // job became eligible, exhausting the site's PHP workers.
143 if (!$this->jobExecutor->acquireGlobalJobsLock(0)) return 0;
144
145 $processed = 0;
146 $deadline = microtime(true) + (float) apply_filters(
147 'sync_basalam_jobs_runner_async_time_limit',
148 self::ASYNC_TIME_LIMIT_SECONDS
149 );
150
151 try {
152 while (microtime(true) < $deadline) {
153 if (!$this->jobManager->hasPendingOrStaleProcessingJobs(self::STALE_PROCESSING_TIMEOUT_SECONDS)) {
154 break;
155 }
156
157 $ranJob = $this->runEligibleJobs();
158
159 if ($ranJob) {
160 $processed++;
161 }
162
163 $delay = $this->secondsUntilNextAllowedRun();
164 if ($delay <= 0.0) {
165 if (!$ranJob) break;
166 continue;
167 }
168
169 if ((microtime(true) + $delay) >= $deadline) {
170 break;
171 }
172
173 usleep((int) ($delay * 1000000));
174 }
175 } finally {
176 $this->jobExecutor->releaseGlobalJobsLock();
177 }
178
179 return $processed;
180 }
181
182 private function runEligibleJobs(): bool
183 {
184 $this->jobManager->ConvertStaleProcessingJobs(self::STALE_PROCESSING_TIMEOUT_SECONDS);
185 $this->discountScheduler->process();
186
187 $circuitBreaker = new CircuitBreaker();
188 if ($circuitBreaker->getState() === CircuitBreaker::STATE_OPEN) {
189 return false;
190 }
191
192 if ($this->jobManager->hasAnyProcessingJob()) {
193 return false;
194 }
195
196 $lastRun = floatval(get_option(self::GLOBAL_RUNNER_LAST_RUN_OPTION, 0));
197 $now = microtime(true);
198
199 if (($now - $lastRun) < $this->getRunThresholdSeconds()) {
200 return false;
201 }
202
203 $sortedJobTypes = $this->jobExecutor->getSortedJobTypes();
204
205 foreach ($sortedJobTypes as $jobType => $jobExecutor) {
206 if (!$this->jobExecutor->canRun($jobType)) {
207 continue;
208 }
209
210 $job = $this->jobManager->getNextEligibleJob($jobType);
211 $processingJob = $this->jobManager->getJob(['job_type' => $jobType, 'status' => 'processing']);
212
213 if (!$job || $processingJob) {
214 continue;
215 }
216
217 update_option(self::GLOBAL_RUNNER_LAST_RUN_OPTION, microtime(true), false);
218
219 $this->jobManager->updateJob(
220 ['status' => 'processing', 'started_at' => time()],
221 ['id' => $job->id]
222 );
223
224 $this->executeJob($job);
225 return true;
226 }
227
228 return false;
229 }
230
231 private function executeJob(object $job): void
232 {
233 $jobType = $job->job_type;
234 $this->jobExecutor->execute($jobType, $job);
235 }
236
237 private function dispatchAsyncRequest(): void
238 {
239 $url = add_query_arg('action', self::ASYNC_ACTION, admin_url('admin-ajax.php'));
240
241 wp_remote_post(esc_url_raw($url), [
242 'timeout' => 0.01,
243 'blocking' => false,
244 'body' => [
245 'action' => self::ASYNC_ACTION,
246 'nonce' => wp_create_nonce(self::ASYNC_ACTION),
247 ],
248 'cookies' => $_COOKIE,
249 'sslverify' => apply_filters('https_local_ssl_verify', false),
250 'headers' => [
251 'X-WP-Async-Request' => self::ASYNC_ACTION,
252 ],
253 ]);
254 }
255
256 private function isCurrentAsyncRequest(): bool
257 {
258 if (!wp_doing_ajax()) return false;
259
260 $action = isset($_REQUEST['action']) ? sanitize_key(wp_unslash($_REQUEST['action'])) : '';
261
262 return $action === self::ASYNC_ACTION;
263 }
264
265 private function secondsUntilNextAllowedRun(): float
266 {
267 $lastRun = floatval(get_option(self::GLOBAL_RUNNER_LAST_RUN_OPTION, 0));
268 $elapsed = microtime(true) - $lastRun;
269
270 return max(0.0, $this->getRunThresholdSeconds() - $elapsed);
271 }
272
273 private function getRunThresholdSeconds(): float
274 {
275 $tasksPerMinute = max(1, intval(Settings::getEffectiveTasksPerMinute()));
276
277 return 60.0 / $tasksPerMinute;
278 }
279 }
280