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 / JobManager.php

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

335 lines 13.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 defined('ABSPATH') || exit;
6
7 class JobManager
8 {
9 private $jobManagerTableName;
10
11 private const ALLOWED_COLUMNS = [
12 'id', 'job_type', 'status', 'payload',
13 'attempts', 'max_attempts', 'retry_after',
14 'started_at', 'created_at', 'failed_at', 'error_message',
15 ];
16
17 function __construct()
18 {
19 global $wpdb;
20 $this->jobManagerTableName = $wpdb->prefix . 'sync_basalam_job_manager';
21 }
22
23 public function createJob($jobType, $status = 'pending', $payload = null, $maxAttempts = 3)
24 {
25 global $wpdb;
26
27 // phpcs:ignore WordPress.DB.DirectDatabaseQuery.DirectQuery, WordPress.DB.DirectDatabaseQuery.NoCaching -- Custom plugin table; no object cache for these operational queries.
28 $result = $wpdb->insert(
29 $this->jobManagerTableName,
30 array(
31 'job_type' => $jobType,
32 'status' => $status,
33 'payload' => $payload,
34 'attempts' => 0,
35 'max_attempts' => $maxAttempts,
36 'created_at' => time(),
37 )
38 );
39
40 if ($result !== false) {
41 do_action('sync_basalam_job_created', $jobType, $status, $payload);
42 }
43
44 return $result;
45 }
46
47 public function getNextEligibleJob(string $jobType): ?object
48 {
49 global $wpdb;
50
51 // phpcs:ignore WordPress.DB.DirectDatabaseQuery.DirectQuery, WordPress.DB.DirectDatabaseQuery.NoCaching, PluginCheck.Security.DirectDB.UnescapedDBParameter -- Custom plugin table; identifier from $wpdb->prefix, not user input.
52 return $wpdb->get_row($wpdb->prepare(
53 "SELECT * FROM {$this->jobManagerTableName}
54 WHERE job_type = %s
55 AND status = 'pending'
56 AND (retry_after IS NULL OR retry_after <= %d)
57 ORDER BY id ASC
58 LIMIT 1",
59 $jobType,
60 time()
61 ));
62 }
63
64 public function hasAnyProcessingJob(): bool
65 {
66 return $this->getCountJobs(['status' => 'processing']) > 0;
67 }
68
69 public function getProductUpdateStatus(): array
70 {
71 $quickJob = $this->getJob(['job_type' => 'sync_basalam_bulk_update_products', 'status' => 'processing'])
72 ?: $this->getJob(['job_type' => 'sync_basalam_bulk_update_products', 'status' => 'pending']);
73 $fullJob = $this->getJob(['job_type' => 'sync_basalam_update_all_products', 'status' => 'processing'])
74 ?: $this->getJob(['job_type' => 'sync_basalam_update_all_products', 'status' => 'pending']);
75 $singleCount = $this->getCountJobs([
76 'job_type' => 'sync_basalam_update_single_product',
77 'status' => ['pending', 'processing'],
78 ]);
79
80 return [
81 'active' => (bool) ($quickJob || $fullJob || $singleCount > 0),
82 'type' => $quickJob ? 'quick' : (($fullJob || $singleCount > 0) ? 'full' : ''),
83 'status' => $quickJob ? $quickJob->status : ($fullJob ? $fullJob->status : ''),
84 'count' => $singleCount,
85 ];
86 }
87
88 public function hasPendingOrStaleProcessingJobs(int $staleProcessingTimeoutSeconds = 120): bool
89 {
90 global $wpdb;
91
92 $now = time();
93 $staleBefore = $now - $staleProcessingTimeoutSeconds;
94
95 // phpcs:ignore WordPress.DB.DirectDatabaseQuery.DirectQuery, WordPress.DB.DirectDatabaseQuery.NoCaching, PluginCheck.Security.DirectDB.UnescapedDBParameter -- Custom plugin table; identifier from $wpdb->prefix, not user input.
96 $result = $wpdb->get_var($wpdb->prepare(
97 "SELECT 1 FROM {$this->jobManagerTableName}
98 WHERE (
99 status = 'pending'
100 AND (retry_after IS NULL OR retry_after <= %d)
101 ) OR (
102 status = 'processing'
103 AND started_at IS NOT NULL
104 AND started_at < %d
105 )
106 LIMIT 1",
107 $now,
108 $staleBefore
109 ));
110
111 return (string) $result === '1';
112 }
113
114 public function getJob($where = array())
115 {
116 global $wpdb;
117
118 if (empty($where)) return null;
119
120 $conditions = [];
121 $values = [];
122
123 foreach ($where as $column => $value) {
124 if (!in_array($column, self::ALLOWED_COLUMNS, true)) {
125 throw new \InvalidArgumentException(esc_html("Invalid column: {$column}"));
126 }
127 $conditions[] = "{$column} = %s";
128 $values[] = $value;
129 }
130
131 $sql = "SELECT * FROM {$this->jobManagerTableName} WHERE " . implode(" AND ", $conditions) . " LIMIT 1";
132
133 // phpcs:ignore WordPress.DB.DirectDatabaseQuery.DirectQuery, WordPress.DB.DirectDatabaseQuery.NoCaching, PluginCheck.Security.DirectDB.UnescapedDBParameter -- Custom plugin table; identifiers from $wpdb->prefix and whitelisted column list, not user input; values are prepared.
134 return $wpdb->get_row($wpdb->prepare($sql, $values));
135 }
136
137 public function getCountJobs($where = array())
138 {
139 global $wpdb;
140
141 if (empty($where)) return 0;
142
143 $conditions = [];
144 $values = [];
145
146 foreach ($where as $column => $value) {
147 if (!in_array($column, self::ALLOWED_COLUMNS, true)) {
148 throw new \InvalidArgumentException(esc_html("Invalid column: {$column}"));
149 }
150 if (is_array($value)) {
151 $placeholders = array_fill(0, count($value), '%s');
152 $conditions[] = "{$column} IN (" . implode(',', $placeholders) . ")";
153 $values = array_merge($values, $value);
154 } else {
155 $conditions[] = "{$column} = %s";
156 $values[] = $value;
157 }
158 }
159
160 $sql = "SELECT COUNT(*) FROM {$this->jobManagerTableName} WHERE " . implode(" AND ", $conditions);
161
162 // phpcs:ignore WordPress.DB.DirectDatabaseQuery.DirectQuery, WordPress.DB.DirectDatabaseQuery.NoCaching, PluginCheck.Security.DirectDB.UnescapedDBParameter -- Custom plugin table; identifiers from $wpdb->prefix and whitelisted column list, not user input; values are prepared.
163 return (int) $wpdb->get_var($wpdb->prepare($sql, $values));
164 }
165
166 public function updateJob($jobData, $where = array())
167 {
168 global $wpdb;
169
170 if (empty($where) || empty($jobData)) return false;
171
172 // phpcs:ignore WordPress.DB.DirectDatabaseQuery.DirectQuery, WordPress.DB.DirectDatabaseQuery.NoCaching -- Custom plugin table; no object cache for these operational queries.
173 return $wpdb->update($this->jobManagerTableName, $jobData, $where);
174 }
175
176 public function deleteJob($where = array())
177 {
178 global $wpdb;
179
180 if (empty($where)) return false;
181
182 // phpcs:ignore WordPress.DB.DirectDatabaseQuery.DirectQuery, WordPress.DB.DirectDatabaseQuery.NoCaching -- Custom plugin table; no object cache for these operational queries.
183 return $wpdb->delete($this->jobManagerTableName, $where);
184 }
185
186 public function ConvertStaleProcessingJobs($timeoutSeconds = 120)
187 {
188 global $wpdb;
189
190 $timeoutTimestamp = time() - $timeoutSeconds;
191
192 // phpcs:ignore WordPress.DB.DirectDatabaseQuery.DirectQuery, WordPress.DB.DirectDatabaseQuery.NoCaching, PluginCheck.Security.DirectDB.UnescapedDBParameter -- Custom plugin table; identifier from $wpdb->prefix, not user input.
193 $wpdb->query(
194 $wpdb->prepare(
195 "UPDATE {$this->jobManagerTableName}
196 SET status = CASE
197 WHEN attempts + 1 >= max_attempts THEN 'failed'
198 ELSE 'pending'
199 END,
200 attempts = attempts + 1,
201 started_at = NULL,
202 failed_at = CASE
203 WHEN attempts + 1 >= max_attempts THEN %d
204 ELSE failed_at
205 END
206 WHERE status = 'processing'
207 AND started_at IS NOT NULL
208 AND started_at < %d",
209 time(),
210 $timeoutTimestamp
211 )
212 );
213 }
214
215 public function hasProductJobInProgress(int $productId, string $jobType): bool
216 {
217 global $wpdb;
218
219 // phpcs:ignore WordPress.DB.DirectDatabaseQuery.DirectQuery, WordPress.DB.DirectDatabaseQuery.NoCaching, PluginCheck.Security.DirectDB.UnescapedDBParameter -- Custom plugin table; identifier from $wpdb->prefix, not user input.
220 $jobs = $wpdb->get_results($wpdb->prepare(
221 "SELECT payload FROM {$this->jobManagerTableName}
222 WHERE job_type = %s
223 AND (status = %s OR status = %s)",
224 $jobType,
225 'pending',
226 'processing'
227 ));
228
229 if (empty($jobs)) {
230 return false;
231 }
232
233 foreach ($jobs as $job) {
234 $payload = json_decode($job->payload, true);
235 $jobProductId = $payload['product_id'] ?? $payload;
236
237 if (intval($jobProductId) === intval($productId)) {
238 return true;
239 }
240 }
241
242 return false;
243 }
244
245 public function retryJob(int $jobId, ?string $errorMessage = null): bool
246 {
247 global $wpdb;
248
249 // phpcs:ignore WordPress.DB.DirectDatabaseQuery.DirectQuery, WordPress.DB.DirectDatabaseQuery.NoCaching, PluginCheck.Security.DirectDB.UnescapedDBParameter -- Custom plugin table; identifier from $wpdb->prefix, not user input.
250 $job = $wpdb->get_row($wpdb->prepare(
251 "SELECT * FROM {$this->jobManagerTableName} WHERE id = %d",
252 $jobId
253 ));
254
255 if (!$job) return false;
256
257 $newAttempts = intval($job->attempts) + 1;
258
259 $errorMessages = [];
260 if (!empty($job->error_message)) {
261 $decoded = json_decode($job->error_message, true);
262 if (json_last_error() === JSON_ERROR_NONE && is_array($decoded)) $errorMessages = $decoded;
263 }
264
265 if ($errorMessage) $errorMessages[$newAttempts] = $errorMessage;
266
267 $encodedErrors = json_encode($errorMessages, JSON_UNESCAPED_UNICODE);
268
269 if ($newAttempts >= intval($job->max_attempts)) {
270 $this->updateJob(
271 [
272 'status' => 'failed',
273 'error_message' => $encodedErrors,
274 'failed_at' => time(),
275 'started_at' => 0,
276 'attempts' => $newAttempts,
277 ],
278 ['id' => $jobId]
279 );
280 return false;
281 }
282
283 // Progressive exponential backoff: 30s, 60s, 120s, 240s, ...
284 $delaySeconds = 30 * (int) pow(2, $newAttempts - 1);
285 $retryAfter = time() + $delaySeconds;
286
287 // Atomic DELETE + INSERT inside a transaction so a crash can't lose the job.
288 // phpcs:ignore WordPress.DB.DirectDatabaseQuery.DirectQuery, WordPress.DB.DirectDatabaseQuery.NoCaching -- Transaction control for atomic job requeue; no object cache applicable.
289 $wpdb->query('START TRANSACTION');
290 try {
291 // phpcs:ignore WordPress.DB.DirectDatabaseQuery.DirectQuery, WordPress.DB.DirectDatabaseQuery.NoCaching -- Custom plugin table; no object cache for these operational queries.
292 $wpdb->delete($this->jobManagerTableName, ['id' => $jobId]);
293
294 // phpcs:ignore WordPress.DB.DirectDatabaseQuery.DirectQuery, WordPress.DB.DirectDatabaseQuery.NoCaching -- Custom plugin table; no object cache for these operational queries.
295 $wpdb->insert(
296 $this->jobManagerTableName,
297 [
298 'job_type' => $job->job_type,
299 'status' => 'pending',
300 'payload' => $job->payload,
301 'attempts' => $newAttempts,
302 'max_attempts' => $job->max_attempts,
303 'error_message' => $encodedErrors,
304 'created_at' => $job->created_at,
305 'retry_after' => $retryAfter,
306 'started_at' => 0,
307 ]
308 );
309
310 // phpcs:ignore WordPress.DB.DirectDatabaseQuery.DirectQuery, WordPress.DB.DirectDatabaseQuery.NoCaching -- Transaction control for atomic job requeue; no object cache applicable.
311 $wpdb->query('COMMIT');
312 } catch (\Exception $e) {
313 // phpcs:ignore WordPress.DB.DirectDatabaseQuery.DirectQuery, WordPress.DB.DirectDatabaseQuery.NoCaching -- Transaction control for atomic job requeue; no object cache applicable.
314 $wpdb->query('ROLLBACK');
315 throw $e;
316 }
317
318 return true;
319 }
320
321 public function failJob(int $jobId, ?string $errorMessage = null): bool
322 {
323 return $this->updateJob(
324 [
325 'status' => 'failed',
326 'error_message' => $errorMessage,
327 'failed_at' => time(),
328 'started_at' => 0,
329 ],
330 ['id' => $jobId]
331 );
332 }
333
334 }
335