From fe26dcb50aa057340fe5eb087f9c3181a8e8b6ef Mon Sep 17 00:00:00 2001 From: Jake Barnby Date: Thu, 17 Apr 2025 22:27:20 +1200 Subject: [PATCH] Add publisher consumer pool adapter --- app/cli.php | 6 +++-- app/http.php | 3 ++- app/init/resources.php | 5 ++-- app/worker.php | 32 +++++++++++------------- composer.json | 6 ++--- composer.lock | 56 ++++++++++++++++++++++++------------------ 6 files changed, 58 insertions(+), 50 deletions(-) diff --git a/app/cli.php b/app/cli.php index f963d32d94..45361c2971 100644 --- a/app/cli.php +++ b/app/cli.php @@ -24,6 +24,7 @@ use Utopia\DSN\DSN; use Utopia\Logger\Log; use Utopia\Platform\Service; use Utopia\Pools\Group; +use Utopia\Queue\Broker\Pool as BrokerPool; use Utopia\Queue\Publisher; use Utopia\Registry\Registry; use Utopia\System\System; @@ -158,6 +159,7 @@ CLI::setResource('getProjectDB', function (Group $pools, Database $dbForPlatform CLI::setResource('getLogsDB', function (Group $pools, Cache $cache) { $database = null; + return function (?Document $project = null) use ($pools, $cache, $database) { if ($database !== null && $project !== null && !$project->isEmpty() && $project->getId() !== 'console') { $database->setTenant($project->getInternalId()); @@ -182,14 +184,14 @@ CLI::setResource('getLogsDB', function (Group $pools, Cache $cache) { }; }, ['pools', 'cache']); -CLI::setResource('queueForStatsUsage', function (Connection $publisher) { +CLI::setResource('queueForStatsUsage', function (Publisher $publisher) { return new StatsUsage($publisher); }, ['publisher']); CLI::setResource('queueForStatsResources', function (Publisher $publisher) { return new StatsResources($publisher); }, ['publisher']); CLI::setResource('publisher', function (Group $pools) { - return $pools->get('publisher')->pop()->getResource(); + return new BrokerPool(publisher: $pools->get('publisher')); }, ['pools']); CLI::setResource('queueForFunctions', function (Publisher $publisher) { return new Func($publisher); diff --git a/app/http.php b/app/http.php index b2af26a631..79f3ebd464 100644 --- a/app/http.php +++ b/app/http.php @@ -463,6 +463,7 @@ $http->on(Constant::EVENT_REQUEST, function (SwooleRequest $swooleRequest, Swool Console::error('[Error] Message: ' . $th->getMessage()); Console::error('[Error] File: ' . $th->getFile()); Console::error('[Error] Line: ' . $th->getLine()); + Console::error('[Error] Trace: ' . $th->getTraceAsString()); $swooleResponse->setStatusCode(500); @@ -484,7 +485,7 @@ $http->on(Constant::EVENT_REQUEST, function (SwooleRequest $swooleRequest, Swool }); // Fetch domains every `DOMAIN_SYNC_TIMER` seconds and update in the memory -$http->on('Task', function () use ($register, $domains) { +$http->on(Constant::EVENT_TASK, function () use ($register, $domains) { $lastSyncUpdate = null; $pools = $register->get('pools'); App::setResource('pools', fn () => $pools); diff --git a/app/init/resources.php b/app/init/resources.php index 4e7d219666..7979e39e64 100644 --- a/app/init/resources.php +++ b/app/init/resources.php @@ -40,6 +40,7 @@ use Utopia\Locale\Locale; use Utopia\Logger\Log; use Utopia\Pools\Group; use Utopia\Queue\Publisher; +use Utopia\Queue\Broker\Pool as BrokerPool; use Utopia\Storage\Device; use Utopia\Storage\Device\AWS; use Utopia\Storage\Device\Backblaze; @@ -74,10 +75,10 @@ App::setResource('localeCodes', function () { // Queues App::setResource('publisher', function (Group $pools) { - return $pools->get('publisher')->pop()->getResource(); + return new BrokerPool(publisher: $pools->get('publisher')); }, ['pools']); App::setResource('consumer', function (Group $pools) { - return $pools->get('consumer')->pop()->getResource(); + return new BrokerPool(consumer: $pools->get('consumer')); }, ['pools']); App::setResource('queueForMessaging', function (Publisher $publisher) { return new Messaging($publisher); diff --git a/app/worker.php b/app/worker.php index 72ad5f3338..79f8fe0d85 100644 --- a/app/worker.php +++ b/app/worker.php @@ -35,6 +35,8 @@ use Utopia\Logger\Log; use Utopia\Logger\Logger; use Utopia\Platform\Service; use Utopia\Pools\Group; +use Utopia\Queue\Adapter\Pool as QueuePool; +use Utopia\Queue\Broker\Pool as BrokerPool; use Utopia\Queue\Message; use Utopia\Queue\Publisher; use Utopia\Queue\Server; @@ -42,7 +44,7 @@ use Utopia\Registry\Registry; use Utopia\System\System; Authorization::disable(); -Runtime::enableCoroutine(SWOOLE_HOOK_ALL); +Runtime::enableCoroutine(); Server::setResource('register', fn () => $register); @@ -239,11 +241,11 @@ Server::setResource('timelimit', function (\Redis $redis) { Server::setResource('log', fn () => new Log()); Server::setResource('publisher', function (Group $pools) { - return $pools->get('publisher')->pop()->getResource(); + return new BrokerPool(publisher: $pools->get('publisher')); }, ['pools']); Server::setResource('consumer', function (Group $pools) { - return $pools->get('consumer')->pop()->getResource(); + return new BrokerPool(consumer: $pools->get('consumer')); }, ['pools']); Server::setResource('queueForStatsUsage', function (Publisher $publisher) { @@ -408,25 +410,21 @@ try { * - _APP_WORKER_PER_CORE The number of worker processes per core (ignored if _APP_WORKERS_NUM is set) * - _APP_QUEUE_NAME The name of the queue to read for database events */ - $platform->init(Service::TYPE_WORKER, [ - 'workersNum' => System::getEnv('_APP_WORKERS_NUM', 1), - 'connection' => $pools->get('consumer')->pop()->getResource(), - 'workerName' => strtolower($workerName) ?? null, - 'queueName' => $queueName - ]); + $platform->init( + type: Service::TYPE_WORKER, + params: ['workerName' => strtolower($workerName) ?? null], + server: new Server(new QueuePool( + $pools->get('consumer'), + workerNum: System::getEnv('_APP_WORKERS_NUM', 1), + queue: $queueName + )) + ); } catch (\Throwable $e) { Console::error($e->getMessage() . ', File: ' . $e->getFile() . ', Line: ' . $e->getLine()); } $worker = $platform->getWorker(); -$worker - ->shutdown() - ->inject('pools') - ->action(function (Group $pools) { - $pools->get('consumer')->reclaim(); - }); - $worker ->error() ->inject('error') @@ -435,8 +433,6 @@ $worker ->inject('pools') ->inject('project') ->action(function (Throwable $error, ?Logger $logger, Log $log, Group $pools, Document $project) use ($worker, $queueName) { - $pools->get('consumer')->reclaim(); - $version = System::getEnv('_APP_VERSION', 'UNKNOWN'); if ($logger) { diff --git a/composer.json b/composer.json index 615d5abd62..5d91aa8116 100644 --- a/composer.json +++ b/composer.json @@ -51,7 +51,7 @@ "utopia-php/cache": "0.13.*", "utopia-php/cli": "0.15.*", "utopia-php/config": "0.2.*", - "utopia-php/database": "dev-feat-pool-init as 0.65.0", + "utopia-php/database": "0.66.*", "utopia-php/domains": "0.5.*", "utopia-php/dsn": "0.2.1", "utopia-php/framework": "0.33.*", @@ -62,10 +62,10 @@ "utopia-php/messaging": "0.16.*", "utopia-php/migration": "0.8.*", "utopia-php/orchestration": "0.9.*", - "utopia-php/platform": "0.7.*", + "utopia-php/platform": "dev-feat-custom-server as 0.7.4", "utopia-php/pools": "0.8.*", "utopia-php/preloader": "0.2.*", - "utopia-php/queue": "0.9.*", + "utopia-php/queue": "dev-feat-pool-adapter as 0.9.1", "utopia-php/registry": "0.5.*", "utopia-php/storage": "0.18.*", "utopia-php/swoole": "0.8.*", diff --git a/composer.lock b/composer.lock index 1d39c1b753..9ad06f5317 100644 --- a/composer.lock +++ b/composer.lock @@ -4,7 +4,7 @@ "Read more about it at https://getcomposer.org/doc/01-basic-usage.md#installing-dependencies", "This file is @generated automatically" ], - "content-hash": "d56d7960edb985f497bcc816ac294652", + "content-hash": "127ac05400882555232b9d9a9d9e5c0e", "packages": [ { "name": "adhocore/jwt", @@ -3498,16 +3498,16 @@ }, { "name": "utopia-php/database", - "version": "dev-feat-pool-init", + "version": "0.66.1", "source": { "type": "git", "url": "https://github.com/utopia-php/database.git", - "reference": "3f03b9cb56f7d2f09bb9c17e5898d3222a285473" + "reference": "155bc1c0ec68e0224661b1e587e01c6451026031" }, "dist": { "type": "zip", - "url": "https://api.github.com/repos/utopia-php/database/zipball/3f03b9cb56f7d2f09bb9c17e5898d3222a285473", - "reference": "3f03b9cb56f7d2f09bb9c17e5898d3222a285473", + "url": "https://api.github.com/repos/utopia-php/database/zipball/155bc1c0ec68e0224661b1e587e01c6451026031", + "reference": "155bc1c0ec68e0224661b1e587e01c6451026031", "shasum": "" }, "require": { @@ -3548,9 +3548,9 @@ ], "support": { "issues": "https://github.com/utopia-php/database/issues", - "source": "https://github.com/utopia-php/database/tree/feat-pool-init" + "source": "https://github.com/utopia-php/database/tree/0.66.1" }, - "time": "2025-04-17T04:27:21+00:00" + "time": "2025-04-17T07:23:15+00:00" }, { "name": "utopia-php/domains", @@ -4058,16 +4058,16 @@ }, { "name": "utopia-php/platform", - "version": "0.7.4", + "version": "dev-feat-custom-server", "source": { "type": "git", "url": "https://github.com/utopia-php/platform.git", - "reference": "a5b93d8177702ec458c3af9137663133c012b71b" + "reference": "1afcdb43643531c82d5beaab30c2cfb5bf40ab52" }, "dist": { "type": "zip", - "url": "https://api.github.com/repos/utopia-php/platform/zipball/a5b93d8177702ec458c3af9137663133c012b71b", - "reference": "a5b93d8177702ec458c3af9137663133c012b71b", + "url": "https://api.github.com/repos/utopia-php/platform/zipball/1afcdb43643531c82d5beaab30c2cfb5bf40ab52", + "reference": "1afcdb43643531c82d5beaab30c2cfb5bf40ab52", "shasum": "" }, "require": { @@ -4102,9 +4102,9 @@ ], "support": { "issues": "https://github.com/utopia-php/platform/issues", - "source": "https://github.com/utopia-php/platform/tree/0.7.4" + "source": "https://github.com/utopia-php/platform/tree/feat-custom-server" }, - "time": "2025-03-13T13:00:12+00:00" + "time": "2025-04-17T07:58:12+00:00" }, { "name": "utopia-php/pools", @@ -4213,16 +4213,16 @@ }, { "name": "utopia-php/queue", - "version": "0.9.1", + "version": "dev-feat-pool-adapter", "source": { "type": "git", "url": "https://github.com/utopia-php/queue.git", - "reference": "32b6f84c55aae761db5a5ae76cc91ca8dbc8bc32" + "reference": "80cc1fb1053950f977f0d027bb4a7f7c91c1e794" }, "dist": { "type": "zip", - "url": "https://api.github.com/repos/utopia-php/queue/zipball/32b6f84c55aae761db5a5ae76cc91ca8dbc8bc32", - "reference": "32b6f84c55aae761db5a5ae76cc91ca8dbc8bc32", + "url": "https://api.github.com/repos/utopia-php/queue/zipball/80cc1fb1053950f977f0d027bb4a7f7c91c1e794", + "reference": "80cc1fb1053950f977f0d027bb4a7f7c91c1e794", "shasum": "" }, "require": { @@ -4231,6 +4231,7 @@ "utopia-php/cli": "0.15.*", "utopia-php/fetch": "0.4.*", "utopia-php/framework": "0.33.*", + "utopia-php/pools": "0.8.*", "utopia-php/telemetry": "0.1.*" }, "require-dev": { @@ -4272,9 +4273,9 @@ ], "support": { "issues": "https://github.com/utopia-php/queue/issues", - "source": "https://github.com/utopia-php/queue/tree/0.9.1" + "source": "https://github.com/utopia-php/queue/tree/feat-pool-adapter" }, - "time": "2025-03-28T19:49:36+00:00" + "time": "2025-04-17T08:43:02+00:00" }, { "name": "utopia-php/registry", @@ -8127,15 +8128,22 @@ ], "aliases": [ { - "package": "utopia-php/database", - "version": "dev-feat-pool-init", - "alias": "0.65.0", - "alias_normalized": "0.65.0.0" + "package": "utopia-php/platform", + "version": "dev-feat-custom-server", + "alias": "0.7.4", + "alias_normalized": "0.7.4.0" + }, + { + "package": "utopia-php/queue", + "version": "dev-feat-pool-adapter", + "alias": "0.9.1", + "alias_normalized": "0.9.1.0" } ], "minimum-stability": "stable", "stability-flags": { - "utopia-php/database": 20 + "utopia-php/platform": 20, + "utopia-php/queue": 20 }, "prefer-stable": false, "prefer-lowest": false,