addressed comments

* renamed getDatabaseDB to getDatabasesDB
* renamed dbForDatabase to dbForDatabases
* removed call_user_func and using newer callable syntax
This commit is contained in:
ArnabChatterjee20k
2025-10-10 20:00:43 +05:30
parent d51486c62a
commit 34e4208b3e
55 changed files with 173 additions and 170 deletions
+1 -4
View File
@@ -43,10 +43,7 @@ _APP_DB_PASS=password
_APP_DB_ROOT_PASS=rootsecretpassword
_APP_DB_ADAPTER_DOCUMENTSDB=mongodb
_APP_DB_HOST_DOCUMENTSDB=mongodb
_APP_DB_HOST_DOCUMENTSDB_PORT=27017
_APP_CONNECTIONS_DATABASE_DOCUMENTSDB=selfhosted_fra1=mongodb://user:password@mongodb:27017/appwrite,selfhosted_fra2=mongodb://user:password@mongodb:27017/appwrite
_APP_DATABASE_DOCUMENTSDB_KEYS=documentsdb_selfhosted_fra1,documentsdb_selfhosted_fra2
_APP_DATABASE_DOCUMENTSDB_OVERRIDE=documentsdb_selfhosted_fra1
_APP_DB_PORT_DOCUMENTSDB=27017
_APP_STORAGE_DEVICE=Local
_APP_STORAGE_S3_ACCESS_KEY=
_APP_STORAGE_S3_SECRET=
+1 -1
View File
@@ -113,7 +113,7 @@ $register->set('pools', function () {
$fallbackForDocumentsDB = 'db_main=' . AppwriteURL::unparse([
'scheme' => 'mongodb',
'host' => System::getEnv('_APP_DB_HOST_DOCUMENTSDB', 'mongodb'),
'port' => System::getEnv('_APP_DB_HOST_DOCUMENTSDB_PORT', '27017'),
'port' => System::getEnv('_APP_DB_PORT_DOCUMENTSDB', '27017'),
'user' => System::getEnv('_APP_DB_USER', ''),
'pass' => System::getEnv('_APP_DB_PASS', ''),
'path' => System::getEnv('_APP_DB_SCHEMA', ''),
+8 -2
View File
@@ -431,11 +431,17 @@ App::setResource('dbForPlatform', function (Group $pools, Cache $cache) {
return $database;
}, ['pools', 'cache']);
App::setResource('getDatabaseDB', function (Group $pools, Cache $cache, Document $project, Request $request, StatsUsage $queueForStatsUsage) {
App::setResource('getDatabasesDB', function (Group $pools, Cache $cache, Document $project, Request $request, StatsUsage $queueForStatsUsage) {
return function (Document $database) use ($pools, $cache, $project, $request, $queueForStatsUsage): Database {
$databaseType = $database->getAttribute('database', '');
$databaseDSN = new DSN($databaseType);
try {
$databaseDSN = new DSN($databaseType);
} catch (\InvalidArgumentException) {
// for old databases migrated through patch script
// databaseType determines the adapter
$databaseDSN = new DSN('mysql://'.$databaseType);
}
try {
$dsn = new DSN($project->getAttribute('database'));
} catch (\InvalidArgumentException) {
+1 -1
View File
@@ -192,7 +192,7 @@ Server::setResource('getLogsDB', function (Group $pools, Cache $cache) {
};
}, ['pools', 'cache']);
Server::setResource('getDatabaseDB', function (Cache $cache, Registry $register, Document $project) {
Server::setResource('getDatabasesDB', function (Cache $cache, Registry $register, Document $project) {
return function (Document $database) use ($cache, $register, $project): Database {
$databaseType = $database->getAttribute('database', '');
$databaseDSN = new DSN($databaseType);
@@ -76,12 +76,12 @@ class Create extends Action
->param('enabled', true, new Boolean(), 'Is collection enabled? When set to \'disabled\', users cannot access the collection but Server SDKs with and API key can still read and write to the collection. No data is lost when this is toggled.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForEvents')
->callback($this->action(...));
}
public function action(string $databaseId, string $collectionId, string $name, ?array $permissions, bool $documentSecurity, bool $enabled, UtopiaResponse $response, Database $dbForProject, callable $getDatabaseDB, Event $queueForEvents): void
public function action(string $databaseId, string $collectionId, string $name, ?array $permissions, bool $documentSecurity, bool $enabled, UtopiaResponse $response, Database $dbForProject, callable $getDatabasesDB, Event $queueForEvents): void
{
$database = Authorization::skip(fn () => $dbForProject->getDocument('databases', $databaseId));
@@ -113,9 +113,9 @@ class Create extends Action
throw new Exception(Exception::DATABASE_NOT_FOUND);
}
$dbForDatabase = call_user_func($getDatabaseDB, $database);
$dbForDatabases = $getDatabasesDB($database);
try {
$dbForDatabase->createCollection(
$dbForDatabases->createCollection(
id: 'database_' . $database->getSequence() . '_collection_' . $collection->getSequence(),
permissions: $permissions,
documentSecurity: $documentSecurity
@@ -62,13 +62,13 @@ class Delete extends Action
->param('collectionId', '', new UID(), 'Collection ID.')
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForDatabase')
->inject('queueForEvents')
->callback($this->action(...));
}
public function action(string $databaseId, string $collectionId, UtopiaResponse $response, Database $dbForProject, callable $getDatabaseDB, EventDatabase $queueForDatabase, Event $queueForEvents): void
public function action(string $databaseId, string $collectionId, UtopiaResponse $response, Database $dbForProject, callable $getDatabasesDB, EventDatabase $queueForDatabase, Event $queueForEvents): void
{
$database = Authorization::skip(fn () => $dbForProject->getDocument('databases', $databaseId));
if ($database->isEmpty()) {
@@ -85,8 +85,8 @@ class Delete extends Action
throw new Exception(Exception::GENERAL_SERVER_ERROR, "Failed to remove $type from DB");
}
$dbForDatabase = call_user_func($getDatabaseDB, $database);
$dbForDatabase->purgeCachedCollection('database_' . $database->getSequence() . '_collection_' . $collection->getSequence());
$dbForDatabases = $getDatabasesDB($database);
$dbForDatabases->purgeCachedCollection('database_' . $database->getSequence() . '_collection_' . $collection->getSequence());
$queueForDatabase
->setType(DATABASE_TYPE_DELETE_COLLECTION)
@@ -81,14 +81,14 @@ class Decrement extends Action
->param('transactionId', null, new UID(), 'Transaction ID for staging the operation.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForEvents')
->inject('queueForStatsUsage')
->inject('plan')
->callback($this->action(...));
}
public function action(string $databaseId, string $collectionId, string $documentId, string $attribute, int|float $value, int|float|null $min, ?string $transactionId, UtopiaResponse $response, Database $dbForProject, callable $getDatabaseDB, Event $queueForEvents, StatsUsage $queueForStatsUsage, array $plan): void
public function action(string $databaseId, string $collectionId, string $documentId, string $attribute, int|float $value, int|float|null $min, ?string $transactionId, UtopiaResponse $response, Database $dbForProject, callable $getDatabasesDB, Event $queueForEvents, StatsUsage $queueForStatsUsage, array $plan): void
{
$isAPIKey = Auth::isAppUser(Authorization::getRoles());
$isPrivilegedUser = Auth::isPrivilegedUser(Authorization::getRoles());
@@ -166,9 +166,9 @@ class Decrement extends Action
return;
}
$dbForDatabase = call_user_func($getDatabaseDB, $database);
$dbForDatabases = $getDatabasesDB($database);
try {
$document = $dbForDatabase->decreaseDocumentAttribute(
$document = $dbForDatabases->decreaseDocumentAttribute(
collection: 'database_' . $database->getSequence() . '_collection_' . $collection->getSequence(),
id: $documentId,
attribute: $attribute,
@@ -81,14 +81,14 @@ class Increment extends Action
->param('transactionId', null, new UID(), 'Transaction ID for staging the operation.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForEvents')
->inject('queueForStatsUsage')
->inject('plan')
->callback($this->action(...));
}
public function action(string $databaseId, string $collectionId, string $documentId, string $attribute, int|float $value, int|float|null $max, ?string $transactionId, UtopiaResponse $response, Database $dbForProject, callable $getDatabaseDB, Event $queueForEvents, StatsUsage $queueForStatsUsage, array $plan): void
public function action(string $databaseId, string $collectionId, string $documentId, string $attribute, int|float $value, int|float|null $max, ?string $transactionId, UtopiaResponse $response, Database $dbForProject, callable $getDatabasesDB, Event $queueForEvents, StatsUsage $queueForStatsUsage, array $plan): void
{
$isAPIKey = Auth::isAppUser(Authorization::getRoles());
$isPrivilegedUser = Auth::isPrivilegedUser(Authorization::getRoles());
@@ -166,9 +166,9 @@ class Increment extends Action
return;
}
$dbForDatabase = call_user_func($getDatabaseDB, $database);
$dbForDatabases = $getDatabasesDB($database);
try {
$document = $dbForDatabase->increaseDocumentAttribute(
$document = $dbForDatabases->increaseDocumentAttribute(
collection: 'database_' . $database->getSequence() . '_collection_' . $collection->getSequence(),
id: $documentId,
attribute: $attribute,
@@ -74,7 +74,7 @@ class Delete extends Action
->param('transactionId', null, new UID(), 'Transaction ID for staging the operation.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForStatsUsage')
->inject('queueForEvents')
->inject('queueForRealtime')
@@ -84,7 +84,7 @@ class Delete extends Action
->callback($this->action(...));
}
public function action(string $databaseId, string $collectionId, array $queries, ?string $transactionId, UtopiaResponse $response, Database $dbForProject, callable $getDatabaseDB, StatsUsage $queueForStatsUsage, Event $queueForEvents, Event $queueForRealtime, Event $queueForFunctions, Event $queueForWebhooks, array $plan): void
public function action(string $databaseId, string $collectionId, array $queries, ?string $transactionId, UtopiaResponse $response, Database $dbForProject, callable $getDatabasesDB, StatsUsage $queueForStatsUsage, Event $queueForEvents, Event $queueForRealtime, Event $queueForFunctions, Event $queueForWebhooks, array $plan): void
{
$database = $dbForProject->getDocument('databases', $databaseId);
if ($database->isEmpty()) {
@@ -159,11 +159,11 @@ class Delete extends Action
return;
}
$dbForDatabase = call_user_func($getDatabaseDB, $database);
$dbForDatabases = $getDatabasesDB($database);
$documents = [];
try {
$modified = $dbForDatabase->deleteDocuments(
$modified = $dbForDatabases->deleteDocuments(
'database_' . $database->getSequence() . '_collection_' . $collection->getSequence(),
$queries,
onNext: function (Document $document) use ($plan, &$documents) {
@@ -78,7 +78,7 @@ class Update extends Action
->param('transactionId', null, new UID(), 'Transaction ID for staging the operation.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForStatsUsage')
->inject('queueForEvents')
->inject('queueForRealtime')
@@ -88,7 +88,7 @@ class Update extends Action
->callback($this->action(...));
}
public function action(string $databaseId, string $collectionId, string|array $data, array $queries, ?string $transactionId, UtopiaResponse $response, Database $dbForProject, callable $getDatabaseDB, StatsUsage $queueForStatsUsage, Event $queueForEvents, Event $queueForRealtime, Event $queueForFunctions, Event $queueForWebhooks, array $plan): void
public function action(string $databaseId, string $collectionId, string|array $data, array $queries, ?string $transactionId, UtopiaResponse $response, Database $dbForProject, callable $getDatabasesDB, StatsUsage $queueForStatsUsage, Event $queueForEvents, Event $queueForRealtime, Event $queueForFunctions, Event $queueForWebhooks, array $plan): void
{
$data = \is_string($data)
? \json_decode($data, true)
@@ -181,12 +181,12 @@ class Update extends Action
return;
}
$dbForDatabase = call_user_func($getDatabaseDB, $database);
$dbForDatabases = $getDatabasesDB($database);
$documents = [];
try {
$modified = $dbForDatabase->withPreserveDates(function () use ($plan, &$documents, $dbForDatabase, $database, $collection, $data, $queries) {
return $dbForDatabase->updateDocuments(
$modified = $dbForDatabases->withPreserveDates(function () use ($plan, &$documents, $dbForDatabases, $database, $collection, $data, $queries) {
return $dbForDatabases->updateDocuments(
'database_' . $database->getSequence() . '_collection_' . $collection->getSequence(),
new Document($data),
$queries,
@@ -76,7 +76,7 @@ class Upsert extends Action
->param('transactionId', null, new UID(), 'Transaction ID for staging the operation.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForStatsUsage')
->inject('queueForEvents')
->inject('queueForRealtime')
@@ -86,7 +86,7 @@ class Upsert extends Action
->callback($this->action(...));
}
public function action(string $databaseId, string $collectionId, array $documents, ?string $transactionId, UtopiaResponse $response, Database $dbForProject, callable $getDatabaseDB, StatsUsage $queueForStatsUsage, Event $queueForEvents, Event $queueForRealtime, Event $queueForFunctions, Event $queueForWebhooks, array $plan): void
public function action(string $databaseId, string $collectionId, array $documents, ?string $transactionId, UtopiaResponse $response, Database $dbForProject, callable $getDatabasesDB, StatsUsage $queueForStatsUsage, Event $queueForEvents, Event $queueForRealtime, Event $queueForFunctions, Event $queueForWebhooks, array $plan): void
{
$database = $dbForProject->getDocument('databases', $databaseId);
if ($database->isEmpty()) {
@@ -158,12 +158,12 @@ class Upsert extends Action
return;
}
$dbForDatabase = call_user_func($getDatabaseDB, $database);
$dbForDatabases = $getDatabasesDB($database);
$upserted = [];
try {
$modified = $dbForDatabase->withPreserveDates(function () use ($dbForDatabase, $database, $collection, $documents, $plan, &$upserted) {
return $dbForDatabase->upsertDocuments(
$modified = $dbForDatabases->withPreserveDates(function () use ($dbForDatabases, $database, $collection, $documents, $plan, &$upserted) {
return $dbForDatabases->upsertDocuments(
'database_' . $database->getSequence() . '_collection_' . $collection->getSequence(),
$documents,
onNext: function (Document $document) use ($plan, &$upserted) {
@@ -126,7 +126,7 @@ class Create extends Action
->param('transactionId', null, new UID(), 'Transaction ID for staging the operation.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('user')
->inject('queueForEvents')
->inject('queueForStatsUsage')
@@ -136,7 +136,7 @@ class Create extends Action
->inject('plan')
->callback($this->action(...));
}
public function action(string $databaseId, string $documentId, string $collectionId, string|array $data, ?array $permissions, ?array $documents, ?string $transactionId, UtopiaResponse $response, Database $dbForProject, callable $getDatabaseDB, Document $user, Event $queueForEvents, StatsUsage $queueForStatsUsage, Event $queueForRealtime, Event $queueForFunctions, Event $queueForWebhooks, array $plan): void
public function action(string $databaseId, string $documentId, string $collectionId, string|array $data, ?array $permissions, ?array $documents, ?string $transactionId, UtopiaResponse $response, Database $dbForProject, callable $getDatabasesDB, Document $user, Event $queueForEvents, StatsUsage $queueForStatsUsage, Event $queueForRealtime, Event $queueForFunctions, Event $queueForWebhooks, array $plan): void
{
$data = \is_string($data)
? \json_decode($data, true)
@@ -435,12 +435,12 @@ class Create extends Action
return;
}
$dbForDatabase = call_user_func($getDatabaseDB, $database);
$dbForDatabases = $getDatabasesDB($database);
try {
$created = [];
$dbForDatabase->withPreserveDates(
function () use (&$created, $dbForDatabase, $database, $collection, $documents) {
$dbForDatabase->createDocuments(
$dbForDatabases->withPreserveDates(
function () use (&$created, $dbForDatabases, $database, $collection, $documents) {
$dbForDatabases->createDocuments(
'database_' . $database->getSequence() . '_collection_' . $collection->getSequence(),
$documents,
onNext: function ($doc) use (&$created) {
@@ -78,7 +78,7 @@ class Delete extends Action
->inject('requestTimestamp')
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForEvents')
->inject('queueForStatsUsage')
->inject('transactionState')
@@ -96,7 +96,7 @@ class Delete extends Action
?\DateTime $requestTimestamp,
UtopiaResponse $response,
Database $dbForProject,
callable $getDatabaseDB,
callable $getDatabasesDB,
Event $queueForEvents,
StatsUsage $queueForStatsUsage,
TransactionState $transactionState,
@@ -117,7 +117,7 @@ class Delete extends Action
throw new Exception($this->getParentNotFoundException());
}
$dbForDatabase = call_user_func($getDatabaseDB, $database);
$dbForDatabases = $getDatabasesDB($database);
// Read permission should not be required for delete
$collectionTableId = 'database_' . $database->getSequence() . '_collection_' . $collection->getSequence();
@@ -125,7 +125,7 @@ class Delete extends Action
// Use transaction-aware document retrieval to see changes from same transaction
$document = $transactionState->getDocument($collectionTableId, $documentId, $transactionId);
} else {
$document = Authorization::skip(fn () => $dbForDatabase->getDocument($collectionTableId, $documentId));
$document = Authorization::skip(fn () => $dbForDatabases->getDocument($collectionTableId, $documentId));
}
if ($document->isEmpty()) {
@@ -187,8 +187,8 @@ class Delete extends Action
}
try {
$dbForDatabase->withRequestTimestamp($requestTimestamp, function () use ($dbForDatabase, $database, $collection, $documentId) {
$dbForDatabase->deleteDocument(
$dbForDatabases->withRequestTimestamp($requestTimestamp, function () use ($dbForDatabases, $database, $collection, $documentId) {
$dbForDatabases->deleteDocument(
'database_' . $database->getSequence() . '_collection_' . $collection->getSequence(),
$documentId
);
@@ -68,14 +68,14 @@ class Get extends Action
->param('transactionId', null, new UID(), 'Transaction ID to read uncommitted changes within the transaction.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForStatsUsage')
->inject('transactionState')
->inject('transactionState')
->callback($this->action(...));
}
public function action(string $databaseId, string $collectionId, string $documentId, array $queries, ?string $transactionId, UtopiaResponse $response, Database $dbForProject, callable $getDatabaseDB, StatsUsage $queueForStatsUsage, TransactionState $transactionState): void
public function action(string $databaseId, string $collectionId, string $documentId, array $queries, ?string $transactionId, UtopiaResponse $response, Database $dbForProject, callable $getDatabasesDB, StatsUsage $queueForStatsUsage, TransactionState $transactionState): void
{
$isAPIKey = Auth::isAppUser(Authorization::getRoles());
$isPrivilegedUser = Auth::isPrivilegedUser(Authorization::getRoles());
@@ -87,7 +87,7 @@ class Get extends Action
$collection = Authorization::skip(fn () => $dbForProject->getDocument('database_' . $database->getSequence(), $collectionId));
$dbForDatabase = call_user_func($getDatabaseDB, $database);
$dbForDatabases = $getDatabasesDB($database);
if ($collection->isEmpty() || (!$collection->getAttribute('enabled', false) && !$isAPIKey && !$isPrivilegedUser)) {
throw new Exception($this->getParentNotFoundException());
}
@@ -108,10 +108,10 @@ class Get extends Action
$document = $transactionState->getDocument($collectionTableId, $documentId, $transactionId, $queries);
} elseif (! empty($selects)) {
// has selects, allow relationship on documents!
$document = $dbForDatabase->getDocument($collectionTableId, $documentId, $queries);
$document = $dbForDatabases->getDocument($collectionTableId, $documentId, $queries);
} else {
// has no selects, disable relationship looping on documents!
$document = $dbForDatabase->skipRelationships(fn () => $dbForDatabase->getDocument($collectionTableId, $documentId, $queries));
$document = $dbForDatabases->skipRelationships(fn () => $dbForDatabases->getDocument($collectionTableId, $documentId, $queries));
}
} catch (QueryException $e) {
throw new Exception(Exception::GENERAL_QUERY_INVALID, $e->getMessage());
@@ -82,7 +82,7 @@ class Update extends Action
->inject('requestTimestamp')
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForEvents')
->inject('queueForStatsUsage')
->inject('transactionState')
@@ -92,7 +92,7 @@ class Update extends Action
->callback($this->action(...));
}
public function action(string $databaseId, string $collectionId, string $documentId, string|array $data, ?array $permissions, ?string $transactionId, ?\DateTime $requestTimestamp, UtopiaResponse $response, Database $dbForProject, callable $getDatabaseDB, Event $queueForEvents, StatsUsage $queueForStatsUsage, TransactionState $transactionState, array $plan): void
public function action(string $databaseId, string $collectionId, string $documentId, string|array $data, ?array $permissions, ?string $transactionId, ?\DateTime $requestTimestamp, UtopiaResponse $response, Database $dbForProject, callable $getDatabasesDB, Event $queueForEvents, StatsUsage $queueForStatsUsage, TransactionState $transactionState, array $plan): void
{
$data = (\is_string($data)) ? \json_decode($data, true) : $data; // Cast to JSON array
@@ -115,7 +115,7 @@ class Update extends Action
throw new Exception($this->getParentNotFoundException());
}
$dbForDatabase = call_user_func($getDatabaseDB, $database);
$dbForDatabases = $getDatabasesDB($database);
// Read permission should not be required for update
/** @var Document $document */
$collectionTableId = 'database_' . $database->getSequence() . '_collection_' . $collection->getSequence();
@@ -124,7 +124,7 @@ class Update extends Action
// Use transaction-aware document retrieval to see changes from same transaction
$document = $transactionState->getDocument($collectionTableId, $documentId, $transactionId);
} else {
$document = Authorization::skip(fn () => $dbForDatabase->getDocument($collectionTableId, $documentId));
$document = Authorization::skip(fn () => $dbForDatabases->getDocument($collectionTableId, $documentId));
}
if ($document->isEmpty()) {
@@ -311,9 +311,9 @@ class Update extends Action
try {
$document = $dbForDatabase->withRequestTimestamp(
$document = $dbForDatabases->withRequestTimestamp(
$requestTimestamp,
fn () => $dbForDatabase->withPreserveDates(fn () => $dbForDatabase->updateDocument(
fn () => $dbForDatabases->withPreserveDates(fn () => $dbForDatabases->updateDocument(
'database_' . $database->getSequence() . '_collection_' . $collection->getSequence(),
$document->getId(),
$newDocument
@@ -86,7 +86,7 @@ class Upsert extends Action
->inject('response')
->inject('user')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForEvents')
->inject('queueForStatsUsage')
->inject('transactionState')
@@ -96,7 +96,7 @@ class Upsert extends Action
->callback($this->action(...));
}
public function action(string $databaseId, string $collectionId, string $documentId, string|array $data, ?array $permissions, ?string $transactionId, ?\DateTime $requestTimestamp, UtopiaResponse $response, Document $user, Database $dbForProject, callable $getDatabaseDB, Event $queueForEvents, StatsUsage $queueForStatsUsage, TransactionState $transactionState, array $plan): void
public function action(string $databaseId, string $collectionId, string $documentId, string|array $data, ?array $permissions, ?string $transactionId, ?\DateTime $requestTimestamp, UtopiaResponse $response, Document $user, Database $dbForProject, callable $getDatabasesDB, Event $queueForEvents, StatsUsage $queueForStatsUsage, TransactionState $transactionState, array $plan): void
{
$data = (\is_string($data)) ? \json_decode($data, true) : $data; // Cast to JSON array
@@ -121,7 +121,7 @@ class Upsert extends Action
throw new Exception($this->getParentNotFoundException());
}
$dbForDatabase = call_user_func($getDatabaseDB, $database);
$dbForDatabases = $getDatabasesDB($database);
$allowedPermissions = [
Database::PERMISSION_READ,
Database::PERMISSION_UPDATE,
@@ -140,7 +140,7 @@ class Upsert extends Action
// Use transaction-aware document retrieval to see changes from same transaction
$oldDocument = $transactionState->getDocument($collectionTableId, $documentId, $transactionId);
} else {
$oldDocument = Authorization::skip(fn () => $dbForDatabase->getDocument($collectionTableId, $documentId));
$oldDocument = Authorization::skip(fn () => $dbForDatabases->getDocument($collectionTableId, $documentId));
}
if ($oldDocument->isEmpty()) {
if (!empty($user->getId())) {
@@ -182,7 +182,7 @@ class Upsert extends Action
$newDocument = new Document($data);
$operations = 0;
$setCollection = (function (Document $collection, Document $document) use ($isAPIKey, $isPrivilegedUser, &$setCollection, $dbForProject, $dbForDatabase, $database, &$operations) {
$setCollection = (function (Document $collection, Document $document) use ($isAPIKey, $isPrivilegedUser, &$setCollection, $dbForProject, $dbForDatabases, $database, &$operations) {
$operations++;
$relationships = \array_filter(
@@ -223,7 +223,7 @@ class Upsert extends Action
if ($relation instanceof Document) {
$relation = $this->removeReadonlyAttributes($relation, $isAPIKey || $isPrivilegedUser);
$oldDocument = Authorization::skip(fn () => $dbForDatabase->getDocument(
$oldDocument = Authorization::skip(fn () => $dbForDatabases->getDocument(
'database_' . $database->getSequence() . '_collection_' . $relatedCollection->getSequence(),
$relation->getId()
));
@@ -321,8 +321,8 @@ class Upsert extends Action
$upserted = [];
try {
$dbForDatabase->withPreserveDates(function () use (&$upserted, $dbForDatabase, $database, $collection, $newDocument) {
return $dbForDatabase->upsertDocuments(
$dbForDatabases->withPreserveDates(function () use (&$upserted, $dbForDatabases, $database, $collection, $newDocument) {
return $dbForDatabases->upsertDocuments(
'database_' . $database->getSequence() . '_collection_' . $collection->getSequence(),
[$newDocument],
onNext: function (Document $document) use (&$upserted) {
@@ -347,7 +347,7 @@ class Upsert extends Action
// For transactions, get the document with transaction changes applied
$upserted[0] = $transactionState->getDocument($collectionTableId, $documentId, $transactionId);
} else {
$upserted[0] = $dbForDatabase->getDocument($collectionTableId, $documentId);
$upserted[0] = $dbForDatabases->getDocument($collectionTableId, $documentId);
}
}
@@ -70,14 +70,14 @@ class XList extends Action
->param('transactionId', null, new UID(), 'Transaction ID to read uncommitted changes within the transaction.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForStatsUsage')
->inject('transactionState')
->inject('transactionState')
->callback($this->action(...));
}
public function action(string $databaseId, string $collectionId, array $queries, ?string $transactionId, UtopiaResponse $response, Database $dbForProject, callable $getDatabaseDB, StatsUsage $queueForStatsUsage, TransactionState $transactionState): void
public function action(string $databaseId, string $collectionId, array $queries, ?string $transactionId, UtopiaResponse $response, Database $dbForProject, callable $getDatabasesDB, StatsUsage $queueForStatsUsage, TransactionState $transactionState): void
{
$isAPIKey = Auth::isAppUser(Authorization::getRoles());
$isPrivilegedUser = Auth::isPrivilegedUser(Authorization::getRoles());
@@ -98,7 +98,7 @@ class XList extends Action
throw new Exception(Exception::GENERAL_QUERY_INVALID, $e->getMessage());
}
$dbForDatabase = call_user_func($getDatabaseDB, $database);
$dbForDatabases = $getDatabasesDB($database);
/**
* Get cursor document if there was a cursor query, we use array_filter and reset for reference $cursor to $queries
*/
@@ -116,7 +116,7 @@ class XList extends Action
$documentId = $cursor->getValue();
$cursorDocument = Authorization::skip(fn () => $dbForDatabase->getDocument('database_' . $database->getSequence() . '_collection_' . $collection->getSequence(), $documentId));
$cursorDocument = Authorization::skip(fn () => $dbForDatabases->getDocument('database_' . $database->getSequence() . '_collection_' . $collection->getSequence(), $documentId));
if ($cursorDocument->isEmpty()) {
$type = ucfirst($this->getContext());
@@ -135,13 +135,13 @@ class XList extends Action
$total = $transactionState->countDocuments($collectionTableId, $transactionId, $queries);
} elseif (! empty($selectQueries)) {
// has selects, allow relationship on documents
$documents = $dbForDatabase->find($collectionTableId, $queries);
$total = $dbForDatabase->count($collectionTableId, $queries, APP_LIMIT_COUNT);
$documents = $dbForDatabases->find($collectionTableId, $queries);
$total = $dbForDatabases->count($collectionTableId, $queries, APP_LIMIT_COUNT);
} else {
// has no selects, disable relationship loading on documents
/* @type Document[] $documents */
$documents = $dbForDatabase->skipRelationships(fn () => $dbForDatabase->find($collectionTableId, $queries));
$total = $dbForDatabase->count($collectionTableId, $queries, APP_LIMIT_COUNT);
$documents = $dbForDatabases->skipRelationships(fn () => $dbForDatabases->find($collectionTableId, $queries));
$total = $dbForDatabases->count($collectionTableId, $queries, APP_LIMIT_COUNT);
}
} catch (OrderException $e) {
$documents = $this->isCollectionsAPI() ? 'documents' : 'rows';
@@ -77,13 +77,13 @@ class Create extends Action
->param('lengths', [], new ArrayList(new Nullable(new Integer()), APP_LIMIT_ARRAY_PARAMS_SIZE), 'Length of index. Maximum of ' . APP_LIMIT_ARRAY_PARAMS_SIZE, optional: true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForDatabase')
->inject('queueForEvents')
->callback($this->action(...));
}
public function action(string $databaseId, string $collectionId, string $key, string $type, array $attributes, array $orders, array $lengths, UtopiaResponse $response, Database $dbForProject, callable $getDatabaseDB, EventDatabase $queueForDatabase, Event $queueForEvents): void
public function action(string $databaseId, string $collectionId, string $key, string $type, array $attributes, array $orders, array $lengths, UtopiaResponse $response, Database $dbForProject, callable $getDatabasesDB, EventDatabase $queueForDatabase, Event $queueForEvents): void
{
$db = Authorization::skip(fn () => $dbForProject->getDocument('databases', $databaseId));
@@ -103,9 +103,9 @@ class Create extends Action
Query::equal('databaseInternalId', [$db->getSequence()])
], 61);
$dbForDatabase = call_user_func($getDatabaseDB, $db);
$dbForDatabases = $getDatabasesDB($db);
$limit = $dbForDatabase->getLimitForIndexes();
$limit = $dbForDatabases->getLimitForIndexes();
if ($count >= $limit) {
throw new Exception($this->getLimitException(), 'Index limit exceeded');
@@ -146,7 +146,7 @@ class Create extends Action
];
$contextType = $this->getParentContext();
if ($dbForDatabase->getAdapter()->getSupportForAttributes()) {
if ($dbForDatabases->getAdapter()->getSupportForAttributes()) {
foreach ($attributes as $i => $attribute) {
// find attribute metadata in collection document
$attributeIndex = \array_search($attribute, array_column($oldAttributes, 'key'));
@@ -210,9 +210,9 @@ class Create extends Action
$supportForSpatialAttributes,
$supportForSpatialIndexNull,
$supportForSpatialIndexOrder,
$dbForDatabase->getAdapter()->getSupportForAttributes(),
$dbForDatabase->getAdapter()->getSupportForMultipleFulltextIndexes(),
$dbForDatabase->getAdapter()->getSupportForIdenticalIndexes(),
$dbForDatabases->getAdapter()->getSupportForAttributes(),
$dbForDatabases->getAdapter()->getSupportForMultipleFulltextIndexes(),
$dbForDatabases->getAdapter()->getSupportForIdenticalIndexes(),
);
if (!$validator->isValid($index)) {
@@ -69,12 +69,12 @@ class Update extends Action
->param('enabled', true, new Boolean(), 'Is collection enabled? When set to \'disabled\', users cannot access the collection but Server SDKs with and API key can still read and write to the collection. No data is lost when this is toggled.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForEvents')
->callback($this->action(...));
}
public function action(string $databaseId, string $collectionId, string $name, ?array $permissions, bool $documentSecurity, bool $enabled, UtopiaResponse $response, Database $dbForProject, callable $getDatabaseDB, Event $queueForEvents): void
public function action(string $databaseId, string $collectionId, string $name, ?array $permissions, bool $documentSecurity, bool $enabled, UtopiaResponse $response, Database $dbForProject, callable $getDatabasesDB, Event $queueForEvents): void
{
$database = Authorization::skip(fn () => $dbForProject->getDocument('databases', $databaseId));
if ($database->isEmpty()) {
@@ -104,8 +104,8 @@ class Update extends Action
->setAttribute('search', \implode(' ', [$collectionId, $name]))
);
$dbForDatabase = call_user_func($getDatabaseDB, $database);
$dbForDatabase->updateCollection('database_' . $database->getSequence() . '_collection_' . $collection->getSequence(), $permissions, $documentSecurity);
$dbForDatabases = $getDatabasesDB($database);
$dbForDatabases->updateCollection('database_' . $database->getSequence() . '_collection_' . $collection->getSequence(), $permissions, $documentSecurity);
$queueForEvents
->setContext('database', $database)
@@ -63,16 +63,16 @@ class Get extends Action
->param('collectionId', '', new UID(), 'Collection ID.')
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->callback($this->action(...));
}
public function action(string $databaseId, string $range, string $collectionId, UtopiaResponse $response, Database $dbForProject, callable $getDatabaseDB): void
public function action(string $databaseId, string $range, string $collectionId, UtopiaResponse $response, Database $dbForProject, callable $getDatabasesDB): void
{
$database = $dbForProject->getDocument('databases', $databaseId);
$collectionDocument = $dbForProject->getDocument('database_' . $database->getSequence(), $collectionId);
$dbForDatabase = call_user_func($getDatabaseDB, $database);
$collection = $dbForDatabase->getCollection('database_' . $database->getSequence() . '_collection_' . $collectionDocument->getSequence());
$dbForDatabases = $getDatabasesDB($database);
$collection = $dbForDatabases->getCollection('database_' . $database->getSequence() . '_collection_' . $collectionDocument->getSequence());
if ($collection->isEmpty()) {
throw new Exception($this->getNotFoundException());
@@ -61,7 +61,7 @@ class Create extends CollectionCreate
->param('enabled', true, new Boolean(), 'Is collection enabled? When set to \'disabled\', users cannot access the collection but Server SDKs with and API key can still read and write to the collection. No data is lost when this is toggled.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForEvents')
->callback($this->action(...));
}
@@ -53,7 +53,7 @@ class Delete extends CollectionDelete
->param('collectionId', '', new UID(), 'Collection ID.')
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForDatabase')
->inject('queueForEvents')
->callback($this->action(...));
@@ -63,7 +63,7 @@ class Decrement extends DecrementDocumentAttribute
->param('transactionId', null, new UID(), 'Transaction ID for staging the operation.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForEvents')
->inject('queueForStatsUsage')
->inject('plan')
@@ -63,7 +63,7 @@ class Increment extends IncrementDocumentAttribute
->param('transactionId', null, new UID(), 'Transaction ID for staging the operation.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForEvents')
->inject('queueForStatsUsage')
->inject('plan')
@@ -59,7 +59,7 @@ class Delete extends DocumentsDelete
->param('transactionId', null, new UID(), 'Transaction ID for staging the operation.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForStatsUsage')
->inject('queueForEvents')
->inject('queueForRealtime')
@@ -61,7 +61,7 @@ class Update extends DocumentsUpdate
->param('transactionId', null, new UID(), 'Transaction ID for staging the operation.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForStatsUsage')
->inject('queueForEvents')
->inject('queueForRealtime')
@@ -61,7 +61,7 @@ class Upsert extends DocumentsUpsert
->param('transactionId', null, new UID(), 'Transaction ID for staging the operation.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForStatsUsage')
->inject('queueForEvents')
->inject('queueForRealtime')
@@ -101,7 +101,7 @@ class Create extends DocumentCreate
->param('transactionId', null, new UID(), 'Transaction ID for staging the operation.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('user')
->inject('queueForEvents')
->inject('queueForStatsUsage')
@@ -65,7 +65,7 @@ class Delete extends DocumentDelete
->inject('requestTimestamp')
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForEvents')
->inject('queueForStatsUsage')
->inject('transactionState')
@@ -55,7 +55,7 @@ class Get extends DocumentGet
->param('transactionId', null, new UID(), 'Transaction ID to read uncommitted changes within the transaction.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForStatsUsage')
->inject('transactionState')
->callback($this->action(...));
@@ -64,7 +64,7 @@ class Update extends DocumentUpdate
->inject('requestTimestamp')
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForEvents')
->inject('queueForStatsUsage')
->inject('transactionState')
@@ -67,7 +67,7 @@ class Upsert extends DocumentUpsert
->inject('response')
->inject('user')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForEvents')
->inject('queueForStatsUsage')
->inject('transactionState')
@@ -54,7 +54,7 @@ class XList extends DocumentXList
->param('transactionId', null, new UID(), 'Transaction ID to read uncommitted changes within the transaction.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForStatsUsage')
->inject('transactionState')
->callback($this->action(...));
@@ -64,7 +64,7 @@ class Create extends IndexCreate
->param('lengths', [], new ArrayList(new Nullable(new Integer()), APP_LIMIT_ARRAY_PARAMS_SIZE), 'Length of index. Maximum of ' . APP_LIMIT_ARRAY_PARAMS_SIZE, optional: true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForDatabase')
->inject('queueForEvents')
->callback($this->action(...));
@@ -60,7 +60,7 @@ class Update extends CollectionUpdate
->param('enabled', true, new Boolean(), 'Is collection enabled? When set to \'disabled\', users cannot access the collection but Server SDKs with and API key can still read and write to the collection. No data is lost when this is toggled.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForEvents')
->callback($this->action(...));
}
@@ -52,7 +52,7 @@ class Get extends CollectionUsageGet
->param('collectionId', '', new UID(), 'Collection ID.')
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->callback($this->action(...));
}
}
@@ -61,7 +61,7 @@ class Create extends CollectionCreate
->param('enabled', true, new Boolean(), 'Is table enabled? When set to \'disabled\', users cannot access the table but Server SDKs with and API key can still read and write to the table. No data is lost when this is toggled.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForEvents')
->callback($this->action(...));
}
@@ -53,7 +53,7 @@ class Delete extends CollectionDelete
->param('tableId', '', new UID(), 'Table ID.')
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForDatabase')
->inject('queueForEvents')
->callback($this->action(...));
@@ -64,7 +64,7 @@ class Create extends IndexCreate
->param('lengths', [], new ArrayList(new Nullable(new Integer()), APP_LIMIT_ARRAY_PARAMS_SIZE), 'Length of index. Maximum of ' . APP_LIMIT_ARRAY_PARAMS_SIZE, optional: true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForDatabase')
->inject('queueForEvents')
->callback($this->action(...));
@@ -59,7 +59,7 @@ class Delete extends DocumentsDelete
->param('transactionId', null, new UID(), 'Transaction ID for staging the operation.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForStatsUsage')
->inject('queueForEvents')
->inject('queueForRealtime')
@@ -61,7 +61,7 @@ class Update extends DocumentsUpdate
->param('transactionId', null, new UID(), 'Transaction ID for staging the operation.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForStatsUsage')
->inject('queueForEvents')
->inject('queueForRealtime')
@@ -61,7 +61,7 @@ class Upsert extends DocumentsUpsert
->param('transactionId', null, new UID(), 'Transaction ID for staging the operation.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForStatsUsage')
->inject('queueForEvents')
->inject('queueForRealtime')
@@ -63,7 +63,7 @@ class Decrement extends DecrementDocumentAttribute
->param('transactionId', null, new UID(), 'Transaction ID for staging the operation.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForEvents')
->inject('queueForStatsUsage')
->inject('plan')
@@ -63,7 +63,7 @@ class Increment extends IncrementDocumentAttribute
->param('transactionId', null, new UID(), 'Transaction ID for staging the operation.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForEvents')
->inject('queueForStatsUsage')
->inject('plan')
@@ -103,7 +103,7 @@ class Create extends DocumentCreate
->param('transactionId', null, new UID(), 'Transaction ID for staging the operation.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('user')
->inject('queueForEvents')
->inject('queueForStatsUsage')
@@ -65,7 +65,7 @@ class Delete extends DocumentDelete
->inject('requestTimestamp')
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForEvents')
->inject('queueForStatsUsage')
->inject('transactionState')
@@ -55,7 +55,7 @@ class Get extends DocumentGet
->param('transactionId', null, new UID(), 'Transaction ID to read uncommitted changes within the transaction.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForStatsUsage')
->inject('transactionState')
->callback($this->action(...));
@@ -64,7 +64,7 @@ class Update extends DocumentUpdate
->inject('requestTimestamp')
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForEvents')
->inject('queueForStatsUsage')
->inject('transactionState')
@@ -67,7 +67,7 @@ class Upsert extends DocumentUpsert
->inject('response')
->inject('user')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForEvents')
->inject('queueForStatsUsage')
->inject('transactionState')
@@ -54,7 +54,7 @@ class XList extends DocumentXList
->param('transactionId', null, new UID(), 'Transaction ID to read uncommitted changes within the transaction.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForStatsUsage')
->inject('transactionState')
->callback($this->action(...));
@@ -60,7 +60,7 @@ class Update extends CollectionUpdate
->param('enabled', true, new Boolean(), 'Is table enabled? When set to \'disabled\', users cannot access the table but Server SDKs with and API key can still read and write to the table. No data is lost when this is toggled.', true)
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForEvents')
->callback($this->action(...));
}
@@ -52,7 +52,7 @@ class Get extends CollectionUsageGet
->param('tableId', '', new UID(), 'Table ID.')
->inject('response')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->callback($this->action(...));
}
}
@@ -36,7 +36,7 @@ class Databases extends Action
->inject('project')
->inject('dbForPlatform')
->inject('dbForProject')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('queueForRealtime')
->inject('log')
->callback($this->action(...));
@@ -52,7 +52,7 @@ class Databases extends Action
* @return void
* @throws \Exception
*/
public function action(Message $message, Document $project, Database $dbForPlatform, Database $dbForProject, callable $getDatabaseDB, Realtime $queueForRealtime, Log $log): void
public function action(Message $message, Document $project, Database $dbForPlatform, Database $dbForProject, callable $getDatabasesDB, Realtime $queueForRealtime, Log $log): void
{
$payload = $message->getPayload() ?? [];
@@ -65,9 +65,9 @@ class Databases extends Action
$collection = new Document($payload['table'] ?? $payload['collection'] ?? []);
$database = new Document($payload['database'] ?? []);
/**
* @var Database $dbForDatabase
* @var Database $dbForDatabases
*/
$dbForDatabase = call_user_func($getDatabaseDB, $database);
$dbForDatabases = $getDatabasesDB($database);
$log->addTag('projectId', $project->getId());
$log->addTag('type', $type);
@@ -78,12 +78,12 @@ class Databases extends Action
$log->addTag('databaseId', $database->getId());
match (\strval($type)) {
DATABASE_TYPE_DELETE_DATABASE => $this->deleteDatabase($database, $dbForProject, $dbForDatabase),
DATABASE_TYPE_DELETE_COLLECTION => $this->deleteCollection($database, $collection, $dbForProject, $dbForDatabase),
DATABASE_TYPE_DELETE_DATABASE => $this->deleteDatabase($database, $dbForProject, $dbForDatabases),
DATABASE_TYPE_DELETE_COLLECTION => $this->deleteCollection($database, $collection, $dbForProject, $dbForDatabases),
DATABASE_TYPE_CREATE_ATTRIBUTE => $this->createAttribute($database, $collection, $document, $project, $dbForPlatform, $dbForProject, $queueForRealtime),
DATABASE_TYPE_DELETE_ATTRIBUTE => $this->deleteAttribute($database, $collection, $document, $project, $dbForPlatform, $dbForProject, $dbForDatabase, $queueForRealtime),
DATABASE_TYPE_CREATE_INDEX => $this->createIndex($database, $collection, $document, $project, $dbForPlatform, $dbForProject, $dbForDatabase, $queueForRealtime),
DATABASE_TYPE_DELETE_INDEX => $this->deleteIndex($database, $collection, $document, $project, $dbForPlatform, $dbForProject, $dbForDatabase, $queueForRealtime),
DATABASE_TYPE_DELETE_ATTRIBUTE => $this->deleteAttribute($database, $collection, $document, $project, $dbForPlatform, $dbForProject, $dbForDatabases, $queueForRealtime),
DATABASE_TYPE_CREATE_INDEX => $this->createIndex($database, $collection, $document, $project, $dbForPlatform, $dbForProject, $dbForDatabases, $queueForRealtime),
DATABASE_TYPE_DELETE_INDEX => $this->deleteIndex($database, $collection, $document, $project, $dbForPlatform, $dbForProject, $dbForDatabases, $queueForRealtime),
default => throw new Exception('No database operation for type: ' . \strval($type)),
};
@@ -236,7 +236,7 @@ class Databases extends Action
* @param Document $project
* @param Database $dbForPlatform
* @param Database $dbForProject
* @param Database $dbForDatabase
* @param Database $dbForDatabases
* @param Realtime $queueForRealtime
* @return void
* @throws Authorization
@@ -244,7 +244,7 @@ class Databases extends Action
* @throws \Exception
* @throws \Throwable
**/
private function deleteAttribute(Document $database, Document $collection, Document $attribute, Document $project, Database $dbForPlatform, Database $dbForDatabase, Database $dbForProject, Realtime $queueForRealtime): void
private function deleteAttribute(Document $database, Document $collection, Document $attribute, Document $project, Database $dbForPlatform, Database $dbForDatabases, Database $dbForProject, Realtime $queueForRealtime): void
{
if ($collection->isEmpty()) {
throw new Exception('Missing collection/table');
@@ -373,7 +373,7 @@ class Databases extends Action
}
if ($exists) { // Delete the duplicate if created, else update in db
$this->deleteIndex($database, $collection, $index, $project, $dbForPlatform, $dbForProject, $dbForDatabase, $queueForRealtime);
$this->deleteIndex($database, $collection, $index, $project, $dbForPlatform, $dbForProject, $dbForDatabases, $queueForRealtime);
} else {
$dbForProject->updateDocument('indexes', $index->getId(), $index);
}
@@ -398,7 +398,7 @@ class Databases extends Action
* @param Document $project
* @param Database $dbForPlatform
* @param Database $dbForProject
* @param Database $dbForDatabase
* @param Database $dbForDatabases
* @param Realtime $queueForRealtime
* @return void
* @throws Authorization
@@ -407,7 +407,7 @@ class Databases extends Action
* @throws DatabaseException
* @throws \Throwable
*/
private function createIndex(Document $database, Document $collection, Document $index, Document $project, Database $dbForPlatform, Database $dbForProject, Database $dbForDatabase, Realtime $queueForRealtime): void
private function createIndex(Document $database, Document $collection, Document $index, Document $project, Database $dbForPlatform, Database $dbForProject, Database $dbForDatabases, Realtime $queueForRealtime): void
{
if ($collection->isEmpty()) {
throw new Exception('Missing collection/table');
@@ -427,7 +427,7 @@ class Databases extends Action
$project = $dbForPlatform->getDocument('projects', $projectId);
try {
if (!$dbForDatabase->createIndex('database_' . $database->getSequence() . '_collection_' . $collection->getSequence(), $key, $type, $attributes, $lengths, $orders)) {
if (!$dbForDatabases->createIndex('database_' . $database->getSequence() . '_collection_' . $collection->getSequence(), $key, $type, $attributes, $lengths, $orders)) {
throw new DatabaseException('Failed to create Index');
}
$dbForProject->updateDocument('indexes', $index->getId(), $index->setAttribute('status', 'available'));
@@ -457,7 +457,7 @@ class Databases extends Action
* @param Document $project
* @param Database $dbForPlatform
* @param Database $dbForProject
* @param Database $dbForDatabase
* @param Database $dbForDatabases
* @param Realtime $queueForRealtime
* @return void
* @throws Authorization
@@ -466,7 +466,7 @@ class Databases extends Action
* @throws DatabaseException
* @throws \Throwable
*/
private function deleteIndex(Document $database, Document $collection, Document $index, Document $project, Database $dbForPlatform, Database $dbForProject, Database $dbForDatabase, Realtime $queueForRealtime): void
private function deleteIndex(Document $database, Document $collection, Document $index, Document $project, Database $dbForPlatform, Database $dbForProject, Database $dbForDatabases, Realtime $queueForRealtime): void
{
if ($collection->isEmpty()) {
throw new Exception('Missing collection/table');
@@ -482,7 +482,7 @@ class Databases extends Action
$project = $dbForPlatform->getDocument('projects', $projectId);
try {
if ($status !== 'failed' && !$dbForDatabase->deleteIndex('database_' . $database->getSequence() . '_collection_' . $collection->getSequence(), $key)) {
if ($status !== 'failed' && !$dbForDatabases->deleteIndex('database_' . $database->getSequence() . '_collection_' . $collection->getSequence(), $key)) {
throw new DatabaseException('Failed to delete index');
}
$dbForProject->deleteDocument('indexes', $index->getId());
@@ -511,14 +511,14 @@ class Databases extends Action
/**
* @param Document $database
* @param Database $dbForProject
* @param Database $dbForDatabase
* @param Database $dbForDatabases
* @return void
* @throws Exception
*/
protected function deleteDatabase(Document $database, Database $dbForProject, Database $dbForDatabase): void
protected function deleteDatabase(Document $database, Database $dbForProject, Database $dbForDatabases): void
{
$this->deleteByGroup('database_' . $database->getSequence(), [], $dbForProject, function ($collection) use ($database, $dbForProject, $dbForDatabase) {
$this->deleteCollection($database, $collection, $dbForProject, $dbForDatabase);
$this->deleteByGroup('database_' . $database->getSequence(), [], $dbForProject, function ($collection) use ($database, $dbForProject, $dbForDatabases) {
$this->deleteCollection($database, $collection, $dbForProject, $dbForDatabases);
});
$dbForProject->deleteCollection('database_' . $database->getSequence());
@@ -528,7 +528,7 @@ class Databases extends Action
* @param Document $database
* @param Document $collection
* @param Database $dbForProject
* @param Database $dbForDatabase
* @param Database $dbForDatabases
* @return void
* @throws Authorization
* @throws Conflict
@@ -537,7 +537,7 @@ class Databases extends Action
* @throws Structure
* @throws Exception
*/
protected function deleteCollection(Document $database, Document $collection, Database $dbForProject, Database $dbForDatabase): void
protected function deleteCollection(Document $database, Document $collection, Database $dbForProject, Database $dbForDatabases): void
{
if ($collection->isEmpty()) {
throw new Exception('Missing collection/table');
@@ -547,7 +547,7 @@ class Databases extends Action
$collectionInternalId = $collection->getSequence();
$databaseInternalId = $database->getSequence();
$dbForDatabase->deleteCollection('database_' . $databaseInternalId . '_collection_' . $collection->getSequence());
$dbForDatabases->deleteCollection('database_' . $databaseInternalId . '_collection_' . $collection->getSequence());
/**
* Related collections relating to current collection
+5 -5
View File
@@ -41,7 +41,7 @@ class Migrations extends Action
/**
* @var callable(string $databaseDSN): Database
*/
protected mixed $getDatabaseDB;
protected mixed $getDatabasesDB;
/**
@@ -72,7 +72,7 @@ class Migrations extends Action
->inject('project')
->inject('dbForProject')
->inject('dbForPlatform')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('logError')
->inject('queueForRealtime')
->inject('deviceForImports')
@@ -82,11 +82,11 @@ class Migrations extends Action
/**
* @throws Exception
*/
public function action(Message $message, Document $project, Database $dbForProject, Database $dbForPlatform, callable $getDatabaseDB, callable $logError, Realtime $queueForRealtime, Device $deviceForImports): void
public function action(Message $message, Document $project, Database $dbForProject, Database $dbForPlatform, callable $getDatabasesDB, callable $logError, Realtime $queueForRealtime, Device $deviceForImports): void
{
$payload = $message->getPayload() ?? [];
$this->deviceForImports = $deviceForImports;
$this->getDatabaseDB = $getDatabaseDB;
$this->getDatabasesDB = $getDatabasesDB;
if (empty($payload)) {
throw new Exception('Missing payload');
@@ -178,7 +178,7 @@ class Migrations extends Action
'http://appwrite/v1',
$apiKey,
$this->dbForProject,
$this->getDatabaseDB,
$this->getDatabasesDB,
Config::getParam('collections', [])['databases']['collections'],
),
default => throw new \Exception('Invalid destination type'),
@@ -47,7 +47,7 @@ class StatsResources extends Action
->inject('project')
->inject('getProjectDB')
->inject('getLogsDB')
->inject('getDatabaseDB')
->inject('getDatabasesDB')
->inject('dbForPlatform')
->inject('logError')
->callback($this->action(...));
@@ -58,12 +58,12 @@ class StatsResources extends Action
* @param Document $project
* @param callable $getProjectDB
* @param callable $getLogsDB
* @param callable $getDatabaseDB
* @param callable $getDatabasesDB
* @return void
* @throws \Utopia\Database\Exception
* @throws Exception
*/
public function action(Message $message, Document $project, callable $getProjectDB, callable $getLogsDB, callable $getDatabaseDB, Database $dbForPlatform, callable $logError): void
public function action(Message $message, Document $project, callable $getProjectDB, callable $getLogsDB, callable $getDatabasesDB, Database $dbForPlatform, callable $logError): void
{
$this->logError = $logError;
@@ -80,13 +80,13 @@ class StatsResources extends Action
$this->documents = [];
$startTime = microtime(true);
$this->countForProject($dbForPlatform, $getLogsDB, $getProjectDB, $getDatabaseDB, $project);
$this->countForProject($dbForPlatform, $getLogsDB, $getProjectDB, $getDatabasesDB, $project);
$endTime = microtime(true);
$executionTime = $endTime - $startTime;
Console::info('Project: ' . $project->getId() . '(' . $project->getSequence() . ') aggregated in ' . $executionTime .' seconds');
}
protected function countForProject(Database $dbForPlatform, callable $getLogsDB, callable $getProjectDB, callable $getDatabaseDB, Document $project): void
protected function countForProject(Database $dbForPlatform, callable $getLogsDB, callable $getProjectDB, callable $getDatabasesDB, Document $project): void
{
Console::info('Begining count for: ' . $project->getId());
@@ -195,7 +195,7 @@ class StatsResources extends Action
}
try {
$this->countForDatabase($dbForProject, $getDatabaseDB, $region);
$this->countForDatabase($dbForProject, $getDatabasesDB, $region);
} catch (Throwable $th) {
call_user_func_array($this->logError, [$th, "StatsResources", "count_for_database_{$project->getId()}"]);
}
@@ -255,21 +255,21 @@ class StatsResources extends Action
$this->createStatsDocuments($region, METRIC_FILES_IMAGES_TRANSFORMED, $totalImageTransformations);
}
protected function countForDatabase(Database $dbForProject, callable $getDatabaseDB, string $region)
protected function countForDatabase(Database $dbForProject, callable $getDatabasesDB, string $region)
{
$totalCollections = 0;
$totalDocuments = 0;
$totalDatabaseStorage = 0;
$this->foreachDocument($dbForProject, 'databases', [], function ($database) use ($dbForProject, $getDatabaseDB, $region, &$totalCollections, &$totalDocuments, &$totalDatabaseStorage) {
$dbForDatabase = call_user_func($getDatabaseDB, $database);
$this->foreachDocument($dbForProject, 'databases', [], function ($database) use ($dbForProject, $getDatabasesDB, $region, &$totalCollections, &$totalDocuments, &$totalDatabaseStorage) {
$dbForDatabases = $getDatabasesDB($database);
$collections = $dbForProject->count('database_' . $database->getSequence());
$metric = str_replace('{databaseInternalId}', $database->getSequence(), METRIC_DATABASE_ID_COLLECTIONS);
$this->createStatsDocuments($region, $metric, $collections);
[$documents, $storage] = $this->countForCollections($dbForProject, $dbForDatabase, $database, $region);
[$documents, $storage] = $this->countForCollections($dbForProject, $dbForDatabases, $database, $region);
$totalDatabaseStorage += $storage;
$totalDocuments += $documents;
@@ -280,17 +280,17 @@ class StatsResources extends Action
$this->createStatsDocuments($region, METRIC_DOCUMENTS, $totalDocuments);
$this->createStatsDocuments($region, METRIC_DATABASES_STORAGE, $totalDatabaseStorage);
}
protected function countForCollections(Database $dbForProject, Database $dbForDatabase, Document $database, string $region): array
protected function countForCollections(Database $dbForProject, Database $dbForDatabases, Document $database, string $region): array
{
$databaseDocuments = 0;
$databaseStorage = 0;
$this->foreachDocument($dbForProject, 'database_' . $database->getSequence(), [], function ($collection) use ($dbForProject, $dbForDatabase, $database, $region, &$databaseStorage, &$databaseDocuments) {
$documents = $dbForDatabase->count('database_' . $database->getSequence() . '_collection_' . $collection->getSequence());
$this->foreachDocument($dbForProject, 'database_' . $database->getSequence(), [], function ($collection) use ($dbForProject, $dbForDatabases, $database, $region, &$databaseStorage, &$databaseDocuments) {
$documents = $dbForDatabases->count('database_' . $database->getSequence() . '_collection_' . $collection->getSequence());
$metric = str_replace(['{databaseInternalId}', '{collectionInternalId}'], [$database->getSequence(), $collection->getSequence()], METRIC_DATABASE_ID_COLLECTION_ID_DOCUMENTS);
$this->createStatsDocuments($region, $metric, $documents);
$databaseDocuments += $documents;
$collectionStorage = $dbForDatabase->getSizeOfCollection('database_' . $database->getSequence() . '_collection_' . $collection->getSequence());
$collectionStorage = $dbForDatabases->getSizeOfCollection('database_' . $database->getSequence() . '_collection_' . $collection->getSequence());
$metric = str_replace(['{databaseInternalId}', '{collectionInternalId}'], [$database->getSequence(), $collection->getSequence()], METRIC_DATABASE_ID_COLLECTION_ID_STORAGE);
$this->createStatsDocuments($region, $metric, $collectionStorage);
$databaseStorage += $collectionStorage;