pools::queue client injection

This commit is contained in:
shimon
2022-11-04 13:01:47 +02:00
parent da5addf2f2
commit 848b997a5f
9 changed files with 65 additions and 102 deletions
+4 -21
View File
@@ -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]]);
+20 -16
View File
@@ -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;
});
+4 -23
View File
@@ -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;
+10 -24
View File
@@ -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);
+4 -4
View File
@@ -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());
+4 -4
View File
@@ -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());
+1 -1
View File
@@ -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",
Generated
+17 -9
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": "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,
+1
View File
@@ -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