Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 13 additions & 0 deletions build/rector-strict.php
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,9 @@
*/

use Rector\DeadCode\Rector\ClassMethod\RemoveDuplicatedReturnSelfDocblockRector;
use Rector\DeadCode\Rector\ClassMethod\RemoveEmptyClassMethodRector;
use Rector\DeadCode\Rector\ClassMethod\RemoveReturnTagIncompatibleWithNativeTypeRector;
use Rector\DeadCode\Rector\ClassMethod\RemoveUnusedPublicMethodParameterRector;
use Rector\Php81\Rector\Property\ReadOnlyPropertyRector;
use Rector\Php82\Rector\Class_\ReadOnlyClassRector;
use Rector\PHPUnit\CodeQuality\Rector\Class_\AddSeeTestAnnotationRector;
Expand Down Expand Up @@ -63,6 +65,10 @@
$nextcloudDir . '/lib/public/SystemReport',
$nextcloudDir . '/lib/private/SystemReport',
$nextcloudDir . '/tests/lib/SystemReport',
$nextcloudDir . '/lib/public/MessageQueue',
$nextcloudDir . '/lib/private/MessageQueue',
$nextcloudDir . '/core/Command/MessageQueue',
$nextcloudDir . '/tests/lib/MessageQueue',
])
->withAutoloadPaths([
// ensure rector properly autoload the public interfaces
Expand Down Expand Up @@ -102,4 +108,11 @@
// non-final type; removing it breaks psalm's
// LessSpecificImplementedReturnType check (psalm-strict).
RemoveDuplicatedReturnSelfDocblockRector::class,
// Message handler fixtures declare the handled message through their parameter type
RemoveUnusedPublicMethodParameterRector::class => [
$nextcloudDir . '/tests/lib/MessageQueue/Fixtures',
],
RemoveEmptyClassMethodRector::class => [
$nextcloudDir . '/tests/lib/MessageQueue/Fixtures',
],
]);
26 changes: 26 additions & 0 deletions config/config.sample.php
Original file line number Diff line number Diff line change
Expand Up @@ -1767,13 +1767,39 @@
* the background jobs which advertise themselves as not time sensitive will be
* delayed during the "working" hours and only run in the 4 hours after the given time.
* This is, e.g., used for activity expiration, suspicious login training, and update checks.
* Messages in the low priority queue of the message queue are delayed the same way.
*
* A value of 1, e.g., will only run these background jobs between 01:00am UTC and 05:00am UTC.
*
* Defaults to ``100`` which disables the feature
*/
'maintenance_window_start' => 1,

/**
* Consume the message queue during cron runs
*
* Messages dispatched by apps are handled after the background jobs of each
* cron run. Large instances can disable this and instead run one or more
* dedicated workers with ``occ message-queue:consume``.
*
* Defaults to ``true``
*/
'message_queue.consume_in_cron' => true,

/**
* Handle messages right after the request
*
* When running under PHP-FPM, messages dispatched during a web request are
* handled by the same PHP process after the response was sent to the
* client, so they don't have to wait for the next cron run. Messages that
* are not handled within a few seconds stay in the queue. Large instances
* running dedicated workers can disable this to keep PHP-FPM processes
* available for requests.
*
* Defaults to ``true``
*/
'message_queue.handle_after_request' => true,

/**
* Log all LDAP requests into a file
*
Expand Down
84 changes: 84 additions & 0 deletions core/Command/MessageQueue/Consume.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,84 @@
<?php

declare(strict_types=1);

// SPDX-FileCopyrightText: 2026 Nextcloud GmbH and Nextcloud contributors
// SPDX-License-Identifier: AGPL-3.0-or-later

namespace OC\Core\Command\MessageQueue;

use OC\Core\Command\Base;
use OC\MessageQueue\Consumer;
use OCP\MessageQueue\Queue;
use OCP\Util;
use Symfony\Component\Console\Input\InputArgument;
use Symfony\Component\Console\Input\InputInterface;
use Symfony\Component\Console\Input\InputOption;
use Symfony\Component\Console\Output\OutputInterface;

final class Consume extends Base {
public function __construct(
private readonly Consumer $consumer,
) {
parent::__construct();
}

#[\Override]
protected function configure(): void {
$this
->setName('message-queue:consume')
->setDescription('Run a worker consuming messages from the message queue')
->addArgument(
'queues',
InputArgument::OPTIONAL | InputArgument::IS_ARRAY,
'Queues to consume, in order of priority (' . implode(', ', array_map(static fn (Queue $queue) => $queue->value, Queue::cases())) . ')',
)
->addOption('time-limit', 't', InputOption::VALUE_REQUIRED, 'Stop after this many seconds')
->addOption('limit', 'l', InputOption::VALUE_REQUIRED, 'Stop after handling this many messages')
->addOption('memory-limit', 'm', InputOption::VALUE_REQUIRED, 'Stop when the memory usage exceeds this limit, e.g. 512M')
->addOption('sleep', null, InputOption::VALUE_REQUIRED, 'Seconds to wait before polling again when all queues are empty', '1')
->addOption('stop-when-empty', null, InputOption::VALUE_NONE, 'Stop as soon as all queues are empty');
}

#[\Override]
protected function execute(InputInterface $input, OutputInterface $output): int {
/** @var list<string> $names */
$names = $input->getArgument('queues');
$queues = [];
foreach ($names as $name) {
$queue = Queue::tryFrom($name);
if ($queue === null) {
$output->writeln('<error>Unknown queue ' . $name . '</error>');
return self::FAILURE;
}

$queues[] = $queue;
}

