Use connection group to reclaim all use connections for request

This commit is contained in:
Jake Barnby
2024-04-09 17:03:28 +12:00
parent a0051eaa81
commit 730e7319df
3 changed files with 83 additions and 65 deletions
+7 -20
View File
@@ -2,6 +2,7 @@
require_once __DIR__ . '/../vendor/autoload.php';
use Appwrite\Utopia\Queue\Connections;
use Appwrite\Utopia\Response;
use Swoole\Process;
use Swoole\Http\Server;
@@ -323,27 +324,13 @@ $http->on('request', function (SwooleRequest $swooleRequest, SwooleResponse $swo
$swooleResponse->end(\json_encode($output));
} finally {
$connectionForConsole = $app->getResource('connectionForConsole');
$connectionForProject = $app->getResource('connectionForProject');
$connectionForQueue = $app->getResource('connectionForQueue');
$connectionsForCache = $app->getResource('connectionsForCache');
/**
* @var Connections $connections
*/
$connections = $app->getResource('connections');
if (!is_null($connectionForConsole)) {
$connectionForConsole->reclaim();
}
if (!is_null($connectionForProject)) {
$connectionForProject->reclaim();
}
if (!is_null($connectionForQueue)) {
$connectionForQueue->reclaim();
}
if (!empty($connectionsForCache)) {
foreach ($connectionsForCache as $connection) {
$connection->reclaim();
}
if (!empty($connections)) {
$connections->reclaim();
}
}
});
+20 -45
View File
@@ -33,6 +33,7 @@ use Appwrite\Network\Validator\Email;
use Appwrite\Network\Validator\Origin;
use Appwrite\OpenSSL\OpenSSL;
use Appwrite\URL\URL as AppwriteURL;
use Appwrite\Utopia\Queue\Connections;
use Utopia\App;
use Utopia\Logger\Logger;
use Utopia\Cache\Adapter\Redis as RedisCache;
@@ -881,17 +882,14 @@ App::setResource('locale', fn() => new Locale(App::getEnv('_APP_LOCALE', 'en')))
App::setResource('localeCodes', function () {
return array_map(fn($locale) => $locale['code'], Config::getParam('locale-codes', []));
});
// Queues
App::setResource('queue', function (Group $pools) {
App::setResource('queue', function (Group $pools, Connections $connections) {
$connection = $pools->get('queue')->pop();
App::setResource('connectionForQueue', function () use ($connection) {
return $connection;
});
$connections->add($connection);
return $connection->getResource();
}, ['pools']);
}, ['pools', 'connections']);
App::setResource('queueForMessaging', function (Connection $queue) {
return new Phone($queue);
}, ['queue']);
@@ -1132,20 +1130,11 @@ App::setResource('console', function () {
]);
}, []);
App::setResource('connectionForProject', function () {
return null;
});
App::setResource('connectionForConsole', function () {
return null;
});
App::setResource('connectionForQueue', function () {
return null;
});
App::setResource('connectionsForCache', function () {
return [];
App::setResource('connections', function () {
return new Connections();
});
App::setResource('dbForProject', function (Group $pools, Database $dbForConsole, Cache $cache, Document $project) {
App::setResource('dbForProject', function (Group $pools, Database $dbForConsole, Cache $cache, Document $project, Connections $connections) {
if ($project->isEmpty() || $project->getId() === 'console') {
return $dbForConsole;
}
@@ -1154,9 +1143,7 @@ App::setResource('dbForProject', function (Group $pools, Database $dbForConsole,
->get($project->getAttribute('database'))
->pop();
App::setResource('connectionForProject', function () use ($connection) {
return $connection;
});
$connections->add($connection);
$dbAdapter = $connection->getResource();
@@ -1169,16 +1156,14 @@ App::setResource('dbForProject', function (Group $pools, Database $dbForConsole,
->setTimeout(APP_DATABASE_TIMEOUT_MILLISECONDS);
return $database;
}, ['pools', 'dbForConsole', 'cache', 'project']);
}, ['pools', 'dbForConsole', 'cache', 'project', 'connections']);
App::setResource('dbForConsole', function (Group $pools, Cache $cache) {
App::setResource('dbForConsole', function (Group $pools, Cache $cache, Connections $connections) {
$connection = $pools
->get('console')
->pop();
App::setResource('connectionForConsole', function () use ($connection) {
return $connection;
});
$connections->add($connection);
$dbAdapter = $connection->getResource();
@@ -1191,12 +1176,12 @@ App::setResource('dbForConsole', function (Group $pools, Cache $cache) {
->setTimeout(APP_DATABASE_TIMEOUT_MILLISECONDS);
return $database;
}, ['pools', 'cache']);
}, ['pools', 'cache', 'connections']);
App::setResource('getProjectDB', function (Group $pools, Database $dbForConsole, $cache) {
App::setResource('getProjectDB', function (Group $pools, Database $dbForConsole, Cache $cache, Connections $connections) {
$databases = []; // TODO: @Meldiron This should probably be responsibility of utopia-php/pools
$getProjectDB = function (Document $project) use ($pools, $dbForConsole, $cache, &$databases) {
$getProjectDB = function (Document $project) use ($pools, $dbForConsole, $cache, $connections, &$databases) {
if ($project->isEmpty() || $project->getId() === 'console') {
return $dbForConsole;
}
@@ -1208,8 +1193,6 @@ App::setResource('getProjectDB', function (Group $pools, Database $dbForConsole,
$database
->setNamespace('_' . $project->getInternalId())
->setMetadata('host', \gethostname())
->setMetadata('project', $project->getId())
->setTimeout(APP_DATABASE_TIMEOUT_MILLISECONDS);
return $database;
@@ -1219,44 +1202,36 @@ App::setResource('getProjectDB', function (Group $pools, Database $dbForConsole,
->get($databaseName)
->pop();
$dbAdapter = $connection->getResource();
$connections->add($connection);
App::setResource('connectionForProject', function () use ($connection) {
return $connection;
}, []);
$dbAdapter = $connection->getResource();
$database = new Database($dbAdapter, $cache);
$database
->setNamespace('_' . $project->getInternalId())
->setMetadata('host', \gethostname())
->setMetadata('project', $project->getId())
->setTimeout(APP_DATABASE_TIMEOUT_MILLISECONDS);
return $database;
};
return $getProjectDB;
}, ['pools', 'dbForConsole', 'cache']);
}, ['pools', 'dbForConsole', 'cache', 'connections']);
App::setResource('cache', function (Group $pools) {
App::setResource('cache', function (Group $pools, Connections $connections) {
$list = Config::getParam('pools-cache', []);
$adapters = [];
$connections = [];
foreach ($list as $value) {
$connection = $pools
->get($value)
->pop();
$connections[] = $connection;
$connections->add($connection);
$adapters[] = $connection->getResource();
}
App::setResource('connectionsForCache', function () use ($connections) {
return $connections;
}, []);
return new Cache(new Sharding($adapters));
}, ['pools']);
+56
View File
@@ -0,0 +1,56 @@
<?php
namespace Appwrite\Utopia\Queue;
use Utopia\Pools\Connection;
class Connections
{
/**
* @var array<Connection>
*/
protected array $connections = [];
/**
* @param Connection $connection
* @return self
*/
public function add(Connection $connection): self
{
$this->connections[$connection->getID()] = $connection;
return $this;
}
/**
* @param string $id
* @return Connection
* @throws \Exception
*/
public function get(string $id): Connection
{
return $this->connections[$id] ?? throw new \Exception("Connection '{$id}' not found");
}
/**
* @param string $id
* @return self
*/
public function remove(string $id): self
{
unset($this->connections[$id]);
return $this;
}
/**
* @return self
* @throws \Exception
*/
public function reclaim(): self
{
foreach ($this->connections as $connection) {
$connection->reclaim();
}
return $this;
}
}