Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 19 additions & 9 deletions src/Sql/MySqlWideTableImporter.php
Original file line number Diff line number Diff line change
Expand Up @@ -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);

Expand All @@ -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,
Expand Down Expand Up @@ -264,36 +266,44 @@ private function storeWeightVariable(
/**
* @param list<mixed> $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'));
$quotedColumns = array_map(
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<mixed> $rows
* @return \Generator<int, list<int|string|null>>
*/
private function caseRows(MySqlWideTableDefinition $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;
}
}

Expand Down
23 changes: 14 additions & 9 deletions src/Sql/PostgreSqlWideTableImporter.php
Original file line number Diff line number Diff line change
Expand Up @@ -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<mixed> $rows
* @return \Generator<int, list<int|string|null>>
*/
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;
}
}

Expand Down
67 changes: 67 additions & 0 deletions src/Sql/PreparedCaseBatch.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,67 @@
<?php

declare(strict_types=1);

namespace OpenStatSpec\Sql;

use OpenStatSpec\Core\DiagnosticCode;
use OpenStatSpec\Core\UnsupportedOperation;
use PDO;
use PDOStatement;

/** @internal Shared bounded sender for already validated, encoded case rows. */
final class PreparedCaseBatch
{
/** @param iterable<list<int|string|null>> $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<int|string|null> $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);
}
}
98 changes: 98 additions & 0 deletions tests/Integration/MySqlFamilySpssRoundTripTestCase.php
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -146,6 +150,100 @@ public function testRealEngineRoundTripsSavAndZsavThroughMySqlFamily(): void
}
}

/** @return iterable<string, array{bool, string}> */
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;
Expand Down
Loading