TYPO3 v15 dev-main snapshot ()
This commit is contained in:
@@ -0,0 +1,358 @@
|
||||
<?php
|
||||
|
||||
declare(strict_types=1);
|
||||
|
||||
/*
|
||||
* This file is part of the TYPO3 CMS project.
|
||||
*
|
||||
* It is free software; you can redistribute it and/or modify it under
|
||||
* the terms of the GNU General Public License, either version 2
|
||||
* of the License, or any later version.
|
||||
*
|
||||
* For the full copyright and license information, please read the
|
||||
* LICENSE.txt file that was distributed with this source code.
|
||||
*
|
||||
* The TYPO3 project - inspiring people to share!
|
||||
*/
|
||||
|
||||
namespace TYPO3\CMS\Core\Command;
|
||||
|
||||
use Psr\Container\ContainerInterface;
|
||||
use Psr\EventDispatcher\EventDispatcherInterface;
|
||||
use Psr\Log\LoggerInterface;
|
||||
use Symfony\Component\Console\Attribute\AsCommand;
|
||||
use Symfony\Component\Console\Command\Command;
|
||||
use Symfony\Component\Console\Completion\CompletionInput;
|
||||
use Symfony\Component\Console\Completion\CompletionSuggestions;
|
||||
use Symfony\Component\Console\Exception\InvalidOptionException;
|
||||
use Symfony\Component\Console\Exception\RuntimeException;
|
||||
use Symfony\Component\Console\Input\InputArgument;
|
||||
use Symfony\Component\Console\Input\InputInterface;
|
||||
use Symfony\Component\Console\Input\InputOption;
|
||||
use Symfony\Component\Console\Logger\ConsoleLogger;
|
||||
use Symfony\Component\Console\Output\ConsoleOutputInterface;
|
||||
use Symfony\Component\Console\Output\OutputInterface;
|
||||
use Symfony\Component\Console\Question\ChoiceQuestion;
|
||||
use Symfony\Component\Console\Style\SymfonyStyle;
|
||||
use Symfony\Component\DependencyInjection\Attribute\Autowire;
|
||||
use Symfony\Component\DependencyInjection\Attribute\AutowireIterator;
|
||||
use Symfony\Component\DependencyInjection\Attribute\AutowireLocator;
|
||||
use Symfony\Component\EventDispatcher\EventSubscriberInterface;
|
||||
use Symfony\Component\Messenger\EventListener\StopWorkerOnFailureLimitListener;
|
||||
use Symfony\Component\Messenger\EventListener\StopWorkerOnMemoryLimitListener;
|
||||
use Symfony\Component\Messenger\EventListener\StopWorkerOnMessageLimitListener;
|
||||
use Symfony\Component\Messenger\EventListener\StopWorkerOnTimeLimitListener;
|
||||
use Symfony\Component\Messenger\MessageBusInterface;
|
||||
use Symfony\Component\Messenger\Transport\Sync\SyncTransport;
|
||||
use Symfony\Component\Messenger\Worker;
|
||||
use TYPO3\CMS\Core\EventDispatcher\ListenerProvider;
|
||||
|
||||
/**
|
||||
* Almost full version of the symfony command with the same name.
|
||||
*/
|
||||
#[AsCommand(name: 'messenger:consume', description: 'Consume messages')]
|
||||
class ConsumeMessagesCommand extends Command
|
||||
{
|
||||
private const int DEFAULT_KEEPALIVE_INTERVAL = 5;
|
||||
|
||||
private ?Worker $worker = null;
|
||||
private ?LoggerInterface $logger = null;
|
||||
private array $receiverNames;
|
||||
|
||||
public function __construct(
|
||||
#[Autowire(service: 'messenger.bus.default')]
|
||||
private readonly MessageBusInterface $messageBus,
|
||||
#[AutowireLocator('messenger.receiver', indexAttribute: 'identifier')]
|
||||
private readonly ContainerInterface $receiverLocator,
|
||||
private readonly EventDispatcherInterface $eventDispatcher,
|
||||
private readonly ListenerProvider $listenerProvider,
|
||||
private readonly ContainerInterface $container,
|
||||
#[AutowireIterator('messenger.receiver', indexAttribute: 'identifier')]
|
||||
iterable $receiverNamesIterator,
|
||||
private readonly array $busIds = [],
|
||||
#[AutowireLocator('messenger.rate_limiter', indexAttribute: 'identifier')]
|
||||
private readonly ?ContainerInterface $rateLimiterLocator = null,
|
||||
private readonly ?array $signals = null,
|
||||
) {
|
||||
$this->receiverNames = array_keys([...$receiverNamesIterator]);
|
||||
parent::__construct();
|
||||
}
|
||||
|
||||
protected function configure(): void
|
||||
{
|
||||
$defaultReceiverName = count($this->receiverNames) === 1 ? current($this->receiverNames) : null;
|
||||
|
||||
$this
|
||||
->setDefinition(
|
||||
[
|
||||
new InputArgument(
|
||||
'receivers',
|
||||
InputArgument::IS_ARRAY,
|
||||
'Names of the receivers/transports to consume in order of priority',
|
||||
$defaultReceiverName ? [$defaultReceiverName] : []
|
||||
),
|
||||
new InputOption('limit', 'l', InputOption::VALUE_REQUIRED, 'Limit the number of received messages'),
|
||||
new InputOption('failure-limit', 'f', InputOption::VALUE_REQUIRED, 'The number of failed messages the worker can consume'),
|
||||
new InputOption('memory-limit', 'm', InputOption::VALUE_REQUIRED, 'The memory limit the worker can consume'),
|
||||
new InputOption('time-limit', 't', InputOption::VALUE_REQUIRED, 'The time limit in seconds the worker can handle new messages'),
|
||||
new InputOption('sleep', null, InputOption::VALUE_REQUIRED, 'Seconds to sleep before asking for new messages after no messages were found', 1),
|
||||
new InputOption('bus', 'b', InputOption::VALUE_REQUIRED, 'Name of the bus to which received messages should be dispatched (if not passed, bus is determined automatically)'),
|
||||
new InputOption('queues', null, InputOption::VALUE_REQUIRED | InputOption::VALUE_IS_ARRAY, 'Limit receivers to only consume from the specified queues'),
|
||||
new InputOption('all', null, InputOption::VALUE_NONE, 'Consume messages from all receivers'),
|
||||
new InputOption('keepalive', null, InputOption::VALUE_OPTIONAL, 'Whether to use the transport\'s keepalive mechanism if implemented', self::DEFAULT_KEEPALIVE_INTERVAL),
|
||||
]
|
||||
)
|
||||
->setHelp(
|
||||
<<<'EOF'
|
||||
The <info>%command.name%</info> command consumes messages and dispatches them to the message bus.
|
||||
|
||||
<info>php %command.full_name% <receiver-name></info>
|
||||
|
||||
To receive from multiple transports, pass each name:
|
||||
|
||||
<info>php %command.full_name% receiver1 receiver2</info>
|
||||
|
||||
Use the --limit option to limit the number of messages received:
|
||||
|
||||
<info>php %command.full_name% <receiver-name> --limit=10</info>
|
||||
|
||||
Use the --failure-limit option to stop the worker when the given number of failed messages is reached:
|
||||
|
||||
<info>php %command.full_name% <receiver-name> --failure-limit=2</info>
|
||||
|
||||
Use the --memory-limit option to stop the worker if it exceeds a given memory usage limit. You can use shorthand byte values [K, M or G]:
|
||||
|
||||
<info>php %command.full_name% <receiver-name> --memory-limit=128M</info>
|
||||
|
||||
Use the --time-limit option to stop the worker when the given time limit (in seconds) is reached.
|
||||
If a message is being handled, the worker will stop after the processing is finished:
|
||||
|
||||
<info>php %command.full_name% <receiver-name> --time-limit=3600</info>
|
||||
|
||||
Use the --bus option to specify the message bus to dispatch received messages
|
||||
to instead of trying to determine it automatically. This is required if the
|
||||
messages didn't originate from Messenger:
|
||||
|
||||
<info>php %command.full_name% <receiver-name> --bus=event_bus</info>
|
||||
|
||||
Use the --queues option to limit a receiver to only certain queues (only supported by some receivers):
|
||||
|
||||
<info>php %command.full_name% <receiver-name> --queues=fasttrack</info>
|
||||
|
||||
Use the --all option to consume from all receivers:
|
||||
|
||||
<info>php %command.full_name% --all</info>
|
||||
EOF
|
||||
)
|
||||
;
|
||||
}
|
||||
|
||||
protected function initialize(InputInterface $input, OutputInterface $output): void
|
||||
{
|
||||
if ($input->hasParameterOption('--keepalive')) {
|
||||
$this->getApplication()->setAlarmInterval((int)($input->getOption('keepalive') ?? self::DEFAULT_KEEPALIVE_INTERVAL));
|
||||
}
|
||||
}
|
||||
|
||||
protected function interact(InputInterface $input, OutputInterface $output): void
|
||||
{
|
||||
$io = new SymfonyStyle($input, $output instanceof ConsoleOutputInterface ? $output->getErrorOutput() : $output);
|
||||
|
||||
if ($input->getOption('all')) {
|
||||
return;
|
||||
}
|
||||
|
||||
if ($this->receiverNames && !$input->getArgument('receivers')) {
|
||||
if (count($this->receiverNames) === 1) {
|
||||
$input->setArgument('receivers', $this->receiverNames);
|
||||
return;
|
||||
}
|
||||
|
||||
$io->block('Which transports/receivers do you want to consume?', null, 'fg=white;bg=blue', ' ', true);
|
||||
|
||||
$io->writeln('Choose which receivers you want to consume messages from in order of priority.');
|
||||
$io->writeln(sprintf('Hint: to consume from multiple, use a list of their names, e.g. <comment>%s</comment>', implode(', ', $this->receiverNames)));
|
||||
|
||||
$question = new ChoiceQuestion('Select receivers to consume:', $this->receiverNames, 0);
|
||||
$question->setMultiselect(true);
|
||||
|
||||
$input->setArgument('receivers', $io->askQuestion($question));
|
||||
}
|
||||
|
||||
if (!$input->getArgument('receivers')) {
|
||||
throw new RuntimeException('Please pass at least one receiver.', 1605305001);
|
||||
}
|
||||
}
|
||||
|
||||
protected function execute(InputInterface $input, OutputInterface $output): int
|
||||
{
|
||||
$this->logger = new ConsoleLogger($output);
|
||||
|
||||
$receivers = [];
|
||||
$rateLimiters = [];
|
||||
$receiverNames = $input->getOption('all') ? $this->receiverNames : $input->getArgument('receivers');
|
||||
foreach ($receiverNames as $receiverName) {
|
||||
if (!$this->receiverLocator->has($receiverName)) {
|
||||
$message = sprintf('The receiver "%s" does not exist.', $receiverName);
|
||||
if ($this->receiverNames) {
|
||||
$message .= sprintf(' Valid receivers are: %s.', implode(', ', $this->receiverNames));
|
||||
}
|
||||
throw new RuntimeException($message, 1605305002);
|
||||
}
|
||||
|
||||
$receiver = $this->receiverLocator->get($receiverName);
|
||||
if ($receiver instanceof SyncTransport) {
|
||||
$idx = array_search($receiverName, $receiverNames);
|
||||
unset($receiverNames[$idx]);
|
||||
|
||||
continue;
|
||||
}
|
||||
|
||||
$receivers[$receiverName] = $receiver;
|
||||
if ($this->rateLimiterLocator?->has($receiverName)) {
|
||||
$rateLimiters[$receiverName] = $this->rateLimiterLocator->get($receiverName);
|
||||
}
|
||||
}
|
||||
|
||||
$stopsWhen = [];
|
||||
if (null !== $limit = $input->getOption('limit')) {
|
||||
if (!is_numeric($limit) || $limit <= 0) {
|
||||
throw new InvalidOptionException(sprintf('Option "limit" must be a positive integer, "%s" passed.', $limit), 1605305003);
|
||||
}
|
||||
|
||||
$stopsWhen[] = "processed {$limit} messages";
|
||||
$this->addSubscriber(new StopWorkerOnMessageLimitListener((int)$limit, $this->logger));
|
||||
}
|
||||
|
||||
if ($failureLimit = $input->getOption('failure-limit')) {
|
||||
$stopsWhen[] = "reached {$failureLimit} failed messages";
|
||||
$this->addSubscriber(new StopWorkerOnFailureLimitListener((int)$failureLimit, $this->logger));
|
||||
}
|
||||
|
||||
if ($memoryLimit = $input->getOption('memory-limit')) {
|
||||
$stopsWhen[] = "exceeded {$memoryLimit} of memory";
|
||||
$this->addSubscriber(new StopWorkerOnMemoryLimitListener($this->convertToBytes($memoryLimit), $this->logger));
|
||||
}
|
||||
|
||||
if (null !== $timeLimit = $input->getOption('time-limit')) {
|
||||
if (!is_numeric($timeLimit) || $timeLimit <= 0) {
|
||||
throw new InvalidOptionException(sprintf('Option "time-limit" must be a positive integer, "%s" passed.', $timeLimit), 1605305004);
|
||||
}
|
||||
|
||||
$stopsWhen[] = "been running for {$timeLimit}s";
|
||||
$this->addSubscriber(new StopWorkerOnTimeLimitListener((int)$timeLimit, $this->logger));
|
||||
}
|
||||
|
||||
$stopsWhen[] = 'received a stop signal via the messenger:stop-workers command';
|
||||
|
||||
$io = new SymfonyStyle($input, $output instanceof ConsoleOutputInterface ? $output->getErrorOutput() : $output);
|
||||
$io->success(sprintf('Consuming messages from transport%s "%s".', count($receivers) > 1 ? 's' : '', implode(', ', $receiverNames)));
|
||||
|
||||
if ($stopsWhen) {
|
||||
$last = array_pop($stopsWhen);
|
||||
$stopsWhen = ($stopsWhen ? implode(', ', $stopsWhen) . ' or ' : '') . $last;
|
||||
$io->comment("The worker will automatically exit once it has {$stopsWhen}.");
|
||||
}
|
||||
|
||||
$io->comment('Quit the worker with CONTROL-C.');
|
||||
|
||||
if ($output->getVerbosity() < OutputInterface::VERBOSITY_VERBOSE) {
|
||||
$io->comment('Re-run the command with a -vv option to see logs about consumed messages.');
|
||||
}
|
||||
|
||||
$this->worker = new Worker($receivers, $this->messageBus, $this->eventDispatcher, $this->logger, $rateLimiters);
|
||||
$options = [
|
||||
'sleep' => $input->getOption('sleep') * 1000000,
|
||||
];
|
||||
if ($queues = $input->getOption('queues')) {
|
||||
$options['queues'] = $queues;
|
||||
}
|
||||
|
||||
try {
|
||||
$this->worker->run($options);
|
||||
} finally {
|
||||
$this->worker = null;
|
||||
}
|
||||
|
||||
return Command::SUCCESS;
|
||||
}
|
||||
|
||||
public function complete(CompletionInput $input, CompletionSuggestions $suggestions): void
|
||||
{
|
||||
if ($input->mustSuggestArgumentValuesFor('receivers')) {
|
||||
$suggestions->suggestValues(array_diff($this->receiverNames, array_diff($input->getArgument('receivers'), [$input->getCompletionValue()])));
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
if ($input->mustSuggestOptionValuesFor('bus')) {
|
||||
$suggestions->suggestValues($this->busIds);
|
||||
}
|
||||
}
|
||||
|
||||
public function getSubscribedSignals(): array
|
||||
{
|
||||
return $this->signals ?? (extension_loaded('pcntl') ? [SIGTERM, SIGINT, SIGQUIT, SIGALRM] : []);
|
||||
}
|
||||
|
||||
public function handleSignal(int $signal, int|false $previousExitCode = 0): int|false
|
||||
{
|
||||
if (!$this->worker) {
|
||||
return false;
|
||||
}
|
||||
|
||||
if (defined('SIGALRM') && $signal === SIGALRM) {
|
||||
$this->logger?->debug('Sending keepalive request.', ['transport_names' => $this->worker->getMetadata()->getTransportNames()]);
|
||||
|
||||
$this->worker->keepalive($this->getApplication()->getAlarmInterval());
|
||||
|
||||
return false;
|
||||
}
|
||||
|
||||
$this->logger?->info('Received signal {signal}.', ['signal' => $signal, 'transport_names' => $this->worker->getMetadata()->getTransportNames()]);
|
||||
|
||||
$this->worker->stop();
|
||||
|
||||
return false;
|
||||
}
|
||||
|
||||
private function convertToBytes(string $memoryLimit): int
|
||||
{
|
||||
$memoryLimit = strtolower($memoryLimit);
|
||||
$max = ltrim($memoryLimit, '+');
|
||||
if (str_starts_with($max, '0x')) {
|
||||
$max = intval($max, 16);
|
||||
} elseif (str_starts_with($max, '0')) {
|
||||
$max = intval($max, 8);
|
||||
} else {
|
||||
$max = (float)$max;
|
||||
}
|
||||
|
||||
switch (substr(rtrim($memoryLimit, 'b'), -1)) {
|
||||
case 't': $max *= 1024;
|
||||
// no break
|
||||
case 'g': $max *= 1024;
|
||||
// no break
|
||||
case 'm': $max *= 1024;
|
||||
// no break
|
||||
case 'k': $max *= 1024;
|
||||
}
|
||||
|
||||
return (int)$max;
|
||||
}
|
||||
|
||||
/**
|
||||
* @todo: This method show be removed when we can add event subscribers dynamically.
|
||||
*/
|
||||
private function addSubscriber(EventSubscriberInterface $subscriber): void
|
||||
{
|
||||
$this->container->set($subscriber::class, $subscriber);
|
||||
foreach ($subscriber->getSubscribedEvents() as $eventName => $params) {
|
||||
$this->listenerProvider->addListener(
|
||||
$eventName,
|
||||
$subscriber::class,
|
||||
is_string($params) ? $params : $params[0],
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user