forked from cakephp/queue
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathProcessor.php
More file actions
129 lines (108 loc) · 4.12 KB
/
Processor.php
File metadata and controls
129 lines (108 loc) · 4.12 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
<?php
declare(strict_types=1);
/**
* CakePHP(tm) : Rapid Development Framework (https://cakephp.org)
* Copyright (c) Cake Software Foundation, Inc. (https://cakefoundation.org/)
*
* Licensed under The MIT License
* For full copyright and license information, please see the LICENSE.txt
* Redistributions of files must retain the above copyright notice.
*
* @copyright Copyright (c) Cake Software Foundation, Inc. (https://cakefoundation.org/)
* @link https://cakephp.org CakePHP(tm) Project
* @since 0.1.0
* @license https://opensource.org/licenses/MIT MIT License
*/
namespace Cake\Queue\Queue;
use Cake\Core\ContainerInterface;
use Cake\Event\EventDispatcherTrait;
use Cake\Queue\Job\Message;
use Enqueue\Consumption\Result;
use Error;
use Interop\Queue\Context;
use Interop\Queue\Message as QueueMessage;
use Interop\Queue\Processor as InteropProcessor;
use Psr\Log\LoggerInterface;
use Psr\Log\NullLogger;
use RuntimeException;
use Throwable;
class Processor implements InteropProcessor
{
use EventDispatcherTrait;
/**
* @var \Psr\Log\LoggerInterface
*/
protected $logger;
/**
* @var \Cake\Core\ContainerInterface|null
*/
protected $container;
/**
* Processor constructor
*
* @param \Psr\Log\LoggerInterface|null $logger Logger instance.
* @param \Cake\Core\ContainerInterface|null $container DI container instance
*/
public function __construct(?LoggerInterface $logger = null, ?ContainerInterface $container = null)
{
$this->logger = $logger ?: new NullLogger();
$this->container = $container;
}
/**
* The method processes messages
*
* @param \Interop\Queue\Message $message Message.
* @param \Interop\Queue\Context $context Context.
* @return string|object with __toString method implemented
*/
public function process(QueueMessage $message, Context $context)
{
$this->dispatchEvent('Processor.message.seen', ['queueMessage' => $message]);
$jobMessage = new Message($message, $context, $this->container);
try {
$jobMessage->getCallable();
} catch (RuntimeException | Error $e) {
$this->logger->debug('Invalid callable for message. Rejecting message from queue.');
$this->dispatchEvent('Processor.message.invalid', ['message' => $jobMessage]);
return InteropProcessor::REJECT;
}
$this->dispatchEvent('Processor.message.start', ['message' => $jobMessage]);
try {
$response = $this->processMessage($jobMessage);
} catch (Throwable $e) {
$message->setProperty('jobException', $e);
$this->logger->debug(sprintf('Message encountered exception: %s', $e->getMessage()));
$this->dispatchEvent('Processor.message.exception', [
'message' => $jobMessage,
'exception' => $e,
]);
return Result::requeue('Exception occurred while processing message');
}
if ($response === InteropProcessor::ACK) {
$this->logger->debug('Message processed sucessfully');
$this->dispatchEvent('Processor.message.success', ['message' => $jobMessage]);
return InteropProcessor::ACK;
}
if ($response === InteropProcessor::REJECT) {
$this->logger->debug('Message processed with rejection');
$this->dispatchEvent('Processor.message.reject', ['message' => $jobMessage]);
return InteropProcessor::REJECT;
}
$this->logger->debug('Message processed with failure, requeuing');
$this->dispatchEvent('Processor.message.failure', ['message' => $jobMessage]);
return InteropProcessor::REQUEUE;
}
/**
* @param \Cake\Queue\Job\Message $message Message.
* @return string|object with __toString method implemented
*/
public function processMessage(Message $message)
{
$callable = $message->getCallable();
$response = $callable($message);
if ($response === null) {
$response = InteropProcessor::ACK;
}
return $response;
}
}