Add publisher consumer pool adapter

This commit is contained in:
Jake Barnby
2025-04-17 22:27:20 +12:00
parent 25a3266301
commit fe26dcb50a
6 changed files with 58 additions and 50 deletions
+4 -2
View File
@@ -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);
+2 -1
View File
@@ -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);
+3 -2
View File
@@ -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);
+14 -18
View File
@@ -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) {
+3 -3
View File
@@ -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.*",
Generated
+32 -24
View File
@@ -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,