| @@ -8,9 +8,19 @@ | ||
| 8 | 8 | defined('ABSPATH') || exit; |
| 9 | 9 | |
| 10 | 10 | class JobsRunner |
| 11 | 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; | |
| 12 | 21 | private const GLOBAL_RUNNER_LAST_RUN_OPTION = 'sync_basalam_jobs_runner_last_run'; |
| 22 | + private const STALE_PROCESSING_TIMEOUT_SECONDS = 120; | |
| 13 | 23 | |
| 14 | 24 | private $jobExecutor; |
| 15 | 25 | private $jobManager; |
| 16 | 26 | private $discountScheduler; |
| @@ -21,9 +31,14 @@ | ||
| 21 | 31 | $jobExecutor, |
| 22 | 32 | $discountScheduler, |
| 23 | 33 | $CheckHttpBlockService |
| 24 | 34 | ) { |
| 25 | - add_action('init', [$this, 'checkAndRunJobs']); | |
| 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 | + | |
| 26 | 41 | $this->jobManager = $jobManager; |
| 27 | 42 | $this->jobExecutor = $jobExecutor; |
| 28 | 43 | $this->discountScheduler = $discountScheduler; |
| 29 | 44 | $this->CheckHttpBlockService = $CheckHttpBlockService; |
| @@ -28,41 +43,162 @@ | ||
| 28 | 43 | $this->discountScheduler = $discountScheduler; |
| 29 | 44 | $this->CheckHttpBlockService = $CheckHttpBlockService; |
| 30 | 45 | } |
| 31 | 46 | |
| 32 | - public function checkAndRunJobs(): void | |
| 47 | + public function maybeDispatchAsyncRequest(): void | |
| 33 | 48 | { |
| 49 | + if ($this->isCurrentAsyncRequest()) return; | |
| 50 | + | |
| 51 | + if ($this->isShutdownCallback()) { | |
| 52 | + $this->finishFastCgiResponse(); | |
| 53 | + $this->repairDatabaseConnection(); | |
| 54 | + } | |
| 55 | + | |
| 34 | 56 | if ($this->CheckHttpBlockService->SyncBasalamHttpBlock()) return; |
| 35 | - if (!$this->jobExecutor->acquireGlobalJobsLock(0)) 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; | |
| 36 | 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 | + | |
| 37 | 128 | try { |
| 38 | - $this->runEligibleJobs(); | |
| 129 | + return $this->runEligibleJobs(); | |
| 39 | 130 | } finally { |
| 40 | 131 | $this->jobExecutor->releaseGlobalJobsLock(); |
| 41 | 132 | } |
| 42 | 133 | } |
| 43 | 134 | |
| 44 | - private function runEligibleJobs(): void | |
| 135 | + private function runAsyncBatch(): int | |
| 45 | 136 | { |
| 46 | - $this->jobManager->ConvertStaleProcessingJobs(120); | |
| 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); | |
| 47 | 185 | $this->discountScheduler->process(); |
| 48 | 186 | |
| 49 | 187 | $circuitBreaker = new CircuitBreaker(); |
| 50 | 188 | if ($circuitBreaker->getState() === CircuitBreaker::STATE_OPEN) { |
| 51 | - return; | |
| 189 | + return false; | |
| 52 | 190 | } |
| 53 | 191 | |
| 54 | 192 | if ($this->jobManager->hasAnyProcessingJob()) { |
| 55 | - return; | |
| 193 | + return false; | |
| 56 | 194 | } |
| 57 | 195 | |
| 58 | - $tasksPerMinute = max(1, intval(Settings::getEffectiveTasksPerMinute())); | |
| 59 | - $thresholdSeconds = 60.0 / $tasksPerMinute; | |
| 60 | 196 | $lastRun = floatval(get_option(self::GLOBAL_RUNNER_LAST_RUN_OPTION, 0)); |
| 61 | 197 | $now = microtime(true); |
| 62 | 198 | |
| 63 | - if (($now - $lastRun) < $thresholdSeconds) { | |
| 64 | - return; | |
| 199 | + if (($now - $lastRun) < $this->getRunThresholdSeconds()) { | |
| 200 | + return false; | |
| 65 | 201 | } |
| 66 | 202 | |
| 67 | 203 | $sortedJobTypes = $this->jobExecutor->getSortedJobTypes(); |
| 68 | 204 | |
| @@ -85,10 +221,12 @@ | ||
| 85 | 221 | ['id' => $job->id] |
| 86 | 222 | ); |
| 87 | 223 | |
| 88 | 224 | $this->executeJob($job); |
| 89 | - break; | |
| 225 | + return true; | |
| 90 | 226 | } |
| 227 | + | |
| 228 | + return false; | |
| 91 | 229 | } |
| 92 | 230 | |
| 93 | 231 | private function executeJob(object $job): void |
| 94 | 232 | { |
| @@ -93,6 +231,49 @@ | ||
| 93 | 231 | private function executeJob(object $job): void |
| 94 | 232 | { |
| 95 | 233 | $jobType = $job->job_type; |
| 96 | 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; | |
| 97 | 278 | } |
| 98 | 279 | } |