Skip to content
Closed
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
11 changes: 11 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -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.
Expand Down
5 changes: 5 additions & 0 deletions lib/src/messages/client_messages.dart
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
49 changes: 44 additions & 5 deletions lib/src/v3/connection.dart
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
}
}
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -958,6 +996,7 @@ class _PgResultStreamSubscription
@override
Future<void> handleMessage(ServerMessage message) async {
switch (message) {
case ParseCompleteMessage():
case BindCompleteMessage():
case NoDataMessage():
// Nothing to do!
Expand Down
21 changes: 15 additions & 6 deletions test/connection_test.dart
Original file line number Diff line number Diff line change
Expand Up @@ -314,23 +314,32 @@ 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 {
orderEnsurer.add(3);
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)));
},
);
},
Expand Down
120 changes: 120 additions & 0 deletions test/single_exchange_test.dart
Original file line number Diff line number Diff line change
@@ -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<ClientMessage> sent;
late Connection connection;

setUp(() async {
sent = [];
connection = await Connection.open(
await server.endpoint(),
settings: ConnectionSettings(
transformer: StreamChannelTransformer<Message, Message>(
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<ClientMessage> throughExecute() =>
sent.takeWhile((message) => message is! ExecuteMessage).followedBy([
sent.whereType<ExecuteMessage>().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<SyncMessage>(),
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<ParseMessage>().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<CloseMessage>(),
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<CloseMessage>(), 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<ParseMessage>().single;
expect(parse.statementName, isNotEmpty);
// Nothing binds it yet, so its exchange ends where it was sent.
expect(sent.whereType<BindMessage>(), 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});
});
});
}