1
0
Fork 0
mirror of synced 2024-10-05 12:43:13 +13:00

Revert "Reclaim only used connections for workers"

This reverts commit fefb60557d.

# Conflicts:
#	app/worker.php
This commit is contained in:
Jake Barnby 2024-04-09 20:50:33 +12:00
parent 9a3f6e7f71
commit 3859f037a1
No known key found for this signature in database
GPG key ID: C437A8CC85B96E9C

View file

@ -16,7 +16,6 @@ use Appwrite\Event\Phone;
use Appwrite\Event\Usage; use Appwrite\Event\Usage;
use Appwrite\Event\UsageDump; use Appwrite\Event\UsageDump;
use Appwrite\Platform\Appwrite; use Appwrite\Platform\Appwrite;
use Appwrite\Utopia\Pools\Connections;
use Swoole\Runtime; use Swoole\Runtime;
use Utopia\App; use Utopia\App;
use Utopia\Cache\Adapter\Sharding; use Utopia\Cache\Adapter\Sharding;
@ -37,32 +36,25 @@ use Utopia\Pools\Group;
use Utopia\Queue\Connection; use Utopia\Queue\Connection;
Authorization::disable(); Authorization::disable();
Runtime::enableCoroutine(); Runtime::enableCoroutine(SWOOLE_HOOK_ALL);
Server::setResource('register', fn () => $register); Server::setResource('register', fn () => $register);
Server::setResource('connections', function () { Server::setResource('dbForConsole', function (Cache $cache, Registry $register) {
return new Connections();
});
Server::setResource('dbForConsole', function (Cache $cache, Registry $register, Connections $connections) {
$pools = $register->get('pools'); $pools = $register->get('pools');
$database = $pools
$connection = $pools
->get('console') ->get('console')
->pop(); ->pop()
->getResource()
$connections->add($connection); ;
$database = $connection->getResource();
$adapter = new Database($database, $cache); $adapter = new Database($database, $cache);
$adapter->setNamespace('_console'); $adapter->setNamespace('_console');
return $adapter; return $adapter;
}, ['cache', 'register', 'connections']); }, ['cache', 'register']);
Server::setResource('dbForProject', function (Cache $cache, Registry $register, Message $message, Database $dbForConsole, Connections $connections) { Server::setResource('dbForProject', function (Cache $cache, Registry $register, Message $message, Database $dbForConsole) {
$payload = $message->getPayload() ?? []; $payload = $message->getPayload() ?? [];
$project = new Document($payload['project'] ?? []); $project = new Document($payload['project'] ?? []);
@ -71,19 +63,16 @@ Server::setResource('dbForProject', function (Cache $cache, Registry $register,
} }
$pools = $register->get('pools'); $pools = $register->get('pools');
$database = $pools
$connection = $pools
->get($project->getAttribute('database')) ->get($project->getAttribute('database'))
->pop(); ->pop()
->getResource()
$database = $connection->getResource(); ;
$connections->add($connection);
$adapter = new Database($database, $cache); $adapter = new Database($database, $cache);
$adapter->setNamespace('_' . $project->getInternalId()); $adapter->setNamespace('_' . $project->getInternalId());
return $adapter; return $adapter;
}, ['cache', 'register', 'message', 'dbForConsole', 'connections']); }, ['cache', 'register', 'message', 'dbForConsole']);
Server::setResource('project', function (Message $message, Database $dbForConsole) { Server::setResource('project', function (Message $message, Database $dbForConsole) {
$payload = $message->getPayload() ?? []; $payload = $message->getPayload() ?? [];
@ -92,14 +81,14 @@ Server::setResource('project', function (Message $message, Database $dbForConsol
if ($project->getId() === 'console') { if ($project->getId() === 'console') {
return $project; return $project;
} }
return $dbForConsole->getDocument('projects', $project->getId()); return $dbForConsole->getDocument('projects', $project->getId());
;
}, ['message', 'dbForConsole']); }, ['message', 'dbForConsole']);
Server::setResource('getProjectDB', function (Group $pools, Database $dbForConsole, Cache $cache, Connections $connections) { Server::setResource('getProjectDB', function (Group $pools, Database $dbForConsole, $cache) {
$databases = []; // TODO: @Meldiron This should probably be responsibility of utopia-php/pools $databases = []; // TODO: @Meldiron This should probably be responsibility of utopia-php/pools
return function (Document $project) use ($pools, $dbForConsole, $cache, $connections, &$databases): Database { return function (Document $project) use ($pools, $dbForConsole, $cache, &$databases): Database {
if ($project->isEmpty() || $project->getId() === 'console') { if ($project->isEmpty() || $project->getId() === 'console') {
return $dbForConsole; return $dbForConsole;
} }
@ -112,13 +101,10 @@ Server::setResource('getProjectDB', function (Group $pools, Database $dbForConso
return $database; return $database;
} }
$dbConnection = $pools $dbAdapter = $pools
->get($databaseName) ->get($databaseName)
->pop(); ->pop()
->getResource();
$dbAdapter = $dbConnection->getResource();
$connections->add($dbConnection);
$database = new Database($dbAdapter, $cache); $database = new Database($dbAdapter, $cache);
@ -128,7 +114,7 @@ Server::setResource('getProjectDB', function (Group $pools, Database $dbForConso
return $database; return $database;
}; };
}, ['pools', 'dbForConsole', 'cache', 'connections']); }, ['pools', 'dbForConsole', 'cache']);
Server::setResource('abuseRetention', function () { Server::setResource('abuseRetention', function () {
return DateTime::addSeconds(new \DateTime(), -1 * App::getEnv('_APP_MAINTENANCE_RETENTION_ABUSE', 86400)); return DateTime::addSeconds(new \DateTime(), -1 * App::getEnv('_APP_MAINTENANCE_RETENTION_ABUSE', 86400));
@ -142,104 +128,82 @@ Server::setResource('executionRetention', function () {
return DateTime::addSeconds(new \DateTime(), -1 * App::getEnv('_APP_MAINTENANCE_RETENTION_EXECUTION', 1209600)); return DateTime::addSeconds(new \DateTime(), -1 * App::getEnv('_APP_MAINTENANCE_RETENTION_EXECUTION', 1209600));
}); });
Server::setResource('cache', function (Registry $register, Connections $connections) { Server::setResource('cache', function (Registry $register) {
$pools = $register->get('pools'); $pools = $register->get('pools');
$list = Config::getParam('pools-cache', []); $list = Config::getParam('pools-cache', []);
$adapters = []; $adapters = [];
foreach ($list as $value) { foreach ($list as $value) {
$connection = $pools $adapters[] = $pools
->get($value) ->get($value)
->pop(); ->pop()
->getResource()
$connections->add($connection); ;
$adapters[] = $connection->getResource();
} }
return new Cache(new Sharding($adapters)); return new Cache(new Sharding($adapters));
}, ['register', 'connections']); }, ['register']);
Server::setResource('log', fn() => new Log()); Server::setResource('log', fn() => new Log());
Server::setResource('queueForUsage', function (Connection $queue) { Server::setResource('queueForUsage', function (Connection $queue) {
return new Usage($queue); return new Usage($queue);
}, ['queue']); }, ['queue']);
Server::setResource('queueForUsageDump', function (Connection $queue) { Server::setResource('queueForUsageDump', function (Connection $queue) {
return new UsageDump($queue); return new UsageDump($queue);
}, ['queue']); }, ['queue']);
Server::setResource('queue', function (Group $pools) { Server::setResource('queue', function (Group $pools) {
return $pools->get('queue')->pop()->getResource(); return $pools->get('queue')->pop()->getResource();
}, ['pools']); }, ['pools']);
Server::setResource('queueForDatabase', function (Connection $queue) { Server::setResource('queueForDatabase', function (Connection $queue) {
return new EventDatabase($queue); return new EventDatabase($queue);
}, ['queue']); }, ['queue']);
Server::setResource('queueForMessaging', function (Connection $queue) { Server::setResource('queueForMessaging', function (Connection $queue) {
return new Phone($queue); return new Phone($queue);
}, ['queue']); }, ['queue']);
Server::setResource('queueForMails', function (Connection $queue) { Server::setResource('queueForMails', function (Connection $queue) {
return new Mail($queue); return new Mail($queue);
}, ['queue']); }, ['queue']);
Server::setResource('queueForBuilds', function (Connection $queue) { Server::setResource('queueForBuilds', function (Connection $queue) {
return new Build($queue); return new Build($queue);
}, ['queue']); }, ['queue']);
Server::setResource('queueForDeletes', function (Connection $queue) { Server::setResource('queueForDeletes', function (Connection $queue) {
return new Delete($queue); return new Delete($queue);
}, ['queue']); }, ['queue']);
Server::setResource('queueForEvents', function (Connection $queue) { Server::setResource('queueForEvents', function (Connection $queue) {
return new Event($queue); return new Event($queue);
}, ['queue']); }, ['queue']);
Server::setResource('queueForAudits', function (Connection $queue) { Server::setResource('queueForAudits', function (Connection $queue) {
return new Audit($queue); return new Audit($queue);
}, ['queue']); }, ['queue']);
Server::setResource('queueForFunctions', function (Connection $queue) { Server::setResource('queueForFunctions', function (Connection $queue) {
return new Func($queue); return new Func($queue);
}, ['queue']); }, ['queue']);
Server::setResource('queueForCertificates', function (Connection $queue) { Server::setResource('queueForCertificates', function (Connection $queue) {
return new Certificate($queue); return new Certificate($queue);
}, ['queue']); }, ['queue']);
Server::setResource('queueForMigrations', function (Connection $queue) { Server::setResource('queueForMigrations', function (Connection $queue) {
return new Migration($queue); return new Migration($queue);
}, ['queue']); }, ['queue']);
Server::setResource('logger', function (Registry $register) { Server::setResource('logger', function (Registry $register) {
return $register->get('logger'); return $register->get('logger');
}, ['register']); }, ['register']);
Server::setResource('pools', function (Registry $register) { Server::setResource('pools', function (Registry $register) {
return $register->get('pools'); return $register->get('pools');
}, ['register']); }, ['register']);
Server::setResource('getFunctionsDevice', function () { Server::setResource('getFunctionsDevice', function () {
return function (string $projectId) { return function (string $projectId) {
return getDevice(APP_STORAGE_FUNCTIONS . '/app-' . $projectId); return getDevice(APP_STORAGE_FUNCTIONS . '/app-' . $projectId);
}; };
}); });
Server::setResource('getFilesDevice', function () { Server::setResource('getFilesDevice', function () {
return function (string $projectId) { return function (string $projectId) {
return getDevice(APP_STORAGE_UPLOADS . '/app-' . $projectId); return getDevice(APP_STORAGE_UPLOADS . '/app-' . $projectId);
}; };
}); });
Server::setResource('getBuildsDevice', function () { Server::setResource('getBuildsDevice', function () {
return function (string $projectId) { return function (string $projectId) {
return getDevice(APP_STORAGE_BUILDS . '/app-' . $projectId); return getDevice(APP_STORAGE_BUILDS . '/app-' . $projectId);
}; };
}); });
Server::setResource('getCacheDevice', function () { Server::setResource('getCacheDevice', function () {
return function (string $projectId) { return function (string $projectId) {
return getDevice(APP_STORAGE_CACHE . '/app-' . $projectId); return getDevice(APP_STORAGE_CACHE . '/app-' . $projectId);
@ -290,9 +254,9 @@ $worker = $platform->getWorker();
$worker $worker
->shutdown() ->shutdown()
->inject('connections') ->inject('pools')
->action(function (Connections $connections) { ->action(function (Group $pools) {
$connections->reclaim(); $pools->reclaim();
}); });
$worker $worker
@ -301,10 +265,7 @@ $worker
->inject('logger') ->inject('logger')
->inject('log') ->inject('log')
->inject('project') ->inject('project')
->inject('connections') ->action(function (Throwable $error, ?Logger $logger, Log $log, Document $project) {
->action(function (Throwable $error, ?Logger $logger, Log $log, Document $project, Connections $connections) {
$connections->reclaim();
$version = App::getEnv('_APP_VERSION', 'UNKNOWN'); $version = App::getEnv('_APP_VERSION', 'UNKNOWN');
if ($error instanceof PDOException) { if ($error instanceof PDOException) {