/** @var mixed $memoryLimit */
$memoryLimit = $input->getOption('memory-limit');
$this->consumer->consume(
queues: $queues ?: Queue::cases(),
timeLimit: $this->getIntOption($input, 'time-limit'),
messageLimit: $this->getIntOption($input, 'limit'),
memoryLimit: is_string($memoryLimit) ? (int)Util::computerFileSize($memoryLimit) : null,
stopWhenEmpty: (bool)$input->getOption('stop-when-empty'),
sleep: (int)$input->getOption('sleep'),
output: $output->isVerbose() ? static fn (string $message) => $output->writeln($message) : null,
);

return self::SUCCESS;
}

#[\Override]
public function cancelOperation(): void {
parent::cancelOperation();
$this->consumer->stop();
}

private function getIntOption(InputInterface $input, string $name): ?int {
/** @var mixed $value */
$value = $input->getOption($name);
return is_string($value) ? (int)$value : null;
}
}
62 changes: 62 additions & 0 deletions core/Migrations/Version36000Date20261004120000.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,62 @@
<?php

declare(strict_types=1);

/**
* SPDX-FileCopyrightText: 2026 Nextcloud GmbH and Nextcloud contributors
* SPDX-License-Identifier: AGPL-3.0-or-later
*/

namespace OC\Core\Migrations;

use Closure;
use OCP\DB\ISchemaWrapper;
use OCP\DB\Types;
use OCP\Migration\Attributes\AddIndex;
use OCP\Migration\Attributes\CreateTable;
use OCP\Migration\Attributes\IndexType;
use OCP\Migration\IOutput;
use OCP\Migration\SimpleMigrationStep;
use Override;

#[CreateTable(
table: 'message_queue',
columns: ['id', 'queue_name', 'message_class', 'body', 'retry_count', 'last_error', 'created_at', 'available_at', 'delivered_at', 'deduplication_hash'],
description: 'New table to store the messages of the message queue',
)]
#[AddIndex(table: 'message_queue', type: IndexType::PRIMARY)]
#[AddIndex(table: 'message_queue', type: IndexType::INDEX, description: 'Allows to fetch the next available message of a queue')]
#[AddIndex(table: 'message_queue', type: IndexType::INDEX, description: 'Allows to find pending duplicates of a message')]
class Version36000Date20261004120000 extends SimpleMigrationStep {
#[Override]
public function changeSchema(IOutput $output, Closure $schemaClosure, array $options): ?ISchemaWrapper {
/** @var ISchemaWrapper $schema */
$schema = $schemaClosure();

if (!$schema->hasTable('message_queue')) {
$table = $schema->createTable('message_queue');
$table->addColumn('id', Types::BIGINT, [
'notnull' => true,
'unsigned' => true,
]);
$table->addColumn('queue_name', Types::STRING, ['notnull' => true, 'length' => 32]);
$table->addColumn('message_class', Types::STRING, ['notnull' => true, 'length' => 255]);
$table->addColumn('body', Types::TEXT, ['notnull' => true]);
$table->addColumn('retry_count', Types::INTEGER, ['notnull' => true, 'default' => 0, 'unsigned' => true]);
$table->addColumn('last_error', Types::TEXT, ['notnull' => false]);
$table->addColumn('created_at', Types::BIGINT, ['notnull' => true, 'unsigned' => true]);
$table->addColumn('available_at', Types::BIGINT, ['notnull' => true, 'unsigned' => true]);
$table->addColumn('delivered_at', Types::BIGINT, ['notnull' => false, 'unsigned' => true]);
$table->addColumn('deduplication_hash', Types::STRING, ['notnull' => false, 'length' => 64]);
$table->setPrimaryKey(['id']);
$table->addIndex(['queue_name', 'available_at'], 'mq_queue_available');
$table->addIndex(['deduplication_hash'], 'mq_dedup_hash');
// Makes sure there is no auto-increment in Oracle
$schema->dropAutoincrementColumn('message_queue', 'id');

return $schema;
}

return null;
}
}
29 changes: 29 additions & 0 deletions core/Service/CronService.php
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
use OC\BackgroundJob\JobClassesRegistry;
use OC\BackgroundJob\JobRuns;
use OC\DB\Connection;
use OC\MessageQueue\Consumer;
use OC\Security\CSRF\TokenStorage\SessionStorage;
use OC\Session\CryptoWrapper;
use OC\Session\Memory;
Expand All @@ -29,6 +30,7 @@
use OCP\ILogger;
use OCP\ISession;
use OCP\ITempManager;
use OCP\MessageQueue\Queue;
use OCP\Util;
use Psr\Log\LoggerInterface;

