Code Coverage |
||||||||||
Lines |
Functions and Methods |
Classes and Traits |
||||||||
| Total | |
98.79% |
163 / 165 |
|
96.88% |
62 / 64 |
CRAP | |
0.00% |
0 / 1 |
| AbstractJob | |
98.79% |
163 / 165 |
|
96.88% |
62 / 64 |
110 | |
0.00% |
0 / 1 |
| __construct | |
100.00% |
4 / 4 |
|
100.00% |
1 / 1 |
3 | |||
| generateJobId | |
100.00% |
2 / 2 |
|
100.00% |
1 / 1 |
1 | |||
| setJobId | |
100.00% |
2 / 2 |
|
100.00% |
1 / 1 |
1 | |||
| getJobId | |
100.00% |
3 / 3 |
|
100.00% |
1 / 1 |
2 | |||
| hasJobId | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| setJobDescription | |
100.00% |
2 / 2 |
|
100.00% |
1 / 1 |
1 | |||
| getJobDescription | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| hasJobDescription | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| getResults | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| hasResults | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| setCallable | |
100.00% |
8 / 8 |
|
100.00% |
1 / 1 |
4 | |||
| setCommand | |
100.00% |
2 / 2 |
|
100.00% |
1 / 1 |
1 | |||
| setExec | |
100.00% |
2 / 2 |
|
100.00% |
1 / 1 |
1 | |||
| getCallable | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| getCommand | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| getExec | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| hasCallable | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| hasCommand | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| hasExec | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| setMaxAttempts | |
100.00% |
2 / 2 |
|
100.00% |
1 / 1 |
1 | |||
| getMaxAttempts | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| hasMaxAttempts | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| isAttemptOnce | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| getAttempts | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| hasAttempts | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| runUntil | |
100.00% |
2 / 2 |
|
100.00% |
1 / 1 |
1 | |||
| hasRunUntil | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| getRunUntil | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| delay | |
100.00% |
7 / 7 |
|
100.00% |
1 / 1 |
4 | |||
| getAvailableAt | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| isAvailable | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
2 | |||
| setTimeout | |
100.00% |
2 / 2 |
|
100.00% |
1 / 1 |
1 | |||
| getTimeout | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| hasTimeout | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| setBackoff | |
100.00% |
2 / 2 |
|
100.00% |
1 / 1 |
1 | |||
| getBackoff | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| hasBackoff | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| getBackoffDelay | |
100.00% |
6 / 6 |
|
100.00% |
1 / 1 |
3 | |||
| isExpired | |
100.00% |
9 / 9 |
|
100.00% |
1 / 1 |
7 | |||
| hasExceededMaxAttempts | |
100.00% |
3 / 3 |
|
100.00% |
1 / 1 |
2 | |||
| isValid | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
2 | |||
| hasNotRun | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
2 | |||
| start | |
100.00% |
2 / 2 |
|
100.00% |
1 / 1 |
1 | |||
| getStarted | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| hasStarted | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| isRunning | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
3 | |||
| complete | |
100.00% |
3 / 3 |
|
100.00% |
1 / 1 |
1 | |||
| getCompleted | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| getDuration | |
100.00% |
2 / 2 |
|
100.00% |
1 / 1 |
3 | |||
| isComplete | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| failed | |
100.00% |
5 / 5 |
|
100.00% |
1 / 1 |
2 | |||
| hasFailed | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| getFailed | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| addFailedMessage | |
100.00% |
3 / 3 |
|
100.00% |
1 / 1 |
1 | |||
| hasFailedMessages | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| getFailedMessages | |
100.00% |
1 / 1 |
|
100.00% |
1 / 1 |
1 | |||
| run | |
100.00% |
8 / 8 |
|
100.00% |
1 / 1 |
5 | |||
| loadCallable | |
90.00% |
9 / 10 |
|
0.00% |
0 / 1 |
4.02 | |||
| runCommand | |
100.00% |
12 / 12 |
|
100.00% |
1 / 1 |
4 | |||
| describeFromCommand | |
100.00% |
5 / 5 |
|
100.00% |
1 / 1 |
5 | |||
| buildExecProcess | |
100.00% |
5 / 5 |
|
100.00% |
1 / 1 |
3 | |||
| runExec | |
100.00% |
4 / 4 |
|
100.00% |
1 / 1 |
1 | |||
| __sleep | |
85.71% |
6 / 7 |
|
0.00% |
0 / 1 |
4.05 | |||
| __wakeup | |
100.00% |
6 / 6 |
|
100.00% |
1 / 1 |
2 | |||
| 1 | <?php |
| 2 | declare(strict_types=1); |
| 3 | /** |
| 4 | * Pop PHP Framework (https://www.popphp.org/) |
| 5 | * |
| 6 | * @link https://github.com/popphp/popphp-framework |
| 7 | * @author Nick Sagona, III <nick@popphp.org> |
| 8 | * @copyright Copyright (c) 2009-2026 Nick Sagona, III |
| 9 | * @license https://www.popphp.org/license New BSD License |
| 10 | */ |
| 11 | |
| 12 | /** |
| 13 | * @namespace |
| 14 | */ |
| 15 | namespace Pop\Queue\Process; |
| 16 | |
| 17 | use Pop\Application; |
| 18 | use Pop\Console\Command\AbstractCommand; |
| 19 | use Pop\Utils\CallableObject; |
| 20 | use Laravel\SerializableClosure\SerializableClosure; |
| 21 | use Symfony\Component\Process\Exception\ProcessFailedException; |
| 22 | use Symfony\Component\Process\Exception\ProcessTimedOutException; |
| 23 | use Symfony\Component\Process\Process; |
| 24 | |
| 25 | /** |
| 26 | * Abstract job class |
| 27 | * |
| 28 | * @category Pop |
| 29 | * @package Pop\Queue |
| 30 | * @author Nick Sagona, III <nick@popphp.org> |
| 31 | * @copyright Copyright (c) 2009-2026 Nick Sagona, III |
| 32 | * @license https://www.popphp.org/license New BSD License |
| 33 | * @version 3.0.0 |
| 34 | */ |
| 35 | abstract class AbstractJob implements JobInterface |
| 36 | { |
| 37 | |
| 38 | /** |
| 39 | * Job ID |
| 40 | * @var ?string |
| 41 | */ |
| 42 | protected ?string $id = null; |
| 43 | |
| 44 | /** |
| 45 | * Job Description |
| 46 | * @var ?string |
| 47 | */ |
| 48 | protected ?string $description = null; |
| 49 | |
| 50 | /** |
| 51 | * Job callable |
| 52 | * @var ?CallableObject |
| 53 | */ |
| 54 | protected ?CallableObject $callable = null; |
| 55 | |
| 56 | /** |
| 57 | * Job application command - an invocation string routed through the |
| 58 | * application (e.g. 'greet Nick'), or an argv-style array of already |
| 59 | * split segments (e.g. ['notify', 'Hello there, world']). The string |
| 60 | * form is split on whitespace by the router, so the array form is the |
| 61 | * only way to pass a value that itself contains spaces. |
| 62 | * @var string|array|null |
| 63 | */ |
| 64 | protected string|array|null $command = null; |
| 65 | |
| 66 | /** |
| 67 | * Job CLI executable command - a shell command string (runs via the |
| 68 | * shell, e.g. 'ls -la | wc -l'), or an argv-style array to run with no |
| 69 | * shell involved at all (e.g. ['ls', '-la'] - the safer form when any |
| 70 | * part of the command isn't a fully-trusted literal, since shell |
| 71 | * metacharacters in an argv element are inert) |
| 72 | * @var string|array|null |
| 73 | */ |
| 74 | protected string|array|null $exec = null; |
| 75 | |
| 76 | /** |
| 77 | * Job started timestamp |
| 78 | * @var ?int |
| 79 | */ |
| 80 | protected ?int $started = null; |
| 81 | |
| 82 | /** |
| 83 | * Job completed timestamp |
| 84 | * @var ?int |
| 85 | */ |
| 86 | protected ?int $completed = null; |
| 87 | |
| 88 | /** |
| 89 | * Job failed timestamp |
| 90 | * @var ?int |
| 91 | */ |
| 92 | protected ?int $failed = null; |
| 93 | |
| 94 | /** |
| 95 | * Job failed messages |
| 96 | * @var array |
| 97 | */ |
| 98 | protected array $failedMessages = []; |
| 99 | |
| 100 | /** |
| 101 | * Max attempts |
| 102 | * @var int |
| 103 | */ |
| 104 | protected int $maxAttempts = 0; |
| 105 | |
| 106 | /** |
| 107 | * Attempts |
| 108 | * @var int |
| 109 | */ |
| 110 | protected int $attempts = 0; |
| 111 | |
| 112 | /** |
| 113 | * Run until property |
| 114 | * @var int|string|null |
| 115 | */ |
| 116 | protected int|string|null $runUntil = null; |
| 117 | |
| 118 | /** |
| 119 | * Serialize closure |
| 120 | * @var ?string |
| 121 | */ |
| 122 | protected ?string $serializedClosure = null; |
| 123 | |
| 124 | /** |
| 125 | * Serialize parameters |
| 126 | * @var ?array |
| 127 | */ |
| 128 | protected ?array $serializedParameters = null; |
| 129 | |
| 130 | /** |
| 131 | * Job results |
| 132 | * @var mixed |
| 133 | */ |
| 134 | protected mixed $results = null; |
| 135 | |
| 136 | /** |
| 137 | * Timestamp before which this job is not eligible for reservation |
| 138 | * @var ?int |
| 139 | */ |
| 140 | protected ?int $availableAt = null; |
| 141 | |
| 142 | /** |
| 143 | * Soft execution timeout, in seconds (only enforced when ext-pcntl is loaded) |
| 144 | * @var ?int |
| 145 | */ |
| 146 | protected ?int $timeout = null; |
| 147 | |
| 148 | /** |
| 149 | * Retry backoff: a fixed delay in seconds, or a per-attempt schedule that |
| 150 | * holds at its last value for further attempts. Null = immediate retry. |
| 151 | * @var int|array|null |
| 152 | */ |
| 153 | protected int|array|null $backoff = null; |
| 154 | |
| 155 | /** |
| 156 | * Constructor |
| 157 | * |
| 158 | * Instantiate the job object |
| 159 | * |
| 160 | * @param mixed $callable |
| 161 | * @param mixed $params |
| 162 | * @param ?string $id |
| 163 | */ |
| 164 | public function __construct(mixed $callable = null, mixed $params = null, ?string $id = null) |
| 165 | { |
| 166 | if ($callable !== null) { |
| 167 | $this->setCallable($callable, $params); |
| 168 | } |
| 169 | if ($id !== null) { |
| 170 | $this->setJobId($id); |
| 171 | } |
| 172 | } |
| 173 | |
| 174 | /** |
| 175 | * Generate job ID |
| 176 | * |
| 177 | * @return string |
| 178 | */ |
| 179 | public function generateJobId(): string |
| 180 | { |
| 181 | $this->id = sha1(uniqid((string)rand()) . time()); |
| 182 | return $this->id; |
| 183 | } |
| 184 | |
| 185 | /** |
| 186 | * Set job ID |
| 187 | * |
| 188 | * @param string $id |
| 189 | * @return AbstractJob |
| 190 | */ |
| 191 | public function setJobId(string $id): AbstractJob |
| 192 | { |
| 193 | $this->id = $id; |
| 194 | return $this; |
| 195 | } |
| 196 | |
| 197 | /** |
| 198 | * Get job ID |
| 199 | * |
| 200 | * @return ?string |
| 201 | */ |
| 202 | public function getJobId(): ?string |
| 203 | { |
| 204 | if (!$this->hasJobId()) { |
| 205 | $this->generateJobId(); |
| 206 | } |
| 207 | return $this->id; |
| 208 | } |
| 209 | |
| 210 | /** |
| 211 | * Has job ID |
| 212 | * |
| 213 | * @return bool |
| 214 | */ |
| 215 | public function hasJobId(): bool |
| 216 | { |
| 217 | return ($this->id !== null); |
| 218 | } |
| 219 | |
| 220 | /** |
| 221 | * Set job description |
| 222 | * |
| 223 | * @param string $description |
| 224 | * @return AbstractJob |
| 225 | */ |
| 226 | public function setJobDescription(string $description): AbstractJob |
| 227 | { |
| 228 | $this->description = $description; |
| 229 | return $this; |
| 230 | } |
| 231 | |
| 232 | /** |
| 233 | * Get job description |
| 234 | * |
| 235 | * @return ?string |
| 236 | */ |
| 237 | public function getJobDescription(): ?string |
| 238 | { |
| 239 | return $this->description; |
| 240 | } |
| 241 | |
| 242 | /** |
| 243 | * Has job description |
| 244 | * |
| 245 | * @return bool |
| 246 | */ |
| 247 | public function hasJobDescription(): bool |
| 248 | { |
| 249 | return ($this->description !== null); |
| 250 | } |
| 251 | |
| 252 | /** |
| 253 | * Get job results |
| 254 | * |
| 255 | * @return mixed |
| 256 | */ |
| 257 | public function getResults(): mixed |
| 258 | { |
| 259 | return $this->results; |
| 260 | } |
| 261 | |
| 262 | /** |
| 263 | * Has job results |
| 264 | * |
| 265 | * @return bool |
| 266 | */ |
| 267 | public function hasResults(): bool |
| 268 | { |
| 269 | return !empty($this->results); |
| 270 | } |
| 271 | |
| 272 | /** |
| 273 | * Set job callable |
| 274 | * |
| 275 | * @param mixed $callable |
| 276 | * @param mixed $params |
| 277 | * @return AbstractJob |
| 278 | */ |
| 279 | public function setCallable(mixed $callable, mixed $params = null): AbstractJob |
| 280 | { |
| 281 | |
| 282 | if (!($callable instanceof CallableObject)) { |
| 283 | $this->callable = new CallableObject($callable, $params); |
| 284 | } else { |
| 285 | $this->callable = $callable; |
| 286 | if ($params !== null) { |
| 287 | if (is_array($params)) { |
| 288 | $this->callable->addParameters($params); |
| 289 | } else { |
| 290 | $this->callable->addParameter($params); |
| 291 | } |
| 292 | } |
| 293 | } |
| 294 | |
| 295 | return $this; |
| 296 | } |
| 297 | |
| 298 | /** |
| 299 | * Set job application command |
| 300 | * |
| 301 | * @param string|array $command |
| 302 | * @return AbstractJob |
| 303 | */ |
| 304 | public function setCommand(string|array $command): AbstractJob |
| 305 | { |
| 306 | $this->command = $command; |
| 307 | return $this; |
| 308 | } |
| 309 | |
| 310 | /** |
| 311 | * Set job CLI executable command |
| 312 | * |
| 313 | * @param string|array $command |
| 314 | * @return AbstractJob |
| 315 | */ |
| 316 | public function setExec(string|array $command): AbstractJob |
| 317 | { |
| 318 | $this->exec = $command; |
| 319 | return $this; |
| 320 | } |
| 321 | |
| 322 | /** |
| 323 | * Get job callable |
| 324 | * |
| 325 | * @return ?CallableObject |
| 326 | */ |
| 327 | public function getCallable(): ?CallableObject |
| 328 | { |
| 329 | return $this->callable; |
| 330 | } |
| 331 | |
| 332 | /** |
| 333 | * Get job application command |
| 334 | * |
| 335 | * @return string|array|null |
| 336 | */ |
| 337 | public function getCommand(): string|array|null |
| 338 | { |
| 339 | return $this->command; |
| 340 | } |
| 341 | |
| 342 | /** |
| 343 | * Get job CLI executable command |
| 344 | * |
| 345 | * @return string|array|null |
| 346 | */ |
| 347 | public function getExec(): string|array|null |
| 348 | { |
| 349 | return $this->exec; |
| 350 | } |
| 351 | |
| 352 | /** |
| 353 | * Has job callable |
| 354 | * |
| 355 | * @return bool |
| 356 | */ |
| 357 | public function hasCallable(): bool |
| 358 | { |
| 359 | return ($this->callable !== null); |
| 360 | } |
| 361 | |
| 362 | /** |
| 363 | * Has job application command |
| 364 | * |
| 365 | * @return bool |
| 366 | */ |
| 367 | public function hasCommand(): bool |
| 368 | { |
| 369 | return ($this->command !== null); |
| 370 | } |
| 371 | |
| 372 | /** |
| 373 | * Has job CLI executable command |
| 374 | * |
| 375 | * @return bool |
| 376 | */ |
| 377 | public function hasExec(): bool |
| 378 | { |
| 379 | return ($this->exec !== null); |
| 380 | } |
| 381 | |
| 382 | /** |
| 383 | * Set max attempts |
| 384 | * |
| 385 | * @param int $maxAttempts |
| 386 | * @return AbstractJob |
| 387 | */ |
| 388 | public function setMaxAttempts(int $maxAttempts): AbstractJob |
| 389 | { |
| 390 | $this->maxAttempts = $maxAttempts; |
| 391 | return $this; |
| 392 | } |
| 393 | |
| 394 | /** |
| 395 | * Get max attempts |
| 396 | * |
| 397 | * @return int |
| 398 | */ |
| 399 | public function getMaxAttempts(): int |
| 400 | { |
| 401 | return $this->maxAttempts; |
| 402 | } |
| 403 | |
| 404 | /** |
| 405 | * Has max attempts |
| 406 | * |
| 407 | * @return bool |
| 408 | */ |
| 409 | public function hasMaxAttempts(): bool |
| 410 | { |
| 411 | return ($this->maxAttempts > 0); |
| 412 | } |
| 413 | |
| 414 | /** |
| 415 | * Is job set for only one max attempt |
| 416 | * |
| 417 | * @return bool |
| 418 | */ |
| 419 | public function isAttemptOnce(): bool |
| 420 | { |
| 421 | return ($this->maxAttempts == 1); |
| 422 | } |
| 423 | |
| 424 | /** |
| 425 | * Get actual attempts |
| 426 | * |
| 427 | * @return int |
| 428 | */ |
| 429 | public function getAttempts(): int |
| 430 | { |
| 431 | return $this->attempts; |
| 432 | } |
| 433 | |
| 434 | /** |
| 435 | * Has actual attempts |
| 436 | * |
| 437 | * @return bool |
| 438 | */ |
| 439 | public function hasAttempts(): bool |
| 440 | { |
| 441 | return ($this->attempts > 0); |
| 442 | } |
| 443 | |
| 444 | /** |
| 445 | * Set the run until property |
| 446 | * |
| 447 | * @param int|string $runUntil |
| 448 | * @return AbstractJob |
| 449 | */ |
| 450 | public function runUntil(int|string $runUntil): AbstractJob |
| 451 | { |
| 452 | $this->runUntil = $runUntil; |
| 453 | return $this; |
| 454 | } |
| 455 | |
| 456 | /** |
| 457 | * Has run until |
| 458 | * |
| 459 | * @return bool |
| 460 | */ |
| 461 | public function hasRunUntil(): bool |
| 462 | { |
| 463 | return ($this->runUntil !== null); |
| 464 | } |
| 465 | |
| 466 | /** |
| 467 | * Get run until value |
| 468 | * |
| 469 | * @return int|string|null |
| 470 | */ |
| 471 | public function getRunUntil(): int|string|null |
| 472 | { |
| 473 | return $this->runUntil; |
| 474 | } |
| 475 | |
| 476 | /** |
| 477 | * Delay job availability |
| 478 | * |
| 479 | * @param int|string $when Seconds from now (int below 1000000000), an absolute |
| 480 | * timestamp (int), or a strtotime()-parseable string |
| 481 | * @throws Exception |
| 482 | * @return AbstractJob |
| 483 | */ |
| 484 | public function delay(int|string $when): AbstractJob |
| 485 | { |
| 486 | if (is_int($when)) { |
| 487 | $this->availableAt = ($when < 1000000000) ? (time() + $when) : $when; |
| 488 | } else { |
| 489 | $timestamp = strtotime($when); |
| 490 | if ($timestamp === false) { |
| 491 | throw new Exception('Error: That delay value is not valid.'); |
| 492 | } |
| 493 | $this->availableAt = $timestamp; |
| 494 | } |
| 495 | |
| 496 | return $this; |
| 497 | } |
| 498 | |
| 499 | /** |
| 500 | * Get available-at timestamp |
| 501 | * |
| 502 | * @return ?int |
| 503 | */ |
| 504 | public function getAvailableAt(): ?int |
| 505 | { |
| 506 | return $this->availableAt; |
| 507 | } |
| 508 | |
| 509 | /** |
| 510 | * Determine if the job is currently available (no delay, or delay has elapsed) |
| 511 | * |
| 512 | * @return bool |
| 513 | */ |
| 514 | public function isAvailable(): bool |
| 515 | { |
| 516 | return ($this->availableAt === null) || (time() >= $this->availableAt); |
| 517 | } |
| 518 | |
| 519 | /** |
| 520 | * Set soft execution timeout |
| 521 | * |
| 522 | * @param int $seconds |
| 523 | * @return AbstractJob |
| 524 | */ |
| 525 | public function setTimeout(int $seconds): AbstractJob |
| 526 | { |
| 527 | $this->timeout = $seconds; |
| 528 | return $this; |
| 529 | } |
| 530 | |
| 531 | /** |
| 532 | * Get soft execution timeout |
| 533 | * |
| 534 | * @return ?int |
| 535 | */ |
| 536 | public function getTimeout(): ?int |
| 537 | { |
| 538 | return $this->timeout; |
| 539 | } |
| 540 | |
| 541 | /** |
| 542 | * Has soft execution timeout |
| 543 | * |
| 544 | * @return bool |
| 545 | */ |
| 546 | public function hasTimeout(): bool |
| 547 | { |
| 548 | return ($this->timeout !== null); |
| 549 | } |
| 550 | |
| 551 | /** |
| 552 | * Set retry backoff (fixed seconds, or a per-attempt schedule) |
| 553 | * |
| 554 | * @param int|array $backoff |
| 555 | * @return AbstractJob |
| 556 | */ |
| 557 | public function setBackoff(int|array $backoff): AbstractJob |
| 558 | { |
| 559 | $this->backoff = $backoff; |
| 560 | return $this; |
| 561 | } |
| 562 | |
| 563 | /** |
| 564 | * Get retry backoff |
| 565 | * |
| 566 | * @return int|array|null |
| 567 | */ |
| 568 | public function getBackoff(): int|array|null |
| 569 | { |
| 570 | return $this->backoff; |
| 571 | } |
| 572 | |
| 573 | /** |
| 574 | * Has retry backoff |
| 575 | * |
| 576 | * @return bool |
| 577 | */ |
| 578 | public function hasBackoff(): bool |
| 579 | { |
| 580 | return ($this->backoff !== null); |
| 581 | } |
| 582 | |
| 583 | /** |
| 584 | * Get the backoff delay, in seconds, for the current attempt count |
| 585 | * |
| 586 | * @return int |
| 587 | */ |
| 588 | public function getBackoffDelay(): int |
| 589 | { |
| 590 | if (empty($this->backoff)) { |
| 591 | return 0; |
| 592 | } |
| 593 | if (is_int($this->backoff)) { |
| 594 | return $this->backoff; |
| 595 | } |
| 596 | |
| 597 | $index = max(min($this->attempts, count($this->backoff)) - 1, 0); |
| 598 | return (int)$this->backoff[$index]; |
| 599 | } |
| 600 | |
| 601 | /** |
| 602 | * Determine if the job has expired |
| 603 | * |
| 604 | * @return bool |
| 605 | */ |
| 606 | public function isExpired(): bool |
| 607 | { |
| 608 | if (!empty($this->runUntil)) { |
| 609 | $runUntil = null; |
| 610 | if (is_string($this->runUntil) && (strtotime($this->runUntil) !== false)) { |
| 611 | $runUntil = strtotime($this->runUntil); |
| 612 | } else if (is_numeric($this->runUntil) && ((string)(int)$this->runUntil == $this->runUntil)) { |
| 613 | $runUntil = $this->runUntil; |
| 614 | } |
| 615 | |
| 616 | if ($runUntil !== null) { |
| 617 | return (time() > $runUntil); |
| 618 | } |
| 619 | } |
| 620 | |
| 621 | return false; |
| 622 | } |
| 623 | |
| 624 | /** |
| 625 | * Determine if the job has exceeded max attempts |
| 626 | * |
| 627 | * @return bool |
| 628 | */ |
| 629 | public function hasExceededMaxAttempts(): bool |
| 630 | { |
| 631 | if ($this->hasMaxAttempts()) { |
| 632 | return ($this->attempts >= $this->maxAttempts); |
| 633 | } |
| 634 | |
| 635 | return false; |
| 636 | } |
| 637 | |
| 638 | /** |
| 639 | * Determine if the job is still valid |
| 640 | * |
| 641 | * Impure on two counts: it reads $attempts, which failed() increments, and |
| 642 | * it compares $runUntil against the current time. The same job can answer |
| 643 | * true and then false without anything being reassigned, so callers must |
| 644 | * not have an earlier answer cached on their behalf. |
| 645 | * |
| 646 | * @phpstan-impure |
| 647 | * @return bool |
| 648 | */ |
| 649 | public function isValid(): bool |
| 650 | { |
| 651 | return ((!$this->isExpired()) && (!$this->hasExceededMaxAttempts())); |
| 652 | } |
| 653 | |
| 654 | /** |
| 655 | * Has job run yet |
| 656 | * |
| 657 | * @return bool |
| 658 | */ |
| 659 | public function hasNotRun(): bool |
| 660 | { |
| 661 | return (($this->started === null) && ($this->completed === null)); |
| 662 | } |
| 663 | |
| 664 | /** |
| 665 | * Start job |
| 666 | * |
| 667 | * @return AbstractJob |
| 668 | */ |
| 669 | public function start(): AbstractJob |
| 670 | { |
| 671 | $this->started = time(); |
| 672 | return $this; |
| 673 | } |
| 674 | |
| 675 | /** |
| 676 | * Get started timestamp |
| 677 | * |
| 678 | * @return ?int |
| 679 | */ |
| 680 | public function getStarted(): ?int |
| 681 | { |
| 682 | return $this->started; |
| 683 | } |
| 684 | |
| 685 | /** |
| 686 | * Has job started |
| 687 | * |
| 688 | * @return bool |
| 689 | */ |
| 690 | public function hasStarted(): bool |
| 691 | { |
| 692 | return ($this->started !== null); |
| 693 | } |
| 694 | |
| 695 | /** |
| 696 | * Is job running and has not completed or failed yet |
| 697 | * |
| 698 | * @return bool |
| 699 | */ |
| 700 | public function isRunning(): bool |
| 701 | { |
| 702 | return (($this->started !== null) && ($this->completed === null) && ($this->failed === null)); |
| 703 | } |
| 704 | |
| 705 | /** |
| 706 | * Complete job |
| 707 | * |
| 708 | * @return AbstractJob |
| 709 | */ |
| 710 | public function complete(): AbstractJob |
| 711 | { |
| 712 | $this->completed = time(); |
| 713 | $this->attempts++; |
| 714 | return $this; |
| 715 | } |
| 716 | |
| 717 | /** |
| 718 | * Get completed timestamp |
| 719 | * |
| 720 | * @return ?int |
| 721 | */ |
| 722 | public function getCompleted(): ?int |
| 723 | { |
| 724 | return $this->completed; |
| 725 | } |
| 726 | |
| 727 | /** |
| 728 | * Get how long the job took to run, in seconds, or null unless it both |
| 729 | * started and completed. Convenience for observability listeners, which |
| 730 | * would otherwise all repeat the same timestamp subtraction. |
| 731 | * |
| 732 | * @return ?int |
| 733 | */ |
| 734 | public function getDuration(): ?int |
| 735 | { |
| 736 | return (($this->started !== null) && ($this->completed !== null)) |
| 737 | ? ($this->completed - $this->started) : null; |
| 738 | } |
| 739 | |
| 740 | /** |
| 741 | * Is job complete |
| 742 | * |
| 743 | * @return bool |
| 744 | */ |
| 745 | public function isComplete(): bool |
| 746 | { |
| 747 | return ($this->completed !== null); |
| 748 | } |
| 749 | |
| 750 | /** |
| 751 | * Set job as failed |
| 752 | * |
| 753 | * @param ?string $message |
| 754 | * @return AbstractJob |
| 755 | */ |
| 756 | public function failed(?string $message = null): AbstractJob |
| 757 | { |
| 758 | $this->failed = time(); |
| 759 | $this->attempts++; |
| 760 | |
| 761 | if ($message !== null) { |
| 762 | $this->addFailedMessage($message); |
| 763 | } |
| 764 | |
| 765 | return $this; |
| 766 | } |
| 767 | |
| 768 | /** |
| 769 | * Has job failed |
| 770 | * |
| 771 | * @return bool |
| 772 | */ |
| 773 | public function hasFailed(): bool |
| 774 | { |
| 775 | return ($this->failed !== null); |
| 776 | } |
| 777 | |
| 778 | /** |
| 779 | * Get failed timestamp |
| 780 | * |
| 781 | * @return ?int |
| 782 | */ |
| 783 | public function getFailed(): ?int |
| 784 | { |
| 785 | return $this->failed; |
| 786 | } |
| 787 | |
| 788 | /** |
| 789 | * Add failed message |
| 790 | * |
| 791 | * @param string $message |
| 792 | * @return AbstractJob |
| 793 | */ |
| 794 | public function addFailedMessage(string $message): AbstractJob |
| 795 | { |
| 796 | $index = $this->failed ?? time(); |
| 797 | $this->failedMessages[$index] = $message; |
| 798 | return $this; |
| 799 | } |
| 800 | |
| 801 | /** |
| 802 | * Has failed messages |
| 803 | * |
| 804 | * @return bool |
| 805 | */ |
| 806 | public function hasFailedMessages(): bool |
| 807 | { |
| 808 | return !empty($this->failedMessages); |
| 809 | } |
| 810 | |
| 811 | /** |
| 812 | * Get failed messages |
| 813 | * |
| 814 | * @return array |
| 815 | */ |
| 816 | public function getFailedMessages(): array |
| 817 | { |
| 818 | return $this->failedMessages; |
| 819 | } |
| 820 | |
| 821 | /** |
| 822 | * Run job |
| 823 | * |
| 824 | * @param ?Application $application |
| 825 | * @return mixed |
| 826 | */ |
| 827 | public function run(?Application $application = null): mixed |
| 828 | { |
| 829 | $this->start(); |
| 830 | |
| 831 | if ($this->hasCallable()) { |
| 832 | return $this->loadCallable($application); |
| 833 | } |
| 834 | if (($this->hasCommand()) && ($application !== null)) { |
| 835 | return $this->runCommand($application); |
| 836 | } |
| 837 | if ($this->hasExec()) { |
| 838 | return $this->runExec(); |
| 839 | } |
| 840 | |
| 841 | return null; |
| 842 | } |
| 843 | |
| 844 | /** |
| 845 | * Load callable |
| 846 | * |
| 847 | * @param ?Application $application |
| 848 | * @throws Exception|\Pop\Utils\Exception|\ReflectionException |
| 849 | * @return mixed |
| 850 | */ |
| 851 | protected function loadCallable(?Application $application = null): mixed |
| 852 | { |
| 853 | if ($this->callable === null) { |
| 854 | throw new Exception('Error: The callable for this job was not set.'); |
| 855 | } |
| 856 | |
| 857 | if ($application !== null) { |
| 858 | if ($this->callable->hasParameters()) { |
| 859 | $parameters = $this->callable->getParameters(); |
| 860 | array_unshift($parameters, $application); |
| 861 | $this->callable->setParameters($parameters); |
| 862 | } else { |
| 863 | $this->callable->addNamedParameter('application', $application); |
| 864 | } |
| 865 | } |
| 866 | |
| 867 | $this->results = $this->callable->call(); |
| 868 | return $this->results; |
| 869 | } |
| 870 | |
| 871 | /** |
| 872 | * Run application command |
| 873 | * |
| 874 | * @param Application $application |
| 875 | * @return mixed |
| 876 | */ |
| 877 | protected function runCommand(Application $application): mixed |
| 878 | { |
| 879 | $router = $application->router(); |
| 880 | $level = ob_get_level(); |
| 881 | $output = ''; |
| 882 | |
| 883 | ob_start(); |
| 884 | |
| 885 | try { |
| 886 | // run(false, ...) - never let an unresolved route call exit() |
| 887 | // and take the whole worker process down with it. |
| 888 | $application->run(false, $this->command); |
| 889 | } finally { |
| 890 | // Unwind to the level we started at, rather than closing exactly |
| 891 | // one buffer: the command may have thrown, or opened a buffer of |
| 892 | // its own and not closed it. Either way a leaked buffer is never |
| 893 | // reclaimed in a long-running worker - it grows without bound and |
| 894 | // silently swallows everything printed afterwards. Inner buffers |
| 895 | // hold the later output, so each unwind prepends to what we have. |
| 896 | while (ob_get_level() > $level) { |
| 897 | $output = ob_get_clean() . $output; |
| 898 | } |
| 899 | } |
| 900 | |
| 901 | // Whether the command resolved is the router's call, not a lookup |
| 902 | // against route definition keys - matching by definition string |
| 903 | // meant a real invocation ('greet Nick') could never be run, only |
| 904 | // the literal route definition ('greet <name>') could. |
| 905 | if (($router === null) || !$router->hasDispatchable()) { |
| 906 | return false; |
| 907 | } |
| 908 | |
| 909 | $this->describeFromCommand($router->getDispatchable()); |
| 910 | |
| 911 | $this->results = array_values(array_filter(explode(PHP_EOL, $output), fn($line) => $line !== '')); |
| 912 | return $this->results; |
| 913 | } |
| 914 | |
| 915 | /** |
| 916 | * Borrow a job description from the dispatched command, if it has one |
| 917 | * and the job was not given one explicitly. A command already documents |
| 918 | * itself, so a queued command job need not be anonymous in the registry. |
| 919 | * |
| 920 | * Guarded by instanceof rather than a hard dependency: pop-console |
| 921 | * arrives transitively through popphp, and instanceof against a missing |
| 922 | * class is simply false rather than an error. |
| 923 | * |
| 924 | * @param mixed $dispatchable |
| 925 | * @return void |
| 926 | */ |
| 927 | protected function describeFromCommand(mixed $dispatchable): void |
| 928 | { |
| 929 | if ($this->hasJobDescription() || !($dispatchable instanceof AbstractCommand)) { |
| 930 | return; |
| 931 | } |
| 932 | |
| 933 | $description = $dispatchable->getHelp() ?: $dispatchable->getName(); |
| 934 | |
| 935 | if (!empty($description)) { |
| 936 | $this->setJobDescription($description); |
| 937 | } |
| 938 | } |
| 939 | |
| 940 | /** |
| 941 | * Build the Process instance for this job's exec command, without |
| 942 | * running it. Split out from runExec() purely so a test can construct |
| 943 | * one and inspect its configured timeout directly, without actually |
| 944 | * executing a command. |
| 945 | * |
| 946 | * @return Process |
| 947 | */ |
| 948 | protected function buildExecProcess(): Process |
| 949 | { |
| 950 | $process = is_array($this->exec) |
| 951 | ? new Process($this->exec) |
| 952 | : Process::fromShellCommandline($this->exec); |
| 953 | |
| 954 | // Process defaults to a 60-second timeout on every instance unless |
| 955 | // explicitly told otherwise - disable it entirely when this job has |
| 956 | // no configured timeout, matching exec()'s old no-timeout-by-default |
| 957 | // behavior. Getting this wrong would silently cap every exec job at |
| 958 | // 60 seconds regardless of what the job actually needed. |
| 959 | $process->setTimeout($this->hasTimeout() ? $this->getTimeout() : null); |
| 960 | |
| 961 | return $process; |
| 962 | } |
| 963 | |
| 964 | /** |
| 965 | * Run CLI executable command |
| 966 | * |
| 967 | * @throws ProcessFailedException|ProcessTimedOutException |
| 968 | * @return mixed |
| 969 | */ |
| 970 | protected function runExec(): mixed |
| 971 | { |
| 972 | $process = $this->buildExecProcess(); |
| 973 | $process->mustRun(); |
| 974 | |
| 975 | $this->results = array_values(array_filter(explode(PHP_EOL, $process->getOutput()), fn($line) => $line !== '')); |
| 976 | return $this->results; |
| 977 | } |
| 978 | |
| 979 | /** |
| 980 | * Sleep magic method |
| 981 | * |
| 982 | * @return array |
| 983 | */ |
| 984 | public function __sleep(): array |
| 985 | { |
| 986 | if (!empty($this->callable) && ($this->callable->getCallable() instanceof \Closure)) { |
| 987 | $serializedClosure = new SerializableClosure($this->callable->getCallable()); |
| 988 | $this->serializedClosure = serialize($serializedClosure); |
| 989 | if ($this->callable->hasParameters()) { |
| 990 | $this->serializedParameters = $this->callable->getParameters(); |
| 991 | } |
| 992 | $this->callable = null; |
| 993 | } |
| 994 | |
| 995 | return array_keys(get_object_vars($this)); |
| 996 | } |
| 997 | |
| 998 | /** |
| 999 | * Wakeup magic method |
| 1000 | * |
| 1001 | * @return void |
| 1002 | */ |
| 1003 | public function __wakeup(): void |
| 1004 | { |
| 1005 | if (!empty($this->serializedClosure)) { |
| 1006 | $serializedClosure = unserialize($this->serializedClosure); |
| 1007 | $callable = $serializedClosure->getClosure(); |
| 1008 | $this->callable = new CallableObject($callable, $this->serializedParameters); |
| 1009 | $this->serializedClosure = null; |
| 1010 | $this->serializedParameters = null; |
| 1011 | } |
| 1012 | } |
| 1013 | |
| 1014 | } |