Skip to content

Commit 5a807d7

Browse files
abnegateclaude
andcommitted
refactor(transaction): own connection recovery in the adapter, drop Swoole PDOProxy dependency
The Swoole PDOProxy keeps its own transaction counter that is incremented on beginTransaction but only decremented on a *successful* commit/rollback, and is never reset on reconnect or on pool checkin (utopia-php/pools does not call its reset() hook). A connection lost mid-transaction therefore leaks the counter and poisons the pooled connection: every later startTransaction trusts the stale counter, rolls back a transaction the real connection no longer holds, and fails with "There is no active transaction". This produced a sustained write outage across all projects on cloud nyc3. Make the library self-sufficient for connection-loss recovery so consumers no longer need to wrap connections in a Swoole PDOProxy: - PDOStatement wraps prepared statements and transparently re-prepares on the reconnected PDO when the connection is lost at execution time, replaying bound params/attributes. Recovery is skipped inside a transaction, where it rethrows so withTransaction can replay the whole transaction from the start. - PDO::prepare() returns the wrapper; prepareNative() re-prepares raw on the reconnected connection. ERRMODE_EXCEPTION is enforced by default. - withTransaction() reconnects on a lost connection before replaying, so the retry runs on a fresh, transaction-less connection. - Transaction state now has a single source of truth (the real PDO via Utopia\Database\PDO::inTransaction()); there is no separate counter to desync. - Replace the Swoole\Database\PDOStatementProxy type hints in the SQL adapters. Stacked on #895. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
1 parent 240b957 commit 5a807d7

7 files changed

Lines changed: 344 additions & 10 deletions

File tree

‎src/Database/Adapter/Postgres.php‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,6 @@
55
use Exception;
66
use PDO;
77
use PDOException;
8-
use Swoole\Database\PDOStatementProxy;
98
use Utopia\Database\Database;
109
use Utopia\Database\Document;
1110
use Utopia\Database\Exception as DatabaseException;
@@ -18,6 +17,7 @@
1817
use Utopia\Database\Exception\Truncate as TruncateException;
1918
use Utopia\Database\Helpers\ID;
2019
use Utopia\Database\Operator;
20+
use Utopia\Database\PDOStatement;
2121
use Utopia\Database\Query;
2222

