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
78 changes: 55 additions & 23 deletions lib/core/database/mysql_connection.dart
Original file line number Diff line number Diff line change
Expand Up @@ -146,6 +146,27 @@ class MysqlConnection {

MySQLConnection? _conn;
Future<void>? _connecting;

/// Statements on this session run one at a time, here. The client would
/// otherwise poll for the previous statement and give up after 10 s, which
/// closed the session for every caller.
Future<void> _statementQueue = Future<void>.value();

/// Test seam: replaces the driver call of [execute].
@visibleForTesting
Future<IResultSet> Function(
String sql,
Map<String, dynamic>? params,
bool iterable,
Duration? timeout,
)? runStatementForTest;

Future<T> _serialized<T>(Future<T> Function() run) {
final previous = _statementQueue;
final turn = Completer<void>();
_statementQueue = turn.future;
return previous.then((_) => run()).whenComplete(() => turn.complete());
}
bool _isConnected = false;
bool _inTransaction = false;

Expand Down Expand Up @@ -441,18 +462,36 @@ class MysqlConnection {
String sql, [
Map<String, dynamic>? params,
bool iterable = false,
]) async {
if (!isConnected || _conn == null) {
throw StateError('Not connected to MySQL');
}
try {
final rs = await _conn!.execute(sql, params, iterable);
_noteTransactionSql(sql);
return rs;
} on TimeoutException {
unawaited(forceClose());
rethrow;
}
]) =>
_run(sql, params, iterable, null);

Future<IResultSet> _run(
String sql,
Map<String, dynamic>? params,
bool iterable,
Duration? timeout,
) {
return _serialized(() async {
final seam = runStatementForTest;
final c = _conn;
if (seam == null && (!isConnected || c == null)) {
throw StateError('Not connected to MySQL');
}
try {
final pending = seam != null
? seam(sql, params, iterable, timeout)
: c!.execute(sql, params, iterable);
final rs = timeout == null ? await pending : await pending.timeout(timeout);
_noteTransactionSql(sql);
return rs;
} on TimeoutException {
// The client cannot cancel a running statement, so the session is out of
// step after a timeout and is closed. Statements queued behind this one
// wait in the queue and are not affected.
unawaited(forceClose());
rethrow;
}
});
}

/// Whether this session has an open `START TRANSACTION` / `BEGIN`.
Expand Down Expand Up @@ -498,22 +537,15 @@ class MysqlConnection {
}
}

/// Runs [execute] with an application-level [timeout] (driver limits still apply).
/// [execute] with [timeout] as the statement's own limit; it starts when the
/// statement starts, not while the statement waits behind another one.
Future<IResultSet> executeWithTimeout(
String sql, {
Duration? timeout,
Map<String, dynamic>? params,
bool iterable = false,
}) async {
final f = execute(sql, params, iterable);
if (timeout == null) return f;
try {
return await f.timeout(timeout);
} on TimeoutException {
unawaited(forceClose());
rethrow;
}
}
}) =>
_run(sql, params, iterable, timeout);

