adopt changes from db-pools

This commit is contained in:
shimon
2022-11-17 18:00:18 +02:00
parent 0abfae8ca1
commit 3f0fcb40d0
14 changed files with 33 additions and 5621 deletions
+1 -1
View File
@@ -322,7 +322,7 @@ RUN mkdir -p /storage/uploads && \
RUN chmod +x /usr/local/bin/doctor && \
chmod +x /usr/local/bin/maintenance && \
chmod +x /usr/local/bin/volume-sync && \
chmod +x /usr/local/bin/sync-edge && \
chmod +x /usr/local/bin/edge-sync && \
chmod +x /usr/local/bin/usage && \
chmod +x /usr/local/bin/install && \
chmod +x /usr/local/bin/migrate && \
+4 -8
View File
@@ -2,15 +2,12 @@
use Ahc\Jwt\JWT;
use Ahc\Jwt\JWTException;
use Appwrite\DSN\DSN;
use Appwrite\Extend\Exception;
use Appwrite\URL\URL as AppwriteURL;
use Appwrite\Utopia\Request;
use Appwrite\Utopia\Response;
use Utopia\App;
use Utopia\Database\Document;
use Utopia\Pools\Group;
use Utopia\Queue\Client as SyncIn;
use Utopia\Queue\Client;
use Utopia\Validator\ArrayList;
use Utopia\Validator\Text;
@@ -37,15 +34,14 @@ App::post('/v1/edge/sync')
->inject('request')
->inject('response')
->inject('pools')
->action(function (array $keys, Request $request, Response $response, Group $pools) {
->action(function (array $keys, Request $request, Response $response, Client $queueForCacheSyncOut) {
if (empty($keys)) {
throw new Exception(Exception::KEY_NOT_FOUND);
}
$client = new SyncIn('syncIn', $pools->get('queue')->pop()->getResource());
$client->enqueue(['value' => ['keys' => $keys]]);
$queueForCacheSyncOut
->enqueue(['value' => ['keys' => $keys]]);
$response->dynamic(new Document([
'keys' => $keys
+12 -7
View File
@@ -40,6 +40,7 @@ use Appwrite\URL\URL as AppwriteURL;
use Appwrite\Usage\Stats;
use Appwrite\Utopia\View;
use Utopia\App;
use Utopia\Queue\Client;
use Utopia\Validator\Range;
use Utopia\Validator\WhiteList;
use Utopia\Database\ID;
@@ -864,6 +865,12 @@ App::setResource('messaging', fn() => new Phone());
App::setResource('queueForFunctions', function (Group $pools) {
return new Func($pools->get('queue')->pop()->getResource());
}, ['pools']);
App::setResource('queueForCacheSyncOut', function (Group $pools) {
return new Client('v1-sync-out', $pools->get('queue')->pop()->getResource());
}, ['pools']);
App::setResource('queueForCacheSyncIn', function (Group $pools) {
return new Client('v1-sync-in', $pools->get('queue')->pop()->getResource());
}, ['pools']);
App::setResource('usage', function ($register) {
return new Stats($register->get('statsd'));
}, ['register']);
@@ -1075,7 +1082,7 @@ App::setResource('dbForConsole', function (Group $pools, Cache $cache) {
return $database;
}, ['pools', 'cache']);
App::setResource('cache', function (Group $pools) {
App::setResource('cache', function (Group $pools, Client $queueForCacheSyncOut) {
$list = Config::getParam('pools-cache', []);
$adapters = [];
@@ -1088,22 +1095,20 @@ App::setResource('cache', function (Group $pools) {
}
$cache = new Cache(new Sharding($adapters));
$client = new SyncOut('syncOut', new QueueRedis(App::getEnv('_APP_REDIS_HOST', 'redis'), App::getEnv('_APP_REDIS_PORT', '6379')));
$cache->on(cache::EVENT_SAVE, function ($key) use ($client) {
$cache->on(cache::EVENT_SAVE, function ($key) use ($queueForCacheSyncOut) {
//Todo fix cache re-invoked
if ($key === 'cache-console:_metadata:users') {
return;
}
$client
$queueForCacheSyncOut
->enqueue(['value' => ['key' => $key]]);
});
$cache->on(cache::EVENT_PURGE, function ($key) use ($client) {
$cache->on(cache::EVENT_PURGE, function ($key) use ($queueForCacheSyncOut) {
if ($key === 'cache-console:_metadata:users') {
return;
}
$client
$queueForCacheSyncOut
->enqueue(['value' => ['key' => $key]]);
});
+4 -2
View File
@@ -10,10 +10,10 @@ use Utopia\Logger\Log;
use Utopia\Queue;
use Utopia\Queue\Message;
global $client;
global $connection;
global $workerNumber;
$adapter = new Queue\Adapter\Swoole($client, $workerNumber, 'syncIn');
$adapter = new Queue\Adapter\Swoole($connection, $workerNumber, 'syncIn');
$server = new Queue\Server($adapter);
$server->job()
@@ -23,10 +23,12 @@ $server->job()
$time = DateTime::now();
$payload = $message->getPayload()['value'];
$cache->setDisableListeners(true);
foreach ($payload['keys'] ?? [] as $key) {
Console::info("[{$time}] Purging {$key}");
$cache->purge($key);
}
$cache->setDisableListeners(false);
});
+1 -4
View File
@@ -17,7 +17,7 @@ use Utopia\Logger\Log;
use Utopia\Queue;
use Utopia\Queue\Message;
global $client;
global $connection;
global $workerNumber;
$regions = array_filter(
@@ -103,9 +103,6 @@ function handle($dbForConsole, $regions, $stack): void
}
}
$adapter = new Queue\Adapter\Swoole($client, $workerNumber, 'syncOut');
$server = new Queue\Server($adapter);
$server->job()
->inject('message')
->action(function (Message $message) use (&$stack, &$failures) {
+3
View File
@@ -0,0 +1,3 @@
#!/bin/sh
php /usr/src/code/app/cli.php edge-sync $@
-3
View File
@@ -1,3 +0,0 @@
#!/bin/sh
php /usr/src/code/app/cli.php sync-edge $@
+1 -1
View File
@@ -1,3 +1,3 @@
#!/bin/sh
php /usr/src/code/app/workers/sync-in.php $@
QUEUE=v1-sync-in php /usr/src/code/app/workers/sync-in.php $@
+1 -1
View File
@@ -1,3 +1,3 @@
#!/bin/sh
php /usr/src/code/app/workers/sync-out.php $@
QUEUE=v1-sync-out php /usr/src/code/app/workers/sync-out.php $@
+2 -4
View File
@@ -46,11 +46,8 @@
"utopia-php/abuse": "0.16.*",
"utopia-php/analytics": "0.2.*",
"utopia-php/audit": "0.17.*",
"utopia-php/cache": "0.8.*",
"utopia-php/cli": "0.14.*",
"utopia-php/audit": "0.15.*",
"utopia-php/cache": "dev-feat-redis-sync as 0.8.1",
"utopia-php/cli": "0.13.*",
"utopia-php/cli": "0.14.*",
"utopia-php/config": "0.2.*",
"utopia-php/database": "0.28.*",
"utopia-php/domains": "1.1.*",
@@ -67,6 +64,7 @@
"utopia-php/storage": "0.11.*",
"utopia-php/swoole": "0.5.*",
"utopia-php/websocket": "0.1.0",
"utopia-php/dsn": "0.1.0",
"resque/php-resque": "1.3.6",
"matomo/device-detector": "6.0.0",
"dragonmantank/cron-expression": "3.3.1",
Generated
-5579
View File
File diff suppressed because it is too large Load Diff
+3 -3
View File
@@ -635,10 +635,10 @@ services:
# - /nfs/config:/data/src
# - /storage/config:/data/dest
appwrite-sync-edge:
entrypoint: sync-edge
appwrite-edge-sync:
entrypoint: edge-sync
<<: *x-logging
container_name: appwrite-sync-edge
container_name: appwrite-edge-sync
image: appwrite-dev
networks:
- appwrite
-3
View File
@@ -2,9 +2,6 @@
namespace Appwrite\Event;
use DateTime;
use Resque;
use ResqueScheduler;
use Utopia\Database\Document;
use Utopia\Queue\Client;
use Utopia\Queue\Connection;
@@ -13,7 +13,7 @@ use Utopia\Queue;
use Utopia\Queue\Client as SyncOut;
$cli
->task('sync-edge')
->task('edge-sync')
->desc('Schedules edge sync tasks')
->action(function () use ($register) {
Console::title('Syncs edges V1');
@@ -23,10 +23,6 @@ $cli
$client = new SyncOut('syncOut', $pools->get('queue')->pop()->getResource());
$database = getConsoleDB();
// Todo fix pdo PDOException
// Table 'appwrite.console__metadata' doesn't exist
sleep(4);
$interval = (int) App::getEnv('_APP_SYNC_EDGE_INTERVAL', '180');
Console::loop(function () use ($interval, $database, $register, $client) {