2323
/**
@@ -2764,12 +2764,12 @@ protected function getOperatorSQL(string $column, Operator $operator, int &$bind
27642764
* Bind operator parameters to statement
27652765
* Override to handle PostgreSQL-specific JSON binding
27662766
*
2767-
* @param \PDOStatement|PDOStatementProxy $stmt
2767+
* @param \PDOStatement|PDOStatement $stmt
27682768
* @param Operator $operator
27692769
* @param int &$bindIndex
27702770
* @return void
27712771
*/
2772-
protected function bindOperatorParams(\PDOStatement|PDOStatementProxy $stmt, Operator $operator, int &$bindIndex): void
2772+
protected function bindOperatorParams(\PDOStatement|PDOStatement $stmt, Operator $operator, int &$bindIndex): void
27732773
{
27742774
$method = $operator->getMethod();
27752775
$values = $operator->getValues();

‎src/Database/Adapter/SQL.php‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,6 @@
44

55
use Exception;
66
use PDOException;
7-
use Swoole\Database\PDOStatementProxy;
87
use Utopia\Database\Adapter;
98
use Utopia\Database\Change;
109
use Utopia\Database\Database;
@@ -16,6 +15,7 @@
1615
use Utopia\Database\Exception\Timeout as TimeoutException;
1716
use Utopia\Database\Exception\Transaction as TransactionException;
1817
use Utopia\Database\Operator;
18+
use Utopia\Database\PDOStatement;
1919
use Utopia\Database\Query;
2020

2121
abstract class SQL extends Adapter
@@ -1978,12 +1978,12 @@ abstract protected function getOperatorSQL(string $column, Operator $operator, i
19781978
/**
19791979
* Bind operator parameters to prepared statement
19801980
*
1981-
* @param \PDOStatement|PDOStatementProxy $stmt
1981+
* @param \PDOStatement|PDOStatement $stmt
19821982
* @param \Utopia\Database\Operator $operator
19831983
* @param int &$bindIndex
19841984
* @return void
19851985
*/
1986-
protected function bindOperatorParams(\PDOStatement|PDOStatementProxy $stmt, Operator $operator, int &$bindIndex): void
1986+
protected function bindOperatorParams(\PDOStatement|PDOStatement $stmt, Operator $operator, int &$bindIndex): void
19871987
{
19881988
$method = $operator->getMethod();
19891989
$values = $operator->getValues();

‎src/Database/Adapter/SQLite.php‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,6 @@
55
use Exception;
66
use PDO;
77
use PDOException;
8-
use Swoole\Database\PDOStatementProxy;
98
use Utopia\Database\Database;
109
use Utopia\Database\Document;
1110
use Utopia\Database\Exception as DatabaseException;
@@ -18,6 +17,7 @@
1817
use Utopia\Database\Exception\Truncate as TruncateException;
1918
use Utopia\Database\Helpers\ID;
2019
use Utopia\Database\Operator;
20+
use Utopia\Database\PDOStatement;
2121
use Utopia\Database\Query;
2222

2323
/**
@@ -2002,12 +2002,12 @@ private function getSupportForMathFunctions(): bool
20022002
* Bind operator parameters to statement
20032003
* Override to handle SQLite-specific operator bindings
20042004
*
2005-
* @param \PDOStatement|PDOStatementProxy $stmt
2005+
* @param \PDOStatement|PDOStatement $stmt
20062006
* @param Operator $operator
20072007
* @param int &$bindIndex
20082008
* @return void
20092009
*/
2010-
protected function bindOperatorParams(\PDOStatement|PDOStatementProxy $stmt, Operator $operator, int &$bindIndex): void
2010+
protected function bindOperatorParams(\PDOStatement|PDOStatement $stmt, Operator $operator, int &$bindIndex): void
20112011
{
20122012
$method = $operator->getMethod();
20132013

‎src/Database/PDO.php‎

Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,8 @@ public function __construct(
2626
protected ?string $password,
2727
protected array $config = []
2828
) {
29+
$this->config[\PDO::ATTR_ERRMODE] ??= \PDO::ERRMODE_EXCEPTION;
30+
2931
$this->pdo = new \PDO(
3032
$this->dsn,
3133
$this->username,
@@ -34,6 +36,42 @@ public function __construct(
3436
);
3537
}
3638

39+
/**
40+
* Prepare a statement, returning a wrapper that transparently re-prepares
41+
* itself on the underlying connection if that connection is lost before the
42+
* statement is executed.
43+
*
44+
* @param array<mixed> $options
45+
* @throws \Throwable
46+
*/
47+
public function prepare(string $query, array $options = []): PDOStatement
48+
{
49+
return new PDOStatement($this, $this->prepareNative($query, $options), $query, $options);
50+
}
51+
52+
/**
53+
* Prepare a raw \PDOStatement on the underlying connection. Used by
54+
* {@see PDOStatement} to re-prepare after a reconnect without re-wrapping
55+
* the result.
56+
*
57+
* A lost connection is not handled here: under emulated prepares this call
58+
* never reaches the server, so the loss surfaces at execution time and is
59+
* recovered by {@see PDOStatement}.
60+
*
61+
* @param array<mixed> $options
62+
* @throws \PDOException
63+
*/
64+
public function prepareNative(string $query, array $options = []): \PDOStatement
65+
{
66+
$statement = $this->pdo->prepare($query, $options);
67+
68+
if ($statement === false) {
69+
throw new \PDOException("Failed to prepare statement: {$query}");
70+
}
71+
72+
return $statement;
73+
}
74+
3775
/**
3876
* @param string $method
3977
* @param array<mixed> $args

‎src/Database/PDOStatement.php‎

Lines changed: 170 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,170 @@
1+
<?php
2+
3+
namespace Utopia\Database;
4+
5+
use Utopia\Console;
6+
7+
/**
8+
* Wraps a \PDOStatement so a connection lost at execution time is recovered
9+
* transparently: the owning PDO reconnects, the statement is re-prepared
10+
* against the fresh connection, previously bound parameters and attributes are
11+
* replayed, and the failed call is retried.
12+
*
13+
* Recovery is skipped while a transaction is open, because the uncommitted
14+
* state lives on the dead connection and cannot survive a reconnect. There the
15+
* call is rethrown so Adapter::withTransaction can roll back and replay the
16+
* whole transaction from the start.
17+
*
18+
* @mixin \PDOStatement
19+
*/
20+
class PDOStatement
21+
{
22+
/**
23+
* @var array<int|string, array{mixed, int}>
24+
*/
25+
private array $values = [];
26+
27+
/**
28+
* @var array<int|string, array{mixed, int, int, mixed}>
29+
*/
30+
private array $params = [];
31+
32+
/**
33+
* @var array<int|string, array{mixed, ?int, ?int, mixed}>
34+
*/
35+
private array $columns = [];
36+
37+
/**
38+
* @var array<int, mixed>
39+
*/
40+
private array $attributes = [];
41+
42+
/**
43+
* @var array<int|string, mixed>|null
44+
*/
45+
private ?array $fetchMode = null;
46+
47+
/**
48+
* @param array<mixed> $options
49+
*/
50+
public function __construct(
51+
private readonly PDO $pdo,
52+
private \PDOStatement $statement,
53+
private readonly string $query,
54+
private readonly array $options = [],
55+
) {
56+
}
57+
58+
public function __get(string $name): mixed
59+
{
60+
return $this->statement->{$name};
61+
}
62+
63+
public function __set(string $name, mixed $value): void
64+
{
65+
$this->statement->{$name} = $value;
66+
}
67+
68+
public function __isset(string $name): bool
69+
{
70+
return isset($this->statement->{$name});
71+
}
72+
73+
public function __unset(string $name): void
74+
{
75+
unset($this->statement->{$name});
76+
}
77+
78+
public function __clone(): void
79+
{
80+
throw new \Error('Trying to clone an uncloneable PDOStatement');
81+
}
82+
83+
/**
84+
* @param array<mixed> $args
85+
* @throws \Throwable
86+
*/
87+
public function __call(string $method, array $args): mixed
88+
{
89+
try {
90+
return $this->statement->{$method}(...$args);
91+
} catch (\Throwable $e) {
92+
if ($this->pdo->inTransaction() || !Connection::hasError($e)) {
93+
throw $e;
94+
}
95+
96+
Console::warning('[Database] ' . $e->getMessage());
97+
Console::warning('[Database] Lost connection detected. Re-preparing statement...');
98+
99+
$this->reprepare();
100+
101+
return $this->statement->{$method}(...$args);
102+
}
103+
}
104+
105+
public function getStatement(): \PDOStatement
106+
{
107+
return $this->statement;
108+
}
109+
110+
public function setAttribute(int $attribute, mixed $value): bool
111+
{
112+
$this->attributes[$attribute] = $value;
113+
114+
return $this->statement->setAttribute($attribute, $value);
115+
}
116+
117+
public function setFetchMode(int $mode, mixed ...$args): bool
118+
{
119+
$this->fetchMode = [$mode, ...$args];
120+
121+
return $this->statement->setFetchMode($mode, ...$args);
122+
}
123+
124+
public function bindValue(int|string $param, mixed $value, int $type = \PDO::PARAM_STR): bool
125+
{
126+
$this->values[$param] = [$value, $type];
127+
128+
return $this->statement->bindValue($param, $value, $type);
129+
}
130+
131+
public function bindParam(int|string $param, mixed &$variable, int $type = \PDO::PARAM_STR, int $maxLength = 0, mixed $driverOptions = null): bool
132+
{
133+
$this->params[$param] = [$variable, $type, $maxLength, $driverOptions];
134+
135+
return $this->statement->bindParam($param, $variable, $type, $maxLength, $driverOptions);
136+
}
137+
138+
public function bindColumn(int|string $column, mixed &$variable, ?int $type = null, ?int $maxLength = null, mixed $driverOptions = null): bool
139+
{
140+
$this->columns[$column] = [$variable, $type, $maxLength, $driverOptions];
141+
142+
return $this->statement->bindColumn($column, $variable, $type, $maxLength, $driverOptions);
143+
}
144+
145+
private function reprepare(): void
146+
{
147+
$this->pdo->reconnect();
148+
$this->statement = $this->pdo->prepareNative($this->query, $this->options);
149+
150+
foreach ($this->attributes as $attribute => $value) {
151+
$this->statement->setAttribute($attribute, $value);
152+
}
153+
154+
if ($this->fetchMode !== null) {
155+
$this->statement->setFetchMode(...$this->fetchMode);
156+
}
157+
158+
foreach ($this->params as $param => [$variable, $type, $maxLength, $driverOptions]) {
159+
$this->statement->bindParam($param, $variable, $type, $maxLength, $driverOptions);
160+
}
161+
162+
foreach ($this->columns as $column => [$variable, $type, $maxLength, $driverOptions]) {
163+
$this->statement->bindColumn($column, $variable, $type, $maxLength, $driverOptions);
164+
}
165+
166+
foreach ($this->values as $param => [$value, $type]) {
167+
$this->statement->bindValue($param, $value, $type);
168+
}
169+
}
170+
}

0 commit comments

Comments
 (0)