Lines
95.05%
250 / 263
Methods
86.44%
51 / 59
Classes
0.00%
0 / 1
| Name | Lines | Methods | CRAP | ||||
|---|---|---|---|---|---|---|---|
| getDefaultQueryGrammar | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | ||
| getDefaultSchemaGrammar | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | ||
| getDefaultSchemaInspector | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | ||
| getDefaultValueCodec | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | ||
| supportsSavepoints | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | ||
| supportsTransactionalDdl | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | ||
| createSavepoint | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | ||
| releaseSavepoint | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | ||
| rollbackToSavepoint | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | ||
| withLock | 100.00% | 2 / 2 | 100.00% | 1 / 1 | 1 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] __construct | 100.00% | 5 / 5 | 100.00% | 1 / 1 | 1 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] assertSql | 100.00% | 2 / 2 | 100.00% | 1 / 1 | 2 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] from | 100.00% | 6 / 6 | 100.00% | 1 / 1 | 2 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] table | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] select | 100.00% | 8 / 8 | 100.00% | 1 / 1 | 3 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] selectColumn | 100.00% | 2 / 2 | 100.00% | 1 / 1 | 1 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] insert | 100.00% | 4 / 4 | 100.00% | 1 / 1 | 1 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] insertGetId | 100.00% | 23 / 23 | 100.00% | 1 / 1 | 8 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] update | 100.00% | 2 / 2 | 100.00% | 1 / 1 | 1 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] delete | 100.00% | 2 / 2 | 100.00% | 1 / 1 | 1 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] cursor | 100.00% | 2 / 2 | 100.00% | 1 / 1 | 1 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] flattenInsertValues | 100.00% | 5 / 5 | 100.00% | 1 / 1 | 2 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] selectSql | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] selectColumnSql | 100.00% | 2 / 2 | 100.00% | 1 / 1 | 1 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] cursorSql | 100.00% | 4 / 4 | 100.00% | 1 / 1 | 2 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] chunkSql | 100.00% | 11 / 11 | 100.00% | 1 / 1 | 6 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] prepareAndExecute | 100.00% | 8 / 8 | 100.00% | 1 / 1 | 3 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] statement | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] affectingStatement | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] bindValues | 100.00% | 11 / 11 | 100.00% | 1 / 1 | 10 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] run | 100.00% | 8 / 8 | 100.00% | 1 / 1 | 3 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] changeRequiresStandaloneTransaction | 0.00% | 0 / 1 | 0.00% | 0 / 1 | 2 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] create | 100.00% | 3 / 3 | 100.00% | 1 / 1 | 2 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] alter | 100.00% | 9 / 9 | 100.00% | 1 / 1 | 5 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] drop | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] apply | 93.33% | 14 / 15 | 0.00% | 0 / 1 | 13.05 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] renameTable | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] renameColumn | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] applyColumnRenames | 100.00% | 2 / 2 | 100.00% | 1 / 1 | 2 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] applyAddColumn | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] applyDropColumn | 100.00% | 5 / 5 | 100.00% | 1 / 1 | 2 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] applyModifyColumn | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] subjectBlueprint | 100.00% | 3 / 3 | 100.00% | 1 / 1 | 2 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] modifyColumn | 100.00% | 2 / 2 | 100.00% | 1 / 1 | 2 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] addForeignKey | 44.44% | 4 / 9 | 0.00% | 0 / 1 | 1.17 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] dropForeignKey | 80.00% | 4 / 5 | 0.00% | 0 / 1 | 1.01 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] addCheck | 66.66% | 4 / 6 | 0.00% | 0 / 1 | 1.04 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] dropCheck | 80.00% | 4 / 5 | 0.00% | 0 / 1 | 1.01 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] rebuildIndexes | 100.00% | 4 / 4 | 100.00% | 1 / 1 | 3 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] transactionLevel | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] beginTransaction | 100.00% | 10 / 10 | 100.00% | 1 / 1 | 3 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] commit | 100.00% | 11 / 11 | 100.00% | 1 / 1 | 4 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] rollBack | 100.00% | 11 / 11 | 100.00% | 1 / 1 | 4 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] reconcileFailedCommit | 100.00% | 2 / 2 | 100.00% | 1 / 1 | 2 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] coroutineId | 88.88% | 8 / 9 | 0.00% | 0 / 1 | 6.05 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] assertSameCoroutine | 100.00% | 8 / 8 | 100.00% | 1 / 1 | 4 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] savepointNameFor | 75.00% | 3 / 4 | 0.00% | 0 / 1 | 1.02 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] transaction | 100.00% | 8 / 8 | 100.00% | 1 / 1 | 3 | ||
| [BlueprintAU\Radiant\Database\Connections\SqlConnection] __destruct | 100.00% | 5 / 5 | 100.00% | 1 / 1 | 3 | ||
| 30 | final class PostgresConnection extends SqlConnection | |
| 31 | { | |
| 32 | /** | |
| 33 | * The default query grammar for this connection. | |
| 34 | * | |
| 35 | * @return Grammar | |
| 36 | */ | |
| 37 | protected function getDefaultQueryGrammar(): Grammar | |
| 38 | { | |
| 39 | return new PostgresGrammar(); | |
| 40 | } | |
| 41 | ||
| 42 | /** | |
| 43 | * The default schema grammar for this connection. | |
| 44 | * | |
| 45 | * @return SchemaGrammar | |
| 46 | */ | |
| 47 | protected function getDefaultSchemaGrammar(): SchemaGrammar | |
| 48 | { | |
| 49 | return new PostgresSchemaGrammar(); | |
| 50 | } | |
| 51 | ||
| 52 | /** | |
| 53 | * The dialect's live-schema reader. | |
| 54 | * | |
| 55 | * @return PostgresSchemaInspector | |
| 56 | */ | |
| 57 | protected function getDefaultSchemaInspector(): SchemaInspector | |
| 58 | { | |
| 59 | return new PostgresSchemaInspector($this->pdo); | |
| 60 | } | |
| 61 | ||
| 62 | /** | |
| 63 | * The default value codec for this connection. | |
| 64 | * | |
| 65 | * Postgres' native `timestamp` stores microseconds, so datetimes are | |
| 66 | * formatted as `Y-m-d H:i:s.u` on the write path. | |
| 67 | * | |
| 68 | * @return ValueCodecInterface | |
| 69 | */ | |
| 70 | protected function getDefaultValueCodec(): ValueCodecInterface | |
| 71 | { | |
| 72 | return new PostgresValueCodec(); | |
| 73 | } | |
| 74 | ||
| 75 | /** | |
| 76 | * Whether this dialect supports savepoints for nested transactions. | |
| 77 | * | |
| 78 | * @return bool | |
| 79 | */ | |
| 80 | protected function supportsSavepoints(): bool | |
| 81 | { | |
| 82 | return true; | |
| 83 | } | |
| 84 | ||
| 85 | /** | |
| 86 | * Postgres DDL is transactional — schema statements roll back with | |
| 87 | * the transaction. | |
| 88 | * | |
| 89 | * @return bool | |
| 90 | */ | |
| 91 | #[Override] | |
| 92 | public function supportsTransactionalDdl(): bool | |
| 93 | { | |
| 94 | return true; | |
| 95 | } | |
| 96 | ||
| 97 | /** | |
| 98 | * Create a named savepoint. | |
| 99 | * | |
| 100 | * @param string $name | |
| 101 | */ | |
| 102 | protected function createSavepoint(string $name): void | |
| 103 | { | |
| 104 | $this->pdo->exec("SAVEPOINT {$name}"); | |
| 105 | } | |
| 106 | ||
| 107 | /** | |
| 108 | * Release a named savepoint. | |
| 109 | * | |
| 110 | * @param string $name | |
| 111 | */ | |
| 112 | protected function releaseSavepoint(string $name): void | |
| 113 | { | |
| 114 | $this->pdo->exec("RELEASE SAVEPOINT {$name}"); | |
| 115 | } | |
| 116 | ||
| 117 | /** | |
| 118 | * Roll back to a named savepoint. | |
| 119 | * | |
| 120 | * @param string $name | |
| 121 | */ | |
| 122 | protected function rollbackToSavepoint(string $name): void | |
| 123 | { | |
| 124 | $this->pdo->exec("ROLLBACK TO SAVEPOINT {$name}"); | |
| 125 | } | |
| 126 | ||
| 127 | /** | |
| 128 | * Run the callback under a Postgres session advisory lock on the given | |
| 129 | * lock domain. | |
| 130 | * | |
| 131 | * @template TReturn | |
| 132 | * | |
| 133 | * @param callable(): TReturn $callback | |
| 134 | * @param string $name | |
| 135 | * @return TReturn | |
| 136 | * @throws \Throwable | |
| 137 | */ | |
| 138 | #[Override] | |
| 139 | public function withLock(callable $callback, string $name): mixed | |
| 140 | { | |
| 141 | return (new \BlueprintAU\Radiant\Database\Locks\PostgresLock($this)) | |
| 142 | ->withLock($callback, $name); | |
| 143 | } | |
| 144 | } |
Inherited from BlueprintAU\Radiant\Database\Connections\SqlConnection
| 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 | } |
| 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 | } |
| 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 | } |
| 126 | final public function table(string $identifier): QueryBuilder | |
| 127 | { | |
| 128 | return new QueryBuilder($this, $identifier); | |
| 129 | } |
| 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 | } |
| 161 | final public function selectColumn(QueryBuilder $query): Collection | |
| 162 | { | |
| 163 | $sql = $this->grammar->compileSelect($query); | |
| 164 | return $this->selectColumnSql($sql, $query->getBindings()); | |
| 165 | } |
| 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 | } |
| 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 | } |
| 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 | } |
| 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 | } |
| 287 | final public function cursor(QueryBuilder $query): \Generator | |
| 288 | { | |
| 289 | $sql = $this->grammar->compileSelect($query); | |
| 290 | return $this->cursorSql($sql, $query->getBindings()); | |
| 291 | } |
| 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 | } |
| 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 | } |
| 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 | } |
| 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 | } |
| 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 | } |
| 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 | } |
| 426 | final public function statement(string $sql, array $bindings = []): void | |
| 427 | { | |
| 428 | $this->run($sql, $bindings, fn() => null); | |
| 429 | } |
| 438 | final public function affectingStatement(string $sql, array $bindings = []): int | |
| 439 | { | |
| 440 | return $this->run($sql, $bindings, fn(\PDOStatement $stmt) => $stmt->rowCount()); | |
| 441 | } |
| 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 | } |
| 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 | } |
| 560 | public function changeRequiresStandaloneTransaction(\BlueprintAU\Radiant\Database\Schema\SchemaChange $change, array $plan = []): bool | |
| 561 | { | |
| 562 | return false; | |
| 563 | } |
| 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 | } |
| 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 | } |
| 608 | final public function drop(string $table): void | |
| 609 | { | |
| 610 | $this->statement($this->schemaGrammar->compileDrop($table)); | |
| 611 | } |
| 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 | } |
| 645 | final public function renameTable(string $from, string $to): void | |
| 646 | { | |
| 647 | $this->statement($this->schemaGrammar->compileRenameTable($from, $to)); | |
| 648 | } |
| 657 | final public function renameColumn(string $table, string $from, string $to): void | |
| 658 | { | |
| 659 | $this->statement($this->schemaGrammar->compileRenameColumn($table, $from, $to)); | |
| 660 | } |
| 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 | } |
| 680 | protected function applyAddColumn(\BlueprintAU\Radiant\Database\Schema\SchemaChange $change): void | |
| 681 | { | |
| 682 | $this->alter(SchemaOperation::AddColumn, $this->subjectBlueprint($change)); | |
| 683 | } |
| 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 | } |
| 711 | protected function applyModifyColumn(\BlueprintAU\Radiant\Database\Schema\SchemaChange $change): void | |
| 712 | { | |
| 713 | $this->modifyColumn($this->subjectBlueprint($change)); | |
| 714 | } |
| 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 | } |
| 739 | public function modifyColumn(Blueprint $blueprint): void | |
| 740 | { | |
| 741 | foreach ($this->schemaGrammar->compileModifyColumn($blueprint) as $sql) { | |
| 742 | $this->statement($sql); | |
| 743 | } | |
| 744 | } |
| 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 | } |
| 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 | } |
| 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 | } |
| 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 | } |
| 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 | } |
| 890 | final public function transactionLevel(): int | |
| 891 | { | |
| 892 | return $this->transactionLevel; | |
| 893 | } |
| 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 | } |
| 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 | } |
| 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 | } |
| 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 | } |
| 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 | } |
| 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 | } |
| 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 | } |
| 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 | } |
| 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 | } |