Lines
94.94%
244 / 257
Functions and Methods
84.61%
44 / 52
Classes and Traits
0.00%
0 / 1
| Name | Lines | Functions and Methods | CRAP | Classes and Traits | ||||||
|---|---|---|---|---|---|---|---|---|---|---|
| SqlConnection | 94.94% | 244 / 257 | 84.61% | 44 / 52 | 129.09 | 0.00% | 0 / 1 | |||
| __construct | 100.00% | 5 / 5 | 100.00% | 1 / 1 | 1 | |||||
| assertSql | 100.00% | 2 / 2 | 100.00% | 1 / 1 | 2 | |||||
| from | 100.00% | 6 / 6 | 100.00% | 1 / 1 | 2 | |||||
| table | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | |||||
| select | 100.00% | 8 / 8 | 100.00% | 1 / 1 | 3 | |||||
| selectColumn | 100.00% | 2 / 2 | 100.00% | 1 / 1 | 1 | |||||
| insert | 100.00% | 4 / 4 | 100.00% | 1 / 1 | 1 | |||||
| insertGetId | 100.00% | 23 / 23 | 100.00% | 1 / 1 | 8 | |||||
| update | 100.00% | 2 / 2 | 100.00% | 1 / 1 | 1 | |||||
| delete | 100.00% | 2 / 2 | 100.00% | 1 / 1 | 1 | |||||
| cursor | 100.00% | 2 / 2 | 100.00% | 1 / 1 | 1 | |||||
| flattenInsertValues | 100.00% | 5 / 5 | 100.00% | 1 / 1 | 2 | |||||
| selectSql | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | |||||
| selectColumnSql | 100.00% | 2 / 2 | 100.00% | 1 / 1 | 1 | |||||
| cursorSql | 100.00% | 4 / 4 | 100.00% | 1 / 1 | 2 | |||||
| chunkSql | 100.00% | 11 / 11 | 100.00% | 1 / 1 | 6 | |||||
| prepareAndExecute | 100.00% | 8 / 8 | 100.00% | 1 / 1 | 3 | |||||
| statement | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | |||||
| affectingStatement | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | |||||
| bindValues | 100.00% | 11 / 11 | 100.00% | 1 / 1 | 10 | |||||
| run | 100.00% | 8 / 8 | 100.00% | 1 / 1 | 3 | |||||
| getDefaultValueCodec | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | |||||
| getDefaultQueryGrammar | n/a | 0 / 0 | n/a | 0 / 0 | 0 | |||||
| getDefaultSchemaGrammar | n/a | 0 / 0 | n/a | 0 / 0 | 0 | |||||
| getDefaultSchemaInspector | n/a | 0 / 0 | n/a | 0 / 0 | 0 | |||||
| supportsTransactionalDdl | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | |||||
| changeRequiresStandaloneTransaction | 0.00% | 0 / 1 | 0.00% | 0 / 1 | 2 | |||||
| create | 100.00% | 3 / 3 | 100.00% | 1 / 1 | 2 | |||||
| alter | 100.00% | 9 / 9 | 100.00% | 1 / 1 | 5 | |||||
| drop | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | |||||
| apply | 93.33% | 14 / 15 | 0.00% | 0 / 1 | 13.05 | |||||
| renameTable | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | |||||
| renameColumn | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | |||||
| applyColumnRenames | 100.00% | 2 / 2 | 100.00% | 1 / 1 | 2 | |||||
| applyAddColumn | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | |||||
| applyDropColumn | 100.00% | 5 / 5 | 100.00% | 1 / 1 | 2 | |||||
| applyModifyColumn | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | |||||
| subjectBlueprint | 100.00% | 3 / 3 | 100.00% | 1 / 1 | 2 | |||||
| modifyColumn | 100.00% | 2 / 2 | 100.00% | 1 / 1 | 2 | |||||
| addForeignKey | 44.44% | 4 / 9 | 0.00% | 0 / 1 | 1.17 | |||||
| dropForeignKey | 80.00% | 4 / 5 | 0.00% | 0 / 1 | 1.01 | |||||
| addCheck | 66.66% | 4 / 6 | 0.00% | 0 / 1 | 1.04 | |||||
| dropCheck | 80.00% | 4 / 5 | 0.00% | 0 / 1 | 1.01 | |||||
| rebuildIndexes | 100.00% | 4 / 4 | 100.00% | 1 / 1 | 3 | |||||
| withLock | 100.00% | 3 / 3 | 100.00% | 1 / 1 | 1 | |||||
| transactionLevel | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | |||||
| supportsSavepoints | n/a | 0 / 0 | n/a | 0 / 0 | 0 | |||||
| createSavepoint | n/a | 0 / 0 | n/a | 0 / 0 | 0 | |||||
| releaseSavepoint | n/a | 0 / 0 | n/a | 0 / 0 | 0 | |||||
| rollbackToSavepoint | n/a | 0 / 0 | n/a | 0 / 0 | 0 | |||||
| beginTransaction | 100.00% | 10 / 10 | 100.00% | 1 / 1 | 3 | |||||
| commit | 100.00% | 11 / 11 | 100.00% | 1 / 1 | 4 | |||||
| rollBack | 100.00% | 11 / 11 | 100.00% | 1 / 1 | 4 | |||||
| reconcileFailedCommit | 100.00% | 2 / 2 | 100.00% | 1 / 1 | 2 | |||||
| coroutineId | 88.88% | 8 / 9 | 0.00% | 0 / 1 | 6.05 | |||||
| assertSameCoroutine | 100.00% | 8 / 8 | 100.00% | 1 / 1 | 4 | |||||
| savepointNameFor | 75.00% | 3 / 4 | 0.00% | 0 / 1 | 1.02 | |||||
| transaction | 100.00% | 8 / 8 | 100.00% | 1 / 1 | 3 | |||||
| __destruct | 100.00% | 5 / 5 | 100.00% | 1 / 1 | 3 | |||||
| 1 | <?php | |
| 2 | ||
| 3 | declare(strict_types=1); | |
| 4 | ||
| 5 | namespace BlueprintAU\Radiant\Database\Connections; | |
| 6 | ||
| 7 | use BlueprintAU\Collections\Collection; | |
| 8 | use BlueprintAU\Radiant\Database\Concerns\DetectsConnectionLoss; | |
| 9 | use BlueprintAU\Radiant\Database\Concerns\NormalizesInsertRows; | |
| 10 | use BlueprintAU\Radiant\Database\Exceptions\QueryException; | |
| 11 | use BlueprintAU\Radiant\Database\Exceptions\UnsupportedFeatureException; | |
| 12 | use BlueprintAU\Radiant\Database\Grammars\Grammar; | |
| 13 | use BlueprintAU\Radiant\Database\Query\Enums\BindingCategory; | |
| 14 | use BlueprintAU\Radiant\Database\Query\QueryBuilder; | |
| 15 | use BlueprintAU\Radiant\Database\Schema\Blueprint; | |
| 16 | use BlueprintAU\Radiant\Database\Schema\Grammars\SchemaGrammar; | |
| 17 | use BlueprintAU\Radiant\Database\Schema\Enums\SchemaOperation; | |
| 18 | use BlueprintAU\Radiant\Database\ValueCodecs\DefaultValueCodec; | |
| 19 | use BlueprintAU\Radiant\Database\ValueCodecs\ValueCodecInterface; | |
| 20 | use Override; | |
| 21 | ||
| 22 | /** | |
| 23 | * A database connection backed by SQL (MySQL, SQLite, Postgres, …). | |
| 24 | * | |
| 25 | * Use it to run queries, raw SQL, transactions and schema changes. You | |
| 26 | * normally get one from a {@see \BlueprintAU\Radiant\Database\DatabaseManager} | |
| 27 | * rather than constructing it yourself. | |
| 28 | * | |
| 29 | * A connection is not safe for concurrent use by multiple coroutines while | |
| 30 | * a transaction is open — use one connection per coroutine in that case. | |
| 31 | * | |
| 32 | * @see ConnectionInterface | |
| 33 | * | |
| 34 | * @template TGrammar of Grammar = Grammar | |
| 35 | * @template TSchemaGrammar of SchemaGrammar = SchemaGrammar | |
| 36 | * @template TSchemaInspector of \BlueprintAU\Radiant\Database\Schema\Inspectors\SchemaInspector = \BlueprintAU\Radiant\Database\Schema\Inspectors\SchemaInspector | |
| 37 | */ | |
| 38 | abstract class SqlConnection implements ConnectionInterface | |
| 39 | { | |
| 40 | use NormalizesInsertRows; | |
| 41 | use DetectsConnectionLoss; | |
| 42 | /** | |
| 43 | * Converts values between PHP types and what the database driver expects. | |
| 44 | * | |
| 45 | * @var ValueCodecInterface | |
| 46 | */ | |
| 47 | public readonly ValueCodecInterface $codec; | |
| 48 | ||
| 49 | /** | |
| 50 | * Compiles query-builder state into dialect SQL. | |
| 51 | * | |
| 52 | * @var TGrammar | |
| 53 | */ | |
| 54 | public readonly Grammar $grammar; | |
| 55 | ||
| 56 | /** | |
| 57 | * Compiles schema definitions into dialect DDL. | |
| 58 | * | |
| 59 | * @var TSchemaGrammar | |
| 60 | */ | |
| 61 | public readonly SchemaGrammar $schemaGrammar; | |
| 62 | ||
| 63 | /** | |
| 64 | * Reads the live schema — the read-side twin of {@see $schemaGrammar}. | |
| 65 | * | |
| 66 | * @var TSchemaInspector | |
| 67 | */ | |
| 68 | public readonly \BlueprintAU\Radiant\Database\Schema\Inspectors\SchemaInspector $schemaInspector; | |
| 69 | ||
| 70 | /** | |
| 71 | * Create a new SQL connection wrapping a PDO instance. | |
| 72 | * | |
| 73 | * @param \PDO $pdo | |
| 74 | */ | |
| 75 | public function __construct(protected \Pdo $pdo) | |
| 76 | { | |
| 77 | $this->pdo->setAttribute(\PDO::ATTR_ERRMODE, \PDO::ERRMODE_EXCEPTION); | |
| 78 | $this->codec = $this->getDefaultValueCodec(); | |
| 79 | $this->grammar = $this->getDefaultQueryGrammar(); | |
| 80 | $this->schemaGrammar = $this->getDefaultSchemaGrammar(); | |
| 81 | $this->schemaInspector = $this->getDefaultSchemaInspector(); | |
| 82 | } | |
| 83 | ||
| 84 | /** | |
| 85 | * Assert the connection speaks SQL — throw when it does not. | |
| 86 | * | |
| 87 | * @param ConnectionInterface $connection | |
| 88 | * @phpstan-assert SqlConnection<Grammar, SchemaGrammar, \BlueprintAU\Radiant\Database\Schema\Inspectors\SchemaInspector> $connection | |
| 89 | */ | |
| 90 | final public static function assertSql(ConnectionInterface $connection): void | |
| 91 | { | |
| 92 | if (!$connection instanceof self) { | |
| 93 | throw new UnsupportedFeatureException('The connection is not a SQL connection.'); | |
| 94 | } | |
| 95 | } | |
| 96 | ||
| 97 | /** | |
| 98 | * Narrow a connection to SQL and to this dialect — return it typed, or | |
| 99 | * throw. | |
| 100 | * | |
| 101 | * @param ConnectionInterface $connection | |
| 102 | * @return static | |
| 103 | * @throws UnsupportedFeatureException | |
| 104 | */ | |
| 105 | final public static function from(ConnectionInterface $connection): static | |
| 106 | { | |
| 107 | if (!$connection instanceof static) { | |
| 108 | // A non-SQL backend gets the generic message; a SQL one gets | |
| 109 | // the dialect-specific one. | |
| 110 | self::assertSql($connection); | |
| 111 | throw new UnsupportedFeatureException( | |
| 112 | 'The connection is SQL, but not a ' . static::class . '.' | |
| 113 | ); | |
| 114 | } | |
| 115 | ||
| 116 | return $connection; | |
| 117 | } | |
| 118 | ||
| 119 | /** | |
| 120 | * Start a fluent query against a table, bound to this connection. | |
| 121 | * | |
| 122 | * @param string $identifier | |
| 123 | * @return QueryBuilder | |
| 124 | */ | |
| 125 | #[Override] | |
| 126 | final public function table(string $identifier): QueryBuilder | |
| 127 | { | |
| 128 | return new QueryBuilder($this, $identifier); | |
| 129 | } | |
| 130 | ||
| 131 | /** | |
| 132 | * Run the query and return the matching rows. | |
| 133 | * | |
| 134 | * @param QueryBuilder $query | |
| 135 | * @return Collection<int,\stdClass> | |
| 136 | * @throws \LogicException | |
| 137 | */ | |
| 138 | #[Override] | |
| 139 | final public function select(QueryBuilder $query): Collection | |
| 140 | { | |
| 141 | if ($query->getLock() !== null && $this->transactionLevel === 0) { | |
| 142 | throw new \LogicException( | |
| 143 | 'Row locks (lockForUpdate/sharedLock) require an open transaction — ' | |
| 144 | . 'outside one, the lock is released at statement end and protects nothing. ' | |
| 145 | . 'Wrap the query in beginTransaction()/transaction().' | |
| 146 | ); | |
| 147 | } | |
| 148 | ||
| 149 | $sql = $this->grammar->compileSelect($query); | |
| 150 | return $this->selectSql($sql, $query->getBindings()); | |
| 151 | } | |
| 152 | ||
| 153 | /** | |
| 154 | * Run the query and return the first selected column's values. | |
| 155 | * | |
| 156 | * @param QueryBuilder $query | |
| 157 | * @return Collection<int, mixed> | |
| 158 | * @throws QueryException | |
| 159 | */ | |
| 160 | #[Override] | |
| 161 | final public function selectColumn(QueryBuilder $query): Collection | |
| 162 | { | |
| 163 | $sql = $this->grammar->compileSelect($query); | |
| 164 | return $this->selectColumnSql($sql, $query->getBindings()); | |
| 165 | } | |
| 166 | ||
| 167 | /** | |
| 168 | * Insert one or more rows into the table. | |
| 169 | * | |
| 170 | * @param QueryBuilder $query | |
| 171 | * @param array<string,mixed>|list<array<string,mixed>> $values | |
| 172 | * @return int | |
| 173 | */ | |
| 174 | #[Override] | |
| 175 | final public function insert(QueryBuilder $query, array $values): int | |
| 176 | { | |
| 177 | $rows = $this->normalizeInsertRows($values); | |
| 178 | ||
| 179 | // Fail fast on ragged rows BEFORE any statement runs — the grammar | |
| 180 | // would reject them at compile, but the binding flattener below | |
| 181 | // would happily emit a mismatched list first. | |
| 182 | $this->assertUniformInsertRows($rows); | |
| 183 | ||
| 184 | $sql = $this->grammar->compileInsert($query, $values); | |
| 185 | return $this->affectingStatement($sql, $this->flattenInsertValues($values)); | |
| 186 | } | |
| 187 | ||
| 188 | /** | |
| 189 | * Insert a single row and return its generated id. | |
| 190 | * | |
| 191 | * @param QueryBuilder $query | |
| 192 | * @param array<string,mixed> $values | |
| 193 | * @return string|int|null | |
| 194 | */ | |
| 195 | #[Override] | |
| 196 | final public function insertGetId(QueryBuilder $query, array $values): string|int|null | |
| 197 | { | |
| 198 | $pk = $query->getInsertIdColumn(); | |
| 199 | ||
| 200 | if ($pk === null) { | |
| 201 | // No key declared — compile the plain insert and report success. | |
| 202 | $this->affectingStatement($this->grammar->compileInsert($query, $values), $this->flattenInsertValues($values)); | |
| 203 | return null; | |
| 204 | } | |
| 205 | ||
| 206 | // The compile-shaped capability contract: the grammar returns the | |
| 207 | // statement PLUS whether that statement yields the key (a RETURNING | |
| 208 | // dialect compiles the clause in; MySQL compiles without it). The | |
| 209 | // connection never probes a boolean — it reads the compile result. | |
| 210 | $compiled = $this->grammar->compileInsertForId($query, $values, $pk); | |
| 211 | $bindings = $this->flattenInsertValues($values); | |
| 212 | ||
| 213 | if ($compiled['returnsKey']) { | |
| 214 | $row = $this->selectSql($compiled['sql'], $bindings)->first(); | |
| 215 | if ($row === null) { | |
| 216 | return null; | |
| 217 | } | |
| 218 | $id = $this->codec->decode($row->{$pk}); | |
| 219 | return is_int($id) || is_string($id) ? $id : null; | |
| 220 | } | |
| 221 | ||
| 222 | $this->statement($compiled['sql'], $bindings); | |
| 223 | // The lastInsertId() fallback is reachable ONLY on dialects without | |
| 224 | // RETURNING (MySQL, and old SQLite) — Postgres' grammar always uses | |
| 225 | // RETURNING, so its sequence-based lastval() hazards never apply | |
| 226 | // here. On MySQL lastInsertId() is connection-scoped and unaffected | |
| 227 | // by concurrent inserts on other connections. PDO always returns a | |
| 228 | // string (or false when there is no generated id); the codec passes | |
| 229 | // strings through untouched, so no decode is needed. | |
| 230 | $id = $this->pdo->lastInsertId(); | |
| 231 | if ($id === false) { | |
| 232 | return null; | |
| 233 | } | |
| 234 | ||
| 235 | // lastInsertId() is only meaningful for an AUTO_INCREMENT/SERIAL | |
| 236 | // column: a caller-declared non-auto-increment PK (UUID, char, or a | |
| 237 | // PK the row value supplies) generates nothing server-side, so the | |
| 238 | // value here is a stale id from an EARLIER insert on this connection | |
| 239 | // (or '0'). Returning it would hand the caller a key that does not | |
| 240 | // identify the row just written — fail fast instead. | |
| 241 | if (!$query->isInsertIdAutoIncrement()) { | |
| 242 | throw new \LogicException( | |
| 243 | "insertGetId() declared key column [{$pk}] is not auto-increment — no id is generated" | |
| 244 | . ' server-side, so lastInsertId() would return a stale value from an earlier insert.' | |
| 245 | . ' Assign the key before inserting and use insert().' | |
| 246 | ); | |
| 247 | } | |
| 248 | ||
| 249 | return $id; | |
| 250 | } | |
| 251 | ||
| 252 | /** | |
| 253 | * Update the rows matching the query's conditions. | |
| 254 | * | |
| 255 | * @param QueryBuilder $query | |
| 256 | * @param array<string,mixed> $values | |
| 257 | * @return int | |
| 258 | */ | |
| 259 | #[Override] | |
| 260 | final public function update(QueryBuilder $query, array $values): int | |
| 261 | { | |
| 262 | $sql = $this->grammar->compileUpdate($query, $values); | |
| 263 | return $this->affectingStatement($sql, array_merge(array_values($values), $query->getBindings([BindingCategory::Join, BindingCategory::Where]))); | |
| 264 | } | |
| 265 | ||
| 266 | /** | |
| 267 | * Delete the rows matching the query's conditions. | |
| 268 | * | |
| 269 | * @param QueryBuilder $query | |
| 270 | * @return int | |
| 271 | */ | |
| 272 | #[Override] | |
| 273 | final public function delete(QueryBuilder $query): int | |
| 274 | { | |
| 275 | $sql = $this->grammar->compileDelete($query); | |
| 276 | return $this->affectingStatement($sql, $query->getBindings([BindingCategory::Join, BindingCategory::Where])); | |
| 277 | } | |
| 278 | ||
| 279 | /** | |
| 280 | * Run the query and yield each matching row as it arrives. | |
| 281 | * | |
| 282 | * @param QueryBuilder $query | |
| 283 | * @return \Generator<int, \stdClass> | |
| 284 | * @throws QueryException | |
| 285 | */ | |
| 286 | #[Override] | |
| 287 | final public function cursor(QueryBuilder $query): \Generator | |
| 288 | { | |
| 289 | $sql = $this->grammar->compileSelect($query); | |
| 290 | return $this->cursorSql($sql, $query->getBindings()); | |
| 291 | } | |
| 292 | ||
| 293 | /** | |
| 294 | * Flatten a single row or a list of rows into one binding list (row-major). | |
| 295 | * | |
| 296 | * @param array<string,mixed>|list<array<string,mixed>> $values | |
| 297 | * @return list<mixed> | |
| 298 | */ | |
| 299 | protected function flattenInsertValues(array $values): array | |
| 300 | { | |
| 301 | $rows = $this->normalizeInsertRows($values); | |
| 302 | $bindings = []; | |
| 303 | foreach ($rows as $row) { | |
| 304 | array_push($bindings, ...array_values($row)); | |
| 305 | } | |
| 306 | return $bindings; | |
| 307 | } | |
| 308 | ||
| 309 | // ---- Running and Binding ---- | |
| 310 | ||
| 311 | ||
| 312 | /** | |
| 313 | * Run a raw SQL query and return every matching row as an object. | |
| 314 | * | |
| 315 | * @param string $sql | |
| 316 | * @param array<string|int, mixed> $bindings | |
| 317 | * @return Collection<int,\stdClass> | |
| 318 | */ | |
| 319 | final public function selectSql(string $sql, array $bindings = []): Collection | |
| 320 | { | |
| 321 | return Collection::make($this->run($sql, $bindings, fn(\PDOStatement $stmt) => $stmt->fetchAll(\PDO::FETCH_OBJ))); | |
| 322 | } | |
| 323 | ||
| 324 | /** | |
| 325 | * Run a raw SQL query and return the first selected column's values. | |
| 326 | * | |
| 327 | * @param string $sql | |
| 328 | * @param array<string|int, mixed> $bindings | |
| 329 | * @return Collection<int, mixed> | |
| 330 | */ | |
| 331 | final public function selectColumnSql(string $sql, array $bindings = []): Collection | |
| 332 | { | |
| 333 | /** @var list<mixed> $column */ | |
| 334 | $column = $this->run($sql, $bindings, fn(\PDOStatement $stmt) => $stmt->fetchAll(\PDO::FETCH_COLUMN, 0)); | |
| 335 | return Collection::make($column); | |
| 336 | } | |
| 337 | ||
| 338 | /** | |
| 339 | * Run a raw SQL query and yield each matching row as it arrives. | |
| 340 | * | |
| 341 | * Consume the generator fully (or let it be garbage collected) before | |
| 342 | * running another query on this connection — an unfinished cursor holds | |
| 343 | * the statement. | |
| 344 | * | |
| 345 | * @param string $sql | |
| 346 | * @param array<string|int, mixed> $bindings | |
| 347 | * @return \Generator<int, \stdClass> | |
| 348 | * @throws QueryException | |
| 349 | */ | |
| 350 | final public function cursorSql(string $sql, array $bindings = []): \Generator | |
| 351 | { | |
| 352 | $stmt = $this->prepareAndExecute($sql, $bindings); | |
| 353 | try { | |
| 354 | while ($row = $stmt->fetch(\PDO::FETCH_OBJ)) { | |
| 355 | yield $row; | |
| 356 | } | |
| 357 | } finally { | |
| 358 | $stmt->closeCursor(); | |
| 359 | } | |
| 360 | } | |
| 361 | ||
| 362 | /** | |
| 363 | * Run a callback over the query's rows in fixed-size chunks. | |
| 364 | * | |
| 365 | * The callback returning `false` (strictly) stops the iteration | |
| 366 | * immediately; any other return value continues. | |
| 367 | * | |
| 368 | * @param string $sql | |
| 369 | * @param array<string|int, mixed> $bindings | |
| 370 | * @param int $size | |
| 371 | * @param callable(list<\stdClass>): mixed $callback | |
| 372 | * @return void | |
| 373 | * @throws \InvalidArgumentException | |
| 374 | * @throws QueryException | |
| 375 | */ | |
| 376 | final public function chunkSql(string $sql, array $bindings, int $size, callable $callback): void | |
| 377 | { | |
| 378 | if ($size < 1) { | |
| 379 | throw new \InvalidArgumentException("Chunk size must be at least 1; got {$size}."); | |
| 380 | } | |
| 381 | $chunk = []; | |
| 382 | foreach ($this->cursorSql($sql, $bindings) as $row) { | |
| 383 | $chunk[] = $row; | |
| 384 | if (count($chunk) === $size) { | |
| 385 | if ($callback($chunk) === false) { | |
| 386 | return; | |
| 387 | } | |
| 388 | $chunk = []; | |
| 389 | } | |
| 390 | } | |
| 391 | if ($chunk !== []) { | |
| 392 | $callback($chunk); | |
| 393 | } | |
| 394 | } | |
| 395 | ||
| 396 | /** | |
| 397 | * Prepare, bind, and execute — returning the statement for callers | |
| 398 | * that manage the cursor themselves. | |
| 399 | * | |
| 400 | * @param string $sql | |
| 401 | * @param array<string|int, mixed> $bindings | |
| 402 | * @return \PDOStatement | |
| 403 | * @throws QueryException | |
| 404 | */ | |
| 405 | private function prepareAndExecute(string $sql, array $bindings): \PDOStatement | |
| 406 | { | |
| 407 | try { | |
| 408 | $stmt = $this->pdo->prepare($sql); | |
| 409 | $this->bindValues($stmt, $bindings); | |
| 410 | $stmt->execute(); | |
| 411 | return $stmt; | |
| 412 | } catch (\PDOException $e) { | |
| 413 | if ($this->isConnectionLoss($e)) { | |
| 414 | $this->stale = true; | |
| 415 | } | |
| 416 | throw new QueryException($sql, $bindings, $e); | |
| 417 | } | |
| 418 | } | |
| 419 | ||
| 420 | /** | |
| 421 | * Run a raw SQL statement that returns no result set. | |
| 422 | * | |
| 423 | * @param string $sql | |
| 424 | * @param array<string|int, mixed> $bindings | |
| 425 | */ | |
| 426 | final public function statement(string $sql, array $bindings = []): void | |
| 427 | { | |
| 428 | $this->run($sql, $bindings, fn() => null); | |
| 429 | } | |
| 430 | ||
| 431 | /** | |
| 432 | * Run a raw SQL statement and return how many rows it affected. | |
| 433 | * | |
| 434 | * @param string $sql | |
| 435 | * @param array<string|int, mixed> $bindings | |
| 436 | * @return int | |
| 437 | */ | |
| 438 | final public function affectingStatement(string $sql, array $bindings = []): int | |
| 439 | { | |
| 440 | return $this->run($sql, $bindings, fn(\PDOStatement $stmt) => $stmt->rowCount()); | |
| 441 | } | |
| 442 | ||
| 443 | /** | |
| 444 | * Bind values to the statement, adapting them through the codec. | |
| 445 | * | |
| 446 | * @param \PDOStatement $stmt | |
| 447 | * @param array<string|int, mixed> $bindings | |
| 448 | * @throws \InvalidArgumentException | |
| 449 | */ | |
| 450 | protected function bindValues(\PDOStatement $stmt, array $bindings): void | |
| 451 | { | |
| 452 | foreach ($bindings as $key => $value) { | |
| 453 | // The write path: bindable value → codec → driver value → bind. | |
| 454 | if (!is_scalar($value) && !$value instanceof \DateTimeInterface && $value !== null) { | |
| 455 | throw new \InvalidArgumentException( | |
| 456 | 'Binding must be a scalar, null, or DateTimeInterface; got ' . get_debug_type($value) . '.' | |
| 457 | ); | |
| 458 | } | |
| 459 | $value = $this->codec->encode($value); | |
| 460 | ||
| 461 | // Bind with an EXPLICIT PDO type. Without one, PDO defaults to | |
| 462 | // PARAM_STR — and an expression like `count(*) > ?` then | |
| 463 | // compares against the string '1', which SQLite evaluates as | |
| 464 | // text-vs-number and always false. Typed columns survive the | |
| 465 | // string bind via column affinity; aggregate expressions have | |
| 466 | // no affinity to save them. (int→PARAM_INT, float→PARAM_STR — | |
| 467 | // PDO has no float type and SQLite compares numerically anyway, | |
| 468 | // bool→PARAM_INT, null→PARAM_NULL.) | |
| 469 | $type = match (true) { | |
| 470 | is_int($value) => \PDO::PARAM_INT, | |
| 471 | is_bool($value) => \PDO::PARAM_INT, | |
| 472 | $value === null => \PDO::PARAM_NULL, | |
| 473 | default => \PDO::PARAM_STR, | |
| 474 | }; | |
| 475 | ||
| 476 | $stmt->bindValue(is_int($key) ? $key + 1 : $key, $value, $type); | |
| 477 | } | |
| 478 | } | |
| 479 | ||
| 480 | /** | |
| 481 | * Prepare, bind, execute, and run the callback — the single run path. | |
| 482 | * | |
| 483 | * @template T | |
| 484 | * | |
| 485 | * @param string $sql | |
| 486 | * @param array<string|int, mixed> $bindings | |
| 487 | * @param callable(\PDOStatement): T $callback | |
| 488 | * @return T | |
| 489 | * @throws QueryException | |
| 490 | */ | |
| 491 | final protected function run(string $sql, array $bindings, callable $callback): mixed | |
| 492 | { | |
| 493 | try { | |
| 494 | $stmt = $this->pdo->prepare($sql); | |
| 495 | $this->bindValues($stmt, $bindings); | |
| 496 | $stmt->execute(); | |
| 497 | return $callback($stmt); | |
| 498 | } catch (\PDOException $e) { | |
| 499 | if ($this->isConnectionLoss($e)) { | |
| 500 | $this->stale = true; | |
| 501 | } | |
| 502 | throw new QueryException($sql, $bindings, $e); | |
| 503 | } | |
| 504 | } | |
| 505 | ||
| 506 | // ---- Encoding and Grammar ---- | |
| 507 | ||
| 508 | /** | |
| 509 | * The default codec for this connection. | |
| 510 | * | |
| 511 | * @return ValueCodecInterface | |
| 512 | */ | |
| 513 | protected function getDefaultValueCodec(): ValueCodecInterface | |
| 514 | { | |
| 515 | return new DefaultValueCodec(); | |
| 516 | } | |
| 517 | ||
| 518 | /** | |
| 519 | * The default query grammar for this connection. | |
| 520 | * | |
| 521 | * @return TGrammar | |
| 522 | */ | |
| 523 | abstract protected function getDefaultQueryGrammar(): Grammar; | |
| 524 | ||
| 525 | /** | |
| 526 | * The default schema grammar for this connection. | |
| 527 | * | |
| 528 | * @return TSchemaGrammar | |
| 529 | */ | |
| 530 | abstract protected function getDefaultSchemaGrammar(): SchemaGrammar; | |
| 531 | ||
| 532 | /** | |
| 533 | * The dialect's live-schema reader. | |
| 534 | * | |
| 535 | * @return TSchemaInspector | |
| 536 | */ | |
| 537 | abstract protected function getDefaultSchemaInspector(): \BlueprintAU\Radiant\Database\Schema\Inspectors\SchemaInspector; | |
| 538 | ||
| 539 | /** | |
| 540 | * Whether this dialect's DDL is transactional. | |
| 541 | * | |
| 542 | * @return bool | |
| 543 | */ | |
| 544 | public function supportsTransactionalDdl(): bool | |
| 545 | { | |
| 546 | return false; | |
| 547 | } | |
| 548 | ||
| 549 | /** | |
| 550 | * Whether applying a change needs a transaction-free connection. | |
| 551 | * | |
| 552 | * Dialects route changes that cannot run inside a transaction through | |
| 553 | * this predicate; the synchronizer consults it before wrapping the | |
| 554 | * apply loop in a transaction. | |
| 555 | * | |
| 556 | * @param \BlueprintAU\Radiant\Database\Schema\SchemaChange $change | |
| 557 | * @param list<\BlueprintAU\Radiant\Database\Schema\SchemaChange> $plan The whole plan, for rename resolution. | |
| 558 | * @return bool | |
| 559 | */ | |
| 560 | public function changeRequiresStandaloneTransaction(\BlueprintAU\Radiant\Database\Schema\SchemaChange $change, array $plan = []): bool | |
| 561 | { | |
| 562 | return false; | |
| 563 | } | |
| 564 | ||
| 565 | // ---- Schema operations (SQL-only) ---- | |
| 566 | ||
| 567 | /** | |
| 568 | * Create a table from a blueprint, plus any indexes declared on it. | |
| 569 | * | |
| 570 | * @param Blueprint $blueprint | |
| 571 | */ | |
| 572 | final public function create(Blueprint $blueprint): void | |
| 573 | { | |
| 574 | $this->statement($this->schemaGrammar->compileCreate($blueprint)); | |
| 575 | ||
| 576 | foreach ($this->schemaGrammar->compileIndexes($blueprint) as $indexSql) { | |
| 577 | $this->statement($indexSql); | |
| 578 | } | |
| 579 | } | |
| 580 | ||
| 581 | /** | |
| 582 | * Alter a table — add or drop columns. | |
| 583 | * | |
| 584 | * @param SchemaOperation $operation | |
| 585 | * @param Blueprint $blueprint | |
| 586 | */ | |
| 587 | final public function alter(SchemaOperation $operation, Blueprint $blueprint): void | |
| 588 | { | |
| 589 | $statements = match ($operation) { | |
| 590 | SchemaOperation::AddColumn => $this->schemaGrammar->compileAddColumns($blueprint), | |
| 591 | SchemaOperation::DropColumn => $this->schemaGrammar->compileDropColumns($blueprint), | |
| 592 | default => throw new \LogicException( | |
| 593 | "Operation [{$operation->value}] is not a column alter; use the " | |
| 594 | . 'dedicated create/drop/rebuildIndexes paths.' | |
| 595 | ), | |
| 596 | }; | |
| 597 | ||
| 598 | foreach ($statements as $sql) { | |
| 599 | $this->statement($sql); | |
| 600 | } | |
| 601 | } | |
| 602 | ||
| 603 | /** | |
| 604 | * Drop a table. | |
| 605 | * | |
| 606 | * @param string $table | |
| 607 | */ | |
| 608 | final public function drop(string $table): void | |
| 609 | { | |
| 610 | $this->statement($this->schemaGrammar->compileDrop($table)); | |
| 611 | } | |
| 612 | ||
| 613 | /** | |
| 614 | * Apply a differ-produced change — dispatches create/alter/drop so the | |
| 615 | * host never writes the match itself. | |
| 616 | * | |
| 617 | * @param \BlueprintAU\Radiant\Database\Schema\SchemaChange $change | |
| 618 | */ | |
| 619 | final public function apply(\BlueprintAU\Radiant\Database\Schema\SchemaChange $change): void | |
| 620 | { | |
| 621 | match ($change->operation) { | |
| 622 | SchemaOperation::CreateTable => $this->create($change->blueprint), | |
| 623 | SchemaOperation::AddColumn => $this->applyAddColumn($change), | |
| 624 | SchemaOperation::DropColumn => $this->applyDropColumn($change), | |
| 625 | SchemaOperation::DropTable => $this->drop($change->table), | |
| 626 | SchemaOperation::AlterIndexes => $this->rebuildIndexes($change->blueprint), | |
| 627 | SchemaOperation::RenameTable => $this->renameTable($change->blueprint->getRenamedFrom() ?? throw new \LogicException( | |
| 628 | "A RenameTable change for [{$change->table}] carries no renamedFrom declaration." | |
| 629 | ), $change->table), | |
| 630 | SchemaOperation::RenameColumn => $this->applyColumnRenames($change->blueprint), | |
| 631 | SchemaOperation::ModifyColumn => $this->applyModifyColumn($change), | |
| 632 | SchemaOperation::AddForeignKey => $this->addForeignKey($change->table, $change->blueprint), | |
| 633 | SchemaOperation::DropForeignKey => $this->dropForeignKey($change->table, $change->blueprint), | |
| 634 | SchemaOperation::AddCheck => $this->addCheck($change->table, $change->blueprint), | |
| 635 | SchemaOperation::DropCheck => $this->dropCheck($change->table, $change->blueprint), | |
| 636 | }; | |
| 637 | } | |
| 638 | ||
| 639 | /** | |
| 640 | * Rename a table. | |
| 641 | * | |
| 642 | * @param string $from | |
| 643 | * @param string $to | |
| 644 | */ | |
| 645 | final public function renameTable(string $from, string $to): void | |
| 646 | { | |
| 647 | $this->statement($this->schemaGrammar->compileRenameTable($from, $to)); | |
| 648 | } | |
| 649 | ||
| 650 | /** | |
| 651 | * Rename a column on a table. | |
| 652 | * | |
| 653 | * @param string $table | |
| 654 | * @param string $from | |
| 655 | * @param string $to | |
| 656 | */ | |
| 657 | final public function renameColumn(string $table, string $from, string $to): void | |
| 658 | { | |
| 659 | $this->statement($this->schemaGrammar->compileRenameColumn($table, $from, $to)); | |
| 660 | } | |
| 661 | ||
| 662 | /** | |
| 663 | * Apply every column rename declared on the blueprint, in declaration | |
| 664 | * order. | |
| 665 | * | |
| 666 | * @param Blueprint $blueprint | |
| 667 | */ | |
| 668 | private function applyColumnRenames(Blueprint $blueprint): void | |
| 669 | { | |
| 670 | foreach ($blueprint->getColumnRenames() as $rename) { | |
| 671 | $this->renameColumn($blueprint->getTable(), $rename['from'], $rename['to']); | |
| 672 | } | |
| 673 | } | |
| 674 | ||
| 675 | /** | |
| 676 | * Apply an AddColumn change in place. | |
| 677 | * | |
| 678 | * @param \BlueprintAU\Radiant\Database\Schema\SchemaChange $change | |
| 679 | */ | |
| 680 | protected function applyAddColumn(\BlueprintAU\Radiant\Database\Schema\SchemaChange $change): void | |
| 681 | { | |
| 682 | $this->alter(SchemaOperation::AddColumn, $this->subjectBlueprint($change)); | |
| 683 | } | |
| 684 | ||
| 685 | /** | |
| 686 | * Apply a DropColumn change. | |
| 687 | * | |
| 688 | * The drop blueprint is built from the change's subject names, falling | |
| 689 | * back to the blueprint's own dropColumn() declarations. | |
| 690 | * | |
| 691 | * @param \BlueprintAU\Radiant\Database\Schema\SchemaChange $change | |
| 692 | */ | |
| 693 | protected function applyDropColumn(\BlueprintAU\Radiant\Database\Schema\SchemaChange $change): void | |
| 694 | { | |
| 695 | $names = $change->subject ?? $change->blueprint->getDropColumns(); | |
| 696 | ||
| 697 | $drop = new Blueprint($change->table); | |
| 698 | ||
| 699 | foreach ($names as $name) { | |
| 700 | $drop = $drop->dropColumn($name); | |
| 701 | } | |
| 702 | ||
| 703 | $this->alter(SchemaOperation::DropColumn, $drop); | |
| 704 | } | |
| 705 | ||
| 706 | /** | |
| 707 | * Apply a ModifyColumn change in place. | |
| 708 | * | |
| 709 | * @param \BlueprintAU\Radiant\Database\Schema\SchemaChange $change | |
| 710 | */ | |
| 711 | protected function applyModifyColumn(\BlueprintAU\Radiant\Database\Schema\SchemaChange $change): void | |
| 712 | { | |
| 713 | $this->modifyColumn($this->subjectBlueprint($change)); | |
| 714 | } | |
| 715 | ||
| 716 | /** | |
| 717 | * The blueprint the in-place dialects compile for a change. | |
| 718 | * | |
| 719 | * The full desired blueprint filtered to the change's subject column | |
| 720 | * names, or the full blueprint when the change acts on the whole table. | |
| 721 | * | |
| 722 | * @param \BlueprintAU\Radiant\Database\Schema\SchemaChange $change | |
| 723 | * @return Blueprint | |
| 724 | */ | |
| 725 | private function subjectBlueprint(\BlueprintAU\Radiant\Database\Schema\SchemaChange $change): Blueprint | |
| 726 | { | |
| 727 | if ($change->subject === null) { | |
| 728 | return $change->blueprint; | |
| 729 | } | |
| 730 | ||
| 731 | return $change->blueprint->onlyColumns($change->subject); | |
| 732 | } | |
| 733 | ||
| 734 | /** | |
| 735 | * Modify one or more columns in place — the content-drift path. | |
| 736 | * | |
| 737 | * @param Blueprint $blueprint | |
| 738 | */ | |
| 739 | public function modifyColumn(Blueprint $blueprint): void | |
| 740 | { | |
| 741 | foreach ($this->schemaGrammar->compileModifyColumn($blueprint) as $sql) { | |
| 742 | $this->statement($sql); | |
| 743 | } | |
| 744 | } | |
| 745 | ||
| 746 | /** | |
| 747 | * Add a foreign-key constraint to an existing table. | |
| 748 | * | |
| 749 | * @param string $table | |
| 750 | * @param Blueprint $blueprint | |
| 751 | */ | |
| 752 | public function addForeignKey(string $table, Blueprint $blueprint): void | |
| 753 | { | |
| 754 | $foreignKeys = $blueprint->getForeignKeys(); | |
| 755 | $first = $foreignKeys[0] ?? throw new \LogicException( | |
| 756 | "An AddForeignKey change for [{$table}] carries no foreign key declaration." | |
| 757 | ); | |
| 758 | ||
| 759 | $this->statement($this->schemaGrammar->compileAddForeignKey( | |
| 760 | $table, | |
| 761 | $first, | |
| 762 | $first['name'], | |
| 763 | )); | |
| 764 | } | |
| 765 | ||
| 766 | /** | |
| 767 | * Drop a foreign-key constraint from an existing table. | |
| 768 | * | |
| 769 | * @param string $table | |
| 770 | * @param Blueprint $blueprint | |
| 771 | */ | |
| 772 | public function dropForeignKey(string $table, Blueprint $blueprint): void | |
| 773 | { | |
| 774 | $names = $blueprint->getDropForeignKeys(); | |
| 775 | $name = $names[0] ?? throw new \LogicException( | |
| 776 | "A DropForeignKey change for [{$table}] carries no constraint name." | |
| 777 | ); | |
| 778 | ||
| 779 | $this->statement($this->schemaGrammar->compileDropForeignKey($table, $name)); | |
| 780 | } | |
| 781 | ||
| 782 | /** | |
| 783 | * Add a CHECK constraint to an existing table. | |
| 784 | * | |
| 785 | * @param string $table | |
| 786 | * @param Blueprint $blueprint | |
| 787 | */ | |
| 788 | public function addCheck(string $table, Blueprint $blueprint): void | |
| 789 | { | |
| 790 | $checks = $blueprint->getChecks(); | |
| 791 | $check = $checks[0] ?? throw new \LogicException( | |
| 792 | "An AddCheck change for [{$table}] carries no CHECK declaration." | |
| 793 | ); | |
| 794 | ||
| 795 | // Every CHECK carries a final name (derived at declaration when | |
| 796 | // omitted) — the name is the drop handle for a later drop. | |
| 797 | $name = $check['name']; | |
| 798 | ||
| 799 | $this->statement($this->schemaGrammar->compileAddCheck($table, $name, $check['expression'])); | |
| 800 | } | |
| 801 | ||
| 802 | /** | |
| 803 | * Drop a CHECK constraint from an existing table. | |
| 804 | * | |
| 805 | * @param string $table | |
| 806 | * @param Blueprint $blueprint | |
| 807 | */ | |
| 808 | public function dropCheck(string $table, Blueprint $blueprint): void | |
| 809 | { | |
| 810 | $names = $blueprint->getDropChecks(); | |
| 811 | $name = $names[0] ?? throw new \LogicException( | |
| 812 | "A DropCheck change for [{$table}] carries no constraint name." | |
| 813 | ); | |
| 814 | ||
| 815 | $this->statement($this->schemaGrammar->compileDropCheck($table, $name)); | |
| 816 | } | |
| 817 | ||
| 818 | /** | |
| 819 | * Rebuild a table's indexes — drop each index named on the blueprint, | |
| 820 | * then re-create it from the blueprint's declaration. | |
| 821 | * | |
| 822 | * @param Blueprint $blueprint | |
| 823 | */ | |
| 824 | final public function rebuildIndexes(Blueprint $blueprint): void | |
| 825 | { | |
| 826 | foreach ($blueprint->getIndexes() as $index) { | |
| 827 | $this->statement($this->schemaGrammar->compileDropIndex($index['name'], $blueprint->getTable())); | |
| 828 | } | |
| 829 | ||
| 830 | foreach ($this->schemaGrammar->compileIndexes($blueprint) as $indexSql) { | |
| 831 | $this->statement($indexSql); | |
| 832 | } | |
| 833 | } | |
| 834 | ||
| 835 | /** | |
| 836 | * Run the callback while holding a cross-process lock taken on this | |
| 837 | * connection. | |
| 838 | * | |
| 839 | * @template TReturn | |
| 840 | * | |
| 841 | * @param callable(): TReturn $callback | |
| 842 | * @param string $name The lock domain — distinct jobs, distinct names. | |
| 843 | * @return TReturn | |
| 844 | * @throws UnsupportedFeatureException | |
| 845 | * @throws \Throwable | |
| 846 | */ | |
| 847 | public function withLock(callable $callback, string $name): mixed | |
| 848 | { | |
| 849 | throw new UnsupportedFeatureException( | |
| 850 | 'This dialect does not provide a native cross-process lock; supply a Lock adapter.', | |
| 851 | ); | |
| 852 | } | |
| 853 | ||
| 854 | // ---- Transactions (depth-counter + savepoints) ---- | |
| 855 | ||
| 856 | /** | |
| 857 | * The current transaction nesting depth. | |
| 858 | * | |
| 859 | * @var int | |
| 860 | */ | |
| 861 | private int $transactionLevel = 0; | |
| 862 | ||
| 863 | /** | |
| 864 | * Monotonic savepoint sequence — makes savepoint names unique. | |
| 865 | * | |
| 866 | * @var int | |
| 867 | */ | |
| 868 | private int $savepointSequence = 0; | |
| 869 | ||
| 870 | /** | |
| 871 | * The savepoint created by the currently-innermost open nested frame, | |
| 872 | * per depth. | |
| 873 | * | |
| 874 | * @var array<int, string> | |
| 875 | */ | |
| 876 | private array $savepointsByLevel = []; | |
| 877 | ||
| 878 | /** | |
| 879 | * The coroutine that opened the current transaction. | |
| 880 | * | |
| 881 | * @var string|null | |
| 882 | */ | |
| 883 | private ?string $transactionOwner = null; | |
| 884 | ||
| 885 | /** | |
| 886 | * The current transaction nesting depth. | |
| 887 | * | |
| 888 | * @return int | |
| 889 | */ | |
| 890 | final public function transactionLevel(): int | |
| 891 | { | |
| 892 | return $this->transactionLevel; | |
| 893 | } | |
| 894 | ||
| 895 | /** | |
| 896 | * Whether this dialect supports savepoints for nested transactions. | |
| 897 | * | |
| 898 | * @return bool | |
| 899 | */ | |
| 900 | abstract protected function supportsSavepoints(): bool; | |
| 901 | ||
| 902 | /** | |
| 903 | * Create a named savepoint. | |
| 904 | * | |
| 905 | * @param string $name | |
| 906 | */ | |
| 907 | abstract protected function createSavepoint(string $name): void; | |
| 908 | ||
| 909 | /** | |
| 910 | * Release a named savepoint. | |
| 911 | * | |
| 912 | * @param string $name | |
| 913 | */ | |
| 914 | abstract protected function releaseSavepoint(string $name): void; | |
| 915 | ||
| 916 | /** | |
| 917 | * Roll back to a named savepoint. | |
| 918 | * | |
| 919 | * @param string $name | |
| 920 | */ | |
| 921 | abstract protected function rollbackToSavepoint(string $name): void; | |
| 922 | ||
| 923 | /** | |
| 924 | * Begin a transaction, nesting via savepoints when supported. | |
| 925 | * | |
| 926 | * @throws \LogicException | |
| 927 | */ | |
| 928 | final public function beginTransaction(): void | |
| 929 | { | |
| 930 | $this->assertSameCoroutine('beginTransaction'); | |
| 931 | $toLevel = $this->transactionLevel + 1; | |
| 932 | if ($toLevel === 1) { | |
| 933 | $this->pdo->beginTransaction(); | |
| 934 | $this->transactionOwner = $this->coroutineId(); | |
| 935 | } elseif ($this->supportsSavepoints()) { | |
| 936 | // Unique per-frame name (depth + sequence): depth alone collides | |
| 937 | // when interleaved coroutine frames nest on one connection. | |
| 938 | $name = 'trans' . $toLevel . '_' . (++$this->savepointSequence); | |
| 939 | $this->createSavepoint($name); | |
| 940 | $this->savepointsByLevel[$toLevel] = $name; | |
| 941 | } | |
| 942 | $this->transactionLevel = $toLevel; | |
| 943 | } | |
| 944 | ||
| 945 | /** | |
| 946 | * Commit the current transaction (or release the innermost savepoint). | |
| 947 | */ | |
| 948 | final public function commit(): void | |
| 949 | { | |
| 950 | $this->assertSameCoroutine('commit'); | |
| 951 | $toLevel = $this->transactionLevel - 1; | |
| 952 | $this->transactionLevel = $toLevel; | |
| 953 | if ($toLevel === 0) { | |
| 954 | try { | |
| 955 | $this->pdo->commit(); | |
| 956 | } catch (\PDOException $e) { | |
| 957 | $this->reconcileFailedCommit(); | |
| 958 | throw $e; | |
| 959 | } | |
| 960 | $this->transactionOwner = null; | |
| 961 | } elseif ($this->supportsSavepoints()) { | |
| 962 | $this->releaseSavepoint($this->savepointNameFor($toLevel + 1)); | |
| 963 | } | |
| 964 | } | |
| 965 | ||
| 966 | /** | |
| 967 | * Roll back the current transaction (or to the innermost savepoint). | |
| 968 | */ | |
| 969 | final public function rollBack(): void | |
| 970 | { | |
| 971 | $this->assertSameCoroutine('rollBack'); | |
| 972 | $toLevel = $this->transactionLevel - 1; | |
| 973 | $this->transactionLevel = $toLevel; | |
| 974 | if ($toLevel === 0) { | |
| 975 | try { | |
| 976 | $this->pdo->rollBack(); | |
| 977 | } catch (\PDOException $e) { | |
| 978 | $this->reconcileFailedCommit(); | |
| 979 | throw $e; | |
| 980 | } | |
| 981 | $this->transactionOwner = null; | |
| 982 | } elseif ($this->supportsSavepoints()) { | |
| 983 | $this->rollbackToSavepoint($this->savepointNameFor($toLevel + 1)); | |
| 984 | } | |
| 985 | } | |
| 986 | ||
| 987 | /** | |
| 988 | * Best-effort clear of a server-side transaction after a failed | |
| 989 | * top-level commit/rollback. | |
| 990 | */ | |
| 991 | private function reconcileFailedCommit(): void | |
| 992 | { | |
| 993 | try { | |
| 994 | $this->pdo->rollBack(); | |
| 995 | } catch (\Throwable) { | |
| 996 | // Nothing more can be done here — the connection is dead or the | |
| 997 | // transaction is already gone; staleness handling takes over. | |
| 998 | } | |
| 999 | } | |
| 1000 | ||
| 1001 | /** | |
| 1002 | * The current coroutine's identity, best-effort. | |
| 1003 | * | |
| 1004 | * @return string | |
| 1005 | */ | |
| 1006 | private function coroutineId(): string | |
| 1007 | { | |
| 1008 | if (\class_exists(\Fiber::class)) { | |
| 1009 | $fiber = \Fiber::getCurrent(); | |
| 1010 | if ($fiber !== null) { | |
| 1011 | return 'fiber:' . (string) \spl_object_id($fiber); | |
| 1012 | } | |
| 1013 | } | |
| 1014 | if (\class_exists('Swoole\Coroutine') | |
| 1015 | && ($cid = \Swoole\Coroutine::getCid()) > 0 | |
| 1016 | ) { | |
| 1017 | return 'swoole:' . (string) $cid; | |
| 1018 | } | |
| 1019 | $pid = getmypid(); | |
| 1020 | return 'proc:' . ($pid === false ? 'unknown' : (string) $pid); | |
| 1021 | } | |
| 1022 | ||
| 1023 | /** | |
| 1024 | * Fail fast when a different coroutine touches an open transaction. | |
| 1025 | * | |
| 1026 | * @param string $operation | |
| 1027 | * @throws \LogicException | |
| 1028 | */ | |
| 1029 | private function assertSameCoroutine(string $operation): void | |
| 1030 | { | |
| 1031 | if ($this->transactionLevel > 0 | |
| 1032 | && $this->transactionOwner !== null | |
| 1033 | && $this->transactionOwner !== $this->coroutineId() | |
| 1034 | ) { | |
| 1035 | throw new \LogicException( | |
| 1036 | "{$operation}() called from a different coroutine than the one that opened" | |
| 1037 | . ' the transaction. Connections are not coroutine-safe while a transaction' | |
| 1038 | . ' is open — use one connection per coroutine.' | |
| 1039 | ); | |
| 1040 | } | |
| 1041 | } | |
| 1042 | ||
| 1043 | /** | |
| 1044 | * The savepoint name the frame at a given depth created, forgetting it. | |
| 1045 | * | |
| 1046 | * @param int $level | |
| 1047 | * @return string | |
| 1048 | */ | |
| 1049 | private function savepointNameFor(int $level): string | |
| 1050 | { | |
| 1051 | $name = $this->savepointsByLevel[$level] | |
| 1052 | ?? 'trans' . $level; | |
| 1053 | ||
| 1054 | unset($this->savepointsByLevel[$level]); | |
| 1055 | ||
| 1056 | return $name; | |
| 1057 | } | |
| 1058 | ||
| 1059 | /** | |
| 1060 | * Run a callback inside a transaction, committing on success and | |
| 1061 | * rolling back on any exception. | |
| 1062 | * | |
| 1063 | * @param callable(SqlConnection): mixed $callback | |
| 1064 | * @return mixed | |
| 1065 | * @throws \Throwable | |
| 1066 | */ | |
| 1067 | final public function transaction(callable $callback): mixed | |
| 1068 | { | |
| 1069 | $this->beginTransaction(); | |
| 1070 | try { | |
| 1071 | $result = $callback($this); | |
| 1072 | $this->commit(); | |
| 1073 | return $result; | |
| 1074 | } catch (\Throwable $e) { | |
| 1075 | try { | |
| 1076 | $this->rollBack(); | |
| 1077 | } catch (\Throwable) { | |
| 1078 | // The original failure is what the caller needs; a failed | |
| 1079 | // rollback is secondary (and usually shares its cause). | |
| 1080 | } | |
| 1081 | throw $e; | |
| 1082 | } | |
| 1083 | } | |
| 1084 | ||
| 1085 | /** | |
| 1086 | * Best-effort rollback of an abandoned transaction at teardown. | |
| 1087 | * | |
| 1088 | * @return void | |
| 1089 | */ | |
| 1090 | public function __destruct() | |
| 1091 | { | |
| 1092 | if ($this->transactionLevel === 0) { | |
| 1093 | return; | |
| 1094 | } | |
| 1095 | try { | |
| 1096 | $this->transactionLevel = 0; | |
| 1097 | $this->pdo->rollBack(); | |
| 1098 | } catch (\Throwable) { | |
| 1099 | // Teardown is best-effort: the connection may already be dead | |
| 1100 | // (which also releases the server-side transaction). | |
| 1101 | } | |
| 1102 | } | |
| 1103 | } |