Merge remote-tracking branch 'origin/multi-region-support' into feat-migration-multi-region

# Conflicts:
#	app/cli.php
#	app/controllers/api/projects.php
#	composer.json
#	composer.lock
This commit is contained in:
Jake Barnby
2024-10-23 18:33:09 +13:00
25 changed files with 249 additions and 105 deletions
+4
View File
@@ -163,15 +163,19 @@ CLI::setResource('getProjectDB', function (Group $pools, Database $dbForConsole,
CLI::setResource('queue', function (Group $pools) {
return $pools->get('queue')->pop()->getResource();
}, ['pools']);
CLI::setResource('queueForFunctions', function (Connection $queue) {
return new Func($queue);
}, ['queue']);
CLI::setResource('queueForDeletes', function (Connection $queue) {
return new Delete($queue);
}, ['queue']);
CLI::setResource('queueForCertificates', function (Connection $queue) {
return new Certificate($queue);
}, ['queue']);
CLI::setResource('logError', function (Registry $register) {
return function (Throwable $error, string $namespace, string $action) use ($register) {
$logger = $register->get('logger');
Submodule
+1
Submodule app/console added at 0959b594b3
+6 -3
View File
@@ -183,7 +183,8 @@ App::post('/v1/functions')
->inject('queueForBuilds')
->inject('dbForConsole')
->inject('gitHub')
->action(function (string $functionId, string $name, string $runtime, array $execute, array $events, string $schedule, int $timeout, bool $enabled, bool $logging, string $entrypoint, string $commands, array $scopes, string $installationId, string $providerRepositoryId, string $providerBranch, bool $providerSilentMode, string $providerRootDirectory, string $templateRepository, string $templateOwner, string $templateRootDirectory, string $templateVersion, string $specification, Request $request, Response $response, Database $dbForProject, Document $project, Document $user, Event $queueForEvents, Build $queueForBuilds, Database $dbForConsole, GitHub $github) use ($redeployVcs) {
->inject('realtimeConnection')
->action(function (string $functionId, string $name, string $runtime, array $execute, array $events, string $schedule, int $timeout, bool $enabled, bool $logging, string $entrypoint, string $commands, array $scopes, string $installationId, string $providerRepositoryId, string $providerBranch, bool $providerSilentMode, string $providerRootDirectory, string $templateRepository, string $templateOwner, string $templateRootDirectory, string $templateVersion, string $specification, Request $request, Response $response, Database $dbForProject, Document $project, Document $user, Event $queueForEvents, Build $queueForBuilds, Database $dbForConsole, GitHub $githubgithub, Callable $realtimeConnection) use ($redeployVcs) {
$functionId = ($functionId == 'unique()') ? ID::unique() : $functionId;
$allowList = \array_filter(\explode(',', System::getEnv('_APP_FUNCTIONS_RUNTIMES', '')));
@@ -249,7 +250,7 @@ App::post('/v1/functions')
$schedule = Authorization::skip(
fn () => $dbForConsole->createDocument('schedules', new Document([
'region' => System::getEnv('_APP_REGION', 'default'), // Todo replace with projects region
'region' => $project->getAttribute('region'),
'resourceType' => 'function',
'resourceId' => $function->getId(),
'resourceInternalId' => $function->getInternalId(),
@@ -374,6 +375,7 @@ App::post('/v1/functions')
project: $project
);
Realtime::send(
redis: $realtimeConnection($queueForEvents->getSourceRegion()),
projectId: 'console',
payload: $rule->getArrayCopy(),
events: $allEvents,
@@ -381,6 +383,7 @@ App::post('/v1/functions')
roles: $target['roles']
);
Realtime::send(
redis: $realtimeConnection($queueForEvents->getSourceRegion()),
projectId: $project->getId(),
payload: $rule->getArrayCopy(),
events: $allEvents,
@@ -1937,7 +1940,7 @@ App::post('/v1/functions/:functionId/executions')
];
$schedule = $dbForConsole->createDocument('schedules', new Document([
'region' => System::getEnv('_APP_REGION', 'default'),
'region' => $project->getAttribute('region'),
'resourceType' => ScheduleExecutions::getSupportedResource(),
'resourceId' => $execution->getId(),
'resourceInternalId' => $execution->getInternalId(),
+5 -5
View File
@@ -2707,7 +2707,7 @@ App::post('/v1/messaging/messages/email')
break;
case MessageStatus::SCHEDULED:
$schedule = $dbForConsole->createDocument('schedules', new Document([
'region' => System::getEnv('_APP_REGION', 'default'),
'region' => $project->getAttribute('region'),
'resourceType' => 'message',
'resourceId' => $message->getId(),
'resourceInternalId' => $message->getInternalId(),
@@ -2823,7 +2823,7 @@ App::post('/v1/messaging/messages/sms')
break;
case MessageStatus::SCHEDULED:
$schedule = $dbForConsole->createDocument('schedules', new Document([
'region' => System::getEnv('_APP_REGION', 'default'),
'region' => $project->getAttribute('region'),
'resourceType' => 'message',
'resourceId' => $message->getId(),
'resourceInternalId' => $message->getInternalId(),
@@ -2999,7 +2999,7 @@ App::post('/v1/messaging/messages/push')
break;
case MessageStatus::SCHEDULED:
$schedule = $dbForConsole->createDocument('schedules', new Document([
'region' => System::getEnv('_APP_REGION', 'default'),
'region' => $project->getAttribute('region'),
'resourceType' => 'message',
'resourceId' => $message->getId(),
'resourceInternalId' => $message->getInternalId(),
@@ -3548,7 +3548,7 @@ App::patch('/v1/messaging/messages/sms/:messageId')
if (\is_null($currentScheduledAt) && !\is_null($scheduledAt)) {
$schedule = $dbForConsole->createDocument('schedules', new Document([
'region' => System::getEnv('_APP_REGION', 'default'),
'region' => $project->getAttribute('region'),
'resourceType' => 'message',
'resourceId' => $message->getId(),
'resourceInternalId' => $message->getInternalId(),
@@ -3712,7 +3712,7 @@ App::patch('/v1/messaging/messages/push/:messageId')
if (\is_null($currentScheduledAt) && !\is_null($scheduledAt)) {
$schedule = $dbForConsole->createDocument('schedules', new Document([
'region' => System::getEnv('_APP_REGION', 'default'),
'region' => $project->getAttribute('region'),
'resourceType' => 'message',
'resourceId' => $message->getId(),
'resourceInternalId' => $message->getInternalId(),
+16 -1
View File
@@ -125,11 +125,20 @@ App::post('/v1/projects')
$databases = Config::getParam('pools-database', []);
$databases = Config::getParam('pools-database', []);
var_dump($databases);
$databaseOverride = System::getEnv('_APP_DATABASE_OVERRIDE');
$index = \array_search($databaseOverride, $databases);
if ($index !== false) {
$dsn = $databases[$index];
} else {
if ($region !== 'default') {
$databases = array_filter($databases, function ($value) use ($region) {
return str_contains($value, $region);
});
}
$dsn = $databases[array_rand($databases)];
}
@@ -192,6 +201,9 @@ App::post('/v1/projects')
$dsn = new DSN('mysql://' . $dsn);
}
var_dump($dsn->getHost());
$adapter = $pools->get($dsn->getHost())->pop()->getResource();
$dbForProject = new Database($adapter, $cache);
$sharedTables = \explode(',', System::getEnv('_APP_DATABASE_SHARED_TABLES', ''));
@@ -263,6 +275,9 @@ App::post('/v1/projects')
continue;
}
$indexes = \array_map(function (array $index) {
return new Document($index);
}, $collection['indexes']);
$attributes = \array_map(fn ($attribute) => new Document($attribute), $collection['attributes']);
$indexes = \array_map(fn (array $index) => new Document($index), $collection['indexes']);
@@ -279,7 +294,7 @@ App::post('/v1/projects')
]));
}
}
}
// }
// Hook allowing instant project mirroring during migration
// Outside of migration, hook is not registered and has no effect
+9 -2
View File
@@ -59,6 +59,12 @@ function router(App $utopia, Database $dbForConsole, callable $getProjectDB, Swo
])
)[0] ?? null;
// var_dump(System::getEnv('_APP_DOMAIN_FUNCTIONS', ''));
// var_dump($host);
// var_dump(APP_HOSTNAME_INTERNAL);
// var_dump($request->getHeader('host'));
// var_dump($request->getHeader('x-forwarded-host'));
if ($route === null) {
if ($host === System::getEnv('_APP_DOMAIN_FUNCTIONS', '')) {
throw new AppwriteException(AppwriteException::GENERAL_ACCESS_FORBIDDEN, 'This domain cannot be used for security reasons. Please use any subdomain instead.');
@@ -461,7 +467,8 @@ App::init()
/*
* Appwrite Router
*/
$host = $request->getHostname() ?? '';
$host = $request->getHostname() ?? '';
$mainDomain = System::getEnv('_APP_DOMAIN', '');
// Only run Router when external domain
if ($host !== $mainDomain) {
@@ -662,6 +669,7 @@ App::init()
) {
throw new AppwriteException(AppwriteException::GENERAL_UNKNOWN_ORIGIN, $originValidator->getDescription());
}
});
App::options()
@@ -1014,7 +1022,6 @@ App::get('/.well-known/acme-challenge/*')
->action(function (Request $request, Response $response) {
$uriChunks = \explode('/', $request->getURI());
$token = $uriChunks[\count($uriChunks) - 1];
$validator = new Text(100, allowList: [
...Text::NUMBERS,
...Text::ALPHABET_LOWER,
+18 -10
View File
@@ -162,6 +162,7 @@ App::init()
->inject('mode')
->inject('team')
->action(function (App $utopia, Request $request, Database $dbForConsole, Document $project, Document $user, ?Document $session, array $servers, string $mode, Document $team) {
$route = $utopia->getRoute();
if ($project->isEmpty()) {
@@ -363,7 +364,6 @@ App::init()
->inject('dbForProject')
->inject('mode')
->action(function (App $utopia, Request $request, Response $response, Document $project, Document $user, Event $queueForEvents, Messaging $queueForMessaging, Audit $queueForAudits, Delete $queueForDeletes, EventDatabase $queueForDatabase, Build $queueForBuilds, Usage $queueForUsage, Database $dbForProject, string $mode) use ($databaseListener) {
$route = $utopia->getRoute();
if (
@@ -461,6 +461,7 @@ App::init()
->on(Database::EVENT_DOCUMENT_DELETE, 'calculate-usage', fn ($event, $document) => $databaseListener($event, $document, $project, $queueForUsage, $dbForProject));
$useCache = $route->getLabel('cache', false);
if ($useCache) {
$key = md5($request->getURI() . '*' . implode('*', $request->getParams()) . '*' . APP_CACHE_BUSTER);
$cacheLog = Authorization::skip(fn () => $dbForProject->getDocument('cache', $key));
@@ -469,8 +470,8 @@ App::init()
);
$timestamp = 60 * 60 * 24 * 30;
$data = $cache->load($key, $timestamp);
if (!empty($data) && !$cacheLog->isEmpty()) {
$timerStart = \microtime(true);
$parts = explode('/', $cacheLog->getAttribute('resourceType'));
$type = $parts[0] ?? null;
@@ -551,6 +552,7 @@ App::shutdown()
->inject('project')
->inject('dbForProject')
->action(function (App $utopia, Request $request, Response $response, Document $project, Database $dbForProject) {
$route = $utopia->getRoute();
$sessionLimit = $project->getAttribute('auths', [])['maxSessions'] ?? APP_LIMIT_USER_SESSIONS_DEFAULT;
$session = $response->getPayload();
$userId = $session['userId'] ?? '';
@@ -595,8 +597,9 @@ App::shutdown()
->inject('queueForFunctions')
->inject('mode')
->inject('dbForConsole')
->action(function (App $utopia, Request $request, Response $response, Document $project, Document $user, Event $queueForEvents, Audit $queueForAudits, Usage $queueForUsage, Delete $queueForDeletes, EventDatabase $queueForDatabase, Build $queueForBuilds, Messaging $queueForMessaging, Database $dbForProject, Func $queueForFunctions, string $mode, Database $dbForConsole) use ($parseLabel) {
->inject('realtimeConnection')
->action(function (App $utopia, Request $request, Response $response, Document $project, Document $user, Event $queueForEvents, Audit $queueForAudits, Usage $queueForUsage, Delete $queueForDeletes, EventDatabase $queueForDatabase, Build $queueForBuilds, Messaging $queueForMessaging, Database $dbForProject, Func $queueForFunctions, string $mode, Database $dbForConsole, callable $realtimeConnection) use ($parseLabel) {
$route = $utopia->getRoute();
$responsePayload = $response->getPayload();
if (!empty($queueForEvents->getEvent())) {
@@ -642,6 +645,7 @@ App::shutdown()
);
Realtime::send(
redis: $realtimeConnection($queueForEvents->getSourceRegion()),
projectId: $target['projectId'] ?? $project->getId(),
payload: $queueForEvents->getRealtimePayload(),
events: $allEvents,
@@ -709,9 +713,11 @@ App::shutdown()
* Cache label
*/
$useCache = $route->getLabel('cache', false);
if ($useCache) {
$resource = $resourceType = null;
$data = $response->getPayload();
if (!empty($data['payload'])) {
$pattern = $route->getLabel('cache.resource', null);
if (!empty($pattern)) {
@@ -728,6 +734,7 @@ App::shutdown()
$cacheLog = Authorization::skip(fn () => $dbForProject->getDocument('cache', $key));
$accessedAt = $cacheLog->getAttribute('accessedAt', '');
$now = DateTime::now();
if ($cacheLog->isEmpty()) {
Authorization::skip(fn () => $dbForProject->createDocument('cache', new Document([
'$id' => $key,
@@ -742,17 +749,18 @@ App::shutdown()
Authorization::skip(fn () => $dbForProject->updateDocument('cache', $cacheLog->getId(), $cacheLog));
}
if ($signature !== $cacheLog->getAttribute('signature')) {
$cache = new Cache(
new Filesystem(APP_STORAGE_CACHE . DIRECTORY_SEPARATOR . 'app-' . $project->getId())
);
$cache = new Cache(
new Filesystem(APP_STORAGE_CACHE . DIRECTORY_SEPARATOR . 'app-' . $project->getId())
);
$timestamp = 60 * 60 * 24 * 30;
$cacheFile = $cache->load($key, $timestamp);
if ($signature !== $cacheLog->getAttribute('signature') || empty($cacheFile)) {
$cache->save($key, $data['payload']);
}
}
}
if ($project->getId() !== 'console') {
if (!Auth::isPrivilegedUser(Authorization::getRoles())) {
$fileSize = 0;
+23 -13
View File
@@ -1519,23 +1519,33 @@ App::setResource('deviceForLocal', function () {
return new Local();
});
App::setResource('deviceForFiles', function ($project) {
return getDevice(APP_STORAGE_UPLOADS . '/app-' . $project->getId());
}, ['project']);
App::setResource('deviceForFiles', function ($project, $connectionString) {
return getDevice(APP_STORAGE_UPLOADS.'/app-'.$project->getId(), $connectionString);
}, ['project', 'connectionString']);
App::setResource('deviceForFunctions', function ($project) {
return getDevice(APP_STORAGE_FUNCTIONS . '/app-' . $project->getId());
}, ['project']);
App::setResource('deviceForFunctions', function ($project, $connectionString) {
return getDevice(APP_STORAGE_FUNCTIONS.'/app-'.$project->getId(), $connectionString);
}, ['project', 'connectionString']);
App::setResource('deviceForBuilds', function ($project) {
return getDevice(APP_STORAGE_BUILDS . '/app-' . $project->getId());
}, ['project']);
App::setResource('deviceForBuilds', function ($project, $connectionString) {
return getDevice(APP_STORAGE_BUILDS.'/app-'.$project->getId(), $connectionString);
}, ['project', 'connectionString']);
function getDevice(string $root, string $connection = ''): Device
App::setResource('connectionString', function () {
return System::getEnv('_APP_CONNECTIONS_STORAGE', '');
});
App::setResource('realtimeConnection',function ($pools) {
return function () use ($pools) {
return $pools->get('pubsub')->pop()->getResource();
};
}, ['pools']);
function getDevice(string $root, string $connectionString = ''): Device
{
$connection = !empty($connection) ? $connection : System::getEnv('_APP_CONNECTIONS_STORAGE', '');
if (!empty($connection)) {
if (! empty($connectionString)) {
$acl = 'private';
$device = Storage::DEVICE_LOCAL;
$accessKey = '';
@@ -1544,7 +1554,7 @@ function getDevice(string $root, string $connection = ''): Device
$region = '';
try {
$dsn = new DSN($connection);
$dsn = new DSN($connectionString);
$device = $dsn->getScheme();
$accessKey = $dsn->getUser() ?? '';
$accessSecret = $dsn->getPassword() ?? '';
+8 -5
View File
@@ -137,14 +137,17 @@ if (!function_exists('getCache')) {
}
}
if (!function_exists('getRealtime')) {
function getRealtime(): Realtime
if (!function_exists("getPubSub")) {
function getPubSub(): \Redis
{
return new Realtime();
global $register;
$pools = $register->get('pools'); /** @var \Utopia\Pools\Group $pools */
return $pools->get('pubsub')->pop()->getResource();
}
}
$realtime = getRealtime();
$realtime = new Realtime();
/**
* Table for statistics across all workers.
@@ -367,7 +370,7 @@ $server->onWorkerStart(function (int $workerId) use ($server, $register, $stats,
}
$start = time();
$redis = $register->get('pools')->get('pubsub')->pop()->getResource(); /** @var Redis $redis */
$redis = getPubSub(); /** @var \Redis $redis */
$redis->setOption(Redis::OPT_READ_TIMEOUT, -1);
if ($redis->ping(true)) {
+15 -10
View File
@@ -262,22 +262,27 @@ Server::setResource('pools', function (Registry $register) {
return $register->get('pools');
}, ['register']);
Server::setResource('deviceForFunctions', function (Document $project) {
return getDevice(APP_STORAGE_FUNCTIONS . '/app-' . $project->getId());
}, ['project']);
Server::setResource('deviceForFunctions', function (Document $project, $connectionString) {
return getDevice(APP_STORAGE_FUNCTIONS.'/app-'.$project->getId(), $connectionString);
}, ['project', 'connectionString']);
Server::setResource('deviceForFiles', function (Document $project) {
return getDevice(APP_STORAGE_UPLOADS . '/app-' . $project->getId());
}, ['project']);
Server::setResource('deviceForFiles', function (Document $project, $connectionString) {
return getDevice(APP_STORAGE_UPLOADS.'/app-'.$project->getId(), $connectionString);
}, ['project', 'connectionString']);
Server::setResource('deviceForBuilds', function (Document $project) {
return getDevice(APP_STORAGE_BUILDS . '/app-' . $project->getId());
}, ['project']);
Server::setResource('deviceForBuilds', function (Document $project, $connectionString) {
return getDevice(APP_STORAGE_BUILDS.'/app-'.$project->getId(), $connectionString);
}, ['project', 'connectionString']);
Server::setResource('deviceForCache', function (Document $project) {
return getDevice(APP_STORAGE_CACHE . '/app-' . $project->getId());
return getDevice(APP_STORAGE_CACHE.'/app-'.$project->getId());
}, ['project']);
Server::setResource('realtimeConnection',function ($pools) {
return function () use ($pools) {
return $pools->get('pubsub')->pop()->getResource();
};
}, ['pools']);
$pools = $register->get('pools');
$platform = new Appwrite();
Generated
+1 -1
View File
@@ -7040,5 +7040,5 @@
"platform-overrides": {
"php": "8.3"
},
"plugin-api-version": "2.6.0"
"plugin-api-version": "2.2.0"
}
+1
View File
@@ -115,6 +115,7 @@ class Build extends Event
$client = new Client($this->queue, $this->connection);
return $client->enqueue([
'sourceRegion' => $this->getSourceRegion(),
'project' => $this->project,
'resource' => $this->resource,
'deployment' => $this->deployment,
+1
View File
@@ -77,6 +77,7 @@ class Certificate extends Event
$client = new Client($this->queue, $this->connection);
return $client->enqueue([
'sourceRegion' => $this->getSourceRegion(),
'project' => $this->project,
'domain' => $this->domain,
'skipRenewCheck' => $this->skipRenewCheck
+2
View File
@@ -108,6 +108,7 @@ class Database extends Event
*/
public function trigger(): string|bool
{
try {
$dsn = new DSN($this->getProject()->getAttribute('database'));
} catch (\InvalidArgumentException) {
@@ -121,6 +122,7 @@ class Database extends Event
try {
$result = $client->enqueue([
'sourceRegion' => $this->getSourceRegion(),
'project' => $this->project,
'user' => $this->user,
'type' => $this->type,
+9
View File
@@ -6,6 +6,7 @@ use InvalidArgumentException;
use Utopia\Database\Document;
use Utopia\Queue\Client;
use Utopia\Queue\Connection;
use Utopia\System\System;
class Event
{
@@ -110,6 +111,13 @@ class Event
return $this->event;
}
public function getSourceRegion(): string
{
return System::getEnv('_APP_REGION', 'default');
}
/**
* Set project for this event.
*
@@ -322,6 +330,7 @@ class Event
$client = new Client($this->queue, $this->connection);
return $client->enqueue([
'sourceRegion' => $this->getSourceRegion(),
'project' => $this->project,
'user' => $this->user,
'userId' => $this->userId,
+1
View File
@@ -222,6 +222,7 @@ class Func extends Event
$events = $this->getEvent() ? Event::generateEvents($this->getEvent(), $this->getParams()) : null;
return $client->enqueue([
'sourceRegion' => $this->getSourceRegion(),
'project' => $this->project,
'user' => $this->user,
'userId' => $this->userId,
+1
View File
@@ -79,6 +79,7 @@ class Migration extends Event
$client = new Client($this->queue, $this->connection);
return $client->enqueue([
'sourceRegion' => $this->getSourceRegion(),
'project' => $this->project,
'user' => $this->user,
'migration' => $this->migration,
+3 -1
View File
@@ -5,6 +5,8 @@ namespace Appwrite\Messaging;
abstract class Adapter
{
abstract public function subscribe(string $projectId, mixed $identifier, array $roles, array $channels): void;
abstract public function unsubscribe(mixed $identifier): void;
abstract public static function send(string $projectId, array $payload, array $events, array $channels, array $roles, array $options): void;
abstract public static function send(\redis $redis, string $projectId, array $payload, array $events, array $channels, array $roles, array $options): void;
}
+4 -4
View File
@@ -122,15 +122,17 @@ class Realtime extends Adapter
/**
* Sends an event to the Realtime Server
* @param \Redis $redis
* @param string $projectId
* @param array $payload
* @param string $event
* @param array $events
* @param array $channels
* @param array $roles
* @param array $options
* @return void
* @throws \RedisException
*/
public static function send(string $projectId, array $payload, array $events, array $channels, array $roles, array $options = []): void
public static function send(\Redis $redis, string $projectId, array $payload, array $events, array $channels, array $roles, array $options = []): void
{
if (empty($channels) || empty($roles) || empty($projectId)) {
return;
@@ -139,8 +141,6 @@ class Realtime extends Adapter
$permissionsChanged = array_key_exists('permissionsChanged', $options) && $options['permissionsChanged'];
$userId = array_key_exists('userId', $options) ? $options['userId'] : null;
$redis = new \Redis(); //TODO: make this part of the constructor
$redis->connect(System::getEnv('_APP_REDIS_HOST', ''), System::getEnv('_APP_REDIS_PORT', ''));
$redis->publish('realtime', json_encode([
'project' => $projectId,
'roles' => $roles,
+1 -1
View File
@@ -731,7 +731,7 @@ class V19 extends Migration
if (empty($document->getAttribute('scheduleId', null))) {
$schedule = $this->consoleDB->createDocument('schedules', new Document([
'region' => System::getEnv('_APP_REGION', 'default'), // Todo replace with projects region
'region' => $project->getAttribute('region'),
'resourceType' => 'function',
'resourceId' => $document->getId(),
'resourceInternalId' => $document->getInternalId(),
+18 -6
View File
@@ -33,6 +33,11 @@ use Utopia\VCS\Adapter\Git\GitHub;
class Builds extends Action
{
/**
* @var mixed|string
*/
protected string $sourceRegion;
public static function getName(): string
{
return 'builds';
@@ -54,7 +59,8 @@ class Builds extends Action
->inject('dbForProject')
->inject('deviceForFunctions')
->inject('log')
->callback(fn ($message, Database $dbForConsole, Event $queueForEvents, Func $queueForFunctions, Usage $usage, Cache $cache, Database $dbForProject, Device $deviceForFunctions, Log $log) => $this->action($message, $dbForConsole, $queueForEvents, $queueForFunctions, $usage, $cache, $dbForProject, $deviceForFunctions, $log));
->inject('realtimeConnection')
->callback(fn ($message, Database $dbForConsole, Event $queueForEvents, Func $queueForFunctions, Usage $usage, Cache $cache, Database $dbForProject, Device $deviceForFunctions, Log $log, Callable $realtimeConnection) => $this->action($message, $dbForConsole, $queueForEvents, $queueForFunctions, $usage, $cache, $dbForProject, $deviceForFunctions, $log, $realtimeConnection));
}
/**
@@ -67,10 +73,11 @@ class Builds extends Action
* @param Database $dbForProject
* @param Device $deviceForFunctions
* @param Log $log
* @param callable $realtimeConnection
* @return void
* @throws \Utopia\Database\Exception
*/
public function action(Message $message, Database $dbForConsole, Event $queueForEvents, Func $queueForFunctions, Usage $queueForUsage, Cache $cache, Database $dbForProject, Device $deviceForFunctions, Log $log): void
public function action(Message $message, Database $dbForConsole, Event $queueForEvents, Func $queueForFunctions, Usage $queueForUsage, Cache $cache, Database $dbForProject, Device $deviceForFunctions, Log $log, Callable $realtimeConnection): void
{
$payload = $message->getPayload() ?? [];
@@ -83,6 +90,7 @@ class Builds extends Action
$resource = new Document($payload['resource'] ?? []);
$deployment = new Document($payload['deployment'] ?? []);
$template = new Document($payload['template'] ?? []);
$this->sourceRegion = $payload['sourceRegion'] ?? 'default';
$log->addTag('projectId', $project->getId());
$log->addTag('type', $type);
@@ -92,7 +100,7 @@ class Builds extends Action
case BUILD_TYPE_RETRY:
Console::info('Creating build for deployment: ' . $deployment->getId());
$github = new GitHub($cache);
$this->buildDeployment($deviceForFunctions, $queueForFunctions, $queueForEvents, $queueForUsage, $dbForConsole, $dbForProject, $github, $project, $resource, $deployment, $template, $log);
$this->buildDeployment($deviceForFunctions, $queueForFunctions, $queueForEvents, $queueForUsage, $dbForConsole, $dbForProject, $github, $project, $resource, $deployment, $template, $log, $realtimeConnection);
break;
default:
@@ -117,7 +125,7 @@ class Builds extends Action
* @throws \Utopia\Database\Exception
* @throws Exception
*/
protected function buildDeployment(Device $deviceForFunctions, Func $queueForFunctions, Event $queueForEvents, Usage $queueForUsage, Database $dbForConsole, Database $dbForProject, GitHub $github, Document $project, Document $function, Document $deployment, Document $template, Log $log): void
protected function buildDeployment(Device $deviceForFunctions, Func $queueForFunctions, Event $queueForEvents, Usage $queueForUsage, Database $dbForConsole, Database $dbForProject, GitHub $github, Document $project, Document $function, Document $deployment, Document $template, Log $log, Callable $realtimeConnection): void
{
$executor = new Executor(System::getEnv('_APP_EXECUTOR_HOST'));
@@ -376,6 +384,7 @@ class Builds extends Action
project: $project
);
Realtime::send(
redis: $realtimeConnection($this->sourceRegion),
projectId: 'console',
payload: $build->getArrayCopy(),
events: $allEvents,
@@ -454,6 +463,7 @@ class Builds extends Action
);
Realtime::send(
redis: $realtimeConnection($this->sourceRegion),
projectId: 'console',
payload: $build->getArrayCopy(),
events: $allEvents,
@@ -552,12 +562,12 @@ class Builds extends Action
$err = $error;
}
}),
Co\go(function () use ($executor, $project, $deployment, &$response, &$build, $dbForProject, $allEvents, &$err, &$isCanceled) {
Co\go(function () use ($realtimeConnection, $executor, $project, $deployment, &$response, &$build, $dbForProject, $allEvents, &$err, &$isCanceled) {
try {
$executor->getLogs(
deploymentId: $deployment->getId(),
projectId: $project->getId(),
callback: function ($logs) use (&$response, &$err, &$build, $dbForProject, $allEvents, $project, &$isCanceled) {
callback: function ($logs) use ($realtimeConnection, &$response, &$err, &$build, $dbForProject, $allEvents, $project, &$isCanceled) {
if ($isCanceled) {
return;
}
@@ -591,6 +601,7 @@ class Builds extends Action
project: $project
);
Realtime::send(
redis: $realtimeConnection($this->sourceRegion),
projectId: 'console',
payload: $build->getArrayCopy(),
events: $allEvents,
@@ -693,6 +704,7 @@ class Builds extends Action
project: $project
);
Realtime::send(
redis: $realtimeConnection($this->sourceRegion),
projectId: 'console',
payload: $build->getArrayCopy(),
events: $allEvents,
+31 -17
View File
@@ -30,6 +30,11 @@ use Utopia\System\System;
class Certificates extends Action
{
/**
* @var mixed|string
*/
protected string $sourceRegion;
public static function getName(): string
{
return 'certificates';
@@ -48,7 +53,8 @@ class Certificates extends Action
->inject('queueForEvents')
->inject('queueForFunctions')
->inject('log')
->callback(fn (Message $message, Database $dbForConsole, Mail $queueForMails, Event $queueForEvents, Func $queueForFunctions, Log $log) => $this->action($message, $dbForConsole, $queueForMails, $queueForEvents, $queueForFunctions, $log));
->inject('realtimeConnection')
->callback(fn (Message $message, Database $dbForConsole, Mail $queueForMails, Event $queueForEvents, Func $queueForFunctions, Log $log, Callable $realtimeConnection) => $this->action($message, $dbForConsole, $queueForMails, $queueForEvents, $queueForFunctions, $log, $realtimeConnection));
}
/**
@@ -58,11 +64,12 @@ class Certificates extends Action
* @param Event $queueForEvents
* @param Func $queueForFunctions
* @param Log $log
* @param callable $realtimeConnection
* @return void
* @throws Throwable
* @throws \Utopia\Database\Exception
*/
public function action(Message $message, Database $dbForConsole, Mail $queueForMails, Event $queueForEvents, Func $queueForFunctions, Log $log): void
public function action(Message $message, Database $dbForConsole, Mail $queueForMails, Event $queueForEvents, Func $queueForFunctions, Log $log, Callable $realtimeConnection): void
{
$payload = $message->getPayload() ?? [];
@@ -73,10 +80,10 @@ class Certificates extends Action
$document = new Document($payload['domain'] ?? []);
$domain = new Domain($document->getAttribute('domain', ''));
$skipRenewCheck = $payload['skipRenewCheck'] ?? false;
$this->sourceRegion = $payload['sourceRegion'] ?? 'default';
$log->addTag('domain', $domain->get());
$this->execute($domain, $dbForConsole, $queueForMails, $queueForEvents, $queueForFunctions, $log, $skipRenewCheck);
$this->execute($domain, $dbForConsole, $queueForMails, $queueForEvents, $queueForFunctions, $log,$realtimeConnection, $skipRenewCheck);
}
/**
@@ -85,12 +92,17 @@ class Certificates extends Action
* @param Mail $queueForMails
* @param Event $queueForEvents
* @param Func $queueForFunctions
* @param Log $log
* @param callable $realtimeConnection
* @param bool $skipRenewCheck
* @return void
* @throws Authorization
* @throws Conflict
* @throws Structure
* @throws Throwable
* @throws \Utopia\Database\Exception
*/
private function execute(Domain $domain, Database $dbForConsole, Mail $queueForMails, Event $queueForEvents, Func $queueForFunctions, Log $log, bool $skipRenewCheck = false): void
protected function execute(Domain $domain, Database $dbForConsole, Mail $queueForMails, Event $queueForEvents, Func $queueForFunctions, Log $log, Callable $realtimeConnection, bool $skipRenewCheck = false): void
{
/**
* 1. Read arguments and validate domain
@@ -162,7 +174,7 @@ class Certificates extends Action
$logs = 'Certificate successfully generated.';
$certificate->setAttribute('logs', \mb_strcut($logs, 0, 1000000));// Limit to 1MB
var_dump($certificate);
// Give certificates to Traefik
$this->applyCertificateFiles($folder, $domain->get(), $letsEncryptData);
@@ -193,7 +205,7 @@ class Certificates extends Action
$certificate->setAttribute('updated', DateTime::now());
// Save all changes we made to certificate document into database
$this->saveCertificateDocument($domain->get(), $certificate, $success, $dbForConsole, $queueForEvents, $queueForFunctions);
$this->saveCertificateDocument($domain->get(), $certificate, $success, $dbForConsole, $queueForEvents, $queueForFunctions, $realtimeConnection);
}
}
@@ -212,7 +224,7 @@ class Certificates extends Action
* @throws Conflict
* @throws Structure
*/
private function saveCertificateDocument(string $domain, Document $certificate, bool $success, Database $dbForConsole, Event $queueForEvents, Func $queueForFunctions): void
protected function saveCertificateDocument(string $domain, Document $certificate, bool $success, Database $dbForConsole, Event $queueForEvents, Func $queueForFunctions, Callable $realtimeConnection): void
{
// Check if update or insert required
$certificateDocument = $dbForConsole->findOne('certificates', [Query::equal('domain', [$domain])]);
@@ -226,7 +238,7 @@ class Certificates extends Action
}
$certificateId = $certificate->getId();
$this->updateDomainDocuments($certificateId, $domain, $success, $dbForConsole, $queueForEvents, $queueForFunctions);
$this->updateDomainDocuments($certificateId, $domain, $success, $dbForConsole, $queueForEvents, $queueForFunctions, $realtimeConnection);
}
/**
@@ -234,7 +246,7 @@ class Certificates extends Action
*
* @return null|string Returns main domain. If null, there is no main domain yet.
*/
private function getMainDomain(): ?string
protected function getMainDomain(): ?string
{
$envDomain = System::getEnv('_APP_DOMAIN', '');
if (!empty($envDomain) && $envDomain !== 'localhost') {
@@ -255,7 +267,7 @@ class Certificates extends Action
* @return void
* @throws Exception
*/
private function validateDomain(Domain $domain, bool $isMainDomain, Log $log): void
protected function validateDomain(Domain $domain, bool $isMainDomain, Log $log): void
{
if (empty($domain->get())) {
throw new Exception('Missing certificate domain.');
@@ -299,7 +311,7 @@ class Certificates extends Action
* @return bool True, if certificate needs to be renewed
* @throws Exception
*/
private function isRenewRequired(string $domain, Log $log): bool
protected function isRenewRequired(string $domain, Log $log): bool
{
$certPath = APP_STORAGE_CERTIFICATES . '/' . $domain . '/cert.pem';
if (\file_exists($certPath)) {
@@ -333,7 +345,7 @@ class Certificates extends Action
* @return array Named array with keys 'stdout' and 'stderr', both string
* @throws Exception
*/
private function issueCertificate(string $folder, string $domain, string $email): array
protected function issueCertificate(string $folder, string $domain, string $email): array
{
$stdout = '';
$stderr = '';
@@ -363,7 +375,7 @@ class Certificates extends Action
* @return string
* @throws \Utopia\Database\Exception
*/
private function getRenewDate(string $domain): string
protected function getRenewDate(string $domain): string
{
$certPath = APP_STORAGE_CERTIFICATES . '/' . $domain . '/cert.pem';
$certData = openssl_x509_parse(file_get_contents($certPath));
@@ -381,7 +393,7 @@ class Certificates extends Action
* @return void
* @throws Exception
*/
private function applyCertificateFiles(string $folder, string $domain, array $letsEncryptData): void
protected function applyCertificateFiles(string $folder, string $domain, array $letsEncryptData): void
{
// Prepare folder in storage for domain
@@ -432,7 +444,7 @@ class Certificates extends Action
* @return void
* @throws Exception
*/
private function notifyError(string $domain, string $errorMessage, int $attempt, Mail $queueForMails): void
protected function notifyError(string $domain, string $errorMessage, int $attempt, Mail $queueForMails): void
{
// Log error into console
Console::warning('Cannot renew domain (' . $domain . ') on attempt no. ' . $attempt . ' certificate: ' . $errorMessage);
@@ -475,7 +487,7 @@ class Certificates extends Action
*
* @return void
*/
private function updateDomainDocuments(string $certificateId, string $domain, bool $success, Database $dbForConsole, Event $queueForEvents, Func $queueForFunctions): void
protected function updateDomainDocuments(string $certificateId, string $domain, bool $success, Database $dbForConsole, Event $queueForEvents, Func $queueForFunctions, Callable $realtimeConnection): void
{
$rule = $dbForConsole->findOne('rules', [
@@ -525,6 +537,7 @@ class Certificates extends Action
project: $project
);
Realtime::send(
redis: $realtimeConnection($this->sourceRegion),
projectId: 'console',
payload: $rule->getArrayCopy(),
events: $allEvents,
@@ -532,6 +545,7 @@ class Certificates extends Action
roles: $target['roles']
);
Realtime::send(
redis: $realtimeConnection($this->sourceRegion),
projectId: $project->getId(),
payload: $rule->getArrayCopy(),
events: $allEvents,
+32 -14
View File
@@ -20,6 +20,11 @@ use Utopia\Queue\Message;
class Databases extends Action
{
/**
* @var array|mixed
*/
protected string $sourceRegion;
public static function getName(): string
{
return 'databases';
@@ -36,7 +41,8 @@ class Databases extends Action
->inject('dbForConsole')
->inject('dbForProject')
->inject('log')
->callback(fn (Message $message, Database $dbForConsole, Database $dbForProject, Log $log) => $this->action($message, $dbForConsole, $dbForProject, $log));
->inject('realtimeConnection')
->callback(fn (Message $message, Database $dbForConsole, Database $dbForProject, Log $log, callable $realtimeConnection) => $this->action($message, $dbForConsole, $dbForProject, $log, $realtimeConnection));
}
/**
@@ -44,10 +50,11 @@ class Databases extends Action
* @param Database $dbForConsole
* @param Database $dbForProject
* @param Log $log
* @param callable $realtimeConnection
* @return void
* @throws \Exception
*/
public function action(Message $message, Database $dbForConsole, Database $dbForProject, Log $log): void
public function action(Message $message, Database $dbForConsole, Database $dbForProject, Log $log, callable $realtimeConnection): void
{
$payload = $message->getPayload() ?? [];
@@ -60,6 +67,7 @@ class Databases extends Action
$collection = new Document($payload['collection'] ?? []);
$document = new Document($payload['document'] ?? []);
$database = new Document($payload['database'] ?? []);
$this->sourceRegion = $payload['sourceRegion'] ?? 'default';
$log->addTag('projectId', $project->getId());
$log->addTag('type', $type);
@@ -73,10 +81,10 @@ class Databases extends Action
match (\strval($type)) {
DATABASE_TYPE_DELETE_DATABASE => $this->deleteDatabase($database, $project, $dbForProject),
DATABASE_TYPE_DELETE_COLLECTION => $this->deleteCollection($database, $collection, $project, $dbForProject),
DATABASE_TYPE_CREATE_ATTRIBUTE => $this->createAttribute($database, $collection, $document, $project, $dbForConsole, $dbForProject),
DATABASE_TYPE_DELETE_ATTRIBUTE => $this->deleteAttribute($database, $collection, $document, $project, $dbForConsole, $dbForProject),
DATABASE_TYPE_CREATE_INDEX => $this->createIndex($database, $collection, $document, $project, $dbForConsole, $dbForProject),
DATABASE_TYPE_DELETE_INDEX => $this->deleteIndex($database, $collection, $document, $project, $dbForConsole, $dbForProject),
DATABASE_TYPE_CREATE_ATTRIBUTE => $this->createAttribute($database, $collection, $document, $project, $dbForConsole, $dbForProject, $realtimeConnection),
DATABASE_TYPE_DELETE_ATTRIBUTE => $this->deleteAttribute($database, $collection, $document, $project, $dbForConsole, $dbForProject, $realtimeConnection),
DATABASE_TYPE_CREATE_INDEX => $this->createIndex($database, $collection, $document, $project, $dbForConsole, $dbForProject, $realtimeConnection),
DATABASE_TYPE_DELETE_INDEX => $this->deleteIndex($database, $collection, $document, $project, $dbForConsole, $dbForProject, $realtimeConnection),
default => throw new \Exception('No database operation for type: ' . \strval($type)),
};
}
@@ -88,12 +96,13 @@ class Databases extends Action
* @param Document $project
* @param Database $dbForConsole
* @param Database $dbForProject
* @param callable $realtimeConnection
* @return void
* @throws Authorization
* @throws Conflict
* @throws \Exception
*/
private function createAttribute(Document $database, Document $collection, Document $attribute, Document $project, Database $dbForConsole, Database $dbForProject): void
private function createAttribute(Document $database, Document $collection, Document $attribute, Document $project, Database $dbForConsole, Database $dbForProject, callable $realtimeConnection): void
{
if ($collection->isEmpty()) {
throw new Exception('Missing collection');
@@ -193,7 +202,7 @@ class Databases extends Action
);
}
} finally {
$this->trigger($database, $collection, $attribute, $project, $projectId, $events);
$this->trigger($database, $collection, $attribute, $project, $projectId, $events, $realtimeConnection);
}
if ($type === Database::VAR_RELATIONSHIP && $options['twoWay']) {
@@ -210,6 +219,7 @@ class Databases extends Action
* @param Document $project
* @param Database $dbForConsole
* @param Database $dbForProject
* @param callable $realtimeConnection
* @return void
* @throws Authorization
* @throws Conflict
@@ -294,7 +304,7 @@ class Databases extends Action
);
}
} finally {
$this->trigger($database, $collection, $attribute, $project, $projectId, $events);
$this->trigger($database, $collection, $attribute, $project, $projectId, $events, $realtimeConnection);
}
// The underlying database removes/rebuilds indexes when attribute is removed
@@ -364,13 +374,14 @@ class Databases extends Action
* @param Document $project
* @param Database $dbForConsole
* @param Database $dbForProject
* @param callable $realtimeConnection
* @return void
* @throws Authorization
* @throws Conflict
* @throws Structure
* @throws DatabaseException
*/
private function createIndex(Document $database, Document $collection, Document $index, Document $project, Database $dbForConsole, Database $dbForProject): void
private function createIndex(Document $database, Document $collection, Document $index, Document $project, Database $dbForConsole, Database $dbForProject, callable $realtimeConnection): void
{
if ($collection->isEmpty()) {
throw new Exception('Missing collection');
@@ -412,7 +423,7 @@ class Databases extends Action
$index->setAttribute('status', 'failed')
);
} finally {
$this->trigger($database, $collection, $index, $project, $projectId, $events);
$this->trigger($database, $collection, $index, $project, $projectId, $events, $realtimeConnection);
}
$dbForProject->purgeCachedDocument('database_' . $database->getInternalId(), $collectionId);
@@ -425,13 +436,14 @@ class Databases extends Action
* @param Document $project
* @param Database $dbForConsole
* @param Database $dbForProject
* @param callable $realtimeConnection
* @return void
* @throws Authorization
* @throws Conflict
* @throws Structure
* @throws DatabaseException
*/
private function deleteIndex(Document $database, Document $collection, Document $index, Document $project, Database $dbForConsole, Database $dbForProject): void
private function deleteIndex(Document $database, Document $collection, Document $index, Document $project, Database $dbForConsole, Database $dbForProject, callable $realtimeConnection): void
{
if ($collection->isEmpty()) {
throw new Exception('Missing collection');
@@ -470,7 +482,7 @@ class Databases extends Action
$index->setAttribute('status', 'stuck')
);
} finally {
$this->trigger($database, $collection, $index, $project, $projectId, $events);
$this->trigger($database, $collection, $index, $project, $projectId, $events, $realtimeConnection);
}
$dbForProject->purgeCachedDocument('database_' . $database->getInternalId(), $collection->getId());
@@ -593,13 +605,18 @@ class Databases extends Action
Console::info("Deleted {$count} document by group in " . ($executionEnd - $executionStart) . " seconds");
}
/**
* @throws \RedisException
* @throws Exception
*/
protected function trigger(
Document $database,
Document $collection,
Document $attribute,
Document $project,
string $projectId,
array $events
array $events,
callable $realtimeConnection
): void {
$target = Realtime::fromPayload(
// Pass first, most verbose event pattern
@@ -608,6 +625,7 @@ class Databases extends Action
project: $project,
);
Realtime::send(
redis: $realtimeConnection($this->sourceRegion),
projectId: 'console',
payload: $attribute->getArrayCopy(),
events: $events,
+16 -2
View File
@@ -28,6 +28,11 @@ use Utopia\System\System;
class Functions extends Action
{
/**
* @var mixed|string
*/
protected string $sourceRegion;
public static function getName(): string
{
return 'functions';
@@ -47,7 +52,8 @@ class Functions extends Action
->inject('queueForEvents')
->inject('queueForUsage')
->inject('log')
->callback(fn (Message $message, Database $dbForProject, Func $queueForFunctions, Event $queueForEvents, Usage $queueForUsage, Log $log) => $this->action($message, $dbForProject, $queueForFunctions, $queueForEvents, $queueForUsage, $log));
->inject('realtimeConnection')
->callback(fn (Message $message, Database $dbForProject, Func $queueForFunctions, Event $queueForEvents, Usage $queueForUsage, Log $log, Callable $realtimeConnection) => $this->action($message, $dbForProject, $queueForFunctions, $queueForEvents, $queueForUsage, $log, $realtimeConnection));
}
/**
@@ -63,7 +69,7 @@ class Functions extends Action
* @throws \Utopia\Database\Exception
* @throws Conflict
*/
public function action(Message $message, Database $dbForProject, Func $queueForFunctions, Event $queueForEvents, Usage $queueForUsage, Log $log): void
public function action(Message $message, Database $dbForProject, Func $queueForFunctions, Event $queueForEvents, Usage $queueForUsage, Log $log, Callable $realtimeConnection): void
{
$payload = $message->getPayload() ?? [];
@@ -109,6 +115,8 @@ class Functions extends Action
return;
}
$this->sourceRegion = $payload['sourceRegion'] ?? 'default';
if ($function->isEmpty() && !empty($functionId)) {
$function = $dbForProject->getDocument('functions', $functionId);
}
@@ -141,6 +149,7 @@ class Functions extends Action
Console::success('Iterating function: ' . $function->getAttribute('name'));
$this->execute(
realtimeConnection: $realtimeConnection,
log: $log,
dbForProject: $dbForProject,
queueForFunctions: $queueForFunctions,
@@ -176,6 +185,7 @@ class Functions extends Action
$execution = new Document($payload['execution'] ?? []);
$user = new Document($payload['user'] ?? []);
$this->execute(
realtimeConnection: $realtimeConnection,
log: $log,
dbForProject: $dbForProject,
queueForFunctions: $queueForFunctions,
@@ -198,6 +208,7 @@ class Functions extends Action
case 'schedule':
$execution = new Document($payload['execution'] ?? []);
$this->execute(
realtimeConnection: $realtimeConnection,
log: $log,
dbForProject: $dbForProject,
queueForFunctions: $queueForFunctions,
@@ -307,6 +318,7 @@ class Functions extends Action
* @throws Conflict
*/
private function execute(
Callable $realtimeConnection,
Log $log,
Database $dbForProject,
Func $queueForFunctions,
@@ -599,6 +611,7 @@ class Functions extends Action
project: $project
);
Realtime::send(
redis: $realtimeConnection($this->sourceRegion),
projectId: 'console',
payload: $execution->getArrayCopy(),
events: $allEvents,
@@ -606,6 +619,7 @@ class Functions extends Action
roles: $target['roles']
);
Realtime::send(
redis: $realtimeConnection($this->sourceRegion),
projectId: $project->getId(),
payload: $execution->getArrayCopy(),
events: $allEvents,
+23 -10
View File
@@ -36,6 +36,10 @@ class Migrations extends Action
protected Database $dbForConsole;
protected Document $project;
/**
* @var string
*/
protected string $sourceRegion;
public static function getName(): string
{
@@ -53,13 +57,14 @@ class Migrations extends Action
->inject('dbForProject')
->inject('dbForConsole')
->inject('log')
->callback(fn (Message $message, Database $dbForProject, Database $dbForConsole, Log $log) => $this->action($message, $dbForProject, $dbForConsole, $log));
->inject('realtimeConnection')
->callback(fn (Message $message, Database $dbForProject, Database $dbForConsole, Log $log, Callable $realtimeConnection) => $this->action($message, $dbForProject, $dbForConsole, $log, $realtimeConnection));
}
/**
* @throws Exception
*/
public function action(Message $message, Database $dbForProject, Database $dbForConsole, Log $log): void
public function action(Message $message, Database $dbForProject, Database $dbForConsole, Log $log, Callable $realtimeConnection): void
{
$payload = $message->getPayload() ?? [];
@@ -75,6 +80,7 @@ class Migrations extends Action
return;
}
$this->sourceRegion = $payload['sourceRegion'] ?? 'default';
$this->dbForProject = $dbForProject;
$this->dbForConsole = $dbForConsole;
$this->project = $project;
@@ -89,7 +95,7 @@ class Migrations extends Action
$log->addTag('migrationId', $migration->getId());
$log->addTag('projectId', $project->getId());
$this->processMigration($migration, $log);
$this->processMigration($migration, $log, $realtimeConnection);
}
/**
@@ -158,7 +164,7 @@ class Migrations extends Action
* @throws \Utopia\Database\Exception
* @throws Exception
*/
protected function updateMigrationDocument(Document $migration, Document $project): Document
protected function updateMigrationDocument(Document $migration, Document $project, Callable $realtimeConnection): Document
{
/** Trigger Realtime */
$allEvents = Event::generateEvents('migrations.[migrationId].update', [
@@ -172,6 +178,7 @@ class Migrations extends Action
);
Realtime::send(
redis: $realtimeConnection($this->sourceRegion),
projectId: 'console',
payload: $migration->getArrayCopy(),
events: $allEvents,
@@ -180,6 +187,7 @@ class Migrations extends Action
);
Realtime::send(
redis: $realtimeConnection($this->sourceRegion),
projectId: $project->getId(),
payload: $migration->getArrayCopy(),
events: $allEvents,
@@ -253,6 +261,11 @@ class Migrations extends Action
}
/**
* @param Document $project
* @param Document $migration
* @param Log $log
* @param callable $realtimeConnection
* @return void
* @throws Authorization
* @throws Conflict
* @throws Restricted
@@ -260,7 +273,7 @@ class Migrations extends Action
* @throws \Utopia\Database\Exception
* @throws Exception
*/
protected function processMigration(Document $migration, Log $log): void
protected function processMigration(Document $migration, Log $log, Callable $realtimeConnection): void
{
$project = $this->project;
$projectDocument = $this->dbForConsole->getDocument('projects', $project->getId());
@@ -286,7 +299,7 @@ class Migrations extends Action
$migration->setAttribute('stage', 'processing');
$migration->setAttribute('status', 'processing');
$this->updateMigrationDocument($migration, $projectDocument);
$this->updateMigrationDocument($migration, $projectDocument, $realtimeConnection);
$log->addTag('type', $migration->getAttribute('source'));
@@ -302,14 +315,14 @@ class Migrations extends Action
/** Start Transfer */
$migration->setAttribute('stage', 'migrating');
$this->updateMigrationDocument($migration, $projectDocument);
$this->updateMigrationDocument($migration, $projectDocument, $realtimeConnection);
$transfer->run(
$migration->getAttribute('resources'),
function () use ($migration, $transfer, $projectDocument) {
$migration->setAttribute('resourceData', json_encode($transfer->getCache()));
$migration->setAttribute('statusCounters', json_encode($transfer->getStatusCounters()));
$this->updateMigrationDocument($migration, $projectDocument);
$this->updateMigrationDocument($migration, $projectDocument, $realtimeConnection);
},
$migration->getAttribute('resourceId'),
$migration->getAttribute('resourceType')
@@ -348,7 +361,7 @@ class Migrations extends Action
$migration->setAttribute('errors', $errorMessages);
$log->addExtra('migrationErrors', json_encode($errorMessages));
$this->updateMigrationDocument($migration, $projectDocument);
$this->updateMigrationDocument($migration, $projectDocument, $realtimeConnection);
return;
}
@@ -389,7 +402,7 @@ class Migrations extends Action
$this->removeAPIKey($tempAPIKey);
}
$this->updateMigrationDocument($migration, $projectDocument);
$this->updateMigrationDocument($migration, $projectDocument, $realtimeConnection);
if ($migration->getAttribute('status', '') === 'failed') {
Console::error('Migration('.$migration->getInternalId().':'.$migration->getId().') failed, Project('.$this->project->getInternalId().':'.$this->project->getId().')');