From e291d81be3d608887d93b0e736ddd23ba5d8d2ad Mon Sep 17 00:00:00 2001 From: shimon Date: Wed, 2 Nov 2022 11:39:16 +0200 Subject: [PATCH] addressing comments --- app/controllers/api/edge.php | 4 +-- app/init.php | 22 +++++++++--- app/tasks/sync-edge.php | 65 ++++++++++++++++++++---------------- app/worker.php | 11 +++++- app/workers/sync-In.php | 52 +++++++++++++++++++++-------- app/workers/sync-out.php | 62 ++++++++++++++++++++++++---------- docker-compose.yml | 2 ++ 7 files changed, 151 insertions(+), 67 deletions(-) diff --git a/app/controllers/api/edge.php b/app/controllers/api/edge.php index 858a39c87d..14bc4fff9b 100644 --- a/app/controllers/api/edge.php +++ b/app/controllers/api/edge.php @@ -8,7 +8,6 @@ use Appwrite\URL\URL as AppwriteURL; use Appwrite\Utopia\Request; use Appwrite\Utopia\Response; use Utopia\App; -use Utopia\CLI\Console; use Utopia\Database\Document; use Utopia\Queue\Client as SyncIn; use Utopia\Queue\Connection\Redis as QueueRedis; @@ -55,13 +54,12 @@ App::post('/v1/edge/sync') $dsns = explode(',', $connection ?? ''); if (empty($dsns)) { - Console::error("No Dsn found"); + throw new Exception(Exception::GENERAL_SERVER_ERROR); } $dsn = explode('=', $dsns[0]); $dsn = $dsn[1] ?? ''; $dsn = new DSN($dsn); - $client = new SyncIn('syncIn', new QueueRedis($dsn->getHost(), $dsn->getPort())); $client->enqueue(['value' => ['keys' => $keys]]); diff --git a/app/init.php b/app/init.php index 4673ccb59f..3c97c1f2d0 100644 --- a/app/init.php +++ b/app/init.php @@ -1028,11 +1028,23 @@ App::setResource('console', function () { ]); }, []); -//App::setResource('queue', function (Group $pools) { -// $pools->get('queue') -// ->pop() -// ->getResource(); -//}, ['pools']); +App::setResource('queue', function () { + + $fallbackForRedis = AppwriteURL::unparse([ + 'scheme' => 'redis', + 'host' => App::getEnv('_APP_REDIS_HOST', 'redis'), + 'port' => App::getEnv('_APP_REDIS_PORT', '6379'), + 'user' => App::getEnv('_APP_REDIS_USER', ''), + 'pass' => App::getEnv('_APP_REDIS_PASS', ''), + ]); + + $connection = App::getEnv('_APP_CONNECTIONS_QUEUE', $fallbackForRedis); + + $dsns = explode(',', $connection ?? ''); + $dsn = explode('=', $dsns[0]); + $dsn = $dsn[1] ?? ''; + return new DSN($dsn); +}, []); App::setResource('dbForProject', function (Group $pools, Database $dbForConsole, Cache $cache, Document $project) { if ($project->isEmpty() || $project->getId() === 'console') { diff --git a/app/tasks/sync-edge.php b/app/tasks/sync-edge.php index 18ad595f02..3310704ffa 100644 --- a/app/tasks/sync-edge.php +++ b/app/tasks/sync-edge.php @@ -9,15 +9,15 @@ use Utopia\App; use Utopia\CLI\Console; use Utopia\Database\DateTime; use Utopia\Database\Query; +use Utopia\Queue; use Utopia\Queue\Client as SyncOut; -use Utopia\Queue\Connection\Redis as QueueRedis; $cli ->task('sync-edge') ->desc('Schedules edge sync tasks') ->action(function () use ($register) { Console::title('Syncs edges V1'); - Console::success(APP_NAME . ' Sync Edge process v1 has started'); + Console::success(APP_NAME . ' Sync failed cache purge process v1 has started'); $fallbackForRedis = AppwriteURL::unparse([ 'scheme' => 'redis', @@ -37,43 +37,52 @@ $cli $dsn = explode('=', $dsns[0]); $dsn = $dsn[1] ?? ''; $dsn = new DSN($dsn); + $redisConnection = new Queue\Connection\Redis($dsn->getHost(), $dsn->getPort()); + $client = new SyncOut('syncOut', $redisConnection); + // Todo fix pdo PDOException - //Table 'appwrite.console__metadata' doesn't exist + // Table 'appwrite.console__metadata' doesn't exist sleep(4); + $interval = (int) App::getEnv('_APP_SYNC_EDGE_INTERVAL', '180'); - Console::loop(function () use ($interval, $register, $dsn) { + Console::loop(function () use ($interval, $register, $client) { + $database = getConsoleDB(); $time = DateTime::now(); - $region = App::getEnv('_APP_REGION', 'default'); - if (App::getEnv('_APP_REGION', 'default') === 'default') { - return; - } + $count = 0; + $chunk = 0; + $limit = 50; + $sum = $limit; - Console::info("[{$time}] Notifying workers with edges tasks every {$interval} seconds"); + Console::success("[{$time}] New task every {$interval} seconds"); - $time = DateTime::now(); - $chunks = $database->find('syncs', [ - Query::equal('region', [$region]), - Query::limit(500) - ]); + while ($sum === $limit) { + $chunk++; - if (count($chunks) > 0) { - $client = new SyncOut('syncOut', new QueueRedis($dsn->getHost(), $dsn->getPort())); - foreach ($chunks as $counter => $chunk) { - Console::info("[{$time}] Sending chunk .$counter. ot of " . count($chunks) . " to {$chunk->getAttribute('target')}"); - $client - ->enqueue([ - 'value' => [ - 'region' => $chunk->getAttribute('target'), - 'keys' => $chunk->getAttribute('keys') - ] - ]); + $results = $database->find('syncs', [ + Query::equal('region', [App::getEnv('_APP_REGION')]), + Query::limit($limit) + ]); - $database->deleteDocument('syncs', $chunk->getId()); + $sum = count($results); + if ($sum > 0) { + foreach ($results as $document) { + Console::info("[{$time}] Enqueueing keys chunk {$count} to {$document->getAttribute('target')}"); + $client + ->enqueue([ + 'value' => [ + 'region' => $document->getAttribute('target'), + 'keys' => $document->getAttribute('keys') + ] + ]); + + $database->deleteDocument('syncs', $document->getId()); + $count++; + } + } else { + Console::info("[{$time}] No cache keys where found."); } - } else { - Console::info("[{$time}] No cache key chunks where found."); } }, $interval); }); diff --git a/app/worker.php b/app/worker.php index d1a68ea278..77f78c3096 100644 --- a/app/worker.php +++ b/app/worker.php @@ -16,6 +16,7 @@ use Utopia\Logger\Log; use Utopia\Logger\Logger; use Utopia\Queue\Server; use Utopia\Registry\Registry; +use Utopia\Queue; global $register; @@ -53,6 +54,10 @@ Server::setResource('cache', function (Registry $register) { return new Cache(new Sharding($adapters)); }, ['register']); +App::setResource('logger', function ($register) { + return $register->get('logger'); +}, ['register']); + Server::setResource('ErrorLog', function (Logger $logger) { return function (Throwable $error, $action) use ($logger) { @@ -114,10 +119,14 @@ $fallbackForRedis = AppwriteURL::unparse([ $connection = App::getEnv('_APP_CONNECTIONS_QUEUE', $fallbackForRedis); $dsns = explode(',', $connection ?? ''); -if (empty($dsns)) { +if (empty($dsns[0])) { Console::error("Dsn not found"); } $dsn = explode('=', $dsns[0]); $dsn = $dsn[1] ?? ''; $dsn = new DSN($dsn); + +$redisConnection = new Queue\Connection\Redis($dsn->getHost(), $dsn->getPort()); +$workerNumber = swoole_cpu_num() * intval(App::getEnv('_APP_WORKER_PER_CORE', 6)); +$workerNumber = 1; diff --git a/app/workers/sync-In.php b/app/workers/sync-In.php index 86f957a0e3..2df457083a 100644 --- a/app/workers/sync-In.php +++ b/app/workers/sync-In.php @@ -2,23 +2,18 @@ require_once __DIR__ . '/../worker.php'; -use Appwrite\Extend\Exception; 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; -if (App::getEnv('_APP_REGION', 'default') === 'default') { - throw new Exception(Exception::GENERAL_SERVER_ERROR); -} +global $redisConnection; +global $workerNumber; -global $register; -global $dsn; - -$connection = new Queue\Connection\Redis($dsn->getHost(), $dsn->getPort()); -$adapter = new Queue\Adapter\Swoole($connection, 2, 'syncIn'); +$adapter = new Queue\Adapter\Swoole($redisConnection, $workerNumber, 'syncIn'); $server = new Queue\Server($adapter); $server->job() @@ -38,10 +33,41 @@ $server->job() $server ->error() ->inject('error') - ->inject('logError') - ->action(function ($error, $logError) { - Console::error($error->getMessage() . ' ' . $error->getFile() . ' ' . $error->getLine()); - call_user_func($logError, $error, 'sync-in-worker'); + ->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("worker"); + $log->setServer(\gethostname()); + $log->setVersion($version); + $log->setType(Log::TYPE_ERROR); + $log->setMessage($error->getMessage()); + $log->setAction('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 diff --git a/app/workers/sync-out.php b/app/workers/sync-out.php index 89fc35d776..fa15114580 100644 --- a/app/workers/sync-out.php +++ b/app/workers/sync-out.php @@ -13,14 +13,12 @@ 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; -if (App::getEnv('_APP_REGION', 'default') === 'default') { - throw new Exception(Exception::GENERAL_SERVER_ERROR); -} - -global $dsn; +global $redisConnection; +global $workerNumber; $regions = array_filter( Config::getParam('regions', []), @@ -85,8 +83,6 @@ function call(string $url, string $token, array $stack): array function handle($dbForConsole, $regions, $stack): void { - global $register; - $jwt = new JWT(App::getEnv('_APP_OPENSSL_KEY_V1'), 'HS256', 600, 10); $token = $jwt->encode([]); @@ -95,7 +91,7 @@ function handle($dbForConsole, $regions, $stack): void $response = call($region['domain'] . '/v1/edge/sync', $token, ['keys' => $stack]); if ($response['status'] !== Response::STATUS_CODE_OK) { Console::error("[{$time}] Request to {$code} has failed"); - try { + $dbForConsole->createDocument('syncs', new Document([ 'region' => App::getEnv('_APP_REGION'), 'target' => $code, @@ -103,16 +99,13 @@ function handle($dbForConsole, $regions, $stack): void 'status' => $response['status'], 'payload' => $response['payload'], ])); - } catch (\Throwable $th) { - $register->get('pools')->reclaim(); - } } } } -$connection = new Queue\Connection\Redis($dsn->getHost(), $dsn->getPort()); -$adapter = new Queue\Adapter\Swoole($connection, 2, 'syncOut'); +$adapter = new Queue\Adapter\Swoole($redisConnection, $workerNumber, 'syncOut'); $server = new Queue\Server($adapter); + $server->job() ->inject('message') ->action(function (Message $message) use (&$stack, &$failures) { @@ -142,10 +135,45 @@ $server->job() $server ->error() ->inject('error') - ->inject('errorLog') - ->action(function ($error, $errorLog) { - Console::error($error->getMessage() . ' ' . $error->getFile() . ' ' . $error->getLine()); - call_user_func($errorLog, $error, 'sync-out-worker'); + ->inject('logger') + ->inject('register') + ->action(function ($error, $logger, $register) { + + $version = App::getEnv('_APP_VERSION', 'UNKNOWN'); + + if ($error instanceof PDOException) { + throw $error; + } + + if ($error->getCode() >= 500 || $error->getCode() === 0) { + $log = new Log(); + + $log->setNamespace("worker"); + $log->setServer(\gethostname()); + $log->setVersion($version); + $log->setType(Log::TYPE_ERROR); + $log->setMessage($error->getMessage()); + $log->setAction('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 diff --git a/docker-compose.yml b/docker-compose.yml index 612b9d7b60..97d85b2413 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -298,6 +298,7 @@ services: - _APP_CONNECTIONS_DB_CONSOLE - _APP_CONNECTIONS_CACHE - _APP_CONNECTIONS_QUEUE + - _APP_WORKER_PER_CORE - _APP_REGION appwrite-worker-sync-in: @@ -322,6 +323,7 @@ services: - _APP_CONNECTIONS_DB_CONSOLE - _APP_CONNECTIONS_CACHE - _APP_CONNECTIONS_QUEUE + - _APP_WORKER_PER_CORE - _APP_REGION appwrite-worker-webhooks: