diff --git a/CHANGELOG.md b/CHANGELOG.md index 744e48bc..3acc98cf 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,16 @@ # Changelog +## 3.5.13-wip + +- Fix queries being split across transactions, which made the driver unusable + through a transaction pooler. Parsing used to be its own Sync-terminated + exchange, so a pooler could hand the server connection to another client + between the parse and the bind: the statement was either gone by the time it + was bound (`26000`), or still there when another client parsed the same + generated name (`42P05`). A one-shot query now parses, binds and executes in + one exchange, under the unnamed statement. This also removes a round trip. +- `ParseMessage.statementName` is now readable. + ## 3.5.12 - Fix `runTx` silently rolling back after `ROLLBACK TO SAVEPOINT` recovery: clear stale `_transactionException` when PostgreSQL confirms a healthy transaction state. diff --git a/lib/src/messages/client_messages.dart b/lib/src/messages/client_messages.dart index e2672b88..3e1c8b97 100644 --- a/lib/src/messages/client_messages.dart +++ b/lib/src/messages/client_messages.dart @@ -124,6 +124,11 @@ class ParseMessage extends ClientMessage { _statementName = statementName, _typeOids = typeOids ?? types?.map((e) => e?.oid).toList() ?? const []; + /// The name the statement is prepared under. Empty for the unnamed + /// statement, which belongs to whichever transaction is using the connection + /// rather than to the connection itself. + String get statementName => _statementName; + @override void applyToBuffer(PgByteDataWriter buffer) { buffer.writeUint8(ClientMessageId.parse); diff --git a/lib/src/v3/connection.dart b/lib/src/v3/connection.dart index 124222eb..c268cf4e 100644 --- a/lib/src/v3/connection.dart +++ b/lib/src/v3/connection.dart @@ -178,13 +178,37 @@ abstract class _PgSessionBase implements Session { await querySubscription.cancel(); } } else { - // The simple query protocol does not support variables. So when we have - // parameters, we need an explicit prepare. - final prepared = await _prepare(description, variables); + // The simple query protocol does not support variables, so this needs the + // extended one. Parsing in its own exchange would end the implicit + // transaction the parse opened, and a pooler in transaction mode hands + // the server connection to another client at that point -- which either + // loses the statement or leaves it behind for someone else to collide + // with. So the parse travels with the bind and the execute, as one + // exchange and one transaction, under the unnamed statement that belongs + // to it. + final stackTrace = StackTrace.current; + final prepared = _PreparedStatement( + description, + '', + this, + Trace.from(stackTrace), + parse: ParseMessage( + description.transformedSql, + statementName: '', + typeOids: _mergeTypeOids(description.parameterTypes, variables), + ), + ); try { + // Nothing to close on the way out: the unnamed statement is replaced + // by the next parse that leaves the name empty, so closing it is an + // exchange that changes nothing -- and every `execute` would make one. return await prepared.run(variables, timeout: timeout); - } finally { + } catch (_) { + // A statement that did not finish can leave the connection with + // messages still to deliver, and this exchange is where they land. + // Closing is beside the point here; having somewhere to arrive is not. await prepared.dispose(); + rethrow; } } } @@ -721,7 +745,18 @@ class _PreparedStatement extends Statement { final Trace _trace; - _PreparedStatement(this._description, this._name, this._session, this._trace); + /// The parse this statement still owes the server, when it was made without + /// sending one. It travels with the first bind so that parsing and binding + /// are one exchange. + final ParseMessage? _parse; + + _PreparedStatement( + this._description, + this._name, + this._session, + this._trace, { + ParseMessage? parse, + }) : _parse = parse; _PgSessionBase get _effectiveSession => _session._connection._activeTransaction ?? _session; @@ -869,6 +904,9 @@ class _PgResultStreamSubscription connection._channel.sink.add( AggregatedClientMessage([ + // A statement that has not been parsed yet parses here, so that the + // parse and the bind cannot be split across two transactions. + ?statement.statement._parse, BindMessage( encodedValues, portalName: _portalName, @@ -958,6 +996,7 @@ class _PgResultStreamSubscription @override Future handleMessage(ServerMessage message) async { switch (message) { + case ParseCompleteMessage(): case BindCompleteMessage(): case NoDataMessage(): // Nothing to do! diff --git a/test/connection_test.dart b/test/connection_test.dart index f29071c6..20c88e01 100644 --- a/test/connection_test.dart +++ b/test/connection_test.dart @@ -314,11 +314,13 @@ void main() { final orderEnsurer = []; // this will emit a query error - conn!.execute('INSERT INTO t (i) VALUES ()').catchError((err) { - orderEnsurer.add(1); - // ignore - return Result(rows: [], affectedRows: 0, schema: ResultSchema([])); - }); + final failing = conn! + .execute('INSERT INTO t (i) VALUES ()') + .catchError((err) { + orderEnsurer.add(1); + // ignore + return Result(rows: [], affectedRows: 0, schema: ResultSchema([])); + }); orderEnsurer.add(2); final res = await conn!.runTx((ctx) async { @@ -326,11 +328,18 @@ void main() { return await ctx.execute('SELECT i FROM t'); }); orderEnsurer.add(4); + // The error is reported when the query's one exchange comes back, + // rather than by a parse of its own that failed before anything else + // was scheduled. Awaiting it says the failure was delivered without + // depending on how many round trips it took to arrive. + await failing; expect(res, [ [1], ]); - expect(orderEnsurer, [2, 1, 3, 4]); + expect(orderEnsurer, containsAll([1, 2, 3, 4])); + expect(orderEnsurer.indexOf(2), lessThan(orderEnsurer.indexOf(3))); + expect(orderEnsurer.indexOf(3), lessThan(orderEnsurer.indexOf(4))); }, ); }, diff --git a/test/single_exchange_test.dart b/test/single_exchange_test.dart new file mode 100644 index 00000000..69393bdc --- /dev/null +++ b/test/single_exchange_test.dart @@ -0,0 +1,120 @@ +import 'dart:async'; + +import 'package:async/async.dart'; +import 'package:postgres/messages.dart'; +import 'package:postgres/postgres.dart'; +import 'package:postgres/src/v3/protocol.dart'; +import 'package:stream_channel/stream_channel.dart'; +import 'package:test/test.dart'; + +import 'docker.dart'; + +/// Everything between two Sync messages is one transaction, so a query that +/// sends a Sync between parsing and binding is two of them. A pooler in +/// transaction mode hands the server connection to another client at the +/// ReadyForQuery a Sync produces, and the bind then arrives somewhere that +/// never saw the parse -- or worse, somewhere still holding another client's +/// statement of the same generated name. +/// +/// So parse, bind and execute travel together. +void main() { + withPostgresServer('single exchange', (server) { + late List sent; + late Connection connection; + + setUp(() async { + sent = []; + connection = await Connection.open( + await server.endpoint(), + settings: ConnectionSettings( + transformer: StreamChannelTransformer( + StreamTransformer.fromHandlers(), + StreamSinkTransformer.fromHandlers( + handleData: (message, sink) { + // Messages travel batched, so what was sent is what the + // batches hold. + if (message is AggregatedClientMessage) { + sent.addAll(message.messages); + } else if (message is ClientMessage) { + sent.add(message); + } + sink.add(message); + }, + ), + ), + ), + ); + // Connecting is not what is under test. + sent.clear(); + }); + + tearDown(() => connection.close()); + + /// What was sent, up to and including the execute. + Iterable throughExecute() => + sent.takeWhile((message) => message is! ExecuteMessage).followedBy([ + sent.whereType().first, + ]); + + test('a parameterised query parses and binds without a sync between', () async { + await connection.execute(r'SELECT $1::int AS value', parameters: [1]); + + expect( + throughExecute().map((message) => message.runtimeType).toList(), + [ParseMessage, BindMessage, DescribeMessage, ExecuteMessage], + ); + }); + + test('a query without parameters is one exchange too', () async { + await connection.execute('SELECT 1 AS value'); + + expect( + throughExecute().whereType(), + isEmpty, + reason: 'a sync before the execute would end the transaction that parsed', + ); + }); + + test('the statement it parses is the unnamed one, which cannot collide', () async { + await connection.execute(r'SELECT $1::int AS value', parameters: [1]); + + expect(sent.whereType().single.statementName, isEmpty); + }); + + test('a one-shot query does not close the statement it never named', () async { + await connection.execute(r'SELECT $1::int AS value', parameters: [1]); + + expect( + sent.whereType(), + isEmpty, + reason: 'the next parse replaces the unnamed statement, so closing it is an exchange that changes nothing', + ); + }); + + test('a named statement is still closed, because its name has to become free again', () async { + final statement = await connection.prepare(r'SELECT $1::int AS value'); + await statement.dispose(); + + expect(sent.whereType(), isNotEmpty); + }); + + test('a statement kept for reuse is still named and still parsed on its own', () async { + final statement = await connection.prepare(r'SELECT $1::int AS value'); + addTearDown(statement.dispose); + + final parse = sent.whereType().single; + expect(parse.statementName, isNotEmpty); + // Nothing binds it yet, so its exchange ends where it was sent. + expect(sent.whereType(), isEmpty); + }); + + test('the results are what was asked for', () async { + final result = await connection.execute( + r'SELECT $1::text AS value, $2::int AS count', + parameters: ['hello', 3], + ); + + expect(result.first.toColumnMap(), {'value': 'hello', 'count': 3}); + }); + }); +}