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
1 change: 1 addition & 0 deletions pkgs/pool/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
7 changes: 3 additions & 4 deletions pkgs/pool/lib/pool.dart
Original file line number Diff line number Diff line change
Expand Up @@ -307,11 +307,8 @@ class Pool {
var completer = Completer<PoolResource>.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;
Expand Down Expand Up @@ -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<void> Function() onRelease) {
if (_released) {
throw StateError('A PoolResource may only be released once.');
Expand Down
52 changes: 49 additions & 3 deletions pkgs/pool/test/pool_test.dart
Original file line number Diff line number Diff line change
Expand Up @@ -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();

Expand All @@ -290,7 +290,7 @@ void main() {
await Future<void>.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',
Expand All @@ -306,7 +306,8 @@ void main() {
await Future<void>.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(
Expand All @@ -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<void>();
resource1.allowRelease(() => completer1.future);
var completer2 = Completer<void>();
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<StateError>()));

var request2Completed = false;
var request2Error = false;
unawaited(requestFuture2.then((_) {
request2Completed = true;
}, onError: (Object e) {
request2Error = true;
}));

await Future<void>.delayed(Duration.zero);

completer1.complete();
await Future<void>.delayed(Duration.zero);

expect(request1Threw, isTrue);
expect(request2Completed, isFalse);
expect(request2Error, isFalse);

completer2.complete();
await Future<void>.delayed(Duration.zero);

expect(request2Completed, isTrue);
expect(request2Error, isFalse);
});
});

group('PoolResource', () {
Expand Down
Loading