diff --git a/app/cli.php b/app/cli.php index c7e5b98157..7c83b63a68 100644 --- a/app/cli.php +++ b/app/cli.php @@ -17,6 +17,7 @@ use Utopia\Database\Database; use Utopia\Database\Document; use Utopia\Logger\Log; use Utopia\Pools\Group; +use Utopia\Queue\Client; use Utopia\Registry\Registry; Authorization::disable(); @@ -114,6 +115,10 @@ CLI::setResource('queueForFunctions', function (Group $pools) { return new Func($pools->get('queue')->pop()->getResource()); }, ['pools']); +CLI::setResource('queueForCacheSyncOut', function (Group $pools) { + return new Client('v1-sync-out', $pools->get('queue')->pop()->getResource()); +}, ['pools']); + CLI::setResource('logError', function (Registry $register) { return function (Throwable $error, string $namespace, string $action) use ($register) { $logger = $register->get('logger'); diff --git a/app/controllers/api/edge.php b/app/controllers/api/edge.php index 396f2f29ef..91b8cf06dc 100644 --- a/app/controllers/api/edge.php +++ b/app/controllers/api/edge.php @@ -33,7 +33,7 @@ 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') - ->inject('pools') + ->inject('queueForCacheSyncOut') ->action(function (array $keys, Request $request, Response $response, Client $queueForCacheSyncOut) { if (empty($keys)) { diff --git a/app/init.php b/app/init.php index 8c0e979e21..a10086c10e 100644 --- a/app/init.php +++ b/app/init.php @@ -595,7 +595,7 @@ $register->set('pools', function () { $dsnUser = $dsn->getUser(); $dsnPass = $dsn->getPassword(); $dsnScheme = $dsn->getScheme(); - $dsnDatabase = $dsn->getDatabase(); + $dsnPath = $dsn->getPath(); if (!in_array($dsnScheme, $schemes)) { throw new Exception(Exception::GENERAL_SERVER_ERROR, "Invalid console database scheme"); @@ -612,9 +612,9 @@ $register->set('pools', function () { switch ($dsnScheme) { case 'mysql': case 'mariadb': - $resource = function () use ($dsnHost, $dsnPort, $dsnUser, $dsnPass, $dsnDatabase) { - return new PDOProxy(function () use ($dsnHost, $dsnPort, $dsnUser, $dsnPass, $dsnDatabase) { - return new PDO("mysql:host={$dsnHost};port={$dsnPort};dbname={$dsnDatabase};charset=utf8mb4", $dsnUser, $dsnPass, array( + $resource = function () use ($dsnHost, $dsnPort, $dsnUser, $dsnPass, $dsnPath) { + return new PDOProxy(function () use ($dsnHost, $dsnPort, $dsnUser, $dsnPass, $dsnPath) { + return new PDO("mysql:host={$dsnHost};port={$dsnPort};dbname={$dsnPath};charset=utf8mb4", $dsnUser, $dsnPass, array( PDO::ATTR_TIMEOUT => 3, // Seconds PDO::ATTR_PERSISTENT => true, PDO::ATTR_DEFAULT_FETCH_MODE => PDO::FETCH_ASSOC, @@ -655,7 +655,7 @@ $register->set('pools', function () { default => null }; - $adapter->setDefaultDatabase($dsn->getDatabase()); + $adapter->setDefaultDatabase($dsn->getPath()); break; case 'pubsub': $adapter = $resource(); @@ -1113,7 +1113,7 @@ App::setResource('cache', function (Group $pools, Client $queueForCacheSyncOut) }); return $cache; -}, ['pools']); +}, ['pools', 'queueForCacheSyncOut']); App::setResource('deviceLocal', function () { return new Local(); diff --git a/docker-compose.yml b/docker-compose.yml index f3dd27de07..8bd0a4d8c8 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -645,6 +645,7 @@ services: volumes: - ./app:/usr/src/code/app - ./src:/usr/src/code/src + - ./vendor/utopia-php/cli:/usr/src/code/vendor/utopia-php/cli depends_on: - mariadb - redis @@ -665,6 +666,7 @@ services: - _APP_CONNECTIONS_CACHE - _APP_CONNECTIONS_QUEUE - _APP_REGION + - _APP_WORKER_PER_CORE appwrite-usage-timeseries: entrypoint: diff --git a/src/Appwrite/Platform/Services/Tasks.php b/src/Appwrite/Platform/Services/Tasks.php index 2968a66b95..915cc34618 100644 --- a/src/Appwrite/Platform/Services/Tasks.php +++ b/src/Appwrite/Platform/Services/Tasks.php @@ -15,6 +15,7 @@ use Appwrite\Platform\Tasks\Usage; use Appwrite\Platform\Tasks\Vars; use Appwrite\Platform\Tasks\Version; use Appwrite\Platform\Tasks\VolumeSync; +use Appwrite\Platform\Tasks\EdgeSync; class Tasks extends Service { @@ -33,6 +34,7 @@ class Tasks extends Service ->addAction(Migrate::getName(), new Migrate()) ->addAction(SDKs::getName(), new SDKs()) ->addAction(VolumeSync::getName(), new VolumeSync()) - ->addAction(Specs::getName(), new Specs()); + ->addAction(Specs::getName(), new Specs()) + ->addAction(EdgeSync::getName(), new EdgeSync()); } } diff --git a/src/Appwrite/Platform/Tasks/EdgeSync.php b/src/Appwrite/Platform/Tasks/EdgeSync.php index e5202a27c5..086477d803 100644 --- a/src/Appwrite/Platform/Tasks/EdgeSync.php +++ b/src/Appwrite/Platform/Tasks/EdgeSync.php @@ -1,30 +1,40 @@ task('edge-sync') - ->desc('Schedules edge sync tasks') - ->action(function () use ($register) { - Console::title('Syncs edges V1'); - Console::success(APP_NAME . ' Sync failed cache purge process v1 has started'); +class EdgeSync extends Action +{ + public static function getName(): string + { + return 'edge-sync'; + } - $pools = $register->get('pools'); - $client = new SyncOut('syncOut', $pools->get('queue')->pop()->getResource()); - $database = getConsoleDB(); + public function __construct() + { + $this + ->desc('Schedules edge sync tasks') + ->inject('pools') + ->inject('dbForConsole') + ->inject('queueForCacheSyncOut') + ->callback(fn (Group $pools, Database $dbForConsole, Client $queueForCacheSyncOut) => $this->action($pools, $dbForConsole, $queueForCacheSyncOut)); + } + + public function action(Group $pools, Database $dbForConsole, Client $queueForCacheSyncOut): void + { + Console::title('Edge-sync V1'); + Console::success(APP_NAME . ' Edge-sync v1 has started'); $interval = (int) App::getEnv('_APP_SYNC_EDGE_INTERVAL', '180'); - Console::loop(function () use ($interval, $database, $register, $client) { + Console::loop(function () use ($interval, $dbForConsole, $queueForCacheSyncOut) { $time = DateTime::now(); $count = 0; @@ -37,7 +47,7 @@ $cli while ($sum === $limit) { $chunk++; - $results = $database->find('syncs', [ + $results = $dbForConsole->find('syncs', [ Query::equal('region', [App::getEnv('_APP_REGION')]), Query::limit($limit) ]); @@ -46,7 +56,7 @@ $cli if ($sum > 0) { foreach ($results as $document) { Console::info("[{$time}] Enqueueing keys chunk {$count} to {$document->getAttribute('target')}"); - $client + $queueForCacheSyncOut ->enqueue([ 'value' => [ 'region' => $document->getAttribute('target'), @@ -54,12 +64,13 @@ $cli ] ]); - $database->deleteDocument('syncs', $document->getId()); + $dbForConsole->deleteDocument('syncs', $document->getId()); $count++; } } else { Console::info("[{$time}] No cache keys where found."); } } - }, $interval); - }); + }, $interval); + } +}