*/ use EventDispatcherTrait; /** * @param \Psr\Log\LoggerInterface $logger Logger instance. * @param \Cake\Core\ContainerInterface|null $container DI container instance */ public function __construct( protected readonly LoggerInterface $logger = new NullLogger(), protected readonly ?ContainerInterface $container = null, ) { } /** * The method processes messages * * @param \Interop\Queue\Message $message Message. * @param \Interop\Queue\Context $context Context. * @return object|string with __toString method implemented */ public function process(QueueMessage $message, Context $context): string|object { $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; } $startTime = microtime(true) * 1000; $this->dispatchEvent('Processor.message.start', ['message' => $jobMessage]); try { $response = $this->executeJob($jobMessage, $message); } catch (Throwable $throwable) { $message->setProperty('jobException', $throwable); $this->logger->debug(sprintf('Message encountered exception: %s', $throwable->getMessage())); $this->dispatchEvent('Processor.message.exception', [ 'message' => $jobMessage, 'exception' => $throwable, 'duration' => (int)((microtime(true) * 1000) - $startTime), ]); return Result::requeue('Exception occurred while processing message'); } $duration = (int)((microtime(true) * 1000) - $startTime); if ($response === InteropProcessor::ACK) { $this->logger->debug('Message processed successfully'); $this->dispatchEvent('Processor.message.success', [ 'message' => $jobMessage, 'duration' => $duration, ]); return InteropProcessor::ACK; } if ($response === InteropProcessor::REJECT) { $this->logger->debug('Message processed with rejection'); $this->dispatchEvent('Processor.message.reject', [ 'message' => $jobMessage, 'duration' => $duration, ]); return InteropProcessor::REJECT; } $this->logger->debug('Message processed with failure, requeuing'); $this->dispatchEvent('Processor.message.failure', [ 'message' => $jobMessage, 'duration' => $duration, ]); return InteropProcessor::REQUEUE; } /** * Execute the job and return the response. * * @param \Cake\Queue\Job\Message $jobMessage Job message wrapper * @param \Interop\Queue\Message $queueMessage Original queue message * @return object|string with __toString method implemented */ protected function executeJob(Message $jobMessage, QueueMessage $queueMessage): string|object { return $this->processMessage($jobMessage); } /** * @param \Cake\Queue\Job\Message $message Message. * @return object|string with __toString method implemented */ public function processMessage(Message $message): string|object { $callable = $message->getCallable(); $response = $callable($message); if ($response === null) { $response = InteropProcessor::ACK; } return $response; } }