adding realtime region sync

This commit is contained in:
shimon
2023-01-02 14:08:30 +02:00
parent f2b6d72511
commit 49f52c8efe
2 changed files with 57 additions and 33 deletions
+24 -7
View File
@@ -2,6 +2,7 @@
require_once __DIR__ . '/../worker.php';
use Appwrite\Messaging\Adapter\Realtime;
use Utopia\App;
use Utopia\Cache\Cache;
use Utopia\CLI\Console;
@@ -12,16 +13,32 @@ $server->job()
->inject('message')
->inject('cache')
->action(function (Message $message, Cache $cache) {
$payload = $message->getPayload();
$type = $payload['type'];
$key = $payload['key'];
$time = DateTime::now();
$cache->setListenersStatus(false);
foreach ($message->getPayload()['keys'] ?? [] as $key) {
Console::log("[{$time}] Purging {$key}");
$cache->purge($key);
switch ($type) {
case 'cache':
$cache->setListenersStatus(false);
Console::log("[{$time}] Purging cache key {$key}");
$cache->purge($key);
$cache->setListenersStatus(true);
break;
case 'realtime':
Console::log("[{$time}] Sending realtime message");
Realtime::send(
projectId: $key['projectId'],
payload: $key['payload'],
events: $key['events'],
channels: $key['channels'],
roles: $key['roles'],
options: $key['options']
);
break;
default:
break;
}
$cache->setListenersStatus(true);
});
$server
+33 -26
View File
@@ -34,12 +34,12 @@ const MAX_CURL_SEND_ATTEMPTS = 4;
/**
* @param string $url
* @param string $token
* @param array $stack
* @param array $payload
* @return array
*/
function call(string $url, string $token, array $stack): array
function call(string $url, string $token, array $payload): array
{
$payload = [];
$ch = curl_init($url);
curl_setopt($ch, CURLOPT_HTTPHEADER, [
'Authorization: Bearer ' . $token,
@@ -48,18 +48,14 @@ function call(string $url, string $token, array $stack): array
curl_setopt($ch, CURLOPT_RETURNTRANSFER, true);
curl_setopt($ch, CURLOPT_TIMEOUT, 5);
curl_setopt($ch, CURLOPT_CUSTOMREQUEST, 'POST');
curl_setopt($ch, CURLOPT_POSTFIELDS, json_encode($stack));
curl_setopt($ch, CURLOPT_POSTFIELDS, json_encode($payload));
for ($attempts = 0; $attempts < MAX_CURL_SEND_ATTEMPTS; $attempts++) {
$response = curl_exec($ch);
$status = curl_getinfo($ch, CURLINFO_HTTP_CODE);
$payload = [
'status' => $status,
'payload' => json_decode($response, true)
];
if ($status === 200) {
return $payload;
return $status;
}
sleep(1);
@@ -67,7 +63,7 @@ function call(string $url, string $token, array $stack): array
curl_close($ch);
return $payload;
return $status;
}
/**
@@ -75,7 +71,7 @@ function call(string $url, string $token, array $stack): array
* @throws Structure
* @throws Exception|\Exception
*/
function handle($dbForConsole, $regions, $stack): void
function handle($dbForConsole, $regions, $payload): void
{
$jwt = new JWT(App::getEnv('_APP_OPENSSL_KEY_V1'), 'HS256', 600, 10);
@@ -83,17 +79,18 @@ function handle($dbForConsole, $regions, $stack): void
foreach ($regions as $code => $region) {
$time = DateTime::now();
$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");
$dbForConsole->createDocument('sync', new Document([
$status = call($region['domain'] . '/v1/edge/sync', $token, ['keys' => $payload]);
if ($status !== Response::STATUS_CODE_OK) {
Console::error("[{$time}] Request to {$code} has failed");
foreach ($payload as $sync) {
$dbForConsole->createDocument('syncs', new Document([
'region' => App::getEnv('_APP_REGION'),
'target' => $code,
'keys' => $stack,
'status' => $response['status'],
'payload' => $response['payload'],
'type' => $sync['type'],
'key' => $sync['key'],
'status' => $status,
]));
}
}
}
}
@@ -103,8 +100,14 @@ $server->job()
->action(function (Message $message) use (&$stack, &$failures) {
$payload = $message->getPayload() ?? [];
$type = $payload['type'] ?? null;
if (!empty($payload['keys'])) {
if (empty($type)) {
return;
}
//Get failed requests
if (!empty($payload['region']) && !empty($payload['keys'])) {
$regions = array_filter(
Config::getParam('regions', []),
fn ($region) => $payload['region'] === $region,
@@ -115,12 +118,15 @@ $server->job()
'regions' => $regions,
'keys' => $payload['keys']
];
return;
}
if (!empty($payload['key'])) {
if (!in_array($payload['key'], $stack['keys'] ?? [])) {
$stack['keys'][] = $payload['key'];
}
$stack['keys'][] = [
'type' => $payload['type'],
'key' => $payload['key'],
];
}
});
@@ -137,6 +143,7 @@ $server
return;
}
//Send failed requests
if (count($failures) > 0) {
$i = 0;
while ($i < count($failures)) {
@@ -147,12 +154,12 @@ $server
}
return;
}
$chunk = array_slice($stack['keys'], 0, CHUNK_MAX_KEYS);
//var_dump($stack['keys']);
$chunk = array_slice($stack['keys'], 0, CHUNK_MAX_KEYS, true);
//var_dump($chunk);
array_splice($stack['keys'], 0, CHUNK_MAX_KEYS);
Console::log("[{$time}] Sending " . count($chunk) . " remains " . count($stack['keys']));
handle($dbForConsole, $stack['regions'], $chunk);
$chunk = [];
});
Console::success("Out [" . App::getEnv('_APP_REGION') . "] edge cache purging worker Started");
});