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
Original file line number Diff line number Diff line change
Expand Up @@ -118,6 +118,9 @@ public static <T, R> AsyncResource<R> transform(
/**
* Returns an {@link AsyncResource} that collects a list of underlying resources into a single lifecycle.
*
* <p>The returned resource becomes ready once every resource in {@code asyncResources} has succeeded, or as soon as
* any of them fails.
*
* <p>Once this method returns, the returned {@link AsyncResource} is the caller's to close, and closing it also
* closes every resource in {@code asyncResources}, so the caller must not close them itself. If this method throws,
* nothing has been taken over and the caller still owns all of them.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@

import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;

/**
Expand All @@ -37,6 +38,12 @@ public class CollectAsyncResource<T> implements AsyncResource<List<T>>
private final AtomicInteger readyCount = new AtomicInteger(0);
private final SettableAsyncResource<List<T>> targetResource = new SettableAsyncResource<>();

/**
* Whether a source has failed. Only the first failure is reported, since {@link #targetResource} can only be
* completed once.
*/
private final AtomicBoolean failed = new AtomicBoolean(false);

/**
* Constructor. Can also be created with {@link AsyncResources#collect(List)}.
*/
Expand All @@ -48,7 +55,7 @@ public class CollectAsyncResource<T> implements AsyncResource<List<T>>
targetResource.set(List.of(), null);
} else {
for (final AsyncResource<T> asyncResource : sourceResources) {
asyncResource.addReadyCallback(this::onOneSourceReady);
asyncResource.addReadyCallback(() -> onOneSourceReady(asyncResource));
}
}
}
Expand Down Expand Up @@ -80,9 +87,23 @@ public void close()
CloseableUtils.closeAndWrapExceptions(closer);
}

private void onOneSourceReady()
private void onOneSourceReady(final AsyncResource<T> sourceResource)
{
try {
sourceResource.get();
}
catch (Throwable e) {
// Fail without waiting for the other sources, but deliberately leave them open: closing them is up to the
// caller, via close().
if (failed.compareAndSet(false, true)) {
targetResource.setException(e);
}
return;
}

if (readyCount.incrementAndGet() == sourceResources.size()) {
// Will only enter this block if all source resources succeeded, because a failed resource
// does not increment readyCount.
try {
final List<T> resources = new ArrayList<>(sourceResources.size());
for (final AsyncResource<T> asyncResource : sourceResources) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,24 @@ public void testCollectOfFromFutureReadyWhenAllComplete()
collected.close();
}

@Test
public void testCollectOfFromFutureFailsWithoutWaitingAndCancelsOthersOnClose()
{
final SettableFuture<String> a = SettableFuture.create();
final SettableFuture<String> b = SettableFuture.create();
final AsyncResource<List<String>> collected =
AsyncResources.collect(List.of(AsyncResources.fromFutureUnmanaged(a), AsyncResources.fromFutureUnmanaged(b)));

final RuntimeException failure = new RuntimeException("boom");
b.setException(failure);
Assertions.assertTrue(collected.isReady(), "one failure must not wait for the other sources");
Assertions.assertSame(failure, Assertions.assertThrows(RuntimeException.class, collected::get));
Assertions.assertFalse(a.isCancelled(), "the other sources are left to close()");

collected.close();
Assertions.assertTrue(a.isCancelled());
}

@Test
public void testFromFutureGetBeforeReadyThrows()
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -96,6 +96,59 @@ public void testSourceFailurePropagates()
collected.close();
}

@Test
public void testSourceFailureFailsWithoutWaitingForOtherSources()
{
final AtomicInteger fired = new AtomicInteger();
final AtomicInteger aCancel = new AtomicInteger();
final AtomicInteger bClose = new AtomicInteger();
final SettableAsyncResource<String> a = new SettableAsyncResource<>();
final SettableAsyncResource<String> b = new SettableAsyncResource<>();
final SettableAsyncResource<String> c = new SettableAsyncResource<>();
a.setCanceler(aCancel::incrementAndGet);
b.set("b", bClose::incrementAndGet);

final AsyncResource<List<String>> collected = AsyncResources.collect(List.of(a, b, c));
collected.addReadyCallback(fired::incrementAndGet);

// a is still pending, but c's failure must not wait for it.
final RuntimeException failure = new RuntimeException("boom");
c.setException(failure);

Assertions.assertTrue(collected.isReady());
Assertions.assertEquals(1, fired.get());
Assertions.assertSame(failure, Assertions.assertThrows(RuntimeException.class, collected::get));
Assertions.assertEquals(0, aCancel.get(), "sources must stay open until the collected resource is closed");
Assertions.assertEquals(0, bClose.get(), "sources must stay open until the collected resource is closed");

collected.close();
Assertions.assertEquals(1, aCancel.get(), "close must cancel the pending source");
Assertions.assertEquals(1, bClose.get(), "close must close the ready source");
}

@Test
public void testAlreadyFailedSourcesFailOnConstructionWithTheFirstFailure()
{
final AtomicInteger cCancel = new AtomicInteger();
final SettableAsyncResource<String> a = new SettableAsyncResource<>();
final SettableAsyncResource<String> b = new SettableAsyncResource<>();
final SettableAsyncResource<String> c = new SettableAsyncResource<>();
final RuntimeException failure = new RuntimeException("boom");
a.setException(failure);
b.setException(new RuntimeException("second"));
c.setCanceler(cCancel::incrementAndGet);

// Both failures fire inline while registering; failing the collected resource twice would throw out of collect.
final AsyncResource<List<String>> collected = AsyncResources.collect(List.of(a, b, c));

Assertions.assertTrue(collected.isReady());
Assertions.assertSame(failure, Assertions.assertThrows(RuntimeException.class, collected::get));
Assertions.assertEquals(0, cCancel.get());

collected.close();
Assertions.assertEquals(1, cCancel.get());
}

@Test
public void testCloseClosesAllReadySources()
{
Expand Down Expand Up @@ -162,8 +215,8 @@ public void testCloseWithAMixOfReadyAndPendingSourcesReportsCancellation()
collected.addReadyCallback(fired::incrementAndGet);
Assertions.assertFalse(collected.isReady(), "one source is still pending");

// The already-ready source counted toward readiness, so closing the pending one brings the internal count to the
// source count and runs the collect body against sources this close just tore down. That must stay harmless.
// Closing the pending source fires its ready callback while this close is still tearing the sources down. That
// must stay harmless.
collected.close();

Assertions.assertEquals(1, fired.get());
Expand Down
Loading