/// Lists user-visible databases (excludes typical system schemas).
Future<List<String>> listDatabases() async {
Expand Down
72 changes: 46 additions & 26 deletions lib/core/database/postgres_connection.dart
Original file line number Diff line number Diff line change
Expand Up @@ -144,6 +144,22 @@ class PostgresConnection {

Connection? _conn;
Future<void>? _connecting;

/// Statements on this session run one at a time, here and not in the driver.
/// The driver starts a statement's timeout before the statement gets its turn,
/// so a statement that waited behind another one could cancel the other one.
Future<void> _statementQueue = Future<void>.value();

/// Test seam: replaces the driver call of [execute].
@visibleForTesting
Future<Result> Function(String sql, Duration? timeout)? runStatementForTest;

Future<T> _serialized<T>(Future<T> Function() run) {
final previous = _statementQueue;
final turn = Completer<void>();
_statementQueue = turn.future;
return previous.then((_) => run()).whenComplete(() => turn.complete());
}
bool _isConnected = false;
bool _inTransaction = false;

Expand Down Expand Up @@ -391,37 +407,41 @@ class PostgresConnection {
}
}

/// Runs SQL on the underlying session. [timeout] overrides
/// [ConnectionSettings.queryTimeout] for this statement (see `postgres`
/// package).
Future<Result> execute(String sql, {Duration? timeout}) async {
if (!isConnected || _conn == null) {
throw StateError('Not connected to PostgreSQL');
}
try {
final result = await _conn!.execute(sql, timeout: timeout);
_inTransaction = applyPostgresTransactionSql(_inTransaction, sql);
return result;
} on TimeoutException {
unawaited(forceClose());
rethrow;
}
/// Runs SQL on the underlying session, one statement at a time. [timeout]
/// overrides [ConnectionSettings.queryTimeout] for this statement and starts
/// when the statement does.
///
/// A server-side cancel (SQLSTATE 57014) arrives as a [PgException], and the
/// session stays usable. Only a bare [TimeoutException] closes it: then the
/// protocol may be out of step.
Future<Result> execute(String sql, {Duration? timeout}) {
return _serialized(() async {
final seam = runStatementForTest;
final c = _conn;
if (seam == null && (!isConnected || c == null)) {
throw StateError('Not connected to PostgreSQL');
}
try {
final result = seam != null
? await seam(sql, timeout)
: await c!.execute(sql, timeout: timeout);
_inTransaction = applyPostgresTransactionSql(_inTransaction, sql);
return result;
} on TimeoutException catch (e) {
if (e is! PgException) unawaited(forceClose());
rethrow;
}
});
}

/// Runs [execute] with an application-level [timeout] (in addition to driver timeout).
/// [execute] with [timeout] as the statement's own limit. The limit starts
/// when the statement starts: time spent waiting behind another statement on
/// this session does not count against it.
Future<Result> executeWithTimeout(
String sql, {
Duration? timeout,
}) async {
final f = execute(sql, timeout: timeout);
if (timeout == null) return f;
try {
return await f.timeout(timeout);
} on TimeoutException {
unawaited(forceClose());
rethrow;
}
}
}) =>
execute(sql, timeout: timeout);

/// Whether the session has an open transaction.
///
Expand Down
15 changes: 7 additions & 8 deletions test/core/database/connection_timeout_protocol_test.dart
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,6 @@ import 'package:mysql_client/mysql_client.dart';
import 'package:querya_desktop/core/database/mysql_connection.dart';
import 'package:querya_desktop/core/database/postgres_connection.dart';
import 'package:querya_desktop/core/database/sqlite_connection.dart';
import 'package:postgres/postgres.dart' as pg;

