diff --git a/src/Appwrite/Platform/Workers/Deletes.php b/src/Appwrite/Platform/Workers/Deletes.php index 25a899bf12..fef4ef5b67 100644 --- a/src/Appwrite/Platform/Workers/Deletes.php +++ b/src/Appwrite/Platform/Workers/Deletes.php @@ -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 diff --git a/src/Appwrite/Platform/Workers/Messaging.php b/src/Appwrite/Platform/Workers/Messaging.php index af7d2027e3..3c34f2f919 100644 --- a/src/Appwrite/Platform/Workers/Messaging.php +++ b/src/Appwrite/Platform/Workers/Messaging.php @@ -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 $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);