From 04f2530fe4657df9867c9a3d681100985cfaa1e0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?T=C3=B5nis=20Ormisson?= Date: Tue, 8 Sep 2026 14:06:56 +0300 Subject: [PATCH 1/2] test: specify bounded server case batches and late-batch rollback --- .../MySqlFamilySpssRoundTripTestCase.php | 98 ++++++++ .../PostgreSqlSpssRoundTripTest.php | 93 +++++++- tests/Sql/MySqlWideTableImporterTest.php | 209 +++++++++++++----- tests/Sql/PostgreSqlWideTableImporterTest.php | 161 ++++++++++++-- 4 files changed, 484 insertions(+), 77 deletions(-) diff --git a/tests/Integration/MySqlFamilySpssRoundTripTestCase.php b/tests/Integration/MySqlFamilySpssRoundTripTestCase.php index 839d0aa..0cf8378 100644 --- a/tests/Integration/MySqlFamilySpssRoundTripTestCase.php +++ b/tests/Integration/MySqlFamilySpssRoundTripTestCase.php @@ -4,10 +4,14 @@ namespace OpenStatSpec\Tests\Integration; +use OpenStatSpec\Sql\Connection; +use OpenStatSpec\Sql\MySqlProfile; +use OpenStatSpec\Sql\MySqlWideTableImporter; use OpenStatSpec\Spss\PhpSpssEngine; use OpenStatSpec\Spss\SpssAdapter; use PDO; use PDOException; +use PHPUnit\Framework\Attributes\DataProvider; use PHPUnit\Framework\TestCase; use RuntimeException; use SPSS\Sav\Alignment; @@ -146,6 +150,100 @@ public function testRealEngineRoundTripsSavAndZsavThroughMySqlFamily(): void } } + /** @return iterable */ + public static function batchImports(): iterable + { + foreach (['native' => false, 'emulated' => true] as $mode => $emulated) { + foreach (['mixed', 'wide NULL', 'escaped packets', 'late failure'] as $shape) { + yield "$mode $shape" => [$emulated, $shape]; + } + } + } + + #[DataProvider('batchImports')] + public function testBoundedCaseImportOnServer(bool $emulated, string $shape): void + { + $pdo = $this->mysql(); + $pdo->setAttribute(PDO::ATTR_EMULATE_PREPARES, $emulated); + $profile = (new Connection($pdo))->profile; + self::assertInstanceOf(MySqlProfile::class, $profile); + $token = bin2hex(random_bytes(6)); + $priorName = 'batch prior ' . $token; + $attemptName = 'batch attempt ' . $token; + $priorTable = $profile->physicalIdentifier('dataset_' . $priorName); + $attemptTable = $profile->physicalIdentifier('dataset_' . $attemptName); + try { + // Keep InnoDB's inline row below a page while exceeding the parameter limit. + $width = $shape === 'wide NULL' ? min(511, $profile->maximumSourceVariables()) : 2; + $variables = $shape === 'wide NULL' + ? array_map(static fn(int $i): array => ['name' => 'v' . $i, 'type' => 'numeric'], range(1, $width)) + : [['name' => 'Score', 'type' => 'numeric'], ['name' => 'Comment', 'type' => 'string']]; + $count = $shape === 'wide NULL' ? 2 * intdiv(65535, $width + 1) + 1 : 513; + $text = ''; + if ($shape === 'escaped packets') { + $budget = $profile->effectiveMaximumStatementBytes($pdo); + self::assertGreaterThan(8, $budget, 'The dedicated service must allow a preflight-valid row.'); + $text = str_repeat("'\\\\", intdiv(min(50000, $budget - 8), 3)); + $count = 41; + } + $rows = []; + for ($i = 0; $i < $count; ++$i) { + $rows[] = $shape === 'wide NULL' ? array_fill(0, $width, null) + : [[0.1, null, 42, 1.0000000000000002, PHP_FLOAT_MAX, 5.0e-324][$i % 6], $shape === 'escaped packets' ? $text : ($i % 2 === 0 ? "õ'\\\\ $i" : '')]; + } + $importer = new MySqlWideTableImporter($pdo, $profile); + $importer->import(['variables' => $variables, 'data' => [$rows[0]]], $priorName); + $prior = $this->rows($pdo, 'SELECT * FROM ' . $this->quote($priorTable), []); + if ($shape === 'late failure') { + // Real server CHECK failure on row 257, confined to this attempt's DDL. + $faultPdo = $this->createMock(PDO::class); + foreach (['getAttribute', 'setAttribute', 'inTransaction', 'beginTransaction', 'commit', 'rollBack', 'query', 'prepare'] as $method) { + $faultPdo->method($method)->willReturnCallback($pdo->$method(...)); + } + $faultPdo->method('exec')->willReturnCallback(function (string $sql) use ($pdo, $attemptTable, $token): int|false { + $result = $pdo->exec($sql); + if (str_starts_with($sql, 'CREATE TABLE ' . $this->quote($attemptTable) . ' ')) { + $pdo->exec('ALTER TABLE ' . $this->quote($attemptTable) . ' ADD CONSTRAINT ' . $this->quote('reject_batch_tail_' . $token) . ' CHECK (__case_ordinal < 257)'); + } + return $result; + }); + try { + (new MySqlWideTableImporter($faultPdo, $profile))->import(['variables' => $variables, 'data' => $rows], $attemptName); + self::fail('The server accepted the forbidden tail row.'); + } catch (PDOException $exception) { + self::assertStringContainsString('reject_batch_tail', $exception->getMessage()); + } + self::assertSame(0, (int) $this->scalar($pdo, 'SELECT COUNT(*) FROM information_schema.TABLES WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = ?', [$attemptTable])); + self::assertSame(0, (int) $this->scalar($pdo, 'SELECT COUNT(*) FROM datasets WHERE dataset_name = ?', [$attemptName])); + self::assertSame(0, (int) $this->scalar($pdo, 'SELECT COUNT(*) FROM variables WHERE dataset_name = ?', [$attemptName])); + } else { + $definition = $importer->import(['variables' => $variables, 'data' => $rows], $attemptName); + $actual = $this->rows($pdo, 'SELECT * FROM ' . $this->quote($definition->tableName) . ' ORDER BY __case_ordinal', []); + self::assertCount(count($rows), $actual); + foreach ($actual as $i => $row) { + self::assertSame($i + 1, (int) array_shift($row)); + foreach (array_values($row) as $j => $value) { + $expected = $rows[$i][$j]; + if (is_int($expected) || is_float($expected)) { + self::assertSame(pack('E', (float) $expected), pack('E', (float) $value)); + } else { + self::assertSame($expected, $value); + } + } + } + } + self::assertSame($prior, $this->rows($pdo, 'SELECT * FROM ' . $this->quote($priorTable), [])); + self::assertFalse($pdo->inTransaction()); + self::assertSame($emulated, $pdo->getAttribute(PDO::ATTR_EMULATE_PREPARES)); + } finally { + if ($pdo->inTransaction()) { + $pdo->rollBack(); + } + $this->cleanup($pdo, $attemptName, $attemptTable); + $this->cleanup($pdo, $priorName, $priorTable); + } + } + abstract protected function serviceName(): string; abstract protected function environmentPrefix(): string; diff --git a/tests/Integration/PostgreSqlSpssRoundTripTest.php b/tests/Integration/PostgreSqlSpssRoundTripTest.php index 355745c..ff2af59 100644 --- a/tests/Integration/PostgreSqlSpssRoundTripTest.php +++ b/tests/Integration/PostgreSqlSpssRoundTripTest.php @@ -9,6 +9,7 @@ use OpenStatSpec\Sql\CatalogOwnership; use OpenStatSpec\Sql\Connection; use OpenStatSpec\Sql\NormativeCatalog; +use OpenStatSpec\Sql\PostgreSqlWideTableImporter; use OpenStatSpec\Spss\PhpSpssEngine; use OpenStatSpec\Transformation\Audit\TransformationAuditMigrator; use OpenStatSpec\Transformation\Execution\InPlaceApplyRequest; @@ -214,7 +215,13 @@ public function testImportFinalizationFailureRollsBackOnlyThisAttempt(bool $jour $pdo->exec('CREATE SCHEMA ' . $this->quote($schema)); try { $pdo->exec('SET search_path TO ' . $this->quote($schema)); - $engine = new FakeSpssEngine($this->fixture('sav')); + $fixture = $this->fixture('sav'); + $engine = new FakeSpssEngine(new Dataset( + new VariableDictionary($fixture->variables()), + array_merge(...array_fill(0, 257, $fixture->rows())), + $fixture->metadata, + $fixture->technicalMetadata, + )); $prior = (new SpssAdapter($pdo, $engine))->import('prior.sav', 'prior'); $before = []; foreach ($this->rows($pdo, "SELECT table_name FROM information_schema.tables WHERE table_schema = current_schema() AND table_type = 'BASE TABLE' AND table_name NOT IN ('operation_catalog', 'operation', 'fidelity_event_catalog', 'fidelity_event')", []) as $table) { @@ -261,6 +268,90 @@ public function testImportFinalizationFailureRollsBackOnlyThisAttempt(bool $jour } } + /** @return iterable */ + public static function batchImports(): iterable + { + foreach (['native' => false, 'emulated' => true] as $mode => $emulated) { + foreach (['mixed', 'wide NULL', 'late failure'] as $shape) { + yield "$mode $shape" => [$emulated, $shape]; + } + } + } + + #[DataProvider('batchImports')] + public function testBoundedCaseImportOnServer(bool $emulated, string $shape): void + { + $pdo = $this->postgres(); + $pdo->setAttribute(PDO::ATTR_EMULATE_PREPARES, $emulated); + $schema = 'case_batches_' . bin2hex(random_bytes(6)); + $searchPath = (string) $this->scalar($pdo, 'SHOW search_path', []); + $pdo->exec('CREATE SCHEMA ' . $this->quote($schema)); + try { + $pdo->exec('SET search_path TO ' . $this->quote($schema)); + $variables = $shape === 'wide NULL' + ? array_map(static fn(int $i): array => ['name' => 'v' . $i, 'type' => 'numeric'], range(1, 1599)) + : [['name' => 'Score', 'type' => 'numeric'], ['name' => 'Comment', 'type' => 'string']]; + $rows = []; + for ($i = 0; $i < ($shape === 'wide NULL' ? 81 : 513); ++$i) { + $rows[] = $shape === 'wide NULL' ? array_fill(0, 1599, null) + : [[0.1, null, 42, 1.0000000000000002, PHP_FLOAT_MAX, 5.0e-324][$i % 6], $i % 2 === 0 ? "õ'\\\\ $i" : '']; + } + $importer = new PostgreSqlWideTableImporter($pdo); + $importer->import(['variables' => $variables, 'data' => [$rows[0]]], 'prior'); + $prior = $this->rows($pdo, 'SELECT * FROM dataset_prior', []); + if ($shape === 'late failure') { + // Delegate to the live connection; only the owned attempt DDL gets a fault constraint. + $faultPdo = $this->createMock(PDO::class); + foreach (['getAttribute', 'setAttribute', 'inTransaction', 'beginTransaction', 'commit', 'rollBack', 'query', 'prepare'] as $method) { + $faultPdo->method($method)->willReturnCallback($pdo->$method(...)); + } + $faultPdo->method('exec')->willReturnCallback(static function (string $sql) use ($pdo): int|false { + if (str_starts_with($sql, 'CREATE TABLE "dataset_attempt" ')) { + $sql = substr($sql, 0, -1) . ', CONSTRAINT reject_batch_tail CHECK (__case_ordinal < 257))'; + } + return $pdo->exec($sql); + }); + $finalized = false; + try { + (new PostgreSqlWideTableImporter($faultPdo))->import(['variables' => $variables, 'data' => $rows], 'attempt', beforeCommit: static function () use (&$finalized): void { + $finalized = true; + }); + self::fail('The server accepted the forbidden tail row.'); + } catch (PDOException $exception) { + self::assertStringContainsString('reject_batch_tail', $exception->getMessage()); + } + self::assertFalse($finalized); + self::assertNull($this->scalar($pdo, "SELECT to_regclass('dataset_attempt')", [])); + self::assertSame(0, (int) $this->scalar($pdo, "SELECT COUNT(*) FROM datasets WHERE dataset_name = 'attempt'", [])); + self::assertSame(0, (int) $this->scalar($pdo, "SELECT COUNT(*) FROM variables WHERE dataset_name = 'attempt'", [])); + } else { + $definition = $importer->import(['variables' => $variables, 'data' => $rows], 'attempt'); + $actual = $this->rows($pdo, 'SELECT * FROM ' . $this->quote($definition->tableName) . ' ORDER BY __case_ordinal', []); + self::assertCount(count($rows), $actual); + foreach ($actual as $i => $row) { + self::assertSame($i + 1, (int) array_shift($row)); + foreach (array_values($row) as $j => $value) { + $expected = $rows[$i][$j]; + if (is_int($expected) || is_float($expected)) { + self::assertSame(pack('E', (float) $expected), pack('E', (float) $value)); + } else { + self::assertSame($expected, $value); + } + } + } + } + self::assertSame($prior, $this->rows($pdo, 'SELECT * FROM dataset_prior', [])); + self::assertFalse($pdo->inTransaction()); + self::assertSame($emulated, $pdo->getAttribute(PDO::ATTR_EMULATE_PREPARES)); + } finally { + if ($pdo->inTransaction()) { + $pdo->rollBack(); + } + $pdo->prepare("SELECT set_config('search_path', ?, false)")->execute([$searchPath]); + $pdo->exec('DROP SCHEMA ' . $this->quote($schema) . ' CASCADE'); + } + } + private function postgres(): PDO { $dsn = getenv('OPENSTATSPEC_PG_DSN'); diff --git a/tests/Sql/MySqlWideTableImporterTest.php b/tests/Sql/MySqlWideTableImporterTest.php index 2f25987..56a0ac0 100644 --- a/tests/Sql/MySqlWideTableImporterTest.php +++ b/tests/Sql/MySqlWideTableImporterTest.php @@ -8,6 +8,7 @@ use OpenStatSpec\Core\UnsupportedOperation; use OpenStatSpec\Sql\DoltProfile; use OpenStatSpec\Sql\MySqlWideTableImporter; +use OpenStatSpec\Sql\MySqlProfile; use PDO; use PDOStatement; use PHPUnit\Framework\Attributes\DataProvider; @@ -40,52 +41,118 @@ public function testRejectsCallerOwnedTransactionBeforeMutation(): void } } - public function testImportsCatalogueAndOrderedRowsAfterMysqlDdl(): void + /** @return iterable}> */ + public static function caseBatches(): iterable + { + foreach (['mysql' => new MySqlProfile(), 'dolt' => new DoltProfile()] as $name => $profile) { + foreach ([0 => [], 1 => [1], 2 => [2], 513 => [256, 256, 1]] as $count => $sizes) { + yield "$name $count rows" => [$profile, $count, 2, 0, 1073741824, $sizes]; + } + $width = $profile->maximumSourceVariables(); + $limit = intdiv(65535, $width + 1); + yield "$name parameter limit includes ordinal" => [$profile, 2 * $limit + 1, $width, 0, 1073741824, [$limit, $limit, 1]]; + yield "$name byte target" => [$profile, 41, 2, 50000, 1073741824, [20, 20, 1]]; + // Real packet 231072 gives a 50000-byte payload, not another halving/reserve. + yield "$name packet full and tail" => [$profile, 5, 2, 20000, 231072, [2, 2, 1]]; + yield "$name preflight boundary singleton" => [$profile, 1, 2, 49992, 231072, [1]]; + yield "$name preflight boundary with following tail" => [$profile, 3, 2, 49992, 231072, [1, 2]]; + } + } + + /** @param positive-int $width + * @param list $batchSizes + */ + #[DataProvider('caseBatches')] + public function testImportsCatalogueAndOrderedRowsAfterMysqlDdl(MySqlProfile $profile, int $count, int $width, int $textBytes, int $packet, array $batchSizes): void { $pdo = $this->createMock(PDO::class); $pdo->method('getAttribute')->with(PDO::ATTR_ERRMODE)->willReturn(PDO::ERRMODE_EXCEPTION); $pdo->method('setAttribute')->with(PDO::ATTR_ERRMODE, PDO::ERRMODE_EXCEPTION)->willReturn(true); $dataset = $this->createMock(PDOStatement::class); $variables = $this->createMock(PDOStatement::class); - $cases = $this->createMock(PDOStatement::class); $variableRows = []; $caseRows = []; + $caseSql = []; + $probes = []; + $ddlStarted = false; + $packetStatement = $this->createMock(PDOStatement::class); + $packetStatement->method('fetchColumn')->willReturn($packet); + $pdo->method('query')->willReturnCallback(static function (string $sql) use ($packetStatement, &$probes, &$ddlStarted): PDOStatement|false { + if ($sql !== 'SELECT @@max_allowed_packet') { + return false; + } + $probes[] = $ddlStarted; + return $packetStatement; + }); + $sourceVariables = $width === 2 ? [ + ['name' => 'Score', 'type' => 'numeric', 'width' => 0], + ['name' => 'Comment', 'type' => 'string', 'width' => 12], + ] : array_map(static fn(int $i): array => ['name' => 'v' . $i, 'type' => 'numeric', 'width' => 0], range(1, $width)); + $rows = []; + $expected = []; + for ($i = 0; $i < $count; ++$i) { + $row = $width === 2 ? [[0.1, null, 42, 1.0000000000000002, PHP_FLOAT_MAX, 5.0e-324][$i % 6], $i % 2 === 0 ? "õ'\\\\" : ''] : array_fill(0, $width, null); + if ($textBytes > 0 && ($textBytes !== 49992 || $i === 0)) { + $row[1] = str_repeat("'\\\\", intdiv($textBytes, 3)); + $row[1] .= str_repeat('x', $textBytes - strlen($row[1])); + } + $rows[] = $row; + $expected[] = array_merge([$i + 1], array_map(static fn($value) => is_float($value) ? sprintf('%.17g', $value) : $value, $row)); + } - $pdo->expects(self::atLeast(18))->method('exec')->willReturn(0); + $pdo->expects(self::atLeast(18))->method('exec')->willReturnCallback(static function () use (&$ddlStarted): int { + $ddlStarted = true; + return 0; + }); $pdo->expects(self::once())->method('beginTransaction')->willReturn(true); $pdo->expects(self::once())->method('commit')->willReturn(true); $pdo->expects(self::never())->method('rollBack'); - $pdo->expects(self::exactly(3))->method('prepare')->willReturnOnConsecutiveCalls($dataset, $variables, $cases); + $pdo->method('prepare')->willReturnCallback(function (string $sql) use ($dataset, $variables, &$caseSql, &$caseRows): PDOStatement { + if (str_starts_with($sql, 'INSERT INTO datasets ')) { + return $dataset; + } + if (str_starts_with($sql, 'INSERT INTO variables ')) { + return $variables; + } + self::assertStringStartsWith('INSERT INTO `dataset_customer_survey` ', $sql); + $caseSql[] = $sql; + $statement = $this->createMock(PDOStatement::class); + $statement->method('execute')->willReturnCallback(static function (array $params) use ($sql, &$caseRows): bool { + $caseRows[] = [$sql, array_values($params)]; + return true; + }); + return $statement; + }); $dataset->expects(self::once())->method('execute')->with(['customer survey', 'dataset_customer_survey'])->willReturn(true); - $variables->expects(self::exactly(2))->method('execute')->willReturnCallback(function (array $row) use (&$variableRows): bool { + $variables->expects(self::exactly($width))->method('execute')->willReturnCallback(function (array $row) use (&$variableRows): bool { $variableRows[] = $row; return true; }); - $cases->expects(self::exactly(2))->method('execute')->willReturnCallback(function (array $row) use (&$caseRows): bool { - $caseRows[] = $row; - - return true; - }); - - $definition = (new MySqlWideTableImporter($pdo))->import([ - 'variables' => [ - ['name' => 'Score', 'type' => 'numeric', 'width' => 0], - ['name' => 'Comment', 'type' => 'string', 'width' => 12], - ], - 'data' => [[0.1, 'blue'], [null, 'green']], + $definition = (new MySqlWideTableImporter($pdo, $profile))->import([ + 'variables' => $sourceVariables, + 'data' => $rows, ], 'customer survey'); self::assertSame('dataset_customer_survey', $definition->tableName); - self::assertSame([ - ['customer survey', 1, 'Score', 'score', 'numeric', 0, 5, 8, 0, 5, 8, 0, null], - ['customer survey', 2, 'Comment', 'comment', 'string', 12, 5, 8, 0, 5, 8, 0, null], - ], $variableRows); - self::assertSame([ - ['value_0' => 1, 'value_1' => '0.10000000000000001', 'value_2' => 'blue'], - ['value_0' => 2, 'value_1' => null, 'value_2' => 'green'], - ], $caseRows); + foreach ($sourceVariables as $i => $variable) { + $column = strtolower($variable['name']); + self::assertSame($column, $definition->columns[$i]['columnName']); + self::assertSame(['customer survey', $i + 1, $variable['name'], $column, $variable['type'], $variable['width'], 5, 8, 0, 5, 8, 0, null], $variableRows[$i]); + } + self::assertSame($batchSizes, array_map(static fn(array $batch): int => intdiv(count($batch[1]), $width + 1), $caseRows), 'Case execution sizes, including the final tail.'); + self::assertSame($count === 0, $caseSql === [], 'Empty imports must not prepare case SQL.'); + self::assertSame(array_fill(0, $count < 2 ? 2 : 3, false), $probes, 'Two preflight probes plus one batch-budget snapshot for multiple rows; all before DDL.'); + $actual = []; + foreach ($caseRows as [$sql, $params]) { + self::assertNotEmpty($params); + self::assertLessThanOrEqual(65535, count($params)); + self::assertSame(count($params), preg_match_all('/\\?|:value_\\d+/', $sql)); + self::assertSame(intdiv(count($params), $width + 1), preg_match_all('/\\([^()]*\\)/', explode(' VALUES ', $sql, 2)[1])); + array_push($actual, ...array_chunk($params, $width + 1)); + } + self::assertSame($expected, $actual); } public function testImportsCoreSpssMetadataThroughMysqlCatalogue(): void @@ -320,45 +387,83 @@ public function testInjectedDoltProfileRejects306VariablesBeforeDdl(): void } } - public function testRolledBackCaseInsertFailureDropsOnlyAttemptPhysicalTable(): void + /** @return iterable */ + public static function lateBatchFailures(): iterable + { + yield 'prepare tail' => ['prepare']; + yield 'execute tail' => ['execute']; + } + + #[DataProvider('lateBatchFailures')] + public function testRolledBackCaseInsertFailureDropsOnlyAttemptPhysicalTable(string $stage): void { $pdo = $this->createMock(PDO::class); - $pdo->method('getAttribute')->with(PDO::ATTR_ERRMODE)->willReturn(PDO::ERRMODE_EXCEPTION); - $pdo->method('setAttribute')->with(PDO::ATTR_ERRMODE, PDO::ERRMODE_EXCEPTION)->willReturn(true); - $dataset = $this->createMock(PDOStatement::class); - $variables = $this->createMock(PDOStatement::class); - $cases = $this->createMock(PDOStatement::class); + $mode = PDO::ERRMODE_SILENT; + $pdo->method('getAttribute')->with(PDO::ATTR_ERRMODE)->willReturn($mode); + $pdo->method('setAttribute')->willReturnCallback(static function (int $attribute, mixed $value) use (&$mode): bool { + self::assertSame(PDO::ATTR_ERRMODE, $attribute); + $mode = $value; + return true; + }); + $pdo->method('inTransaction')->willReturnOnConsecutiveCalls(false, true); + $pdo->expects(self::once())->method('beginTransaction')->willReturn(true); + $commits = $rollbacks = $prepares = $executes = $written = 0; + $pdo->method('commit')->willReturnCallback(static function () use (&$commits): bool { + ++$commits; + return true; + }); + $pdo->method('rollBack')->willReturnCallback(static function () use (&$rollbacks): bool { + ++$rollbacks; + return true; + }); $executedSql = []; - - $pdo->expects(self::exactly(22))->method('exec')->willReturnCallback(function (string $sql) use (&$executedSql): int { + $pdo->method('exec')->willReturnCallback(static function (string $sql) use (&$executedSql): int { $executedSql[] = $sql; - return 0; }); - $pdo->expects(self::once())->method('beginTransaction')->willReturn(true); - $pdo->expects(self::exactly(2))->method('inTransaction')->willReturnOnConsecutiveCalls(false, true); - $pdo->expects(self::once())->method('rollBack')->willReturn(true); - $pdo->expects(self::never())->method('commit'); - $pdo->expects(self::exactly(3))->method('prepare')->willReturnOnConsecutiveCalls( - $dataset, - $variables, - $cases, - ); - - $dataset->expects(self::once())->method('execute')->willReturn(true); + $dataset = $this->createMock(PDOStatement::class); + $variables = $this->createMock(PDOStatement::class); + $dataset->expects(self::once())->method('execute')->with(['customer survey', 'dataset_customer_survey'])->willReturn(true); $variables->expects(self::once())->method('execute')->willReturn(true); - $cases->expects(self::once())->method('execute')->willThrowException(new RuntimeException('insert failed')); - - $this->expectExceptionMessage('insert failed'); + $injected = new RuntimeException('second batch ' . $stage); + $pdo->method('prepare')->willReturnCallback(function (string $sql) use ($dataset, $variables, $stage, $injected, &$prepares, &$executes, &$written): PDOStatement { + if (str_starts_with($sql, 'INSERT INTO datasets ')) { + return $dataset; + } + if (str_starts_with($sql, 'INSERT INTO variables ')) { + return $variables; + } + self::assertStringStartsWith('INSERT INTO `dataset_customer_survey` ', $sql); + if (++$prepares === 2 && $stage === 'prepare') { + throw $injected; + } + $statement = $this->createMock(PDOStatement::class); + $statement->method('execute')->willReturnCallback(static function (array $params) use ($stage, $injected, &$executes, &$written): bool { + if (++$executes === 2 && $stage === 'execute') { + throw $injected; + } + $written += intdiv(count($params), 2); + return true; + }); + return $statement; + }); + $caught = null; try { (new MySqlWideTableImporter($pdo))->import([ 'variables' => [['name' => 'Score', 'type' => 'numeric']], - 'data' => [[1.0]], + 'data' => array_fill(0, 257, [0.1]), ], 'customer survey'); - } finally { - self::assertStringStartsWith('DROP TABLE IF EXISTS ', $executedSql[21] ?? ''); - self::assertStringContainsString('dataset_customer_survey', $executedSql[21] ?? ''); + } catch (RuntimeException $exception) { + $caught = $exception; } + self::assertSame($injected, $caught, 'The second case batch must reach the injected failure.'); + self::assertSame(256, $written, 'One full batch must succeed before the tail fails.'); + self::assertSame(0, $commits); + self::assertSame(1, $rollbacks); + self::assertSame(PDO::ERRMODE_SILENT, $mode); + self::assertCount(22, $executedSql); + self::assertSame(['DROP TABLE IF EXISTS `dataset_customer_survey`'], array_values(array_filter($executedSql, static fn(string $sql): bool => str_starts_with($sql, 'DROP ')))); + self::assertSame('DROP TABLE IF EXISTS `dataset_customer_survey`', $executedSql[21]); } #[DataProvider('nonFiniteValues')] diff --git a/tests/Sql/PostgreSqlWideTableImporterTest.php b/tests/Sql/PostgreSqlWideTableImporterTest.php index 67a0bc8..17f9250 100644 --- a/tests/Sql/PostgreSqlWideTableImporterTest.php +++ b/tests/Sql/PostgreSqlWideTableImporterTest.php @@ -70,50 +70,163 @@ public function testCreatesCatalogAndStrictWideTableInOneTransaction(): void self::assertStringContainsString('DOUBLE PRECISION NULL', $definition->createSql); self::assertStringContainsString('TEXT NOT NULL', $definition->createSql); } - public function testImportsCatalogueAndOrderedRowsThroughPdoTransaction(): void + /** @return iterable}> */ + public static function caseBatches(): iterable + { + yield 'empty' => [0, 2, 0, []]; + yield 'singleton' => [1, 2, 0, [1]]; + yield 'two fitting rows' => [2, 2, 0, [2]]; + yield 'row limit and tail' => [513, 2, 0, [256, 256, 1]]; + yield 'parameter limit includes ordinal' => [81, 1599, 0, [40, 40, 1]]; + yield 'byte target and tail' => [5, 2, 400000, [2, 2, 1]]; + yield 'oversized valid row and following tail' => [3, 2, 1048577, [1, 2]]; + } + + /** @param positive-int $width + * @param list $batchSizes + */ + #[DataProvider('caseBatches')] + public function testImportsCatalogueAndOrderedRowsThroughPdoTransaction(int $count, int $width, int $textBytes, array $batchSizes): void { $pdo = $this->createMock(PDO::class); $pdo->method('getAttribute')->with(PDO::ATTR_ERRMODE)->willReturn(PDO::ERRMODE_EXCEPTION); $pdo->method('setAttribute')->with(PDO::ATTR_ERRMODE, PDO::ERRMODE_EXCEPTION)->willReturn(true); $dataset = $this->createMock(PDOStatement::class); $variables = $this->createMock(PDOStatement::class); - $cases = $this->createMock(PDOStatement::class); $variableRows = []; $caseRows = []; + $caseSql = []; + $sourceVariables = $width === 2 ? [ + ['name' => 'Score', 'type' => 'numeric', 'width' => 0], + ['name' => 'Comment', 'type' => 'string', 'width' => 12], + ] : array_map(static fn(int $i): array => ['name' => 'v' . $i, 'type' => 'numeric', 'width' => 0], range(1, $width)); + $rows = []; + $expected = []; + for ($i = 0; $i < $count; ++$i) { + $row = $width === 2 ? [[0.1, null, 42, 1.0000000000000002, PHP_FLOAT_MAX, 5.0e-324][$i % 6], $i % 2 === 0 ? "õ'\\\\" : ''] : array_fill(0, $width, null); + if ($textBytes > 0 && ($textBytes <= 1048576 || $i === 0)) { + $row[1] = str_repeat('x', $textBytes); + } + $rows[] = $row; + $expected[] = array_merge([$i + 1], array_map(static fn($value) => is_float($value) ? sprintf('%.17g', $value) : $value, $row)); + } $pdo->expects(self::once())->method('beginTransaction')->willReturn(true); $pdo->expects(self::once())->method('commit')->willReturn(true); $pdo->expects(self::never())->method('rollBack'); $pdo->expects(self::atLeast(18))->method('exec')->willReturn(0); - $pdo->expects(self::exactly(3))->method('prepare')->willReturnOnConsecutiveCalls($dataset, $variables, $cases); + $pdo->method('prepare')->willReturnCallback(function (string $sql) use ($dataset, $variables, &$caseSql, &$caseRows): PDOStatement { + if (str_starts_with($sql, 'INSERT INTO datasets ')) { + return $dataset; + } + if (str_starts_with($sql, 'INSERT INTO variables ')) { + return $variables; + } + self::assertStringStartsWith('INSERT INTO "dataset_customer_survey" ', $sql); + $caseSql[] = $sql; + $statement = $this->createMock(PDOStatement::class); + $statement->method('execute')->willReturnCallback(static function (array $params) use ($sql, &$caseRows): bool { + $caseRows[] = [$sql, array_values($params)]; + return true; + }); + return $statement; + }); $dataset->expects(self::once())->method('execute')->with(['customer survey', 'dataset_customer_survey'])->willReturn(true); - $variables->expects(self::exactly(2))->method('execute')->willReturnCallback(function ($row) use (&$variableRows): bool { + $variables->expects(self::exactly($width))->method('execute')->willReturnCallback(function ($row) use (&$variableRows): bool { $variableRows[] = $row; return true; }); - $cases->expects(self::exactly(2))->method('execute')->willReturnCallback(function ($row) use (&$caseRows): bool { - $caseRows[] = $row; - return true; - }); - $definition = (new PostgreSqlWideTableImporter($pdo))->import([ - 'variables' => [ - ['name' => 'Score', 'type' => 'numeric', 'width' => 0], - ['name' => 'Comment', 'type' => 'string', 'width' => 12], - ], - 'data' => [[0.1, 'blue'], [null, 'green']], + 'variables' => $sourceVariables, + 'data' => $rows, ], 'customer survey'); - self::assertSame('score', $definition->columns[0]['columnName']); - self::assertSame('comment', $definition->columns[1]['columnName']); - self::assertSame([ - ['customer survey', 1, 'Score', 'score', 'numeric', 0, 5, 8, 0, 5, 8, 0, null], - ['customer survey', 2, 'Comment', 'comment', 'string', 12, 5, 8, 0, 5, 8, 0, null], - ], $variableRows); - self::assertSame([ - ['value_0' => 1, 'value_1' => '0.10000000000000001', 'value_2' => 'blue'], - ['value_0' => 2, 'value_1' => null, 'value_2' => 'green'], - ], $caseRows); + foreach ($sourceVariables as $i => $variable) { + $column = strtolower($variable['name']); + self::assertSame($column, $definition->columns[$i]['columnName']); + self::assertSame(['customer survey', $i + 1, $variable['name'], $column, $variable['type'], $variable['width'], 5, 8, 0, 5, 8, 0, null], $variableRows[$i]); + } + self::assertSame($batchSizes, array_map(static fn(array $batch): int => intdiv(count($batch[1]), $width + 1), $caseRows), 'Case execution sizes, including the final tail.'); + self::assertSame($count === 0, $caseSql === [], 'Empty imports must not prepare case SQL.'); + $actual = []; + foreach ($caseRows as [$sql, $params]) { + self::assertNotEmpty($params); + self::assertLessThanOrEqual(65535, count($params)); + self::assertSame(count($params), preg_match_all('/\?|:value_\d+/', $sql)); + self::assertSame(intdiv(count($params), $width + 1), preg_match_all('/\\([^()]*\\)/', explode(' VALUES ', $sql, 2)[1])); + array_push($actual, ...array_chunk($params, $width + 1)); + } + self::assertSame($expected, $actual); + } + + /** @return iterable */ + public static function lateBatchFailures(): iterable + { + yield 'prepare tail' => ['prepare']; + yield 'execute tail' => ['execute']; + } + + #[DataProvider('lateBatchFailures')] + public function testSecondBatchFailureRollsBackWithoutFinalization(string $stage): void + { + $pdo = $this->createMock(PDO::class); + $mode = PDO::ERRMODE_SILENT; + $pdo->method('getAttribute')->with(PDO::ATTR_ERRMODE)->willReturn($mode); + $pdo->method('setAttribute')->willReturnCallback(static function (int $attribute, mixed $value) use (&$mode): bool { + self::assertSame(PDO::ATTR_ERRMODE, $attribute); + $mode = $value; + return true; + }); + $pdo->method('inTransaction')->willReturnOnConsecutiveCalls(false, true); + $pdo->expects(self::once())->method('beginTransaction')->willReturn(true); + $commits = $rollbacks = $prepares = $executes = $written = 0; + $pdo->method('commit')->willReturnCallback(static function () use (&$commits): bool { + ++$commits; + return true; + }); + $pdo->method('rollBack')->willReturnCallback(static function () use (&$rollbacks): bool { + ++$rollbacks; + return true; + }); + $pdo->method('exec')->willReturn(0); + $metadata = $this->createMock(PDOStatement::class); + $metadata->method('execute')->willReturn(true); + $injected = new \RuntimeException('second batch ' . $stage); + $pdo->method('prepare')->willReturnCallback(function (string $sql) use ($metadata, $stage, $injected, &$prepares, &$executes, &$written): PDOStatement { + if (!str_starts_with($sql, 'INSERT INTO "dataset_attempt" ')) { + return $metadata; + } + if (++$prepares === 2 && $stage === 'prepare') { + throw $injected; + } + $statement = $this->createMock(PDOStatement::class); + $statement->method('execute')->willReturnCallback(static function (array $params) use ($stage, $injected, &$executes, &$written): bool { + if (++$executes === 2 && $stage === 'execute') { + throw $injected; + } + $written += intdiv(count($params), 2); + return true; + }); + return $statement; + }); + $finalized = false; + $caught = null; + try { + (new PostgreSqlWideTableImporter($pdo))->import([ + 'variables' => [['name' => 'Score', 'type' => 'numeric']], + 'data' => array_fill(0, 257, [0.1]), + ], 'attempt', beforeCommit: static function () use (&$finalized): void { + $finalized = true; + }); + } catch (\RuntimeException $exception) { + $caught = $exception; + } + self::assertSame($injected, $caught, 'The second case batch must reach the injected failure.'); + self::assertSame(256, $written, 'One full batch must succeed before the tail fails.'); + self::assertSame(0, $commits); + self::assertSame(1, $rollbacks); + self::assertFalse($finalized); + self::assertSame(PDO::ERRMODE_SILENT, $mode); } public function testImportsFileLabelDocumentsAndTechnicalMetadata(): void From 6a753ce9953e0df6812e67d47bc837a18a9ea99e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?T=C3=B5nis=20Ormisson?= Date: Tue, 8 Sep 2026 14:21:38 +0300 Subject: [PATCH 2/2] perf: batch server case inserts within packet and parameter budgets --- src/Sql/MySqlWideTableImporter.php | 28 +++++++---- src/Sql/PostgreSqlWideTableImporter.php | 23 +++++---- src/Sql/PreparedCaseBatch.php | 67 +++++++++++++++++++++++++ 3 files changed, 100 insertions(+), 18 deletions(-) create mode 100644 src/Sql/PreparedCaseBatch.php diff --git a/src/Sql/MySqlWideTableImporter.php b/src/Sql/MySqlWideTableImporter.php index 6c100c6..93784f9 100644 --- a/src/Sql/MySqlWideTableImporter.php +++ b/src/Sql/MySqlWideTableImporter.php @@ -65,6 +65,8 @@ public function import( $definition = $schema->wideTableDefinition($datasetName, $variables); $v3Metadata = $this->assertSourceMetadata($source, $variables, $definition); + $payloadBudget = count($rows) > 1 ? $this->profile->effectiveMaximumStatementBytes($this->pdo) : null; + $schema->createCatalog(); $this->pdo->exec($definition->createSql); @@ -82,7 +84,7 @@ public function import( if ($v3Metadata !== null) { (new SqliteV3MetadataImporter($this->pdo))->storeValidated($datasetName, $v3Metadata); } - $this->insertCases($definition, $rows); + $this->insertCases($definition, $rows, $payloadBudget); if ($sourcePath !== "" || $verifiedSourceSha256 !== null) { $datasetId = (new NormativeCatalog($this->pdo))->storeImportedDataset( $datasetName, @@ -264,7 +266,7 @@ private function storeWeightVariable( /** * @param list $rows */ - private function insertCases(MySqlWideTableDefinition $definition, array $rows): void + private function insertCases(MySqlWideTableDefinition $definition, array $rows, ?int $payloadBudget): void { $quote = chr(96); $columns = array_merge(['__case_ordinal'], array_column($definition->columns, 'columnName')); @@ -272,13 +274,21 @@ private function insertCases(MySqlWideTableDefinition $definition, array $rows): static fn(string $column): string => chr(96) . str_replace(chr(96), chr(96) . chr(96), $column) . chr(96), $columns, ); - $parameters = array_map(static fn(int $index): string => ':value_' . $index, array_keys($columns)); - $statement = $this->requiredStatement( + PreparedCaseBatch::send( + $this->pdo, 'INSERT INTO ' . $quote . str_replace($quote, $quote . $quote, $definition->tableName) . $quote . ' (' - . implode(', ', $quotedColumns) . ') VALUES (' . implode(', ', $parameters) . ')', - 'wide-table data', + . implode(', ', $quotedColumns) . ') VALUES ', + $this->caseRows($definition, $rows), + $payloadBudget, ); + } + /** + * @param list $rows + * @return \Generator> + */ + private function caseRows(MySqlWideTableDefinition $definition, array $rows): \Generator + { foreach ($rows as $caseOrdinal => $row) { if (!is_array($row)) { throw new UnsupportedOperation( @@ -286,14 +296,14 @@ private function insertCases(MySqlWideTableDefinition $definition, array $rows): 'Every SPSS case must be an ordered value list or source-name map.', ); } - $values = ['value_0' => $caseOrdinal + 1]; + $values = [$caseOrdinal + 1]; foreach ($definition->columns as $index => $column) { $value = array_key_exists($column['sourceName'], $row) ? $row[$column['sourceName']] : ($row[$index] ?? null); - $values['value_' . ($index + 1)] = $this->caseValue($value, $column['storageKind']); + $values[] = $this->caseValue($value, $column['storageKind']); } - $statement->execute($values); + yield $values; } } diff --git a/src/Sql/PostgreSqlWideTableImporter.php b/src/Sql/PostgreSqlWideTableImporter.php index 0a78dbb..47defcc 100644 --- a/src/Sql/PostgreSqlWideTableImporter.php +++ b/src/Sql/PostgreSqlWideTableImporter.php @@ -408,24 +408,29 @@ private function insertCases(PostgreSqlWideTableDefinition $definition, array $r { $columns = array_merge(['__case_ordinal'], array_column($definition->columns, 'columnName')); $quoted = array_map(fn(string $column): string => '"' . str_replace('"', '""', $column) . '"', $columns); - $parameters = array_map(static fn(int $index): string => ':value_' . $index, array_keys($columns)); - $statement = $this->pdo->prepare( - 'INSERT INTO "' . str_replace('"', '""', $definition->tableName) . '" (' . implode(', ', $quoted) . ') VALUES (' . implode(', ', $parameters) . ')', + PreparedCaseBatch::send( + $this->pdo, + 'INSERT INTO "' . str_replace('"', '""', $definition->tableName) . '" (' . implode(', ', $quoted) . ') VALUES ', + $this->caseRows($definition, $rows), ); - if ($statement === false) { - throw new UnsupportedOperation(DiagnosticCode::InvalidSourceDataset, 'The PostgreSQL profile could not prepare a required data statement.'); - } + } + /** + * @param list $rows + * @return \Generator> + */ + private function caseRows(PostgreSqlWideTableDefinition $definition, array $rows): \Generator + { foreach ($rows as $caseOrdinal => $row) { if (!is_array($row)) { throw new UnsupportedOperation(DiagnosticCode::InvalidSourceDataset, 'Every SPSS case must be an ordered value list or source-name map.'); } - $values = ['value_0' => $caseOrdinal + 1]; + $values = [$caseOrdinal + 1]; foreach ($definition->columns as $index => $column) { $value = array_key_exists($column['sourceName'], $row) ? $row[$column['sourceName']] : ($row[$index] ?? null); - $values['value_' . ($index + 1)] = $this->caseValue($value, $column['storageKind']); + $values[] = $this->caseValue($value, $column['storageKind']); } - $statement->execute($values); + yield $values; } } diff --git a/src/Sql/PreparedCaseBatch.php b/src/Sql/PreparedCaseBatch.php new file mode 100644 index 0000000..2c7b8dd --- /dev/null +++ b/src/Sql/PreparedCaseBatch.php @@ -0,0 +1,67 @@ +> $rows */ + public static function send(PDO $pdo, string $prefix, iterable $rows, ?int $payloadBudget = null): void + { + $byteLimit = min(1_048_576, $payloadBudget ?? 1_048_576); + $values = []; + $count = 0; + $bytes = 0; + $statement = null; + $lastCount = 0; + $tuple = ''; + + foreach ($rows as $row) { + $columns = count($row); + if ($tuple === '') { + $tuple = '(' . implode(', ', array_fill(0, $columns, '?')) . ')'; + } + // 16 bytes/value cover length/type framing and placeholder/quote + // syntax; NULL needs only eight, including its marker and separator. + // MySQL's budget already reserves the prefix and worst-case escaping. + $rowBytes = 4; + foreach ($row as $value) { + $rowBytes += $value === null ? 8 : 16 + strlen((string) $value); + } + if ($count > 0 && ($count >= min(256, intdiv(65_535, $columns)) || $bytes + $rowBytes > $byteLimit)) { + self::flush($pdo, $prefix, $tuple, $values, $count, $statement, $lastCount); + $values = []; + $count = $bytes = 0; + } + // Grouping is not acceptance: preflight-valid oversized rows go alone. + array_push($values, ...$row); + ++$count; + $bytes += $rowBytes; + } + self::flush($pdo, $prefix, $tuple, $values, $count, $statement, $lastCount); + } + + /** @param list $values */ + private static function flush(PDO $pdo, string $prefix, string $tuple, array $values, int $count, ?PDOStatement &$statement, int &$lastCount): void + { + if ($count === 0) { + return; + } + if ($statement === null || $count !== $lastCount) { + $prepared = $pdo->prepare($prefix . implode(', ', array_fill(0, $count, $tuple))); + if ($prepared === false) { + throw new UnsupportedOperation(DiagnosticCode::InvalidSourceDataset, 'Could not prepare a required data statement.'); + } + $statement = $prepared; + $lastCount = $count; + } + $statement->execute($values); + } +}