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
17 changes: 16 additions & 1 deletion lib/core/database/mysql_connection.dart
Original file line number Diff line number Diff line change
Expand Up @@ -145,6 +145,7 @@ class MysqlConnection {
String? get connectionString => _connectionString;

MySQLConnection? _conn;
Future<void>? _connecting;
bool _isConnected = false;
bool _inTransaction = false;

Expand All @@ -164,7 +165,21 @@ class MysqlConnection {
return '`${id.replaceAll('`', '``')}`';
}

Future<void> connect({int connectTimeoutMs = 10000}) async {
/// Single-flight: a caller that arrives while an attempt is in progress waits
/// for that attempt. Without this, two callers opened two sockets and two SSH
/// tunnels, and a failing attempt dropped the session the other one had made.
Future<void> connect({int connectTimeoutMs = 10000}) {
if (_isConnected && _conn != null) return Future<void>.value();
final inFlight = _connecting;
if (inFlight != null) return inFlight;
final attempt = _connectOnce(connectTimeoutMs);
_connecting = attempt;
return attempt.whenComplete(() {
if (identical(_connecting, attempt)) _connecting = null;
});
}

Future<void> _connectOnce(int connectTimeoutMs) async {
if (_isConnected && _conn != null) return;

var effectivePassword = _password;
Expand Down
26 changes: 18 additions & 8 deletions lib/core/database/mysql_connection_pool.dart
Original file line number Diff line number Diff line change
Expand Up @@ -77,10 +77,7 @@ class MysqlConnectionPool {
entry.idleTimer?.cancel();
entry.idleTimer = null;
entry.refs++;
if (!entry.connection.isConnected) {
await entry.connection.connect();
await entry.connection.setSessionReadOnly(mode.isReadOnlySession);
}
if (!entry.connection.isConnected) await _reconnect(entry, mode);
return MysqlLease._(this, entry, entry.connection);
}

Expand Down Expand Up @@ -113,10 +110,7 @@ class MysqlConnectionPool {
entry.idleTimer?.cancel();
entry.idleTimer = null;
entry.refs++;
if (!entry.connection.isConnected) {
await entry.connection.connect();
await entry.connection.setSessionReadOnly(mode.isReadOnlySession);
}
if (!entry.connection.isConnected) await _reconnect(entry, mode);
return MysqlLease._(this, entry, entry.connection);
}

Expand All @@ -143,6 +137,21 @@ class MysqlConnectionPool {
/// Releases one lease on [entry]. A lease taken before an [interrupt] points
/// at an entry that is no longer in the pool: it must not touch the entry
/// that replaced it under the same key, so it does nothing.
/// Brings a dropped entry back with one reconnect and one session setting,
/// even when several callers ask at once; they all wait for the same attempt.
Future<void> _reconnect(_PoolEntry entry, MysqlSessionMode mode) {
final inFlight = entry.reconnecting;
if (inFlight != null) return inFlight;
final attempt = () async {
await entry.connection.connect();
await entry.connection.setSessionReadOnly(mode.isReadOnlySession);
}();
entry.reconnecting = attempt;
return attempt.whenComplete(() {
if (identical(entry.reconnecting, attempt)) entry.reconnecting = null;
});
}

void _release(_PoolEntry entry) {
if (!identical(_pool[entry.key], entry)) return;
entry.refs--;
Expand Down Expand Up @@ -203,6 +212,7 @@ class _PoolEntry {

final MysqlConnection connection;
final String key;
Future<void>? reconnecting;
int refs = 0;
Timer? idleTimer;
DateTime lastUsed;
Expand Down
17 changes: 16 additions & 1 deletion lib/core/database/postgres_connection.dart
Original file line number Diff line number Diff line change
Expand Up @@ -143,6 +143,7 @@ class PostgresConnection {
String? get connectionString => _connectionString;

Connection? _conn;
Future<void>? _connecting;
bool _isConnected = false;
bool _inTransaction = false;

Expand Down Expand Up @@ -201,7 +202,21 @@ class PostgresConnection {
/// [openFromUrl] already parses `sslmode`, `connect_timeout`, `query_timeout`
/// from the URI. If `sslmode` is omitted, we fall back to [useSSL] so the
/// form checkbox still applies; otherwise libpq-style URLs drive TLS mode.
Future<void> connect() async {
/// Single-flight: a caller that arrives while an attempt is in progress waits
/// for that attempt. Without this, two callers opened two sockets and two SSH
/// tunnels, and a failing attempt dropped the session the other one had made.
Future<void> connect() {
if (_isConnected && _conn != null) return Future<void>.value();
final inFlight = _connecting;
if (inFlight != null) return inFlight;
final attempt = _connectOnce();
_connecting = attempt;
return attempt.whenComplete(() {
if (identical(_connecting, attempt)) _connecting = null;
});
}

Future<void> _connectOnce() async {
if (_isConnected && _conn != null) return;

var effectivePassword = _password;
Expand Down
24 changes: 17 additions & 7 deletions lib/core/database/postgres_connection_pool.dart
Original file line number Diff line number Diff line change
Expand Up @@ -89,10 +89,7 @@ class PostgresConnectionPool {
entry.idleTimer?.cancel();
entry.idleTimer = null;
entry.refs++;
if (!entry.connection.isConnected) {
await entry.connection.connect();
await entry.connection.setSessionReadOnly(mode.isReadOnlySession);
}
if (!entry.connection.isConnected) await _reconnect(entry, mode);
return PgLease._(this, entry, entry.connection);
}

Expand Down Expand Up @@ -125,11 +122,23 @@ class PostgresConnectionPool {
entry.idleTimer?.cancel();
entry.idleTimer = null;
entry.refs++;
if (!entry.connection.isConnected) {
if (!entry.connection.isConnected) await _reconnect(entry, mode);
return PgLease._(this, entry, entry.connection);
}

/// Brings a dropped entry back with one reconnect and one session setting,
/// even when several callers ask at once; they all wait for the same attempt.
Future<void> _reconnect(_PoolEntry entry, PgSessionMode mode) {
final inFlight = entry.reconnecting;
if (inFlight != null) return inFlight;
final attempt = () async {
await entry.connection.connect();
await entry.connection.setSessionReadOnly(mode.isReadOnlySession);
}
return PgLease._(this, entry, entry.connection);
}();
entry.reconnecting = attempt;
return attempt.whenComplete(() {
if (identical(entry.reconnecting, attempt)) entry.reconnecting = null;
});
}

/// Drops idle LRU slots until there is room for one more key.
Expand Down Expand Up @@ -222,6 +231,7 @@ class _PoolEntry {

final PostgresConnection connection;
final String key;
Future<void>? reconnecting;
int refs = 0;
Timer? idleTimer;
DateTime lastUsed;
Expand Down
32 changes: 32 additions & 0 deletions test/core/database/mysql_connection_pool_test.dart
Original file line number Diff line number Diff line change
Expand Up @@ -234,4 +234,36 @@ void main() {
expect(replacement.disconnectCount, 1);
});
});

group('concurrent reconnect (#1214)', () {
test('two callers on a dropped slot share one reconnect and one setting',
() async {
final pool = MysqlConnectionPool(
createAndConnect: (row, {required database, required mode}) async {
final c = FakeMysqlConnection();
await c.connect();
return c;
},
);
final r = _row();
final first = await pool.acquire(r,
database: 'app', mode: MysqlSessionMode.readOnly);
final conn = first.connection as FakeMysqlConnection;
await conn.forceClose();

final leases = await Future.wait([
pool.acquire(r, database: 'app', mode: MysqlSessionMode.readOnly),
pool.acquire(r, database: 'app', mode: MysqlSessionMode.readOnly),
]);

expect(conn.connectCount, 2, reason: 'initial connect plus one reconnect');
expect(conn.setReadOnlyCount, 1);
expect(identical(leases[0].connection, conn), isTrue);
expect(identical(leases[1].connection, conn), isTrue);
first.release();
for (final l in leases) {
l.release();
}
});
});
}
34 changes: 34 additions & 0 deletions test/core/database/postgres_connection_pool_test.dart
Original file line number Diff line number Diff line change
Expand Up @@ -593,4 +593,38 @@ void main() {
expect(replacement.disconnectCount, 1);
});
});

group('concurrent reconnect (#1214)', () {
test('two callers on a dropped slot share one reconnect and one setting',
() async {
final created = <FakePostgresConnection>[];
final pool = PostgresConnectionPool(
createAndConnect: (row, {required database, required mode}) async {
final c = FakePostgresConnection();
await c.connect();
created.add(c);
return c;
},
);
final r = _row();
final first = await pool.acquire(r,
database: 'postgres', mode: PgSessionMode.readOnly);
final conn = first.connection as FakePostgresConnection;
await conn.forceClose();

final leases = await Future.wait([
pool.acquire(r, database: 'postgres', mode: PgSessionMode.readOnly),
pool.acquire(r, database: 'postgres', mode: PgSessionMode.readOnly),
]);

expect(conn.connectCount, 2, reason: 'initial connect plus one reconnect');
expect(conn.setReadOnlyCount, 1);
expect(identical(leases[0].connection, conn), isTrue);
expect(identical(leases[1].connection, conn), isTrue);
first.release();
for (final l in leases) {
l.release();
}
});
});
}
Loading