From 848b997a5fc44e589ea645c0759d356bcdd0beab Mon Sep 17 00:00:00 2001 From: shimon Date: Fri, 4 Nov 2022 13:01:47 +0200 Subject: [PATCH] pools::queue client injection --- app/controllers/api/edge.php | 25 ++++--------------------- app/init.php | 36 ++++++++++++++++++++---------------- app/tasks/sync-edge.php | 27 ++++----------------------- app/worker.php | 34 ++++++++++------------------------ app/workers/sync-In.php | 8 ++++---- app/workers/sync-out.php | 8 ++++---- composer.json | 2 +- composer.lock | 26 +++++++++++++++++--------- docker-compose.yml | 1 + 9 files changed, 65 insertions(+), 102 deletions(-) diff --git a/app/controllers/api/edge.php b/app/controllers/api/edge.php index 14bc4fff9b..353befb83e 100644 --- a/app/controllers/api/edge.php +++ b/app/controllers/api/edge.php @@ -9,8 +9,8 @@ use Appwrite\Utopia\Request; use Appwrite\Utopia\Response; use Utopia\App; use Utopia\Database\Document; +use Utopia\Pools\Group; use Utopia\Queue\Client as SyncIn; -use Utopia\Queue\Connection\Redis as QueueRedis; use Utopia\Validator\ArrayList; use Utopia\Validator\Text; @@ -36,31 +36,14 @@ App::post('/v1/edge/sync') ->param('keys', '', new ArrayList(new Text(100), 1000), 'Cache keys. an array containing alphanumerical cache keys') ->inject('request') ->inject('response') - ->action(function (array $keys, Request $request, Response $response) { + ->inject('pools') + ->action(function (array $keys, Request $request, Response $response, Group $pools) { if (empty($keys)) { throw new Exception(Exception::KEY_NOT_FOUND); } - $fallbackForRedis = AppwriteURL::unparse([ - 'scheme' => 'redis', - 'host' => App::getEnv('_APP_REDIS_HOST', 'redis'), - 'port' => App::getEnv('_APP_REDIS_PORT', '6379'), - 'user' => App::getEnv('_APP_REDIS_USER', ''), - 'pass' => App::getEnv('_APP_REDIS_PASS', ''), - ]); - - $connection = App::getEnv('_APP_CONNECTIONS_QUEUE', $fallbackForRedis); - $dsns = explode(',', $connection ?? ''); - - if (empty($dsns)) { - throw new Exception(Exception::GENERAL_SERVER_ERROR); - } - - $dsn = explode('=', $dsns[0]); - $dsn = $dsn[1] ?? ''; - $dsn = new DSN($dsn); - $client = new SyncIn('syncIn', new QueueRedis($dsn->getHost(), $dsn->getPort())); + $client = new SyncIn('syncIn', $pools->get('queue')->pop()->getResource()); $client->enqueue(['value' => ['keys' => $keys]]); diff --git a/app/init.php b/app/init.php index 3c97c1f2d0..4c38024e5a 100644 --- a/app/init.php +++ b/app/init.php @@ -77,6 +77,7 @@ use Ahc\Jwt\JWTException; use MaxMind\Db\Reader; use PHPMailer\PHPMailer\PHPMailer; use Swoole\Database\PDOProxy; +use Utopia\Queue; const APP_NAME = 'Appwrite'; const APP_DOMAIN = 'appwrite.io'; @@ -526,30 +527,35 @@ $register->set('pools', function () { 'dsns' => App::getEnv('_APP_CONNECTIONS_DB_CONSOLE', $fallbackForDB), 'multiple' => false, 'schemes' => ['mariadb', 'mysql'], + 'useResource' => true, ], 'database' => [ 'type' => 'database', 'dsns' => App::getEnv('_APP_CONNECTIONS_DB_PROJECT', $fallbackForDB), 'multiple' => true, 'schemes' => ['mariadb', 'mysql'], + 'useResource' => true, + ], + 'queue' => [ + 'type' => 'queue', + 'dsns' => App::getEnv('_APP_CONNECTIONS_QUEUE', $fallbackForRedis), + 'multiple' => false, + 'schemes' => ['redis'], + 'useResource' => false, ], -// 'queue' => [ -// 'type' => 'queue', -// 'dsns' => App::getEnv('_APP_CONNECTIONS_QUEUE', $fallbackForRedis), -// 'multiple' => false, -// 'schemes' => ['redis'], -// ], 'pubsub' => [ 'type' => 'pubsub', 'dsns' => App::getEnv('_APP_CONNECTIONS_PUBSUB', $fallbackForRedis), 'multiple' => false, 'schemes' => ['redis'], + 'useResource' => true, ], 'cache' => [ 'type' => 'cache', 'dsns' => App::getEnv('_APP_CONNECTIONS_CACHE', $fallbackForRedis), 'multiple' => true, 'schemes' => ['redis'], + 'useResource' => true, ], ]; @@ -558,6 +564,7 @@ $register->set('pools', function () { $dsns = $connection['dsns'] ?? ''; $multipe = $connection['multiple'] ?? false; $schemes = $connection['schemes'] ?? []; + $useResource = $connection['useResource'] ?? true; $config = []; $dsns = explode(',', $connection['dsns'] ?? ''); @@ -580,7 +587,7 @@ $register->set('pools', function () { $dsnScheme = $dsn->getScheme(); $dsnDatabase = $dsn->getDatabase(); - if (!in_array($dsnScheme, $schemes)) { + if (!in_array($dsnScheme, $schemes) && $useResource) { throw new Exception(Exception::GENERAL_SERVER_ERROR, "Invalid console database scheme"); } @@ -643,9 +650,12 @@ $register->set('pools', function () { case 'pubsub': break; $adapter = $resource(); -// case 'queue': -// $adapter = $resource(); -// break; + case 'queue': + $adapter = match ($dsn->getScheme()) { + 'redis' => new Queue\Connection\Redis($dsn->getHost(), $dsn->getPort()), + default => 'bla' + }; + break; case 'cache': $adapter = match ($dsn->getScheme()) { 'redis' => new RedisCache($resource()), @@ -667,12 +677,6 @@ $register->set('pools', function () { Config::setParam('pools-' . $key, $config); } - try { - $group->fill(); - } catch (\Throwable $th) { - Console::error('Connection failure: ' . $th->getMessage()); - } - return $group; }); diff --git a/app/tasks/sync-edge.php b/app/tasks/sync-edge.php index 3310704ffa..102d5ca2a2 100644 --- a/app/tasks/sync-edge.php +++ b/app/tasks/sync-edge.php @@ -19,36 +19,17 @@ $cli Console::title('Syncs edges V1'); Console::success(APP_NAME . ' Sync failed cache purge process v1 has started'); - $fallbackForRedis = AppwriteURL::unparse([ - 'scheme' => 'redis', - 'host' => App::getEnv('_APP_REDIS_HOST', 'redis'), - 'port' => App::getEnv('_APP_REDIS_PORT', '6379'), - 'user' => App::getEnv('_APP_REDIS_USER', ''), - 'pass' => App::getEnv('_APP_REDIS_PASS', ''), - ]); - - $connection = App::getEnv('_APP_CONNECTIONS_QUEUE', $fallbackForRedis); - $dsns = explode(',', $connection ?? ''); - - if (empty($dsns)) { - Console::error("No Dsn found"); - } - - $dsn = explode('=', $dsns[0]); - $dsn = $dsn[1] ?? ''; - $dsn = new DSN($dsn); - $redisConnection = new Queue\Connection\Redis($dsn->getHost(), $dsn->getPort()); - $client = new SyncOut('syncOut', $redisConnection); - + $pools = $register->get('pools'); + $client = new SyncOut('syncOut', $pools->get('queue')->pop()->getResource()); + $database = getConsoleDB(); // Todo fix pdo PDOException // Table 'appwrite.console__metadata' doesn't exist sleep(4); $interval = (int) App::getEnv('_APP_SYNC_EDGE_INTERVAL', '180'); - Console::loop(function () use ($interval, $register, $client) { + Console::loop(function () use ($interval, $database, $register, $client) { - $database = getConsoleDB(); $time = DateTime::now(); $count = 0; $chunk = 0; diff --git a/app/worker.php b/app/worker.php index f705636bf8..d2668f6417 100644 --- a/app/worker.php +++ b/app/worker.php @@ -2,9 +2,7 @@ require_once __DIR__ . '/init.php'; -use Appwrite\DSN\DSN; -use Appwrite\Extend\Exception; -use Appwrite\URL\URL as AppwriteURL; + use Swoole\Runtime; use Utopia\App; use Utopia\Cache\Adapter\Sharding; @@ -13,11 +11,11 @@ use Utopia\Config\Config; use Utopia\Database\Database; use Utopia\Queue\Server; use Utopia\Registry\Registry; -use Utopia\Queue; + global $register; -Runtime::enableCoroutine(SWOOLE_HOOK_ALL); + Server::setResource('register', fn() => $register); @@ -56,26 +54,14 @@ App::setResource('logger', function ($register) { }, ['register']); -// Todo better job to inject the client as a resource -$fallbackForRedis = AppwriteURL::unparse([ - 'scheme' => 'redis', - 'host' => App::getEnv('_APP_REDIS_HOST', 'redis'), - 'port' => App::getEnv('_APP_REDIS_PORT', '6379'), - 'user' => App::getEnv('_APP_REDIS_USER', ''), - 'pass' => App::getEnv('_APP_REDIS_PASS', ''), -]); +$pools = $register->get('pools'); +$client = $pools + ->get('queue') + ->pop() + ->getResource(); -$connection = App::getEnv('_APP_CONNECTIONS_QUEUE', $fallbackForRedis); -$dsns = explode(',', $connection ?? ''); -if (empty($dsns)) { - throw new Exception(Exception::GENERAL_SERVER_ERROR); -} - -$dsn = explode('=', $dsns[0]); -$dsn = $dsn[1] ?? ''; -$dsn = new DSN($dsn); - -$redisConnection = new Queue\Connection\Redis($dsn->getHost(), $dsn->getPort()); $workerNumber = swoole_cpu_num() * intval(App::getEnv('_APP_WORKER_PER_CORE', 6)); $workerNumber = 1; + +Runtime::enableCoroutine(SWOOLE_HOOK_ALL); diff --git a/app/workers/sync-In.php b/app/workers/sync-In.php index 2df457083a..e359a954e4 100644 --- a/app/workers/sync-In.php +++ b/app/workers/sync-In.php @@ -10,10 +10,10 @@ use Utopia\Logger\Log; use Utopia\Queue; use Utopia\Queue\Message; -global $redisConnection; +global $client; global $workerNumber; -$adapter = new Queue\Adapter\Swoole($redisConnection, $workerNumber, 'syncIn'); +$adapter = new Queue\Adapter\Swoole($client, $workerNumber, 'syncIn'); $server = new Queue\Server($adapter); $server->job() @@ -44,12 +44,12 @@ $server if ($error->getCode() >= 500 || $error->getCode() === 0) { $log = new Log(); - $log->setNamespace("worker"); + $log->setNamespace("appwrite-worker"); $log->setServer(\gethostname()); $log->setVersion($version); $log->setType(Log::TYPE_ERROR); $log->setMessage($error->getMessage()); - $log->setAction('worker-sync-out'); + $log->setAction('appwrite-worker-sync-out'); $log->addTag('verboseType', get_class($error)); $log->addTag('code', $error->getCode()); $log->addExtra('file', $error->getFile()); diff --git a/app/workers/sync-out.php b/app/workers/sync-out.php index 8c36cbe808..564e39e0d7 100644 --- a/app/workers/sync-out.php +++ b/app/workers/sync-out.php @@ -17,7 +17,7 @@ use Utopia\Logger\Log; use Utopia\Queue; use Utopia\Queue\Message; -global $redisConnection; +global $client; global $workerNumber; $regions = array_filter( @@ -103,7 +103,7 @@ function handle($dbForConsole, $regions, $stack): void } } -$adapter = new Queue\Adapter\Swoole($redisConnection, $workerNumber, 'syncOut'); +$adapter = new Queue\Adapter\Swoole($client, $workerNumber, 'syncOut'); $server = new Queue\Server($adapter); $server->job() @@ -149,12 +149,12 @@ $server if ($error->getCode() >= 500 || $error->getCode() === 0) { $log = new Log(); - $log->setNamespace("worker"); + $log->setNamespace("appwrite-worker"); $log->setServer(\gethostname()); $log->setVersion($version); $log->setType(Log::TYPE_ERROR); $log->setMessage($error->getMessage()); - $log->setAction('worker-sync-out'); + $log->setAction('appwrite-worker-sync-out'); $log->addTag('verboseType', get_class($error)); $log->addTag('code', $error->getCode()); $log->addExtra('file', $error->getFile()); diff --git a/composer.json b/composer.json index 110e44ca9d..d2fef924fe 100644 --- a/composer.json +++ b/composer.json @@ -62,7 +62,7 @@ "utopia-php/image": "0.5.*", "utopia-php/orchestration": "0.6.*", "utopia-php/queue": "0.4.0", - "utopia-php/pools": "0.1.*", + "utopia-php/pools": "dev-upgrade-cli as 0.2.0", "resque/php-resque": "1.3.6", "matomo/device-detector": "6.0.0", "dragonmantank/cron-expression": "3.3.1", diff --git a/composer.lock b/composer.lock index 2aa790ead4..d39044963f 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": "a02c3502dca5a3a9f0f283e06e11e30e", + "content-hash": "88fa95b378e538da3ff62fe81bf76bc1", "packages": [ { "name": "adhocore/jwt", @@ -2431,23 +2431,24 @@ }, { "name": "utopia-php/pools", - "version": "0.1.0", + "version": "dev-upgrade-cli", "source": { "type": "git", "url": "https://github.com/utopia-php/pools.git", - "reference": "5a467a569a80aefc846a97dc195b4adc2fd71805" + "reference": "88a2c1ed2badbfdf2787ce0a12def2c988fc1097" }, "dist": { "type": "zip", - "url": "https://api.github.com/repos/utopia-php/pools/zipball/5a467a569a80aefc846a97dc195b4adc2fd71805", - "reference": "5a467a569a80aefc846a97dc195b4adc2fd71805", + "url": "https://api.github.com/repos/utopia-php/pools/zipball/88a2c1ed2badbfdf2787ce0a12def2c988fc1097", + "reference": "88a2c1ed2badbfdf2787ce0a12def2c988fc1097", "shasum": "" }, "require": { "ext-mongodb": "*", "ext-pdo": "*", "ext-redis": "*", - "php": ">=8.0" + "php": ">=8.0", + "utopia-php/cli": "0.13.*" }, "require-dev": { "phpunit/phpunit": "^9.4", @@ -2478,9 +2479,9 @@ ], "support": { "issues": "https://github.com/utopia-php/pools/issues", - "source": "https://github.com/utopia-php/pools/tree/0.1.0" + "source": "https://github.com/utopia-php/pools/tree/upgrade-cli" }, - "time": "2022-10-11T19:31:07+00:00" + "time": "2022-11-04T08:33:04+00:00" }, { "name": "utopia-php/preloader", @@ -5469,12 +5470,19 @@ "version": "dev-feat-update-cache-lib", "alias": "0.26.1", "alias_normalized": "0.26.1.0" + }, + { + "package": "utopia-php/pools", + "version": "dev-upgrade-cli", + "alias": "0.2.0", + "alias_normalized": "0.2.0.0" } ], "minimum-stability": "stable", "stability-flags": { "utopia-php/cache": 20, - "utopia-php/database": 20 + "utopia-php/database": 20, + "utopia-php/pools": 20 }, "prefer-stable": false, "prefer-lowest": false, diff --git a/docker-compose.yml b/docker-compose.yml index 97d85b2413..df87436ae8 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -280,6 +280,7 @@ services: volumes: - ./app:/usr/src/code/app - ./src:/usr/src/code/src + - ./vendor/utopia-php/pools:/usr/src/code/vendor/utopia-php/pools depends_on: - mariadb - redis