diff --git a/pkgs/pool/CHANGELOG.md b/pkgs/pool/CHANGELOG.md index 0b19c40f8b..e721ef715c 100644 --- a/pkgs/pool/CHANGELOG.md +++ b/pkgs/pool/CHANGELOG.md @@ -2,6 +2,7 @@ * Fix `Pool.forEach` to ensure all workers complete before the stream is closed, even on error or cancellation. * 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 e47b853a63..f9bb3a1414 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 ce651893c9..f769c0785d 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( @@ -315,6 +316,51 @@ void main() { 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', () {