Expand All @@ -54,6 +56,7 @@ public function __construct(
private readonly JobRuns $jobRuns,
private readonly JobClassesRegistry $jobClassesRegistry,
private readonly ISetupManager $setupManager,
private readonly Consumer $messageQueueConsumer,
private readonly bool $isCLI,
) {
}
Expand Down Expand Up @@ -271,6 +274,11 @@ private function runCli(string $appMode, ?array $jobClasses): void {
}
}

if ($jobClasses === null) {
$queues = $onlyTimeSensitive ? [Queue::High, Queue::Default] : Queue::cases();
$this->consumeMessages($queues, max(60, $endTime - time()));
}

// Makes sure last error isn't caught by shutdown function
error_clear_last();
}
Expand All @@ -287,6 +295,27 @@ private function runWeb(string $appMode): void {
$job->start($this->jobList);
$this->jobList->setLastJob($job);
}
$this->consumeMessages(Queue::cases(), 10);
}
}

/**
* @param list<Queue> $queues
*/
private function consumeMessages(array $queues, int $timeLimit): void {
if (!$this->config->getSystemValueBool('message_queue.consume_in_cron', true)) {
return;
}

try {
$this->messageQueueConsumer->consume(
queues: $queues,
timeLimit: $timeLimit,
stopWhenEmpty: true,
output: $this->verboseCallback,
);
} catch (\Throwable $e) {
$this->logger->error('Error while consuming the message queue: ' . $e->getMessage(), ['app' => 'cron', 'exception' => $e]);
}
}

Expand Down
1 change: 1 addition & 0 deletions core/register_command.php
Original file line number Diff line number Diff line change
Expand Up @@ -156,6 +156,7 @@
$application->addCommand(Server::get(Delete::class));
$application->addCommand(Server::get(JobWorker::class));
$application->addCommand(Server::get(RunningJobs::class));
$application->addCommand(Server::get(Command\MessageQueue\Consume::class));
$application->addCommand(Server::get(JobsHistory::class));

