|
18 | 18 |
|
19 | 19 | use Cake\Core\ContainerInterface; |
20 | 20 | use Cake\Queue\Job\Message; |
21 | | -use Enqueue\Consumption\Result; |
22 | | -use Error; |
23 | | -use Interop\Queue\Context; |
24 | 21 | use Interop\Queue\Message as QueueMessage; |
25 | | -use Interop\Queue\Processor as InteropProcessor; |
26 | 22 | use Psr\Log\LoggerInterface; |
27 | 23 | use RuntimeException; |
28 | | -use Throwable; |
29 | 24 |
|
30 | 25 | /** |
31 | 26 | * Subprocess processor that executes jobs in isolated PHP processes. |
@@ -70,76 +65,18 @@ public function __construct( |
70 | 65 | } |
71 | 66 |
|
72 | 67 | /** |
73 | | - * Process a message in a subprocess. |
74 | | - * Overrides parent to execute in subprocess, but reuses parent's event dispatching. |
| 68 | + * Execute the job in a subprocess. |
75 | 69 | * |
76 | | - * @param \Interop\Queue\Message $message Message. |
77 | | - * @param \Interop\Queue\Context $context Context. |
| 70 | + * @param \Cake\Queue\Job\Message $jobMessage Job message wrapper |
| 71 | + * @param \Interop\Queue\Message $queueMessage Original queue message |
78 | 72 | * @return object|string with __toString method implemented |
79 | 73 | */ |
80 | | - public function process(QueueMessage $message, Context $context): string|object |
| 74 | + protected function executeJob(Message $jobMessage, QueueMessage $queueMessage): string|object |
81 | 75 | { |
82 | | - $this->dispatchEvent('Processor.message.seen', ['queueMessage' => $message]); |
| 76 | + $jobData = $this->prepareJobData($queueMessage); |
| 77 | + $subprocessResult = $this->executeInSubprocess($jobData); |
83 | 78 |
|
84 | | - $jobMessage = new Message($message, $context, $this->container); |
85 | | - try { |
86 | | - $jobMessage->getCallable(); |
87 | | - } catch (RuntimeException | Error $e) { |
88 | | - $this->logger->debug('Invalid callable for message. Rejecting message from queue.'); |
89 | | - $this->dispatchEvent('Processor.message.invalid', ['message' => $jobMessage]); |
90 | | - |
91 | | - return InteropProcessor::REJECT; |
92 | | - } |
93 | | - |
94 | | - $startTime = microtime(true) * 1000; |
95 | | - $this->dispatchEvent('Processor.message.start', ['message' => $jobMessage]); |
96 | | - |
97 | | - try { |
98 | | - $jobData = $this->prepareJobData($message); |
99 | | - $subprocessResult = $this->executeInSubprocess($jobData); |
100 | | - $response = $this->handleSubprocessResult($subprocessResult, $message); |
101 | | - } catch (Throwable $throwable) { |
102 | | - $message->setProperty('jobException', $throwable); |
103 | | - |
104 | | - $this->logger->debug(sprintf('Message encountered exception: %s', $throwable->getMessage())); |
105 | | - $this->dispatchEvent('Processor.message.exception', [ |
106 | | - 'message' => $jobMessage, |
107 | | - 'exception' => $throwable, |
108 | | - 'duration' => (int)((microtime(true) * 1000) - $startTime), |
109 | | - ]); |
110 | | - |
111 | | - return Result::requeue('Exception occurred while processing message'); |
112 | | - } |
113 | | - |
114 | | - $duration = (int)((microtime(true) * 1000) - $startTime); |
115 | | - |
116 | | - if ($response === InteropProcessor::ACK) { |
117 | | - $this->logger->debug('Message processed successfully'); |
118 | | - $this->dispatchEvent('Processor.message.success', [ |
119 | | - 'message' => $jobMessage, |
120 | | - 'duration' => $duration, |
121 | | - ]); |
122 | | - |
123 | | - return InteropProcessor::ACK; |
124 | | - } |
125 | | - |
126 | | - if ($response === InteropProcessor::REJECT) { |
127 | | - $this->logger->debug('Message processed with rejection'); |
128 | | - $this->dispatchEvent('Processor.message.reject', [ |
129 | | - 'message' => $jobMessage, |
130 | | - 'duration' => $duration, |
131 | | - ]); |
132 | | - |
133 | | - return InteropProcessor::REJECT; |
134 | | - } |
135 | | - |
136 | | - $this->logger->debug('Message processed with failure, requeuing'); |
137 | | - $this->dispatchEvent('Processor.message.failure', [ |
138 | | - 'message' => $jobMessage, |
139 | | - 'duration' => $duration, |
140 | | - ]); |
141 | | - |
142 | | - return InteropProcessor::REQUEUE; |
| 79 | + return $this->handleSubprocessResult($subprocessResult, $queueMessage); |
143 | 80 | } |
144 | 81 |
|
145 | 82 | /** |
|
0 commit comments