diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 9e3e6fcd81..6395aa656b 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -321,8 +321,10 @@ These are the current metrics we collect usage stats for: | databases | Total number of databases per project | | collections | Total number of collections per project | | {databaseInternalId}.collections | Total number of collections per database| +| {databaseInternalId}.storage | Sum of database storage (in bytes) | | documents | Total number of documents per project | | {databaseInternalId}.{collectionInternalId}.documents | Total number of documents per collection | +| {databaseInternalId}.{collectionInternalId}.storage | Sum of database storage used by the collection (in bytes) | | buckets | Total number of buckets per project | | files | Total number of files per project | | {bucketInternalId}.files.storage | Sum of files.storage per bucket (in bytes) | diff --git a/app/controllers/shared/api.php b/app/controllers/shared/api.php index 2b0013db29..b134ebf23d 100644 --- a/app/controllers/shared/api.php +++ b/app/controllers/shared/api.php @@ -55,7 +55,7 @@ $parseLabel = function (string $label, array $responsePayload, array $requestPar return $label; }; -$databaseListener = function (string $event, Document $document, Document $project, Usage $queueForUsage, Database $dbForProject) { +$databaseListener = function (string $event, Document $document, Document $project, Usage $queueForUsage, EventDatabase $queueForDatabase, Database $dbForProject) { $value = 1; if ($event === Database::EVENT_DOCUMENT_DELETE) { @@ -109,6 +109,16 @@ $databaseListener = function (string $event, Document $document, Document $proje ->addMetric(METRIC_DOCUMENTS, $value) // per project ->addMetric(str_replace('{databaseInternalId}', $databaseInternalId, METRIC_DATABASE_ID_DOCUMENTS), $value) // per database ->addMetric(str_replace(['{databaseInternalId}', '{collectionInternalId}'], [$databaseInternalId, $collectionInternalId], METRIC_DATABASE_ID_COLLECTION_ID_DOCUMENTS), $value); // per collection + + $queueForDatabase + ->setType(DATABASE_TYPE_CALCULATE_STORAGE_USAGE) + ->setPayload([ + 'collectionInternalId' => $collectionInternalId, + 'databaseInternalId' => $databaseInternalId, + ]); + + $queueForDatabase->trigger(); + break; case $document->getCollection() === 'buckets': //buckets $queueForUsage @@ -406,8 +416,8 @@ App::init() $queueForMessaging->setProject($project); $dbForProject - ->on(Database::EVENT_DOCUMENT_CREATE, 'calculate-usage', fn ($event, $document) => $databaseListener($event, $document, $project, $queueForUsage, $dbForProject)) - ->on(Database::EVENT_DOCUMENT_DELETE, 'calculate-usage', fn ($event, $document) => $databaseListener($event, $document, $project, $queueForUsage, $dbForProject)); + ->on(Database::EVENT_DOCUMENT_CREATE, 'calculate-usage', fn ($event, $document) => $databaseListener($event, $document, $project, $queueForUsage, $queueForDatabase, $dbForProject)) + ->on(Database::EVENT_DOCUMENT_DELETE, 'calculate-usage', fn ($event, $document) => $databaseListener($event, $document, $project, $queueForUsage, $queueForDatabase, $dbForProject)); $useCache = $route->getLabel('cache', false); if ($useCache) { diff --git a/app/init.php b/app/init.php index a86156c750..65493cb95c 100644 --- a/app/init.php +++ b/app/init.php @@ -154,6 +154,7 @@ const DATABASE_TYPE_DELETE_ATTRIBUTE = 'deleteAttribute'; const DATABASE_TYPE_DELETE_INDEX = 'deleteIndex'; const DATABASE_TYPE_DELETE_COLLECTION = 'deleteCollection'; const DATABASE_TYPE_DELETE_DATABASE = 'deleteDatabase'; +const DATABASE_TYPE_CALCULATE_STORAGE_USAGE = 'calculateStorageUsage'; // Build Worker Types const BUILD_TYPE_DEPLOYMENT = 'deployment'; @@ -217,9 +218,11 @@ const METRIC_SESSIONS = 'sessions'; const METRIC_DATABASES = 'databases'; const METRIC_COLLECTIONS = 'collections'; const METRIC_DATABASE_ID_COLLECTIONS = '{databaseInternalId}.collections'; +Const METRIC_DATABASE_ID_STORAGE = '{databaseInternalId}.storage'; const METRIC_DOCUMENTS = 'documents'; const METRIC_DATABASE_ID_DOCUMENTS = '{databaseInternalId}.documents'; const METRIC_DATABASE_ID_COLLECTION_ID_DOCUMENTS = '{databaseInternalId}.{collectionInternalId}.documents'; +const METRIC_DATABASE_ID_COLLECTION_ID_STORAGE = '{databaseInternalId}.{collectionInternalId}.storage'; const METRIC_BUCKETS = 'buckets'; const METRIC_FILES = 'files'; const METRIC_FILES_STORAGE = 'files.storage'; diff --git a/src/Appwrite/Event/Database.php b/src/Appwrite/Event/Database.php index f9eb7d9a7d..05e80633a3 100644 --- a/src/Appwrite/Event/Database.php +++ b/src/Appwrite/Event/Database.php @@ -119,6 +119,15 @@ class Database extends Event $client = new Client($this->queue, $this->connection); + if ($this->type === DATABASE_TYPE_CALCULATE_STORAGE_USAGE) { + $result = $client->enqueue(array_merge([ + 'project' => $this->project, + 'type' => $this->type, + 'user' => $this->user + ], $this->payload)); + return $result; + } + try { $result = $client->enqueue([ 'project' => $this->project, diff --git a/src/Appwrite/Platform/Workers/Databases.php b/src/Appwrite/Platform/Workers/Databases.php index 56f5f012e8..f0b854e6af 100644 --- a/src/Appwrite/Platform/Workers/Databases.php +++ b/src/Appwrite/Platform/Workers/Databases.php @@ -3,6 +3,7 @@ namespace Appwrite\Platform\Workers; use Appwrite\Event\Event; +use Appwrite\Event\Usage; use Appwrite\Messaging\Adapter\Realtime; use Exception; use Utopia\Audit\Audit; @@ -36,19 +37,21 @@ class Databases extends Action ->inject('message') ->inject('dbForConsole') ->inject('dbForProject') + ->inject('queueForUsage') ->inject('log') - ->callback(fn (Message $message, Database $dbForConsole, Database $dbForProject, Log $log) => $this->action($message, $dbForConsole, $dbForProject, $log)); + ->callback(fn (Message $message, Database $dbForConsole, Database $dbForProject, Usage $queueForUsage, Log $log) => $this->action($message, $dbForConsole, $dbForProject, $queueForUsage, $log)); } /** * @param Message $message * @param Database $dbForConsole * @param Database $dbForProject + * @param Usage $queueForUsage * @param Log $log * @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, Usage $queueForUsage, Log $log): void { $payload = $message->getPayload() ?? []; @@ -65,11 +68,13 @@ class Databases extends Action $log->addTag('projectId', $project->getId()); $log->addTag('type', $type); - if ($database->isEmpty()) { + if ($database->isEmpty() && $type !== DATABASE_TYPE_CALCULATE_STORAGE_USAGE) { throw new Exception('Missing database'); } - $log->addTag('databaseId', $database->getId()); + if (!$database->isEmpty()) { + $log->addTag('databaseId', $database->getId()); + } match (\strval($type)) { DATABASE_TYPE_DELETE_DATABASE => $this->deleteDatabase($database, $project, $dbForProject), @@ -78,6 +83,8 @@ class Databases extends Action 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_CALCULATE_STORAGE_USAGE => $this->calculateStorageUsage($payload, $project, $dbForProject, $queueForUsage), + default => throw new \Exception('No database operation for type: ' . \strval($type)), }; } @@ -612,6 +619,56 @@ class Databases extends Action Console::info("Deleted {$count} document by group in " . ($executionEnd - $executionStart) . " seconds"); } + /** + * @param Document $database + * @param Document $collection + * @param Database $dbForConsole + * @param Usage $queueForUsage + * @return void + * @throws Exception + * @throws Authorization + * @throws DatabaseException + */ + private function calculateStorageUsage(array $payload, Document $project, Database $dbForProject, Usage $queueForUsage): void + { + if (!isset($payload['databaseInternalId'])) { + throw new Exception('Missing Database'); + } + + if (!isset($payload['collectionInternalId'])) { + throw new Exception('Missing Collection'); + } + + $databaseInternalId = $payload['databaseInternalId']; + $collectionInternalId = $payload['collectionInternalId']; + + // Calculate storage usage for collection + $collectionStorageUsage = $dbForProject->getSizeOfCollection('database_'. $databaseInternalId . '_collection_' . $collectionInternalId); + + //TODO: Optimize using reduce + // Calculate storage usage for database + $databsaeStorageUsage = 0; + + $collections = $dbForProject->find('database_' . $databaseInternalId); + foreach ($collections as $collection) { + $databsaeStorageUsage += $dbForProject->getSizeOfCollection('database_' . $databaseInternalId . '_collection_' . $collection->getInternalId()); + } + + $queueForUsage + ->addMetric(str_replace('{databaseInternalId}', $databaseInternalId, METRIC_DATABASE_ID_STORAGE), $databsaeStorageUsage) + ->addMetric(str_replace([ + '{databaseInternalId}', + '{collectionInternalId}' + ], [ + $databaseInternalId, + $collectionInternalId + ], METRIC_DATABASE_ID_COLLECTION_ID_STORAGE), $collectionStorageUsage); + + Console::info('Calculated storage usage for database: ' . $databaseInternalId . ' and collection: ' . $collectionInternalId); + $queueForUsage->setProject($project); + $queueForUsage->trigger(); + } + protected function trigger( Document $database, Document $collection,