diff --git a/lib/core/database/mysql_connection.dart b/lib/core/database/mysql_connection.dart index 7efc5b6b..2927b94e 100644 --- a/lib/core/database/mysql_connection.dart +++ b/lib/core/database/mysql_connection.dart @@ -145,6 +145,7 @@ class MysqlConnection { String? get connectionString => _connectionString; MySQLConnection? _conn; + Future? _connecting; bool _isConnected = false; bool _inTransaction = false; @@ -164,7 +165,21 @@ class MysqlConnection { return '`${id.replaceAll('`', '``')}`'; } - Future 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 connect({int connectTimeoutMs = 10000}) { + if (_isConnected && _conn != null) return Future.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 _connectOnce(int connectTimeoutMs) async { if (_isConnected && _conn != null) return; var effectivePassword = _password; diff --git a/lib/core/database/mysql_connection_pool.dart b/lib/core/database/mysql_connection_pool.dart index 59a2d34e..85df2667 100644 --- a/lib/core/database/mysql_connection_pool.dart +++ b/lib/core/database/mysql_connection_pool.dart @@ -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); } @@ -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); } @@ -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 _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--; @@ -203,6 +212,7 @@ class _PoolEntry { final MysqlConnection connection; final String key; + Future? reconnecting; int refs = 0; Timer? idleTimer; DateTime lastUsed; diff --git a/lib/core/database/postgres_connection.dart b/lib/core/database/postgres_connection.dart index 91833c2f..68bd4a6a 100644 --- a/lib/core/database/postgres_connection.dart +++ b/lib/core/database/postgres_connection.dart @@ -143,6 +143,7 @@ class PostgresConnection { String? get connectionString => _connectionString; Connection? _conn; + Future? _connecting; bool _isConnected = false; bool _inTransaction = false; @@ -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 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 connect() { + if (_isConnected && _conn != null) return Future.value(); + final inFlight = _connecting; + if (inFlight != null) return inFlight; + final attempt = _connectOnce(); + _connecting = attempt; + return attempt.whenComplete(() { + if (identical(_connecting, attempt)) _connecting = null; + }); + } + + Future _connectOnce() async { if (_isConnected && _conn != null) return; var effectivePassword = _password; diff --git a/lib/core/database/postgres_connection_pool.dart b/lib/core/database/postgres_connection_pool.dart index c9017256..1602cdf9 100644 --- a/lib/core/database/postgres_connection_pool.dart +++ b/lib/core/database/postgres_connection_pool.dart @@ -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); } @@ -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 _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. @@ -222,6 +231,7 @@ class _PoolEntry { final PostgresConnection connection; final String key; + Future? reconnecting; int refs = 0; Timer? idleTimer; DateTime lastUsed; diff --git a/test/core/database/mysql_connection_pool_test.dart b/test/core/database/mysql_connection_pool_test.dart index b9785055..3d5e82d9 100644 --- a/test/core/database/mysql_connection_pool_test.dart +++ b/test/core/database/mysql_connection_pool_test.dart @@ -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(); + } + }); + }); } diff --git a/test/core/database/postgres_connection_pool_test.dart b/test/core/database/postgres_connection_pool_test.dart index ed669581..483b2d92 100644 --- a/test/core/database/postgres_connection_pool_test.dart +++ b/test/core/database/postgres_connection_pool_test.dart @@ -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 = []; + 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(); + } + }); + }); }