addressing comments

This commit is contained in:
shimon
2022-11-02 11:39:16 +02:00
parent 6842286a3a
commit e291d81be3
7 changed files with 151 additions and 67 deletions
+1 -3
View File
@@ -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]]);
+17 -5
View File
@@ -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') {
+37 -28
View File
@@ -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);
});
+10 -1
View File
@@ -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;
+39 -13
View File
@@ -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
+45 -17
View File
@@ -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
+2
View File
@@ -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: