Move database work to sn-refactors

This commit is contained in:
Bradley Schofield
2024-07-01 14:30:52 +09:00
parent 1a860e846f
commit 4a0b4f6141
5 changed files with 88 additions and 7 deletions
+2
View File
@@ -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) |
+13 -3
View File
@@ -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) {
+3
View File
@@ -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';
+9
View File
@@ -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,
+61 -4
View File
@@ -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,