-
Notifications
You must be signed in to change notification settings - Fork 5
Expand file tree
/
Copy pathAdapter.php
More file actions
529 lines (474 loc) · 19 KB
/
Copy pathAdapter.php
File metadata and controls
529 lines (474 loc) · 19 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
<?php
namespace Utopia\Queue;
use Utopia\DI\Container;
abstract class Adapter
{
protected const int RECEIVE_TIMEOUT = 2;
/**
* Pause before asking again after the broker failed to answer, so an
* unreachable broker is retried at a steady rate rather than in a tight loop.
*/
protected const int RECEIVE_BACKOFF = 1;
/**
* Active queue for the sequential / single-loop hot path. Concurrent
* multi-queue loops pass Queue explicitly via {@see nextMessageFrom()} /
* {@see processFrom()} so they do not race this property. Bound by
* consume() / run() before the first receive.
*/
public Queue $queue;
protected ?Container $context = null;
protected bool $stopped = false;
public Consumer $consumer;
/**
* @var callable(string): Consumer
*/
protected $consumerFactory;
protected bool $sharedConsumer = false;
/**
* How often a running worker sweeps its consumers for idle upkeep.
*
* Sized against the shortest server-side idle deadline a broker is likely
* to be holding: NATS pings every 120s and closes after two go unanswered,
* so a sweep every 30s leaves several chances to answer before that runs
* out. Cheap enough to be unconditional — a sweep with nothing to do is a
* method_exists() per consumer.
*/
protected const float MAINTENANCE_INTERVAL = 30.0;
/**
* Prefer a callable factory so each consume loop gets its own receive
* connection. A bare Consumer is OK for single-queue only.
*
* @param Consumer|callable $consumer Consumer instance, `(string $queue): Consumer`,
* or a zero-arg factory that returns a Consumer
* @param int $workerNum Process/worker count for pool adapters (Swoole/Workerman)
* @param string $namespace Broker key prefix shared by every job on this adapter
*/
public function __construct(
Consumer|callable $consumer,
public int $workerNum,
public string $namespace = 'utopia-queue',
protected Container $resources = new Container(),
) {
if ($consumer instanceof Consumer) {
$this->consumer = $consumer;
$this->consumerFactory = static fn (string $queue): Consumer => $consumer;
$this->sharedConsumer = true;
} else {
$this->consumerFactory = self::normalizeFactory($consumer);
$this->consumer = ($this->consumerFactory)('');
$this->sharedConsumer = false;
}
}
/**
* Invoke the adapter's consumer factory for a queue.
*/
public function createConsumer(string $queue = ''): Consumer
{
return ($this->consumerFactory)($queue);
}
/**
* True when the adapter was constructed with a bare shared Consumer.
*/
public function sharesConsumer(): bool
{
return $this->sharedConsumer;
}
/**
* @return callable(string): Consumer
*/
protected static function normalizeFactory(callable $factory): callable
{
$closure = $factory instanceof \Closure ? $factory : \Closure::fromCallable($factory);
$reflection = new \ReflectionFunction($closure);
if ($reflection->getNumberOfRequiredParameters() === 0) {
return static fn (string $queue): Consumer => $factory();
}
return $closure;
}
/**
* Starts the Server.
*/
abstract public function start(): self;
/**
* Stops the Server.
*/
abstract public function stop(): self;
/** @phpstan-impure stop() flips this from a signal handler mid-consume(). */
protected function isStopped(): bool
{
return $this->stopped;
}
/**
* @param callable(Message): void $messageCallback
* @param callable(Message): void $successCallback
* @param callable(?Message, \Throwable): void $errorCallback Receives null when
* the failure was in obtaining a message rather than handling one.
* @param array<int, array{queue: Queue, coroutines: int, prefetch?: int, consumer?: Consumer}> $queues
* Queue identity and concurrency come from Server::job(); sequential
* adapters run specs one after another, Swoole runs independent loops.
*/
public function consume(
callable $messageCallback,
callable $successCallback,
callable $errorCallback,
array $queues,
): void {
$this->stopped = false;
if ($queues === []) {
throw new \LogicException('At least one queue is required');
}
foreach ($queues as $spec) {
$this->run(
$spec['queue'],
$spec['coroutines'],
$messageCallback,
$successCallback,
$errorCallback,
$spec['consumer'] ?? $this->consumer,
$spec['prefetch'] ?? $spec['coroutines'],
);
}
}
/**
* One-queue loop. `$coroutines` and `$prefetch` are accepted for adapter
* parity; the sequential fallback processes one message at a time.
*
* Binds `$this->queue` / `$this->consumer` for the duration so the hot
* path matches pre-multi-queue (no per-message queue/consumer args).
*
* @param callable(Message): void $messageCallback
* @param callable(Message): void $successCallback
* @param callable(?Message, \Throwable): void $errorCallback
*/
protected function run(
Queue $queue,
int $coroutines,
callable $messageCallback,
callable $successCallback,
callable $errorCallback,
Consumer $consumer,
?int $prefetch = null,
): void {
// Sequential adapters receive and acknowledge one message at a time.
unset($coroutines, $prefetch);
$previousConsumer = $this->consumer;
$this->queue = $queue;
$this->consumer = $consumer;
try {
while (!$this->isStopped()) {
$message = $this->nextMessage($errorCallback);
if (!$message instanceof Message) {
continue;
}
$this->context = new Container($this->resources());
$this->process($message, $messageCallback, $successCallback, $errorCallback);
}
} finally {
$this->consumer = $previousConsumer;
}
}
/**
* Sweep the consumers this adapter owns for idle upkeep.
*
* A broker holding connections nobody is using still has servers on the
* other end counting silence. The consume loop's own traffic keeps its
* receive connection alive, but a pooled broker's spare slots get no
* traffic at all, and are reaped on a timer with nothing watching.
*
* Deliberately maintain() and not tick(): a tick reads the socket and needs
* the caller to hold the resource exclusively, which a running receive loop
* does not allow. maintain() is the pool-level sweep — it touches only what
* is idle, so it is safe to call while the loop is mid-receive. Consumers
* that expose neither are skipped, which is why this can run unconditionally.
*
* Never throws: upkeep failing must not take the worker down with it.
*
* @param callable(?Message, \Throwable): void|null $errorCallback
*/
public function maintain(?callable $errorCallback = null): void
{
foreach ($this->maintenanceTargets() as $target) {
$sweep = [$target, 'maintain'];
if (!\is_callable($sweep)) {
continue;
}
try {
$sweep();
} catch (\Throwable $error) {
// Reported rather than swallowed: a pool that cannot keep its
// idle connections alive will hand the next caller a dead one,
// and that is exactly the silent loss this is here to prevent.
if ($errorCallback === null) {
continue;
}
try {
$errorCallback(null, $error);
} catch (\Throwable) {
}
}
}
}
/**
* Consumers eligible for a maintenance sweep. Adapters that hold more than
* the one bound consumer override this to include them.
*
* @return list<Consumer>
*/
protected function maintenanceTargets(): array
{
return [$this->consumer];
}
/**
* Never throws: a broker that cannot be reached is reported to
* $errorCallback and retried after RECEIVE_BACKOFF. Losing the worker to a
* transient outage is worse than waiting for the broker to come back.
*
* $errorCallback takes a nullable message for exactly this case — the
* failure is in obtaining one, so there is none to report alongside it.
*
* @param callable(?Message, \Throwable): void $errorCallback
*/
protected function nextMessage(callable $errorCallback): ?Message
{
return $this->nextMessageFrom($errorCallback, $this->queue, $this->consumer);
}
/**
* Concurrent multi-queue variant: queue/consumer are explicit so loops do
* not race {@see $queue} / {@see $consumer}.
*
* @param callable(?Message, \Throwable): void $errorCallback
*/
protected function nextMessageFrom(callable $errorCallback, Queue $queue, Consumer $consumer): ?Message
{
return $this->nextBatchFrom($errorCallback, $queue, $consumer, 1)[0] ?? null;
}
/**
* Claim up to $max messages at once.
*
* @param callable(?Message, \Throwable): void $errorCallback
* @return list<Message>
*/
protected function nextBatchFrom(callable $errorCallback, Queue $queue, Consumer $consumer, int $max): array
{
try {
return $consumer->receive($queue, static::RECEIVE_TIMEOUT, $max);
} catch (\Throwable $error) {
try {
$errorCallback(null, $error);
} catch (\Throwable $reportFailure) {
$this->reportUnreported($error, $reportFailure);
}
sleep(static::RECEIVE_BACKOFF);
return [];
}
}
/**
* Never throws: a failed handler is rejected and reported to $errorCallback;
* a failing reject or callback is swallowed rather than left to escape (and
* be lost on a coroutine).
*/
protected function process(
Message $message,
callable $messageCallback,
callable $successCallback,
callable $errorCallback,
): void {
$this->processFrom($message, $messageCallback, $successCallback, $errorCallback, $this->queue, $this->consumer);
}
/**
* Concurrent multi-queue variant of {@see process()}.
*
* The three phases are separated rather than sharing one try, because only
* the first of them means the work failed. Committing and the success hook
* run after the handler has already succeeded, and routing their failures
* to reject() gives the message back to the broker after the job is done:
* the handler runs a second time, or — at the delivery ceiling — a job that
* worked is dead-lettered as though it never had.
*
* Contract change for $errorCallback. It used to fire only for work that
* had failed, so "reported" and "will be retried" were the same statement.
* It now also fires for a failed commit and a throwing success hook, and in
* both of those the handler has already run to completion:
*
* - handler threw — the work did not happen; the message is rejected
* and will be retried, unless the handler threw
* {@see PermanentFailure} and it is dead-lettered.
* - commit threw — the work happened; nothing is rejected, and the
* broker may still redeliver on its own deadline.
* - success hook threw — the work happened and is acked; nothing will
* redeliver it.
*
* A callback that treats every report as a failed job will over-count and,
* where it drives alerting or compensation, act on work that succeeded.
* Implementations that need to tell them apart should key off the phase
* rather than the presence of a report.
*/
protected function processFrom(
Message $message,
callable $messageCallback,
callable $successCallback,
callable $errorCallback,
Queue $queue,
Consumer $consumer,
): void {
try {
$this->runPhases($message, $messageCallback, $successCallback, $errorCallback, $queue, $consumer);
} finally {
// The phases return early on failure, so the container is dropped
// here rather than at the end of any one of them.
$this->releaseContext();
}
}
/**
* The three phases themselves. Separated from processFrom() only so the
* per-message container is released on every exit path.
*/
private function runPhases(
Message $message,
callable $messageCallback,
callable $successCallback,
callable $errorCallback,
Queue $queue,
Consumer $consumer,
): void {
try {
$this->withAckExtension($consumer, $queue, $message, static function () use ($messageCallback, $message): void {
$messageCallback($message);
});
} catch (\Throwable $error) {
// A handler that knows the work can never succeed says so by throwing
// PermanentFailure, and the verdict has to be on the message before it
// is rejected: reject() is where the broker decides between another
// attempt and the dead letter, and it runs here — ahead of the error
// report below, which is the only other place a host sees the failure.
//
// A type error is the same verdict without the handler saying it: the
// payload and the signature disagree, and they will disagree identically
// on every delivery. Staging spent maxDeliver=5 over ~22 minutes of
// backoff on 701 of them, each holding one of 60 ack slots, while the
// worker delivered nothing else for hours. Not \Error at large --
// OutOfMemoryError and the stack overflow say the host was short at that
// moment, which is what the redelivery budget is for.
if (
$error instanceof PermanentFailure
|| $error instanceof \TypeError // ArgumentCountError extends this
|| $error instanceof \ValueError
) {
$message->terminal();
}
// The work did not happen, so hand the message back to be retried
// (or, for a terminal verdict, to be dead-lettered now).
try {
$consumer->reject($queue, $message);
} catch (\Throwable) {
}
$this->report($errorCallback, $error, $message);
return;
}
try {
$consumer->commit($queue, $message);
} catch (\Throwable $error) {
// A transient ack failure over completed work. Not rejected: the
// job is done, and the broker will redeliver on its own deadline if
// the ack genuinely never landed — a duplicate the handler can
// guard against, where a NAK here is a duplicate guaranteed.
$this->report($errorCallback, $error, $message);
return;
}
try {
$successCallback($message);
} catch (\Throwable $error) {
// Bookkeeping after the message is acked and gone. There is nothing
// left to reject, and re-running the job would not fix a shutdown
// hook, so this is reported and no more.
$this->report($errorCallback, $error, $message);
}
}
/**
* Run the handler, keeping the broker's delivery deadline extended for as
* long as it takes.
*
* The default is to just run it: extending needs a scheduler to run
* alongside the handler, which only a coroutine adapter has. Adapters that
* have one override this, and where they do, the broker's ack deadline
* stops being a ceiling on how long a job may run.
*
* @param \Closure(): void $work
*/
protected function withAckExtension(Consumer $consumer, Queue $queue, Message $message, \Closure $work): void
{
$work();
}
/**
* Report a failure, with the last-resort trace when reporting fails too.
*
* @param callable(?Message, \Throwable): void $errorCallback
*/
private function report(callable $errorCallback, \Throwable $error, Message $message): void
{
try {
$errorCallback($message, $error);
} catch (\Throwable $reportFailure) {
$this->reportUnreported($error, $reportFailure, $message);
}
}
/**
* Drop the per-message container after commit/reject and outcome callbacks.
* The next message (or process shutdown) must not inherit the previous
* message's resolved graph.
*/
protected function releaseContext(): void
{
$this->context = null;
}
/**
* Last-resort trace for a failure whose reporting hook also failed.
*
* A hook typically needs resources of its own — a database handle to
* resolve the message's project, say — so the very outages that fail a
* message also fail the report of it, and the message is then rejected
* with nothing written anywhere. Production lost whole batches this way,
* visible only as messages appearing on the failed list. Stderr is the one
* sink that needs nothing to be working.
*/
protected function reportUnreported(\Throwable $error, \Throwable $reportFailure, ?Message $message = null): void
{
try {
fwrite($this->trace(), \sprintf(
"[queue] %s failed and its error report failed too: %s (%s:%d) | report: %s\n",
$message instanceof Message ? "message {$message->getPid()}" : 'receive',
$error->getMessage(),
$error->getFile(),
$error->getLine(),
$reportFailure->getMessage(),
));
} catch (\Throwable) {
}
}
/**
* Where {@see self::reportUnreported()} writes. Overridable so a caller can
* route the trace somewhere it will be retained, and so it can be asserted.
*
* @return resource
*/
protected function trace(): mixed
{
return \defined('STDERR') ? STDERR : fopen('php://stderr', 'w');
}
public function resources(): Container
{
return $this->resources;
}
public function context(): Container
{
return $this->context ??= new Container($this->resources());
}
/**
* Is called when a Worker starts.
*/
abstract public function workerStart(callable $callback): self;
/**
* Is called when a Worker stops.
*/
abstract public function workerStop(callable $callback): self;
}