(perf): Parallelize delete and messaging worker operations with Promise::map

This commit is contained in:
Jake Barnby
2026-03-25 14:27:14 +13:00
parent 33566f2052
commit ed780f58de
2 changed files with 221 additions and 267 deletions
+192 -234
View File
@@ -10,6 +10,7 @@ use Appwrite\Extend\Exception;
use Executor\Executor;
use Throwable;
use Utopia\Abuse\Adapters\TimeLimit\Database as AbuseDatabase;
use Utopia\Async\Promise;
use Utopia\Audit\Adapter\SQL;
use Utopia\Audit\Audit;
use Utopia\Cache\Adapter\Filesystem;
@@ -32,8 +33,6 @@ use Utopia\Queue\Message;
use Utopia\Storage\Device;
use Utopia\System\System;
use function Swoole\Coroutine\batch;
class Deletes extends Action
{
protected array $selects = ['$sequence', '$id', '$collection', '$permissions', '$updatedAt'];
@@ -635,56 +634,52 @@ class Deletes extends Action
}
});
// Delete Platforms
$this->deleteByGroup('platforms', [
Query::equal('projectInternalId', [$projectInternalId]),
Query::orderAsc()
], $dbForPlatform);
// Delete project and function rules
$this->deleteByGroup('rules', [
Query::equal('projectInternalId', [$projectInternalId]),
Query::orderAsc()
], $dbForPlatform, function (Document $document) use ($dbForPlatform, $certificates) {
$this->deleteRule($dbForPlatform, $document, $certificates);
});
// Delete Keys
$this->deleteByGroup('keys', [
Query::equal('resourceType', ['projects']),
Query::equal('resourceInternalId', [$projectInternalId]),
Query::orderAsc()
], $dbForPlatform);
// Delete Webhooks
$this->deleteByGroup('webhooks', [
Query::equal('projectInternalId', [$projectInternalId]),
Query::orderAsc()
], $dbForPlatform);
// Delete VCS Installations
$this->deleteByGroup('installations', [
Query::equal('projectInternalId', [$projectInternalId]),
Query::orderAsc()
], $dbForPlatform);
// Delete VCS Repositories
$this->deleteByGroup('repositories', [
Query::equal('projectInternalId', [$projectInternalId]),
Query::orderAsc()
], $dbForPlatform);
// Delete VCS comments
$this->deleteByGroup('vcsComments', [
Query::equal('projectInternalId', [$projectInternalId]),
Query::orderAsc()
], $dbForPlatform);
// Delete Schedules
$this->deleteByGroup('schedules', [
Query::equal('projectId', [$projectId]),
Query::orderAsc()
], $dbForPlatform);
// Delete platform-level resources concurrently
Promise::map([
// Delete Platforms
fn () => $this->deleteByGroup('platforms', [
Query::equal('projectInternalId', [$projectInternalId]),
Query::orderAsc()
], $dbForPlatform),
// Delete project and function rules
fn () => $this->deleteByGroup('rules', [
Query::equal('projectInternalId', [$projectInternalId]),
Query::orderAsc()
], $dbForPlatform, function (Document $document) use ($dbForPlatform, $certificates) {
$this->deleteRule($dbForPlatform, $document, $certificates);
}),
// Delete Keys
fn () => $this->deleteByGroup('keys', [
Query::equal('resourceType', ['projects']),
Query::equal('resourceInternalId', [$projectInternalId]),
Query::orderAsc()
], $dbForPlatform),
// Delete Webhooks
fn () => $this->deleteByGroup('webhooks', [
Query::equal('projectInternalId', [$projectInternalId]),
Query::orderAsc()
], $dbForPlatform),
// Delete VCS Installations
fn () => $this->deleteByGroup('installations', [
Query::equal('projectInternalId', [$projectInternalId]),
Query::orderAsc()
], $dbForPlatform),
// Delete VCS Repositories
fn () => $this->deleteByGroup('repositories', [
Query::equal('projectInternalId', [$projectInternalId]),
Query::orderAsc()
], $dbForPlatform),
// Delete VCS comments
fn () => $this->deleteByGroup('vcsComments', [
Query::equal('projectInternalId', [$projectInternalId]),
Query::orderAsc()
], $dbForPlatform),
// Delete Schedules
fn () => $this->deleteByGroup('schedules', [
Query::equal('projectId', [$projectId]),
Query::orderAsc()
], $dbForPlatform),
])->await();
// Delete metadata table
if ($projectTables) {
@@ -755,48 +750,46 @@ class Deletes extends Action
$userInternalId = $document->getSequence();
$dbForProject = $getProjectDB($project);
// Delete all sessions of this user from the sessions table and update the sessions field of the user record
$this->deleteByGroup('sessions', [
Query::equal('userInternalId', [$userInternalId]),
Query::orderAsc()
], $dbForProject);
if ($project->getId() === 'console') {
// Delete Keys
$this->deleteByGroup('keys', [
Query::equal('resourceInternalId', [$userInternalId]),
Query::equal('resourceType', ['users']),
Query::orderAsc()
], $dbForProject);
}
$dbForProject->purgeCachedDocument('users', $userId);
// Delete Memberships and decrement team membership counts
$this->deleteByGroup('memberships', [
Query::equal('userInternalId', [$userInternalId]),
Query::orderAsc()
], $dbForProject, function (Document $document) use ($dbForProject) {
if ($document->getAttribute('confirm')) { // Count only confirmed members
$teamId = $document->getAttribute('teamId');
$team = $dbForProject->getDocument('teams', $teamId);
if (!$team->isEmpty()) {
$dbForProject->decreaseDocumentAttribute('teams', $teamId, 'total', 1, 0);
// Delete user-related resources concurrently
Promise::map([
// Delete all sessions of this user
fn () => $this->deleteByGroup('sessions', [
Query::equal('userInternalId', [$userInternalId]),
Query::orderAsc()
], $dbForProject),
// Delete Keys (console project only)
fn () => $project->getId() === 'console'
? $this->deleteByGroup('keys', [
Query::equal('resourceInternalId', [$userInternalId]),
Query::equal('resourceType', ['users']),
Query::orderAsc()
], $dbForProject)
: null,
// Delete Memberships and decrement team membership counts
fn () => $this->deleteByGroup('memberships', [
Query::equal('userInternalId', [$userInternalId]),
Query::orderAsc()
], $dbForProject, function (Document $document) use ($dbForProject) {
if ($document->getAttribute('confirm')) { // Count only confirmed members
$teamId = $document->getAttribute('teamId');
$team = $dbForProject->getDocument('teams', $teamId);
if (!$team->isEmpty()) {
$dbForProject->decreaseDocumentAttribute('teams', $teamId, 'total', 1, 0);
}
}
}
});
// Delete tokens
$this->deleteByGroup('tokens', [
Query::equal('userInternalId', [$userInternalId]),
Query::orderAsc()
], $dbForProject);
// Delete identities
Identities::delete($dbForProject, Query::equal('userInternalId', [$userInternalId]));
// Delete targets
Targets::delete($dbForProject, Query::equal('userInternalId', [$userInternalId]));
}),
// Delete tokens
fn () => $this->deleteByGroup('tokens', [
Query::equal('userInternalId', [$userInternalId]),
Query::orderAsc()
], $dbForProject),
// Delete identities
fn () => Identities::delete($dbForProject, Query::equal('userInternalId', [$userInternalId])),
// Delete targets
fn () => Targets::delete($dbForProject, Query::equal('userInternalId', [$userInternalId])),
])->await();
}
/**
@@ -815,7 +808,7 @@ class Deletes extends Action
// Delete Executions
$this->deleteByGroup('executions', [
Query::select([...$this->selects, '$createdAt']),
Query::lessThan('$createdAt', $datetime),
Query::createdBefore($datetime),
Query::orderDesc('$createdAt'),
Query::orderDesc(),
], $dbForProject);
@@ -866,7 +859,7 @@ class Deletes extends Action
Query::select([...$this->selects, '$createdAt']),
Query::equal('resourceInternalId', [$resourceInternalId]),
Query::equal('resourceType', [$resourceType]),
Query::lessThan('$createdAt', $cutoffTime),
Query::createdBefore($cutoffTime),
Query::orderDesc('$createdAt'),
Query::orderDesc(),
], $dbForProject);
@@ -889,10 +882,10 @@ class Deletes extends Action
};
/* perform processing in parallel */
batch([
Promise::map([
fn () => $processResource(RESOURCE_TYPE_SITES),
fn () => $processResource(RESOURCE_TYPE_FUNCTIONS),
]);
])->await();
}
}
@@ -911,7 +904,7 @@ class Deletes extends Action
// Delete Sessions
$this->deleteByGroup('sessions', [
Query::select([...$this->selects, '$createdAt']),
Query::lessThan('$createdAt', $expired),
Query::createdBefore($expired),
Query::orderDesc('$createdAt'),
Query::orderDesc(),
], $dbForProject);
@@ -1003,73 +996,58 @@ class Deletes extends Action
$siteId = $document->getId();
$siteInternalId = $document->getSequence();
/**
* Delete rules for site
*/
Console::info("Deleting rules for site " . $siteId);
$this->deleteByGroup('rules', [
Query::equal('type', ['deployment']),
Query::equal('deploymentResourceType', ['site']),
Query::equal('deploymentResourceInternalId', [$siteInternalId]),
Query::equal('projectInternalId', [$project->getSequence()])
], $dbForPlatform, function (Document $document) use ($dbForPlatform, $certificates) {
$this->deleteRule($dbForPlatform, $document, $certificates);
});
/**
* Delete Variables
*/
Console::info("Deleting variables for site " . $siteId);
$this->deleteByGroup('variables', [
Query::equal('resourceType', ['site']),
Query::equal('resourceInternalId', [$siteInternalId])
], $dbForProject);
/**
* Delete Deployments
*/
Console::info("Deleting deployments for site " . $siteId);
// Delete site resources concurrently
Console::info("Deleting resources for site " . $siteId);
$deploymentInternalIds = [];
$deploymentIds = [];
$this->deleteByGroup('deployments', [
Query::equal('resourceInternalId', [$siteInternalId]),
Query::equal('resourceType', ['sites']),
Query::orderAsc()
], $dbForProject, function (Document $document) use ($project, $certificates, $deviceForSites, $deviceForBuilds, $deviceForFiles, $dbForPlatform, &$deploymentInternalIds) {
$deploymentInternalIds[] = $document->getSequence();
$deploymentIds[] = $document->getId();
$this->deleteBuildFiles($deviceForBuilds, $document);
$this->deleteDeploymentFiles($deviceForSites, $document);
$this->deleteDeploymentScreenshots($deviceForFiles, $dbForPlatform, $document);
});
/**
* Delete Logs
*/
Console::info("Deleting logs for site " . $siteId);
$this->deleteByGroup('executions', [
Query::select($this->selects),
Query::equal('resourceInternalId', [$siteInternalId]),
Query::equal('resourceType', ['sites']),
Query::orderAsc()
], $dbForProject);
/**
* Delete VCS Repositories and VCS Comments
*/
Console::info("Deleting VCS repositories and comments linked to site " . $siteId);
$this->deleteByGroup('repositories', [
Query::equal('projectInternalId', [$project->getSequence()]),
Query::equal('resourceInternalId', [$siteInternalId]),
Query::equal('resourceType', ['site']),
], $dbForPlatform, function (Document $document) use ($dbForPlatform) {
$providerRepositoryId = $document->getAttribute('providerRepositoryId', '');
$projectInternalId = $document->getAttribute('projectInternalId', '');
$this->deleteByGroup('vcsComments', [
Query::equal('providerRepositoryId', [$providerRepositoryId]),
Query::equal('projectInternalId', [$projectInternalId]),
], $dbForPlatform);
});
Promise::map([
// Delete rules for site
fn () => $this->deleteByGroup('rules', [
Query::equal('type', ['deployment']),
Query::equal('deploymentResourceType', ['site']),
Query::equal('deploymentResourceInternalId', [$siteInternalId]),
Query::equal('projectInternalId', [$project->getSequence()])
], $dbForPlatform, function (Document $document) use ($dbForPlatform, $certificates) {
$this->deleteRule($dbForPlatform, $document, $certificates);
}),
// Delete Variables
fn () => $this->deleteByGroup('variables', [
Query::equal('resourceType', ['site']),
Query::equal('resourceInternalId', [$siteInternalId])
], $dbForProject),
// Delete Deployments
fn () => $this->deleteByGroup('deployments', [
Query::equal('resourceInternalId', [$siteInternalId]),
Query::equal('resourceType', ['sites']),
Query::orderAsc()
], $dbForProject, function (Document $document) use ($project, $certificates, $deviceForSites, $deviceForBuilds, $deviceForFiles, $dbForPlatform, &$deploymentInternalIds) {
$deploymentInternalIds[] = $document->getSequence();
$deploymentIds[] = $document->getId();
$this->deleteBuildFiles($deviceForBuilds, $document);
$this->deleteDeploymentFiles($deviceForSites, $document);
$this->deleteDeploymentScreenshots($deviceForFiles, $dbForPlatform, $document);
}),
// Delete Logs
fn () => $this->deleteByGroup('executions', [
Query::select($this->selects),
Query::equal('resourceInternalId', [$siteInternalId]),
Query::equal('resourceType', ['sites']),
Query::orderAsc()
], $dbForProject),
// Delete VCS Repositories and VCS Comments
fn () => $this->deleteByGroup('repositories', [
Query::equal('projectInternalId', [$project->getSequence()]),
Query::equal('resourceInternalId', [$siteInternalId]),
Query::equal('resourceType', ['site']),
], $dbForPlatform, function (Document $document) use ($dbForPlatform) {
$providerRepositoryId = $document->getAttribute('providerRepositoryId', '');
$projectInternalId = $document->getAttribute('projectInternalId', '');
$this->deleteByGroup('vcsComments', [
Query::equal('providerRepositoryId', [$providerRepositoryId]),
Query::equal('projectInternalId', [$projectInternalId]),
], $dbForPlatform);
}),
])->await();
}
/**
@@ -1089,82 +1067,62 @@ class Deletes extends Action
$functionId = $document->getId();
$functionInternalId = $document->getSequence();
/**
* Delete rules
*/
Console::info("Deleting rules for function " . $functionId);
$this->deleteByGroup('rules', [
Query::equal('type', ['deployment']),
Query::equal('deploymentResourceType', ['function']),
Query::equal('deploymentResourceInternalId', [$functionInternalId]),
Query::equal('projectInternalId', [$project->getSequence()]),
Query::orderAsc()
], $dbForPlatform, function (Document $document) use ($project, $dbForPlatform, $certificates) {
$this->deleteRule($dbForPlatform, $document, $certificates);
});
/**
* Delete Variables
*/
Console::info("Deleting variables for function " . $functionId);
$this->deleteByGroup('variables', [
Query::equal('resourceInternalId', [$functionInternalId]),
Query::equal('resourceType', ['function']),
Query::orderAsc()
], $dbForProject);
/**
* Delete Deployments
*/
Console::info("Deleting deployments for function " . $functionId);
// Delete function resources concurrently
Console::info("Deleting resources for function " . $functionId);
$deploymentInternalIds = [];
$this->deleteByGroup('deployments', [
Query::equal('resourceInternalId', [$functionInternalId]),
Query::equal('resourceType', ['functions']),
Query::orderAsc()
], $dbForProject, function (Document $document) use ($dbForPlatform, $project, $certificates, $deviceForFunctions, $deviceForBuilds, &$deploymentInternalIds) {
$deploymentInternalIds[] = $document->getSequence();
$this->deleteDeploymentFiles($deviceForFunctions, $document);
$this->deleteBuildFiles($deviceForBuilds, $document);
});
/**
* Delete Executions
*/
Console::info("Deleting executions for function " . $functionId);
$this->deleteByGroup('executions', [
Query::select($this->selects),
Query::equal('resourceInternalId', [$functionInternalId]),
Query::equal('resourceType', ['functions']),
Query::orderAsc()
], $dbForProject);
/**
* Delete VCS Repositories and VCS Comments
*/
Console::info("Deleting VCS repositories and comments linked to function " . $functionId);
$this->deleteByGroup('repositories', [
Query::equal('projectInternalId', [$project->getSequence()]),
Query::equal('resourceInternalId', [$functionInternalId]),
Query::equal('resourceType', ['function']),
Query::orderAsc()
], $dbForPlatform, function (Document $document) use ($dbForPlatform) {
$providerRepositoryId = $document->getAttribute('providerRepositoryId', '');
$projectInternalId = $document->getAttribute('projectInternalId', '');
$this->deleteByGroup('vcsComments', [
Query::equal('providerRepositoryId', [$providerRepositoryId]),
Query::equal('projectInternalId', [$projectInternalId]),
Promise::map([
// Delete rules
fn () => $this->deleteByGroup('rules', [
Query::equal('type', ['deployment']),
Query::equal('deploymentResourceType', ['function']),
Query::equal('deploymentResourceInternalId', [$functionInternalId]),
Query::equal('projectInternalId', [$project->getSequence()]),
Query::orderAsc()
], $dbForPlatform);
});
], $dbForPlatform, function (Document $document) use ($project, $dbForPlatform, $certificates) {
$this->deleteRule($dbForPlatform, $document, $certificates);
}),
// Delete Variables
fn () => $this->deleteByGroup('variables', [
Query::equal('resourceInternalId', [$functionInternalId]),
Query::equal('resourceType', ['function']),
Query::orderAsc()
], $dbForProject),
// Delete Deployments
fn () => $this->deleteByGroup('deployments', [
Query::equal('resourceInternalId', [$functionInternalId]),
Query::equal('resourceType', ['functions']),
Query::orderAsc()
], $dbForProject, function (Document $document) use ($dbForPlatform, $project, $certificates, $deviceForFunctions, $deviceForBuilds, &$deploymentInternalIds) {
$deploymentInternalIds[] = $document->getSequence();
$this->deleteDeploymentFiles($deviceForFunctions, $document);
$this->deleteBuildFiles($deviceForBuilds, $document);
}),
// Delete Executions
fn () => $this->deleteByGroup('executions', [
Query::select($this->selects),
Query::equal('resourceInternalId', [$functionInternalId]),
Query::equal('resourceType', ['functions']),
Query::orderAsc()
], $dbForProject),
// Delete VCS Repositories and VCS Comments
fn () => $this->deleteByGroup('repositories', [
Query::equal('projectInternalId', [$project->getSequence()]),
Query::equal('resourceInternalId', [$functionInternalId]),
Query::equal('resourceType', ['function']),
Query::orderAsc()
], $dbForPlatform, function (Document $document) use ($dbForPlatform) {
$providerRepositoryId = $document->getAttribute('providerRepositoryId', '');
$projectInternalId = $document->getAttribute('projectInternalId', '');
/**
* Request executor to delete all deployment containers
*/
Console::info("Requesting executor to delete all deployment containers for function " . $functionId);
$this->deleteRuntimes($getProjectDB, $document, $project, $executor);
$this->deleteByGroup('vcsComments', [
Query::equal('providerRepositoryId', [$providerRepositoryId]),
Query::equal('projectInternalId', [$projectInternalId]),
Query::orderAsc()
], $dbForPlatform);
}),
// Request executor to delete all deployment containers
fn () => $this->deleteRuntimes($getProjectDB, $document, $project, $executor),
])->await();
}
private function deleteDeploymentScreenshots(Device $deviceForFiles, Database $dbForPlatform, Document $deployment): void
+29 -33
View File
@@ -9,6 +9,7 @@ use Appwrite\Usage\Context as UsageContext;
use libphonenumber\NumberParseException;
use libphonenumber\PhoneNumberUtil;
use Swoole\Runtime;
use Utopia\Async\Promise;
use Utopia\Config\Config;
use Utopia\Database\Database;
use Utopia\Database\DateTime;
@@ -48,8 +49,6 @@ use Utopia\Storage\Device\Local;
use Utopia\Storage\Storage;
use Utopia\System\System;
use function Swoole\Coroutine\batch;
class Messaging extends Action
{
private ?Local $localDevice = null;
@@ -147,44 +146,41 @@ class Messaging extends Action
*/
$allTargets = [];
if (\count($topicIds) > 0) {
$topics = $dbForProject->find('topics', [
// Fetch topics, users, and targets concurrently
$results = Promise::map([
'topics' => fn () => \count($topicIds) > 0 ? $dbForProject->find('topics', [
Query::equal('$id', $topicIds),
Query::limit(\count($topicIds)),
]);
foreach ($topics as $topic) {
$targets = \array_filter($topic->getAttribute('targets'), function (Document $target) use ($providerType) {
return $target->getAttribute('providerType') === $providerType;
});
\array_push($allTargets, ...$targets);
}
}
if (\count($userIds) > 0) {
$users = $dbForProject->find('users', [
]) : [],
'users' => fn () => \count($userIds) > 0 ? $dbForProject->find('users', [
Query::equal('$id', $userIds),
Query::limit(\count($userIds)),
]);
foreach ($users as $user) {
$targets = \array_filter($user->getAttribute('targets'), function (Document $target) use ($providerType) {
return $target->getAttribute('providerType') === $providerType;
});
\array_push($allTargets, ...$targets);
}
}
if (\count($targetIds) > 0) {
$targets = $dbForProject->find('targets', [
]) : [],
'targets' => fn () => \count($targetIds) > 0 ? $dbForProject->find('targets', [
Query::equal('$id', $targetIds),
Query::equal('providerType', [$providerType]),
Query::limit(\count($targetIds)),
]);
]) : [],
])->await();
foreach ($results['topics'] as $topic) {
$targets = \array_filter($topic->getAttribute('targets'), function (Document $target) use ($providerType) {
return $target->getAttribute('providerType') === $providerType;
});
\array_push($allTargets, ...$targets);
}
foreach ($results['users'] as $user) {
$targets = \array_filter($user->getAttribute('targets'), function (Document $target) use ($providerType) {
return $target->getAttribute('providerType') === $providerType;
});
\array_push($allTargets, ...$targets);
}
\array_push($allTargets, ...$results['targets']);
if (empty($allTargets)) {
$dbForProject->updateDocument('messages', $message->getId(), $message->setAttributes([
'status' => MessageStatus::FAILED,
@@ -241,7 +237,7 @@ class Messaging extends Action
/**
* @var array<array> $results
*/
$results = batch(\array_map(function ($providerId) use ($identifiers, &$providers, $default, $message, $dbForProject, $deviceForFiles, $project, $publisherForUsage) {
$results = Promise::map(\array_map(function ($providerId) use ($identifiers, &$providers, $default, $message, $dbForProject, $deviceForFiles, $project, $publisherForUsage) {
return function () use ($providerId, $identifiers, &$providers, $default, $message, $dbForProject, $deviceForFiles, $project, $publisherForUsage) {
if (\array_key_exists($providerId, $providers)) {
$provider = $providers[$providerId];
@@ -269,7 +265,7 @@ class Messaging extends Action
$adapter->getMaxMessagesPerRequest()
);
return batch(\array_map(function ($batch) use ($message, $provider, $adapter, $dbForProject, $deviceForFiles, $project, $publisherForUsage) {
return Promise::map(\array_map(function ($batch) use ($message, $provider, $adapter, $dbForProject, $deviceForFiles, $project, $publisherForUsage) {
return function () use ($batch, $message, $provider, $adapter, $dbForProject, $deviceForFiles, $project, $publisherForUsage) {
$deliveredTotal = 0;
$deliveryErrors = [];
@@ -333,9 +329,9 @@ class Messaging extends Action
];
}
};
}, $batches));
}, $batches))->await();
};
}, \array_keys($identifiers)));
}, \array_keys($identifiers)))->await();
$results = \array_merge(...$results);