| 38 | | abstract class SqlConnection implements ConnectionInterface |
| 39 | | { |
| 40 | | use NormalizesInsertRows; |
| 41 | | use DetectsConnectionLoss; |
| 42 | | |
| 43 | | |
| 44 | | |
| 45 | | |
| 46 | | |
| 47 | | public readonly ValueCodecInterface $codec; |
| 48 | | |
| 49 | | |
| 50 | | |
| 51 | | |
| 52 | | |
| 53 | | |
| 54 | | public readonly Grammar $grammar; |
| 55 | | |
| 56 | | |
| 57 | | |
| 58 | | |
| 59 | | |
| 60 | | |
| 61 | | public readonly SchemaGrammar $schemaGrammar; |
| 62 | | |
| 63 | | |
| 64 | | |
| 65 | | |
| 66 | | |
| 67 | | |
| 68 | | public readonly \BlueprintAU\Radiant\Database\Schema\Inspectors\SchemaInspector $schemaInspector; |
| 69 | | |
| 70 | | |
| 71 | | |
| 72 | | |
| 73 | | |
| 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 | | |
| 86 | | |
| 87 | | |
| 88 | | |
| 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 | | |
| 99 | | |
| 100 | | |
| 101 | | |
| 102 | | |
| 103 | | |
| 104 | | |
| 105 | | final public static function from(ConnectionInterface $connection): static |
| 106 | | { |
| 107 | | if (!$connection instanceof static) { |
| 108 | | |
| 109 | | |
| 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 | | |
| 121 | | |
| 122 | | |
| 123 | | |
| 124 | | |
| 125 | | #[Override] |
| 126 | | final public function table(string $identifier): QueryBuilder |
| 127 | | { |
| 128 | | return new QueryBuilder($this, $identifier); |
| 129 | | } |
| 130 | | |
| 131 | | |
| 132 | | |
| 133 | | |
| 134 | | |
| 135 | | |
| 136 | | |
| 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 | | |
| 155 | | |
| 156 | | |
| 157 | | |
| 158 | | |
| 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 | | |
| 169 | | |
| 170 | | |
| 171 | | |
| 172 | | |
| 173 | | |
| 174 | | #[Override] |
| 175 | | final public function insert(QueryBuilder $query, array $values): int |
| 176 | | { |
| 177 | | $rows = $this->normalizeInsertRows($values); |
| 178 | | |
| 179 | | |
| 180 | | |
| 181 | | |
| 182 | | $this->assertUniformInsertRows($rows); |
| 183 | | |
| 184 | | $sql = $this->grammar->compileInsert($query, $values); |
| 185 | | return $this->affectingStatement($sql, $this->flattenInsertValues($values)); |
| 186 | | } |
| 187 | | |
| 188 | | |
| 189 | | |
| 190 | | |
| 191 | | |
| 192 | | |
| 193 | | |
| 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 | | |
| 202 | | $this->affectingStatement($this->grammar->compileInsert($query, $values), $this->flattenInsertValues($values)); |
| 203 | | return null; |
| 204 | | } |
| 205 | | |
| 206 | | |
| 207 | | |
| 208 | | |
| 209 | | |
| 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 | | |
| 224 | | |
| 225 | | |
| 226 | | |
| 227 | | |
| 228 | | |
| 229 | | |
| 230 | | $id = $this->pdo->lastInsertId(); |
| 231 | | if ($id === false) { |
| 232 | | return null; |
| 233 | | } |
| 234 | | |
| 235 | | |
| 236 | | |
| 237 | | |
| 238 | | |
| 239 | | |
| 240 | | |
| 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 | | |
| 254 | | |
| 255 | | |
| 256 | | |
| 257 | | |
| 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 | | |
| 268 | | |
| 269 | | |
| 270 | | |
| 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 | | |
| 281 | | |
| 282 | | |
| 283 | | |
| 284 | | |
| 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 | | |
| 295 | | |
| 296 | | |
| 297 | | |
| 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 | | |
| 310 | | |
| 311 | | |
| 312 | | |
| 313 | | |
| 314 | | |
| 315 | | |
| 316 | | |
| 317 | | |
| 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 | | |
| 326 | | |
| 327 | | |
| 328 | | |
| 329 | | |
| 330 | | |
| 331 | | final public function selectColumnSql(string $sql, array $bindings = []): Collection |
| 332 | | { |
| 333 | | |
| 334 | | $column = $this->run($sql, $bindings, fn(\PDOStatement $stmt) => $stmt->fetchAll(\PDO::FETCH_COLUMN, 0)); |
| 335 | | return Collection::make($column); |
| 336 | | } |
| 337 | | |
| 338 | | |
| 339 | | |
| 340 | | |
| 341 | | |
| 342 | | |
| 343 | | |
| 344 | | |
| 345 | | |
| 346 | | |
| 347 | | |
| 348 | | |
| 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 | | |
| 364 | | |
| 365 | | |
| 366 | | |
| 367 | | |
| 368 | | |
| 369 | | |
| 370 | | |
| 371 | | |
| 372 | | |
| 373 | | |
| 374 | | |
| 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 | | |
| 398 | | |
| 399 | | |
| 400 | | |
| 401 | | |
| 402 | | |
| 403 | | |
| 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 | | |
| 422 | | |
| 423 | | |
| 424 | | |
| 425 | | |
| 426 | | final public function statement(string $sql, array $bindings = []): void |
| 427 | | { |
| 428 | | $this->run($sql, $bindings, fn() => null); |
| 429 | | } |
| 430 | | |
| 431 | | |
| 432 | | |
| 433 | | |
| 434 | | |
| 435 | | |
| 436 | | |
| 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 | | |
| 445 | | |
| 446 | | |
| 447 | | |
| 448 | | |
| 449 | | |
| 450 | | protected function bindValues(\PDOStatement $stmt, array $bindings): void |
| 451 | | { |
| 452 | | foreach ($bindings as $key => $value) { |
| 453 | | |
| 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 | | |
| 462 | | |
| 463 | | |
| 464 | | |
| 465 | | |
| 466 | | |
| 467 | | |
| 468 | | |
| 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 | | |
| 482 | | |
| 483 | | |
| 484 | | |
| 485 | | |
| 486 | | |
| 487 | | |
| 488 | | |
| 489 | | |
| 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 | | |
| 507 | | |
| 508 | | |
| 509 | | |
| 510 | | |
| 511 | | |
| 512 | | |
| 513 | | protected function getDefaultValueCodec(): ValueCodecInterface |
| 514 | | { |
| 515 | | return new DefaultValueCodec(); |
| 516 | | } |
| 517 | | |
| 518 | | |
| 519 | | |
| 520 | | |
| 521 | | |
| 522 | | |
| 523 | | abstract protected function getDefaultQueryGrammar(): Grammar; |
| 524 | | |
| 525 | | |
| 526 | | |
| 527 | | |
| 528 | | |
| 529 | | |
| 530 | | abstract protected function getDefaultSchemaGrammar(): SchemaGrammar; |
| 531 | | |
| 532 | | |
| 533 | | |
| 534 | | |
| 535 | | |
| 536 | | |
| 537 | | abstract protected function getDefaultSchemaInspector(): \BlueprintAU\Radiant\Database\Schema\Inspectors\SchemaInspector; |
| 538 | | |
| 539 | | |
| 540 | | |
| 541 | | |
| 542 | | |
| 543 | | |
| 544 | | public function supportsTransactionalDdl(): bool |
| 545 | | { |
| 546 | | return false; |
| 547 | | } |
| 548 | | |
| 549 | | |
| 550 | | |
| 551 | | |
| 552 | | |
| 553 | | |
| 554 | | |
| 555 | | |
| 556 | | |
| 557 | | |
| 558 | | |
| 559 | | |
| 560 | | public function changeRequiresStandaloneTransaction(\BlueprintAU\Radiant\Database\Schema\SchemaChange $change, array $plan = []): bool |
| 561 | | { |
| 562 | | return false; |
| 563 | | } |
| 564 | | |
| 565 | | |
| 566 | | |
| 567 | | |
| 568 | | |
| 569 | | |
| 570 | | |
| 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 | | |
| 583 | | |
| 584 | | |
| 585 | | |
| 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 | | |
| 605 | | |
| 606 | | |
| 607 | | |
| 608 | | final public function drop(string $table): void |
| 609 | | { |
| 610 | | $this->statement($this->schemaGrammar->compileDrop($table)); |
| 611 | | } |
| 612 | | |
| 613 | | |
| 614 | | |
| 615 | | |
| 616 | | |
| 617 | | |
| 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 | | |
| 641 | | |
| 642 | | |
| 643 | | |
| 644 | | |
| 645 | | final public function renameTable(string $from, string $to): void |
| 646 | | { |
| 647 | | $this->statement($this->schemaGrammar->compileRenameTable($from, $to)); |
| 648 | | } |
| 649 | | |
| 650 | | |
| 651 | | |
| 652 | | |
| 653 | | |
| 654 | | |
| 655 | | |
| 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 | | |
| 664 | | |
| 665 | | |
| 666 | | |
| 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 | | |
| 677 | | |
| 678 | | |
| 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 | | |
| 687 | | |
| 688 | | |
| 689 | | |
| 690 | | |
| 691 | | |
| 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 | | |
| 708 | | |
| 709 | | |
| 710 | | |
| 711 | | protected function applyModifyColumn(\BlueprintAU\Radiant\Database\Schema\SchemaChange $change): void |
| 712 | | { |
| 713 | | $this->modifyColumn($this->subjectBlueprint($change)); |
| 714 | | } |
| 715 | | |
| 716 | | |
| 717 | | |
| 718 | | |
| 719 | | |
| 720 | | |
| 721 | | |
| 722 | | |
| 723 | | |
| 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 | | |
| 736 | | |
| 737 | | |
| 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 | | |
| 748 | | |
| 749 | | |
| 750 | | |
| 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 | | |
| 768 | | |
| 769 | | |
| 770 | | |
| 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 | | |
| 784 | | |
| 785 | | |
| 786 | | |
| 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 | | |
| 796 | | |
| 797 | | $name = $check['name']; |
| 798 | | |
| 799 | | $this->statement($this->schemaGrammar->compileAddCheck($table, $name, $check['expression'])); |
| 800 | | } |
| 801 | | |
| 802 | | |
| 803 | | |
| 804 | | |
| 805 | | |
| 806 | | |
| 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 | | |
| 820 | | |
| 821 | | |
| 822 | | |
| 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 | | |
| 837 | | |
| 838 | | |
| 839 | | |
| 840 | | |
| 841 | | |
| 842 | | |
| 843 | | |
| 844 | | |
| 845 | | |
| 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 | | |
| 855 | | |
| 856 | | |
| 857 | | |
| 858 | | |
| 859 | | |
| 860 | | |
| 861 | | private int $transactionLevel = 0; |
| 862 | | |
| 863 | | |
| 864 | | |
| 865 | | |
| 866 | | |
| 867 | | |
| 868 | | private int $savepointSequence = 0; |
| 869 | | |
| 870 | | |
| 871 | | |
| 872 | | |
| 873 | | |
| 874 | | |
| 875 | | |
| 876 | | private array $savepointsByLevel = []; |
| 877 | | |
| 878 | | |
| 879 | | |
| 880 | | |
| 881 | | |
| 882 | | |
| 883 | | private ?string $transactionOwner = null; |
| 884 | | |
| 885 | | |
| 886 | | |
| 887 | | |
| 888 | | |
| 889 | | |
| 890 | | final public function transactionLevel(): int |
| 891 | | { |
| 892 | | return $this->transactionLevel; |
| 893 | | } |
| 894 | | |
| 895 | | |
| 896 | | |
| 897 | | |
| 898 | | |
| 899 | | |
| 900 | | abstract protected function supportsSavepoints(): bool; |
| 901 | | |
| 902 | | |
| 903 | | |
| 904 | | |
| 905 | | |
| 906 | | |
| 907 | | abstract protected function createSavepoint(string $name): void; |
| 908 | | |
| 909 | | |
| 910 | | |
| 911 | | |
| 912 | | |
| 913 | | |
| 914 | | abstract protected function releaseSavepoint(string $name): void; |
| 915 | | |
| 916 | | |
| 917 | | |
| 918 | | |
| 919 | | |
| 920 | | |
| 921 | | abstract protected function rollbackToSavepoint(string $name): void; |
| 922 | | |
| 923 | | |
| 924 | | |
| 925 | | |
| 926 | | |
| 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 | | |
| 937 | | |
| 938 | | $name = 'trans' . $toLevel . '_' . (++$this->savepointSequence); |
| 939 | | $this->createSavepoint($name); |
| 940 | | $this->savepointsByLevel[$toLevel] = $name; |
| 941 | | } |
| 942 | | $this->transactionLevel = $toLevel; |
| 943 | | } |
| 944 | | |
| 945 | | |
| 946 | | |
| 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 | | |
| 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 | | |
| 989 | | |
| 990 | | |
| 991 | | private function reconcileFailedCommit(): void |
| 992 | | { |
| 993 | | try { |
| 994 | | $this->pdo->rollBack(); |
| 995 | | } catch (\Throwable) { |
| 996 | | |
| 997 | | |
| 998 | | } |
| 999 | | } |
| 1000 | | |
| 1001 | | |
| 1002 | | |
| 1003 | | |
| 1004 | | |
| 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 | | |
| 1025 | | |
| 1026 | | |
| 1027 | | |
| 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 | | |
| 1045 | | |
| 1046 | | |
| 1047 | | |
| 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 | | |
| 1061 | | |
| 1062 | | |
| 1063 | | |
| 1064 | | |
| 1065 | | |
| 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 | | |
| 1079 | | |
| 1080 | | } |
| 1081 | | throw $e; |
| 1082 | | } |
| 1083 | | } |
| 1084 | | |
| 1085 | | |
| 1086 | | |
| 1087 | | |
| 1088 | | |
| 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 | | |
| 1100 | | |
| 1101 | | } |
| 1102 | | } |
| 1103 | | } |