diff --git a/app/http.php b/app/http.php index 05e989b330..91da4e8cab 100644 --- a/app/http.php +++ b/app/http.php @@ -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(); } } }); diff --git a/app/init.php b/app/init.php index 5844cafa5c..1830fd1763 100644 --- a/app/init.php +++ b/app/init.php @@ -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']); diff --git a/src/Appwrite/Utopia/Queue/Connections.php b/src/Appwrite/Utopia/Queue/Connections.php new file mode 100644 index 0000000000..e3d5740b97 --- /dev/null +++ b/src/Appwrite/Utopia/Queue/Connections.php @@ -0,0 +1,56 @@ + + */ + 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; + } +}