From b746d7a345eadc35ee88a98afc3c5ff0a4d58075 Mon Sep 17 00:00:00 2001 From: kevmoo Date: Tue, 7 Apr 2026 13:45:16 -0700 Subject: [PATCH 1/2] Add tests for PoolResource and close, and fix allowRelease leak --- pkgs/pool/CHANGELOG.md | 1 + pkgs/pool/lib/pool.dart | 14 ++-- pkgs/pool/test/pool_test.dart | 126 ++++++++++++++++++++++++++++++++++ 3 files changed, 135 insertions(+), 6 deletions(-) diff --git a/pkgs/pool/CHANGELOG.md b/pkgs/pool/CHANGELOG.md index 964260196..4f3c31e50 100644 --- a/pkgs/pool/CHANGELOG.md +++ b/pkgs/pool/CHANGELOG.md @@ -1,6 +1,7 @@ ## 1.5.3-wip * Added an example. +* Fix a resource leak when `PoolResource.allowRelease` callback throws. ## 1.5.2 diff --git a/pkgs/pool/lib/pool.dart b/pkgs/pool/lib/pool.dart index 70e9df158..36346497c 100644 --- a/pkgs/pool/lib/pool.dart +++ b/pkgs/pool/lib/pool.dart @@ -27,7 +27,7 @@ class Pool { /// allocated. /// /// See [PoolResource.allowRelease]. - final _onReleaseCallbacks = Queue(); + final _onReleaseCallbacks = Queue Function()>(); /// Completers that will be completed once `onRelease` callbacks are done /// running. @@ -275,7 +275,7 @@ class Pool { /// If there are any pending requests, this will fire the oldest one after /// running [onRelease]. - void _onResourceReleaseAllowed(void Function() onRelease) { + void _onResourceReleaseAllowed(FutureOr Function() onRelease) { _resetTimer(); if (_requestedResources.isNotEmpty) { @@ -297,15 +297,17 @@ class Pool { /// /// Futures returned by [_runOnRelease] always complete in the order they were /// created, even if earlier [onRelease] callbacks take longer to run. - Future _runOnRelease(void Function() onRelease) { + Future _runOnRelease(FutureOr Function() onRelease) { + var completer = Completer.sync(); + _onReleaseCompleters.add(completer); + Future.sync(onRelease).then((value) { _onReleaseCompleters.removeFirst().complete(PoolResource._(this)); - }).catchError((Object error, StackTrace stackTrace) { + }, onError: (Object error, StackTrace stackTrace) { _onReleaseCompleters.removeFirst().completeError(error, stackTrace); + _onResourceReleased(); }); - var completer = Completer.sync(); - _onReleaseCompleters.add(completer); return completer.future; } diff --git a/pkgs/pool/test/pool_test.dart b/pkgs/pool/test/pool_test.dart index 23f073de7..1086148b4 100644 --- a/pkgs/pool/test/pool_test.dart +++ b/pkgs/pool/test/pool_test.dart @@ -277,6 +277,119 @@ void main() { await pool.request(); }); + + test('request() throws if allowRelease callback throws', () async { + var pool = Pool(1); + var resource = await pool.request(); + + var requestFuture = pool.request(); + + var completer = Completer(); + resource.allowRelease(() => completer.future); + + await Future.delayed(Duration.zero); + completer.completeError('oh no!'); + + await expectLater(requestFuture, throwsA('oh no!')); + }); + + test('request() does not leak resources when allowRelease throws', + () async { + var pool = Pool(1); + var resource = await pool.request(); + + var requestFuture = pool.request(); + + var completer = Completer(); + resource.allowRelease(() => completer.future); + + await Future.delayed(Duration.zero); + completer.completeError('oh no!'); + + await expectLater(requestFuture, throwsA('oh no!')); + + // Without the fix, this will hang because the slot is leaked. + var nextRequest = pool.request().timeout( + const Duration(milliseconds: 100), + onTimeout: () => throw TimeoutException('Leaked!')); + + await expectLater(nextRequest, completes); + }); + + test('throwing in request listener does not corrupt state', () async { + var pool = Pool(2); + var resource1 = await pool.request(); + var resource2 = await pool.request(); + + var completer1 = Completer(); + resource1.allowRelease(() => completer1.future); + var completer2 = Completer(); + resource2.allowRelease(() => completer2.future); + + var requestFuture1 = pool.request(); + var requestFuture2 = pool.request(); + + var request1Threw = false; + var requestFuture1WithListener = requestFuture1.then((_) { + request1Threw = true; + throw StateError('Listener 1 threw!'); + }); + + expect(requestFuture1WithListener, throwsA(isA())); + + var request2Completed = false; + var request2Error = false; + unawaited(requestFuture2.then((_) { + request2Completed = true; + }, onError: (Object e) { + request2Error = true; + })); + + await Future.delayed(Duration.zero); + + completer1.complete(); + await Future.delayed(Duration.zero); + + expect(request1Threw, isTrue); + expect(request2Completed, isFalse); + expect(request2Error, isFalse); + + completer2.complete(); + await Future.delayed(Duration.zero); + + expect(request2Completed, isTrue); + expect(request2Error, isFalse); + }); + }); + + group('PoolResource', () { + test('release() throws StateError if called twice', () async { + var pool = Pool(1); + var resource = await pool.request(); + resource.release(); + expect(resource.release, throwsStateError); + }); + + test('allowRelease() throws StateError if called twice', () async { + var pool = Pool(1); + var resource = await pool.request(); + resource.allowRelease(() {}); + expect(() => resource.allowRelease(() {}), throwsStateError); + }); + + test('allowRelease() throws if called after release()', () async { + var pool = Pool(1); + var resource = await pool.request(); + resource.release(); + expect(() => resource.allowRelease(() {}), throwsStateError); + }); + + test('release() throws if called after allowRelease()', () async { + var pool = Pool(1); + var resource = await pool.request(); + resource.allowRelease(() {}); + expect(resource.release, throwsStateError); + }); }); test("done doesn't complete without close", () async { @@ -296,6 +409,19 @@ void main() { expect(() => pool.withResource(() {}), throwsStateError); }); + test('can be called multiple times', () async { + var pool = Pool(1); + var resource = await pool.request(); + + var closeFuture1 = pool.close(); + var closeFuture2 = pool.close(); + + expect(closeFuture1, equals(closeFuture2)); + + resource.release(); + await closeFuture1; + }); + test('pending requests are fulfilled', () async { var pool = Pool(1); var resource1 = await pool.request(); From 724eb3354fed3123d7670d3f96583139f2f7c8c1 Mon Sep 17 00:00:00 2001 From: kevmoo Date: Thu, 23 Jul 2026 08:32:20 -0700 Subject: [PATCH 2/2] Ignore errors from allowRelease callback --- pkgs/pool/lib/pool.dart | 7 +++---- pkgs/pool/test/pool_test.dart | 7 ++++--- 2 files changed, 7 insertions(+), 7 deletions(-) diff --git a/pkgs/pool/lib/pool.dart b/pkgs/pool/lib/pool.dart index e47b853a6..f9bb3a141 100644 --- a/pkgs/pool/lib/pool.dart +++ b/pkgs/pool/lib/pool.dart @@ -307,11 +307,8 @@ class Pool { var completer = Completer.sync(); _onReleaseCompleters.add(completer); - Future.sync(onRelease).then((value) { + Future.sync(onRelease).catchError((_) {}).then((_) { _onReleaseCompleters.removeFirst().complete(PoolResource._(this)); - }).onError((Object error, StackTrace stackTrace) { - _onReleaseCompleters.removeFirst().completeError(error, stackTrace); - _onResourceReleased(); }); return completer.future; @@ -378,6 +375,8 @@ class PoolResource { /// This is useful when a resource's main function is complete, but it may /// produce additional information later on. For example, an isolate's task /// may be complete, but it could still emit asynchronous errors. + /// + /// Any errors thrown by [onRelease] or the future it returns are ignored. void allowRelease(FutureOr Function() onRelease) { if (_released) { throw StateError('A PoolResource may only be released once.'); diff --git a/pkgs/pool/test/pool_test.dart b/pkgs/pool/test/pool_test.dart index 0bca4c240..f769c0785 100644 --- a/pkgs/pool/test/pool_test.dart +++ b/pkgs/pool/test/pool_test.dart @@ -278,7 +278,7 @@ void main() { await pool.request(); }); - test('request() throws if allowRelease callback throws', () async { + test('request() completes if allowRelease callback throws', () async { var pool = Pool(1); var resource = await pool.request(); @@ -290,7 +290,7 @@ void main() { await Future.delayed(Duration.zero); completer.completeError('oh no!'); - await expectLater(requestFuture, throwsA('oh no!')); + await expectLater(requestFuture, completes); }); test('request() does not leak resources when allowRelease throws', @@ -306,7 +306,8 @@ void main() { await Future.delayed(Duration.zero); completer.completeError('oh no!'); - await expectLater(requestFuture, throwsA('oh no!')); + var resource2 = await requestFuture; + resource2.release(); // Without the fix, this will hang because the slot is leaked. var nextRequest = pool.request().timeout(