class FakeSlowMysqlConnection extends MysqlConnection {
FakeSlowMysqlConnection({super.id = 1})
Expand Down Expand Up @@ -53,12 +52,6 @@ class FakeSlowPostgresConnection extends PostgresConnection {
@override
bool get isConnected => _connected;

@override
Future<pg.Result> execute(String sql, {Duration? timeout}) async {
await Future.delayed(const Duration(seconds: 10));
throw Exception('should not reach here');
}

@override
Future<void> forceClose() async {
forceCloseCount++;
Expand Down Expand Up @@ -100,6 +93,9 @@ void main() {
test('MysqlConnection.executeWithTimeout force-closes on timeout', () async {
final conn = FakeSlowMysqlConnection();
expect(conn.isConnected, isTrue);
// The statement never answers: the timeout is the only way out.
conn.runStatementForTest = (sql, params, iterable, timeout) =>
Completer<IResultSet>().future;

try {
await conn.executeWithTimeout(
Expand All @@ -116,9 +112,12 @@ void main() {
expect(conn.isConnected, isFalse);
});

test('PostgresConnection.executeWithTimeout force-closes on TimeoutException', () async {
test('PostgresConnection.executeWithTimeout force-closes on a bare TimeoutException', () async {
final conn = FakeSlowPostgresConnection();
expect(conn.isConnected, isTrue);
// The driver's own timer fires as a bare TimeoutException.
conn.runStatementForTest = (sql, timeout) =>
Future.error(TimeoutException('driver timer', timeout));

try {
await conn.executeWithTimeout(
Expand Down
135 changes: 135 additions & 0 deletions test/core/database/statement_queue_test.dart
Original file line number Diff line number Diff line change
@@ -0,0 +1,135 @@
import 'dart:async';

import 'package:flutter_test/flutter_test.dart';
import 'package:mysql_client/mysql_client.dart';
import 'package:postgres/postgres.dart' show PgException;
import 'package:querya_desktop/core/database/mysql_connection.dart';
import 'package:querya_desktop/core/database/postgres_connection.dart';

/// Records force-closes instead of touching a socket.
class _ProbePostgres extends PostgresConnection {
_ProbePostgres()
: super(
id: 1,
name: 'probe',
host: 'localhost',
port: 5432,
database: 'postgres',
);

int forced = 0;

@override
Future<void> forceClose() async {
forced++;
}
}

class _ProbeMysql extends MysqlConnection {
_ProbeMysql()
: super(
id: 1,
name: 'probe',
host: 'localhost',
port: 3306,
database: 'testdb',
);

int forced = 0;

@override
Future<void> forceClose() async {
forced++;
}
}

void main() {
group('PostgreSQL statements on one session (#1216)', () {
test('a server-side cancel is a PgException and keeps the session',
() async {
final conn = _ProbePostgres();
conn.runStatementForTest = (sql, timeout) =>
Future.error(PgException('canceling statement due to user request'));

await expectLater(
conn.execute('SELECT 1'),
throwsA(isA<PgException>()),
);
await Future<void>.delayed(Duration.zero);
expect(conn.forced, 0);
});

test('a bare TimeoutException closes the session', () async {
final conn = _ProbePostgres();
conn.runStatementForTest = (sql, timeout) =>
Future.error(TimeoutException('driver timer', timeout));

await expectLater(
conn.execute('SELECT 1', timeout: const Duration(seconds: 1)),
throwsA(isA<TimeoutException>()),
);
await Future<void>.delayed(Duration.zero);
expect(conn.forced, 1);
});

test('a statement waits for the one before it on the same session',
() async {
final conn = _ProbePostgres();
final log = <String>[];
final gate = Completer<void>();
conn.runStatementForTest = (sql, timeout) async {
log.add('start $sql');
if (sql == 'A') await gate.future;
log.add('end $sql');
throw StateError('no result in this probe');
};

final a = conn.execute('A').then<void>((_) {}, onError: (_) {});
final b = conn.execute('B').then<void>((_) {}, onError: (_) {});
await Future<void>.delayed(Duration.zero);
expect(log, ['start A'], reason: 'B must not start while A runs');

gate.complete();
await Future.wait([a, b]);
expect(log, ['start A', 'end A', 'start B', 'end B']);
});
});

group('MySQL statements on one session (#1216)', () {
test('a statement waits for the one before it on the same session',
() async {
final conn = _ProbeMysql();
final log = <String>[];
final gate = Completer<void>();
conn.runStatementForTest = (sql, params, iterable, timeout) async {
log.add('start $sql');
if (sql == 'A') await gate.future;
log.add('end $sql');
throw StateError('no result in this probe');
};

final a = conn.execute('A').then<void>((_) {}, onError: (_) {});
final b = conn.execute('B').then<void>((_) {}, onError: (_) {});
await Future<void>.delayed(Duration.zero);
expect(log, ['start A']);

gate.complete();
await Future.wait([a, b]);
expect(log, ['start A', 'end A', 'start B', 'end B']);
});

test('a timed-out statement closes the session it ran on', () async {
final conn = _ProbeMysql();
conn.runStatementForTest = (sql, params, iterable, timeout) =>
Completer<IResultSet>().future;

await expectLater(
conn.executeWithTimeout('SELECT SLEEP(100)',
timeout: const Duration(milliseconds: 20)),
throwsA(isA<TimeoutException>()),
);
await Future<void>.delayed(Duration.zero);
expect(conn.forced, 1);
});
});
}
Loading