$application->addCommand(Server::get(Test::class));
Expand Down
17 changes: 17 additions & 0 deletions lib/composer/composer/autoload_classmap.php
Original file line number Diff line number Diff line change
Expand Up @@ -784,6 +784,13 @@
'OCP\\Mail\\Provider\\IProvider' => $baseDir . '/lib/public/Mail/Provider/IProvider.php',
'OCP\\Mail\\Provider\\IService' => $baseDir . '/lib/public/Mail/Provider/IService.php',
'OCP\\Mail\\Provider\\Message' => $baseDir . '/lib/public/Mail/Provider/Message.php',
'OCP\\MessageQueue\\Attribute\\AsMessage' => $baseDir . '/lib/public/MessageQueue/Attribute/AsMessage.php',
'OCP\\MessageQueue\\Attribute\\AsMessageHandler' => $baseDir . '/lib/public/MessageQueue/Attribute/AsMessageHandler.php',
'OCP\\MessageQueue\\Exception\\InvalidMessageException' => $baseDir . '/lib/public/MessageQueue/Exception/InvalidMessageException.php',
'OCP\\MessageQueue\\Exception\\RecoverableMessageException' => $baseDir . '/lib/public/MessageQueue/Exception/RecoverableMessageException.php',
'OCP\\MessageQueue\\Exception\\UnrecoverableMessageException' => $baseDir . '/lib/public/MessageQueue/Exception/UnrecoverableMessageException.php',
'OCP\\MessageQueue\\IMessageBus' => $baseDir . '/lib/public/MessageQueue/IMessageBus.php',
'OCP\\MessageQueue\\Queue' => $baseDir . '/lib/public/MessageQueue/Queue.php',
'OCP\\Migration\\Attributes\\AddColumn' => $baseDir . '/lib/public/Migration/Attributes/AddColumn.php',
'OCP\\Migration\\Attributes\\AddIndex' => $baseDir . '/lib/public/Migration/Attributes/AddIndex.php',
'OCP\\Migration\\Attributes\\ColumnMigrationAttribute' => $baseDir . '/lib/public/Migration/Attributes/ColumnMigrationAttribute.php',
Expand Down Expand Up @@ -1534,6 +1541,7 @@
'OC\\Core\\Command\\Memcache\\DistributedGet' => $baseDir . '/core/Command/Memcache/DistributedGet.php',
'OC\\Core\\Command\\Memcache\\DistributedSet' => $baseDir . '/core/Command/Memcache/DistributedSet.php',
'OC\\Core\\Command\\Memcache\\RedisCommand' => $baseDir . '/core/Command/Memcache/RedisCommand.php',
'OC\\Core\\Command\\MessageQueue\\Consume' => $baseDir . '/core/Command/MessageQueue/Consume.php',
'OC\\Core\\Command\\OCM\\ActivateKey' => $baseDir . '/core/Command/OCM/ActivateKey.php',
'OC\\Core\\Command\\OCM\\ListKeys' => $baseDir . '/core/Command/OCM/ListKeys.php',
'OC\\Core\\Command\\OCM\\RetireKey' => $baseDir . '/core/Command/OCM/RetireKey.php',
Expand Down Expand Up @@ -1743,6 +1751,7 @@
'OC\\Core\\Migrations\\Version34000Date20260518163022' => $baseDir . '/core/Migrations/Version34000Date20260518163022.php',
'OC\\Core\\Migrations\\Version34000Date20260521110333' => $baseDir . '/core/Migrations/Version34000Date20260521110333.php',
'OC\\Core\\Migrations\\Version35000Date20260527162338' => $baseDir . '/core/Migrations/Version35000Date20260527162338.php',
'OC\\Core\\Migrations\\Version36000Date20261004120000' => $baseDir . '/core/Migrations/Version36000Date20261004120000.php',
'OC\\Core\\Notification\\CoreNotifier' => $baseDir . '/core/Notification/CoreNotifier.php',
'OC\\Core\\ResponseDefinitions' => $baseDir . '/core/ResponseDefinitions.php',
'OC\\Core\\Service\\CronService' => $baseDir . '/core/Service/CronService.php',
Expand Down Expand Up @@ -2098,6 +2107,14 @@
'OC\\Memcache\\Redis' => $baseDir . '/lib/private/Memcache/Redis.php',
'OC\\Memcache\\WithLocalCache' => $baseDir . '/lib/private/Memcache/WithLocalCache.php',
'OC\\MemoryInfo' => $baseDir . '/lib/private/MemoryInfo.php',
'OC\\MessageQueue\\AfterRequestHandler' => $baseDir . '/lib/private/MessageQueue/AfterRequestHandler.php',
'OC\\MessageQueue\\Consumer' => $baseDir . '/lib/private/MessageQueue/Consumer.php',
'OC\\MessageQueue\\HandlerRegistry' => $baseDir . '/lib/private/MessageQueue/HandlerRegistry.php',
'OC\\MessageQueue\\MessageBus' => $baseDir . '/lib/private/MessageQueue/MessageBus.php',
'OC\\MessageQueue\\MessageMetadataReader' => $baseDir . '/lib/private/MessageQueue/MessageMetadataReader.php',
'OC\\MessageQueue\\MessageNormalizer' => $baseDir . '/lib/private/MessageQueue/MessageNormalizer.php',
'OC\\MessageQueue\\MessageStore' => $baseDir . '/lib/private/MessageQueue/MessageStore.php',
'OC\\MessageQueue\\QueuedMessage' => $baseDir . '/lib/private/MessageQueue/QueuedMessage.php',
'OC\\Migration\\BackgroundRepair' => $baseDir . '/lib/private/Migration/BackgroundRepair.php',
'OC\\Migration\\ConsoleOutput' => $baseDir . '/lib/private/Migration/ConsoleOutput.php',
'OC\\Migration\\Exceptions\\AttributeException' => $baseDir . '/lib/private/Migration/Exceptions/AttributeException.php',
Expand Down
Loading
Loading