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
36 changes: 21 additions & 15 deletions lib/core/database/mysql_connection_pool.dart
Original file line number Diff line number Diff line change
Expand Up @@ -29,18 +29,18 @@ typedef MysqlPoolConnectionFactory = Future<MysqlConnection> 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;

void release() {
if (_released) return;
_released = true;
_pool._release(_key);
_pool._release(_entry);
}
}

Expand Down Expand Up @@ -81,15 +81,15 @@ 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 {
await _creationLock.createIfAbsent(k, () async {
_evictIfNeededBeforeNewSlot();
final conn =
await createAndConnect(row, database: database, mode: mode);
_pool[k] = _PoolEntry(conn);
_pool[k] = _PoolEntry(conn, k);
return conn;
});
} on StateError {
Expand Down Expand Up @@ -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() {
Expand All @@ -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);
});
}

Expand Down Expand Up @@ -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;
Expand Down
36 changes: 21 additions & 15 deletions lib/core/database/postgres_connection_pool.dart
Original file line number Diff line number Diff line change
Expand Up @@ -33,10 +33,10 @@ typedef PostgresPoolConnectionFactory = Future<PostgresConnection> 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;
Expand All @@ -45,7 +45,7 @@ class PgLease {
void release() {
if (_released) return;
_released = true;
_pool._release(_key);
_pool._release(_entry);
}
}

Expand Down Expand Up @@ -93,15 +93,15 @@ 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 {
await _creationLock.createIfAbsent(k, () async {
_evictIfNeededBeforeNewSlot();
final conn =
await createAndConnect(row, database: database, mode: mode);
_pool[k] = _PoolEntry(conn);
_pool[k] = _PoolEntry(conn, k);
return conn;
});
} on StateError {
Expand Down Expand Up @@ -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.
Expand All @@ -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);
});
}

Expand Down Expand Up @@ -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;
Expand Down
36 changes: 21 additions & 15 deletions lib/core/database/sqlite_connection_pool.dart
Original file line number Diff line number Diff line change
Expand Up @@ -29,24 +29,25 @@ typedef SqlitePoolConnectionFactory = Future<SqliteConnection> 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;

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;
Expand Down Expand Up @@ -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 {
Expand All @@ -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() {
Expand All @@ -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);
});
}
}
Expand Down
34 changes: 34 additions & 0 deletions test/core/database/mysql_connection_pool_test.dart
Original file line number Diff line number Diff line change
Expand Up @@ -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 = <FakeMysqlConnection>[];
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<void>.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<void>.delayed(const Duration(milliseconds: 100));
expect(replacement.disconnectCount, 1);
});
});
}
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 @@ -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 = <FakePostgresConnection>[];
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<void>.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<void>.delayed(const Duration(milliseconds: 100));
expect(replacement.disconnectCount, 1);
});
});
}
32 changes: 32 additions & 0 deletions test/core/database/sqlite_connection_pool_test.dart
Original file line number Diff line number Diff line change
Expand Up @@ -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 = <FakeSqliteConnection>[];
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<void>.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<void>.delayed(const Duration(milliseconds: 100));
expect(replacement.disconnectCount, 1);
});
});
}
Loading