diff --git a/lib/core/database/mysql_connection_pool.dart b/lib/core/database/mysql_connection_pool.dart index c4670dbc..59a2d34e 100644 --- a/lib/core/database/mysql_connection_pool.dart +++ b/lib/core/database/mysql_connection_pool.dart @@ -29,10 +29,10 @@ typedef MysqlPoolConnectionFactory = Future Function( /// Lease for a pooled [MysqlConnection]. Call [release] when the UI is done. class MysqlLease { - MysqlLease._(this._pool, this._key, this.connection); + MysqlLease._(this._pool, this._entry, this.connection); final MysqlConnectionPool _pool; - final String _key; + final _PoolEntry _entry; final MysqlConnection connection; bool _released = false; @@ -40,7 +40,7 @@ class MysqlLease { void release() { if (_released) return; _released = true; - _pool._release(_key); + _pool._release(_entry); } } @@ -81,7 +81,7 @@ class MysqlConnectionPool { await entry.connection.connect(); await entry.connection.setSessionReadOnly(mode.isReadOnlySession); } - return MysqlLease._(this, k, entry.connection); + return MysqlLease._(this, entry, entry.connection); } try { @@ -89,7 +89,7 @@ class MysqlConnectionPool { _evictIfNeededBeforeNewSlot(); final conn = await createAndConnect(row, database: database, mode: mode); - _pool[k] = _PoolEntry(conn); + _pool[k] = _PoolEntry(conn, k); return conn; }); } on StateError { @@ -117,7 +117,7 @@ class MysqlConnectionPool { await entry.connection.connect(); await entry.connection.setSessionReadOnly(mode.isReadOnlySession); } - return MysqlLease._(this, k, entry.connection); + return MysqlLease._(this, entry, entry.connection); } void _evictIfNeededBeforeNewSlot() { @@ -140,18 +140,23 @@ class MysqlConnectionPool { unawaited(entry.connection.forceClose()); } - void _release(String k) { - final entry = _pool[k]; - if (entry == null) return; + /// 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. + void _release(_PoolEntry entry) { + if (!identical(_pool[entry.key], entry)) return; entry.refs--; + assert( + entry.refs >= 0, + 'pool entry ${entry.key} released more often than leased', + ); if (entry.refs > 0) return; entry.idleTimer?.cancel(); entry.idleTimer = Timer(idleDisposeDelay, () { - final e = _pool[k]; - if (e == null || e.refs > 0) return; - e.idleTimer = null; - unawaited(e.connection.disconnect()); - _pool.remove(k); + if (!identical(_pool[entry.key], entry) || entry.refs > 0) return; + entry.idleTimer = null; + unawaited(entry.connection.disconnect()); + _pool.remove(entry.key); }); } @@ -194,9 +199,10 @@ class MysqlConnectionPool { } class _PoolEntry { - _PoolEntry(this.connection) : lastUsed = DateTime.now(); + _PoolEntry(this.connection, this.key) : lastUsed = DateTime.now(); final MysqlConnection connection; + final String key; int refs = 0; Timer? idleTimer; DateTime lastUsed; diff --git a/lib/core/database/postgres_connection_pool.dart b/lib/core/database/postgres_connection_pool.dart index cb60357e..c9017256 100644 --- a/lib/core/database/postgres_connection_pool.dart +++ b/lib/core/database/postgres_connection_pool.dart @@ -33,10 +33,10 @@ typedef PostgresPoolConnectionFactory = Future Function( /// Lease for a pooled [PostgresConnection]. Call [release] when the UI is done /// (typically in [State.dispose]). class PgLease { - PgLease._(this._pool, this._key, this.connection); + PgLease._(this._pool, this._entry, this.connection); final PostgresConnectionPool _pool; - final String _key; + final _PoolEntry _entry; final PostgresConnection connection; bool _released = false; @@ -45,7 +45,7 @@ class PgLease { void release() { if (_released) return; _released = true; - _pool._release(_key); + _pool._release(_entry); } } @@ -93,7 +93,7 @@ class PostgresConnectionPool { await entry.connection.connect(); await entry.connection.setSessionReadOnly(mode.isReadOnlySession); } - return PgLease._(this, k, entry.connection); + return PgLease._(this, entry, entry.connection); } try { @@ -101,7 +101,7 @@ class PostgresConnectionPool { _evictIfNeededBeforeNewSlot(); final conn = await createAndConnect(row, database: database, mode: mode); - _pool[k] = _PoolEntry(conn); + _pool[k] = _PoolEntry(conn, k); return conn; }); } on StateError { @@ -129,7 +129,7 @@ class PostgresConnectionPool { await entry.connection.connect(); await entry.connection.setSessionReadOnly(mode.isReadOnlySession); } - return PgLease._(this, k, entry.connection); + return PgLease._(this, entry, entry.connection); } /// Drops idle LRU slots until there is room for one more key. @@ -153,18 +153,23 @@ class PostgresConnectionPool { unawaited(entry.connection.forceClose()); } - void _release(String k) { - final entry = _pool[k]; - if (entry == null) return; + /// 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. + void _release(_PoolEntry entry) { + if (!identical(_pool[entry.key], entry)) return; entry.refs--; + assert( + entry.refs >= 0, + 'pool entry ${entry.key} released more often than leased', + ); if (entry.refs > 0) return; entry.idleTimer?.cancel(); entry.idleTimer = Timer(idleDisposeDelay, () { - final e = _pool[k]; - if (e == null || e.refs > 0) return; - e.idleTimer = null; - unawaited(e.connection.disconnect()); - _pool.remove(k); + if (!identical(_pool[entry.key], entry) || entry.refs > 0) return; + entry.idleTimer = null; + unawaited(entry.connection.disconnect()); + _pool.remove(entry.key); }); } @@ -213,9 +218,10 @@ class PostgresConnectionPool { } class _PoolEntry { - _PoolEntry(this.connection) : lastUsed = DateTime.now(); + _PoolEntry(this.connection, this.key) : lastUsed = DateTime.now(); final PostgresConnection connection; + final String key; int refs = 0; Timer? idleTimer; DateTime lastUsed; diff --git a/lib/core/database/sqlite_connection_pool.dart b/lib/core/database/sqlite_connection_pool.dart index 52543087..a99bc08e 100644 --- a/lib/core/database/sqlite_connection_pool.dart +++ b/lib/core/database/sqlite_connection_pool.dart @@ -29,10 +29,10 @@ typedef SqlitePoolConnectionFactory = Future Function( /// Lease for a pooled [SqliteConnection]. Call [release] when the UI is done. class SqliteLease { - SqliteLease._(this._pool, this._key, this.connection); + SqliteLease._(this._pool, this._entry, this.connection); final SqliteConnectionPool _pool; - final String _key; + final _PoolEntry _entry; final SqliteConnection connection; bool _released = false; @@ -40,13 +40,14 @@ class SqliteLease { void release() { if (_released) return; _released = true; - _pool._release(_key); + _pool._release(_entry); } } class _PoolEntry { - _PoolEntry(this.connection); + _PoolEntry(this.connection, this.key); final SqliteConnection connection; + final String key; int refs = 0; DateTime lastUsed = DateTime.now(); Timer? idleTimer; @@ -90,14 +91,14 @@ class SqliteConnectionPool { if (!entry.connection.isConnected) { await entry.connection.connect(); } - return SqliteLease._(this, k, entry.connection); + return SqliteLease._(this, entry, entry.connection); } try { await _creationLock.createIfAbsent(k, () async { _evictIfNeededBeforeNewSlot(); final conn = await createAndConnect(row, mode: mode); - _pool[k] = _PoolEntry(conn); + _pool[k] = _PoolEntry(conn, k); return conn; }); } on StateError { @@ -124,7 +125,7 @@ class SqliteConnectionPool { if (!entry.connection.isConnected) { await entry.connection.connect(); } - return SqliteLease._(this, k, entry.connection); + return SqliteLease._(this, entry, entry.connection); } void _evictIfNeededBeforeNewSlot() { @@ -147,19 +148,24 @@ class SqliteConnectionPool { unawaited(entry.connection.forceClose()); } - void _release(String key) { - final entry = _pool[key]; - if (entry == null) return; + /// 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. + void _release(_PoolEntry entry) { + if (!identical(_pool[entry.key], entry)) return; entry.refs--; + assert( + entry.refs >= 0, + 'pool entry ${entry.key} released more often than leased', + ); if (entry.refs <= 0) { entry.refs = 0; entry.idleTimer?.cancel(); entry.idleTimer = Timer(idleDisposeDelay, () { - final e = _pool[key]; - if (e == null || e.refs > 0) return; - e.idleTimer = null; - unawaited(e.connection.disconnect()); - _pool.remove(key); + if (!identical(_pool[entry.key], entry) || entry.refs > 0) return; + entry.idleTimer = null; + unawaited(entry.connection.disconnect()); + _pool.remove(entry.key); }); } } diff --git a/test/core/database/mysql_connection_pool_test.dart b/test/core/database/mysql_connection_pool_test.dart index cdafe6a9..b9785055 100644 --- a/test/core/database/mysql_connection_pool_test.dart +++ b/test/core/database/mysql_connection_pool_test.dart @@ -200,4 +200,38 @@ void main() { grid.release(); }); }); + + group('lease after interrupt (#1213)', () { + test('releasing a lease from before an interrupt keeps the replacement', + () async { + final created = []; + final pool = MysqlConnectionPool( + idleDisposeDelay: const Duration(milliseconds: 20), + createAndConnect: (row, {required database, required mode}) async { + final c = FakeMysqlConnection(); + await c.connect(); + created.add(c); + return c; + }, + ); + final r = _row(); + final before = await pool.acquire(r, + database: 'app', mode: MysqlSessionMode.readOnly); + pool.interrupt(r, database: 'app', mode: MysqlSessionMode.readOnly); + final after = await pool.acquire(r, + database: 'app', mode: MysqlSessionMode.readOnly); + + before.release(); + await Future.delayed(const Duration(milliseconds: 100)); + + final replacement = created.last; + expect(identical(after.connection, replacement), isTrue); + expect(replacement.isConnected, isTrue); + expect(replacement.disconnectCount, 0); + + after.release(); + await Future.delayed(const Duration(milliseconds: 100)); + expect(replacement.disconnectCount, 1); + }); + }); } diff --git a/test/core/database/postgres_connection_pool_test.dart b/test/core/database/postgres_connection_pool_test.dart index b439ac63..ed669581 100644 --- a/test/core/database/postgres_connection_pool_test.dart +++ b/test/core/database/postgres_connection_pool_test.dart @@ -559,4 +559,38 @@ void main() { ); }); }); + + group('lease after interrupt (#1213)', () { + test('releasing a lease from before an interrupt keeps the replacement', + () async { + final created = []; + final pool = PostgresConnectionPool( + idleDisposeDelay: const Duration(milliseconds: 20), + createAndConnect: (row, {required database, required mode}) async { + final c = FakePostgresConnection(); + await c.connect(); + created.add(c); + return c; + }, + ); + final r = _row(); + final before = await pool.acquire(r, + database: 'postgres', mode: PgSessionMode.readOnly); + pool.interrupt(r, database: 'postgres', mode: PgSessionMode.readOnly); + final after = await pool.acquire(r, + database: 'postgres', mode: PgSessionMode.readOnly); + + before.release(); + await Future.delayed(const Duration(milliseconds: 100)); + + final replacement = created.last; + expect(identical(after.connection, replacement), isTrue); + expect(replacement.isConnected, isTrue); + expect(replacement.disconnectCount, 0); + + after.release(); + await Future.delayed(const Duration(milliseconds: 100)); + expect(replacement.disconnectCount, 1); + }); + }); } diff --git a/test/core/database/sqlite_connection_pool_test.dart b/test/core/database/sqlite_connection_pool_test.dart index c476f29f..954e3370 100644 --- a/test/core/database/sqlite_connection_pool_test.dart +++ b/test/core/database/sqlite_connection_pool_test.dart @@ -210,4 +210,36 @@ void main() { grid.release(); }); }); + + group('lease after interrupt (#1213)', () { + test('releasing a lease from before an interrupt keeps the replacement', + () async { + final created = []; + final pool = SqliteConnectionPool( + idleDisposeDelay: const Duration(milliseconds: 20), + createAndConnect: (row, {required mode}) async { + final c = FakeSqliteConnection(); + await c.connect(); + created.add(c); + return c; + }, + ); + final r = _row(); + final before = await pool.acquire(r, mode: SqliteSessionMode.readOnly); + pool.interrupt(r, mode: SqliteSessionMode.readOnly); + final after = await pool.acquire(r, mode: SqliteSessionMode.readOnly); + + before.release(); + await Future.delayed(const Duration(milliseconds: 100)); + + final replacement = created.last; + expect(identical(after.connection, replacement), isTrue); + expect(replacement.isConnected, isTrue); + expect(replacement.disconnectCount, 0); + + after.release(); + await Future.delayed(const Duration(milliseconds: 100)); + expect(replacement.disconnectCount, 1); + }); + }); }