Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
84 changes: 60 additions & 24 deletions src/Database/Adapter/Postgres.php
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
use Utopia\Database\Exception\Unique as UniqueException;
use Utopia\Database\Helpers\ID;
use Utopia\Database\Operator;
use Utopia\Database\PDO as PDOWrapper;
use Utopia\Database\Query;

/**
Expand All @@ -32,6 +33,14 @@ class Postgres extends SQL
{
public const MAX_IDENTIFIER_NAME = 63;

/**
* Last statement_timeout sent on each physical connection, shared by every
* adapter using that connection and dropped together with it.
*
* @var \WeakMap<object, int>|null
*/
private static ?\WeakMap $timeouts = null;

/**
* @inheritDoc
*/
Expand Down Expand Up @@ -66,25 +75,53 @@ protected function execute(mixed $stmt): bool
{
$pdo = $this->getPDO();

// Choose the right SET command based on transaction state
$sql = $this->inTransaction === 0
? "SET statement_timeout = '{$this->timeout}ms'"
: "SET LOCAL statement_timeout = '{$this->timeout}ms'";
// A persistent connection is shared by every PDO opened with the same DSN and
// outlives them, so its timeout is set and reset around each statement instead
if ($pdo->getAttribute(PDO::ATTR_PERSISTENT)) {
$sql = $this->inTransaction === 0
? "SET statement_timeout = '{$this->timeout}ms'"
: "SET LOCAL statement_timeout = '{$this->timeout}ms'";

// Apply timeout
$pdo->exec($sql);
$pdo->exec($sql);

try {
return $stmt->execute();
} finally {
// Only reset the global timeout when not in a transaction
if ($this->inTransaction === 0) {
$pdo->exec("RESET statement_timeout");
try {
return $stmt->execute();
} finally {
if ($this->inTransaction === 0) {
$pdo->exec("RESET statement_timeout");
}
}
}

$connection = $pdo instanceof PDOWrapper ? $pdo->getConnection() : $pdo;
$timeouts = self::$timeouts ??= new \WeakMap();

// Only send the timeout when it differs from what the connection already has
if (($timeouts[$connection] ?? null) !== $this->timeout) {
// Another adapter sharing the connection may have a transaction open, and
// its rollback would undo a session SET
if ($this->inTransaction === 0 && !$pdo->inTransaction()) {
$pdo->exec("SET statement_timeout = '{$this->timeout}ms'");
$timeouts[$connection] = $this->timeout;
} else {
// SET LOCAL ends with the transaction, so the session value is unknown until set again
$pdo->exec("SET LOCAL statement_timeout = '{$this->timeout}ms'");
unset($timeouts[$connection]);
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}

return $stmt->execute();
}

public function reconnect(): void
{
$pdo = $this->getPDO();
if (self::$timeouts !== null) {
unset(self::$timeouts[$pdo instanceof PDOWrapper ? $pdo->getConnection() : $pdo]);
}

parent::reconnect();
}

/**
* Returns Max Execution Time
Expand Down Expand Up @@ -124,14 +161,13 @@ public function create(string $name): bool
$sql = "CREATE SCHEMA \"{$name}\"";
$sql = $this->trigger(Database::EVENT_DATABASE_CREATE, $sql);

$dbCreation = $this->getPDO()
->prepare($sql)
->execute();
$dbCreation = $this->execute($this->getPDO()
->prepare($sql));

// Enable extensions
$this->getPDO()->prepare('CREATE EXTENSION IF NOT EXISTS postgis')->execute();
$this->getPDO()->prepare('CREATE EXTENSION IF NOT EXISTS vector')->execute();
$this->getPDO()->prepare('CREATE EXTENSION IF NOT EXISTS pg_trgm')->execute();
$this->execute($this->getPDO()->prepare('CREATE EXTENSION IF NOT EXISTS postgis'));
$this->execute($this->getPDO()->prepare('CREATE EXTENSION IF NOT EXISTS vector'));
$this->execute($this->getPDO()->prepare('CREATE EXTENSION IF NOT EXISTS pg_trgm'));

$collation = "
CREATE COLLATION IF NOT EXISTS utf8_ci_ai (
Expand All @@ -140,7 +176,7 @@ public function create(string $name): bool
deterministic = false
)
";
$this->getPDO()->prepare($collation)->execute();
$this->execute($this->getPDO()->prepare($collation));
return $dbCreation;
}

Expand All @@ -159,7 +195,7 @@ public function delete(string $name): bool
$sql = "DROP SCHEMA IF EXISTS \"{$name}\" CASCADE";
$sql = $this->trigger(Database::EVENT_DATABASE_DELETE, $sql);

return $this->getPDO()->prepare($sql)->execute();
return $this->execute($this->getPDO()->prepare($sql));
}

/**
Expand Down Expand Up @@ -283,9 +319,9 @@ public function createCollection(string $name, array $attributes = [], array $in
$permissions = $this->trigger(Database::EVENT_COLLECTION_CREATE, $permissions);

try {
$this->getPDO()->prepare($collection)->execute();
$this->execute($this->getPDO()->prepare($collection));

$this->getPDO()->prepare($permissions)->execute();
$this->execute($this->getPDO()->prepare($permissions));

foreach ($indexes as $index) {
$indexId = $this->filter($index->getId());
Expand Down Expand Up @@ -413,7 +449,7 @@ public function deleteCollection(string $id): bool
$sql = $this->trigger(Database::EVENT_COLLECTION_DELETE, $sql);

try {
return $this->getPDO()->prepare($sql)->execute();
return $this->execute($this->getPDO()->prepare($sql));
} catch (PDOException $e) {
throw $this->processException($e);
}
Expand Down Expand Up @@ -923,7 +959,7 @@ public function createIndex(string $collection, string $id, string $type, array
$sql = $this->trigger(Database::EVENT_INDEX_CREATE, $sql);

try {
return $this->getPDO()->prepare($sql)->execute();
return $this->execute($this->getPDO()->prepare($sql));
} catch (PDOException $e) {
throw $this->processException($e);
}
Expand Down
44 changes: 20 additions & 24 deletions src/Database/Adapter/SQL.php
Original file line number Diff line number Diff line change
Expand Up @@ -215,7 +215,7 @@ public function exists(string $database, ?string $collection = null): bool
}

try {
$stmt->execute();
$this->execute($stmt);
$document = $stmt->fetchAll();
$stmt->closeCursor();
} catch (PDOException $e) {
Expand Down Expand Up @@ -266,9 +266,8 @@ public function createAttribute(string $collection, string $id, string $type, in
$sql = $this->trigger(Database::EVENT_ATTRIBUTE_CREATE, $sql);

try {
return $this->getPDO()
->prepare($sql)
->execute();
return $this->execute($this->getPDO()
->prepare($sql));
} catch (PDOException $e) {
throw $this->processException($e);
}
Expand Down Expand Up @@ -303,9 +302,8 @@ public function createAttributes(string $collection, array $attributes): bool
$sql = $this->trigger(Database::EVENT_ATTRIBUTE_CREATE, $sql);

try {
return $this->getPDO()
->prepare($sql)
->execute();
return $this->execute($this->getPDO()
->prepare($sql));
} catch (PDOException $e) {
throw $this->processException($e);
}
Expand All @@ -332,9 +330,8 @@ public function renameAttribute(string $collection, string $old, string $new): b
$sql = $this->trigger(Database::EVENT_ATTRIBUTE_UPDATE, $sql);

try {
return $this->getPDO()
->prepare($sql)
->execute();
return $this->execute($this->getPDO()
->prepare($sql));
} catch (PDOException $e) {
throw $this->processException($e);
}
Expand All @@ -357,9 +354,8 @@ public function deleteAttribute(string $collection, string $id, bool $array = fa
$sql = $this->trigger(Database::EVENT_ATTRIBUTE_DELETE, $sql);

try {
return $this->getPDO()
->prepare($sql)
->execute();
return $this->execute($this->getPDO()
->prepare($sql));
} catch (PDOException $e) {
throw $this->processException($e);
}
Expand Down Expand Up @@ -608,7 +604,7 @@ public function updateDocuments(Document $collection, Document $updates, array $
}

try {
$stmt->execute();
$this->execute($stmt);
} catch (PDOException $e) {
throw $this->processException($e);
}
Expand Down Expand Up @@ -644,7 +640,7 @@ public function updateDocuments(Document $collection, Document $updates, array $
$permissionsStmt->bindValue(':_tenant', $this->tenant);
}

$permissionsStmt->execute();
$this->execute($permissionsStmt);
$permissions = $permissionsStmt->fetchAll();
$permissionsStmt->closeCursor();

Expand Down Expand Up @@ -744,7 +740,7 @@ public function updateDocuments(Document $collection, Document $updates, array $
if ($this->sharedTables) {
$stmtRemovePermissions->bindValue(':_tenant', $this->tenant);
}
$stmtRemovePermissions->execute();
$this->execute($stmtRemovePermissions);
}

if (!empty($addQuery)) {
Expand All @@ -770,7 +766,7 @@ public function updateDocuments(Document $collection, Document $updates, array $
$stmtAddPermissions->bindValue(':_tenant', $this->tenant);
}

$stmtAddPermissions->execute();
$this->execute($stmtAddPermissions);
}
}

Expand Down Expand Up @@ -815,7 +811,7 @@ public function deleteDocuments(string $collection, array $sequences, array $per
$stmt->bindValue(':_tenant', $this->tenant);
}

if (!$stmt->execute()) {
if (!$this->execute($stmt)) {
throw new DatabaseException('Failed to delete documents');
}

Expand All @@ -838,7 +834,7 @@ public function deleteDocuments(string $collection, array $sequences, array $per
$stmtPermissions->bindValue(':_tenant', $this->tenant);
}

if (!$stmtPermissions->execute()) {
if (!$this->execute($stmtPermissions)) {
throw new DatabaseException('Failed to delete permissions');
}
}
Expand Down Expand Up @@ -904,7 +900,7 @@ public function getSequences(string $collection, array $documents): array
$stmt->bindValue($key, $value);
}

$stmt->execute();
$this->execute($stmt);
$sequences = $stmt->fetchAll(\PDO::FETCH_KEY_PAIR); // Fetch as [documentId => sequence]
$stmt->closeCursor();

Expand Down Expand Up @@ -2686,7 +2682,7 @@ public function upsertDocuments(
}

$stmt = $this->getUpsertStatement($name, $columns, $batchKeys, $regularAttributes, $bindValues, $attribute, []);
$stmt->execute();
$this->execute($stmt);
$stmt->closeCursor();
} else {
$groups = [];
Expand Down Expand Up @@ -2822,7 +2818,7 @@ public function upsertDocuments(
$operators
);

$stmt->execute();
$this->execute($stmt);
$stmt->closeCursor();
}
}
Expand Down Expand Up @@ -2888,7 +2884,7 @@ public function upsertDocuments(
foreach ($removeBindValues as $key => $value) {
$stmtRemovePermissions->bindValue($key, $value, $this->getPDOType($value));
}
$stmtRemovePermissions->execute();
$this->execute($stmtRemovePermissions);
}

if (!empty($addQueries)) {
Expand All @@ -2901,7 +2897,7 @@ public function upsertDocuments(
foreach ($addBindValues as $key => $value) {
$stmtAddPermissions->bindValue($key, $value, $this->getPDOType($value));
}
$stmtAddPermissions->execute();
$this->execute($stmtAddPermissions);
}
} catch (PDOException $e) {
throw $this->processException($e);
Expand Down
10 changes: 10 additions & 0 deletions src/Database/PDO.php
Original file line number Diff line number Diff line change
Expand Up @@ -115,6 +115,16 @@ public function __call(string $method, array $args): mixed
}
}

/**
* Get the underlying connection, which is replaced on reconnect
*
* @return \PDO
*/
public function getConnection(): \PDO
{
return $this->pdo;
}

/**
* Create a new connection to the database
*
Expand Down
9 changes: 3 additions & 6 deletions tests/unit/SQLGetDocumentTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -85,13 +85,10 @@ public function testUsesPostgresExecuteHook(): void
$pdo->expects($this->once())
->method('prepare')
->willReturn($statement);
$pdo->expects($this->exactly(2))
$pdo->expects($this->once())
->method('exec')
->withConsecutive(
["SET statement_timeout = '25ms'"],
['RESET statement_timeout']
)
->willReturnOnConsecutiveCalls(0, 0);
->with("SET statement_timeout = '25ms'")
->willReturn(0);

$adapter = new Postgres($pdo);
$adapter->setDatabase('database');
Expand Down
Loading