diff --git a/app/controllers/api/edge.php b/app/controllers/api/edge.php index 91b8cf06dc..b04c26972e 100644 --- a/app/controllers/api/edge.php +++ b/app/controllers/api/edge.php @@ -33,15 +33,17 @@ 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('queueForCacheSyncOut') - ->action(function (array $keys, Request $request, Response $response, Client $queueForCacheSyncOut) { + ->inject('queueForCacheSyncIn') + ->action(function (array $keys, Request $request, Response $response, Client $queueForCacheSyncIn) { if (empty($keys)) { throw new Exception(Exception::KEY_NOT_FOUND); } - $queueForCacheSyncOut - ->enqueue(['value' => ['keys' => $keys]]); + $queueForCacheSyncIn + ->enqueue([ + 'keys' => $keys + ]); $response->dynamic(new Document([ 'keys' => $keys diff --git a/app/init.php b/app/init.php index a10086c10e..2833130957 100644 --- a/app/init.php +++ b/app/init.php @@ -1101,7 +1101,9 @@ App::setResource('cache', function (Group $pools, Client $queueForCacheSyncOut) return; } $queueForCacheSyncOut - ->enqueue(['value' => ['key' => $key]]); + ->enqueue([ + 'key' => $key + ]); }); $cache->on(cache::EVENT_PURGE, function ($key) use ($queueForCacheSyncOut) { @@ -1109,7 +1111,9 @@ App::setResource('cache', function (Group $pools, Client $queueForCacheSyncOut) return; } $queueForCacheSyncOut - ->enqueue(['value' => ['key' => $key]]); + ->enqueue([ + 'key' => $key + ]); }); return $cache; diff --git a/app/worker.php b/app/worker.php index 42a5f92439..050d352caa 100644 --- a/app/worker.php +++ b/app/worker.php @@ -101,6 +101,8 @@ if (empty(App::getEnv('QUEUE'))) { throw new Exception('Please configure "QUEUE" environemnt variable.'); } + +$workerNumber =1; $adapter = new Swoole($connection, $workerNumber, App::getEnv('QUEUE')); $server = new Server($adapter); diff --git a/app/workers/sync-In.php b/app/workers/sync-In.php index 60ccea9a23..4eab3be337 100644 --- a/app/workers/sync-In.php +++ b/app/workers/sync-In.php @@ -6,72 +6,27 @@ use Utopia\App; use Utopia\Cache\Cache; use Utopia\CLI\Console; use Utopia\Database\DateTime; -use Utopia\Logger\Log; -use Utopia\Queue; use Utopia\Queue\Message; global $connection; global $workerNumber; -$adapter = new Queue\Adapter\Swoole($connection, $workerNumber, 'syncIn'); -$server = new Queue\Server($adapter); - $server->job() ->inject('message') ->inject('cache') ->action(function (Message $message, Cache $cache) { $time = DateTime::now(); - $payload = $message->getPayload()['value']; + $cache->setDisableListeners(true); - foreach ($payload['keys'] ?? [] as $key) { - Console::info("[{$time}] Purging {$key}"); + foreach ($message->getPayload()['keys'] ?? [] as $key) { + Console::log("[{$time}] Purging {$key}"); $cache->purge($key); } $cache->setDisableListeners(false); }); -$server - ->error() - ->inject('error') - ->inject('logger') - ->action(function ($error, $logger) { - $version = App::getEnv('_APP_VERSION', 'UNKNOWN'); - - if ($error instanceof PDOException) { - throw $error; - } - - if ($error->getCode() >= 500 || $error->getCode() === 0) { - $log = new Log(); - - $log->setNamespace("appwrite-worker"); - $log->setServer(\gethostname()); - $log->setVersion($version); - $log->setType(Log::TYPE_ERROR); - $log->setMessage($error->getMessage()); - $log->setAction('appwrite-worker-sync-out'); - $log->addTag('verboseType', get_class($error)); - $log->addTag('code', $error->getCode()); - $log->addExtra('file', $error->getFile()); - $log->addExtra('line', $error->getLine()); - $log->addExtra('trace', $error->getTraceAsString()); - $log->addExtra('detailedTrace', $error->getTrace()); - $log->addExtra('roles', \Utopia\Database\Validator\Authorization::$roles); - - $isProduction = App::getEnv('_APP_ENV', 'development') === 'production'; - $log->setEnvironment($isProduction ? Log::ENVIRONMENT_PRODUCTION : Log::ENVIRONMENT_STAGING); - - $logger->addLog($log); - } - - Console::error('[Error] Type: ' . get_class($error)); - Console::error('[Error] Message: ' . $error->getMessage()); - Console::error('[Error] File: ' . $error->getFile()); - Console::error('[Error] Line: ' . $error->getLine()); - }); - $server ->workerStart() ->action(function () { diff --git a/app/workers/sync-out.php b/app/workers/sync-out.php index 3458ca5b38..12f8db3aa0 100644 --- a/app/workers/sync-out.php +++ b/app/workers/sync-out.php @@ -13,8 +13,6 @@ use Utopia\Database\DateTime; use Utopia\Database\Document; use Utopia\Database\Exception\Authorization; use Utopia\Database\Exception\Structure; -use Utopia\Logger\Log; -use Utopia\Queue; use Utopia\Queue\Message; global $connection; @@ -107,7 +105,7 @@ $server->job() ->inject('message') ->action(function (Message $message) use (&$stack, &$failures) { - $payload = $message->getPayload()['value'] ?? []; + $payload = $message->getPayload() ?? []; if (!empty($payload['keys'])) { $regions = array_filter( @@ -129,51 +127,6 @@ $server->job() } }); -$server - ->error() - ->inject('error') - ->inject('logger') - ->inject('register') - ->action(function ($error, $logger, $register) { - - // Todo better job of abstracting the error log - $version = App::getEnv('_APP_VERSION', 'UNKNOWN'); - - if ($error instanceof PDOException) { - throw $error; - } - - if ($error->getCode() >= 500 || $error->getCode() === 0) { - $log = new Log(); - - $log->setNamespace("appwrite-worker"); - $log->setServer(\gethostname()); - $log->setVersion($version); - $log->setType(Log::TYPE_ERROR); - $log->setMessage($error->getMessage()); - $log->setAction('appwrite-worker-sync-out'); - $log->addTag('verboseType', get_class($error)); - $log->addTag('code', $error->getCode()); - $log->addExtra('file', $error->getFile()); - $log->addExtra('line', $error->getLine()); - $log->addExtra('trace', $error->getTraceAsString()); - $log->addExtra('detailedTrace', $error->getTrace()); - $log->addExtra('roles', \Utopia\Database\Validator\Authorization::$roles); - - $isProduction = App::getEnv('_APP_ENV', 'development') === 'production'; - $log->setEnvironment($isProduction ? Log::ENVIRONMENT_PRODUCTION : Log::ENVIRONMENT_STAGING); - - $logger->addLog($log); - } - - Console::error('[Error] Type: ' . get_class($error)); - Console::error('[Error] Message: ' . $error->getMessage()); - Console::error('[Error] File: ' . $error->getFile()); - Console::error('[Error] Line: ' . $error->getLine()); - - $register->get('pools')->reclaim(); - }); - $server ->workerStart() ->inject('dbForConsole') @@ -200,7 +153,7 @@ $server $chunk = array_slice($stack['keys'], 0, CHUNK_MAX_KEYS); array_splice($stack['keys'], 0, CHUNK_MAX_KEYS); - Console::info("[{$time}] Sending " . count($chunk) . " remains " . count($stack['keys'])); + Console::log("[{$time}] Sending " . count($chunk) . " remains " . count($stack['keys'])); handle($dbForConsole, $stack['regions'], $chunk); $chunk = []; }); diff --git a/src/Appwrite/Platform/Tasks/EdgeSync.php b/src/Appwrite/Platform/Tasks/EdgeSync.php index 086477d803..db2ba03924 100644 --- a/src/Appwrite/Platform/Tasks/EdgeSync.php +++ b/src/Appwrite/Platform/Tasks/EdgeSync.php @@ -58,10 +58,8 @@ class EdgeSync extends Action Console::info("[{$time}] Enqueueing keys chunk {$count} to {$document->getAttribute('target')}"); $queueForCacheSyncOut ->enqueue([ - 'value' => [ 'region' => $document->getAttribute('target'), 'keys' => $document->getAttribute('keys') - ] ]); $dbForConsole->deleteDocument('syncs', $document->getId());