-
Notifications
You must be signed in to change notification settings - Fork 21
Expand file tree
/
Copy pathProcessor.php
More file actions
127 lines (106 loc) · 4.04 KB
/
Processor.php
File metadata and controls
127 lines (106 loc) · 4.04 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
<?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 Queue\Queue;
use Cake\Event\EventDispatcherTrait;
use Cake\Log\LogTrait;
use Exception;
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 Queue\Job\Message;
class Processor implements InteropProcessor
{
use EventDispatcherTrait;
use LogTrait;
/**
* @var \Psr\Log\LoggerInterface
*/
protected $logger;
/**
* Processor constructor
*
* @param \Psr\Log\LoggerInterface $logger Logger instance.
*/
public function __construct(?LoggerInterface $logger = null)
{
$this->logger = $logger ?: new NullLogger();
}
/**
* The method processes messages
*
* @param \Interop\Queue\Message $queueMessage Message.
* @param \Interop\Queue\Context $context Context.
* @return string|object with __toString method implemented
*/
public function process(QueueMessage $queueMessage, Context $context)
{
$this->dispatchEvent('Processor.message.seen', ['queueMessage' => $queueMessage]);
$message = new Message($queueMessage, $context);
if (!is_callable($message->getCallable())) {
$this->logger->debug('Invalid callable for message. Rejecting message from queue.');
$this->dispatchEvent('Processor.message.invalid', ['message' => $message]);
return InteropProcessor::REJECT;
}
$this->dispatchEvent('Processor.message.start', ['message' => $message]);
try {
$response = $this->processMessage($message);
} catch (Exception $e) {
$this->logger->debug(sprintf('Message encountered exception: %s', $e->getMessage()));
$this->dispatchEvent('Processor.message.exception', [
'message' => $message,
'exception' => $e,
]);
return InteropProcessor::REQUEUE;
}
if ($response === InteropProcessor::ACK) {
$this->logger->debug('Message processed sucessfully');
$this->dispatchEvent('Processor.message.success', ['message' => $message]);
return InteropProcessor::ACK;
}
if ($response === InteropProcessor::REJECT) {
$this->logger->debug('Message processed with rejection');
$this->dispatchEvent('Processor.message.reject', ['message' => $message]);
return InteropProcessor::REJECT;
}
$this->logger->debug('Message processed with failure, requeuing');
$this->dispatchEvent('Processor.message.failure', ['message' => $message]);
return InteropProcessor::REQUEUE;
}
/**
* @param \Queue\Job\Message $message Message.
* @return string|object with __toString method implemented
*/
public function processMessage(Message $message)
{
$callable = $message->getCallable();
$response = InteropProcessor::REQUEUE;
if (is_array($callable) && count($callable) == 2) {
$className = $callable[0];
$methodName = $callable[1];
$instance = new $className();
$response = $instance->$methodName($message);
} elseif (is_string($callable)) {
/** @psalm-suppress InvalidArgument */
$response = call_user_func($callable, $message);
}
if ($response === null) {
$response = InteropProcessor::ACK;
}
return $response;
}
}