8157522: Performance improvements to CompletableFuture

Reviewed-by: martin, psandoz, rriggs, plevart, dfuchs
This commit is contained in:
Doug Lea 2016-07-15 13:59:58 -07:00
parent 7fa43fa58f
commit a09ddefd05
2 changed files with 1003 additions and 634 deletions

View File

@ -57,6 +57,7 @@ import java.util.concurrent.ExecutionException;
import java.util.concurrent.Executor;
import java.util.concurrent.ForkJoinPool;
import java.util.concurrent.ForkJoinTask;
import java.util.concurrent.RejectedExecutionException;
import java.util.concurrent.TimeoutException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
@ -486,62 +487,68 @@ public class CompletableFutureTest extends JSR166TestCase {
class FailingSupplier extends CheckedAction
implements Supplier<Integer>
{
FailingSupplier(ExecutionMode m) { super(m); }
final CFException ex;
FailingSupplier(ExecutionMode m) { super(m); ex = new CFException(); }
public Integer get() {
invoked();
throw new CFException();
throw ex;
}
}
class FailingConsumer extends CheckedIntegerAction
implements Consumer<Integer>
{
FailingConsumer(ExecutionMode m) { super(m); }
final CFException ex;
FailingConsumer(ExecutionMode m) { super(m); ex = new CFException(); }
public void accept(Integer x) {
invoked();
value = x;
throw new CFException();
throw ex;
}
}
class FailingBiConsumer extends CheckedIntegerAction
implements BiConsumer<Integer, Integer>
{
FailingBiConsumer(ExecutionMode m) { super(m); }
final CFException ex;
FailingBiConsumer(ExecutionMode m) { super(m); ex = new CFException(); }
public void accept(Integer x, Integer y) {
invoked();
value = subtract(x, y);
throw new CFException();
throw ex;
}
}
class FailingFunction extends CheckedIntegerAction
implements Function<Integer, Integer>
{
FailingFunction(ExecutionMode m) { super(m); }
final CFException ex;
FailingFunction(ExecutionMode m) { super(m); ex = new CFException(); }
public Integer apply(Integer x) {
invoked();
value = x;
throw new CFException();
throw ex;
}
}
class FailingBiFunction extends CheckedIntegerAction
implements BiFunction<Integer, Integer, Integer>
{
FailingBiFunction(ExecutionMode m) { super(m); }
final CFException ex;
FailingBiFunction(ExecutionMode m) { super(m); ex = new CFException(); }
public Integer apply(Integer x, Integer y) {
invoked();
value = subtract(x, y);
throw new CFException();
throw ex;
}
}
class FailingRunnable extends CheckedAction implements Runnable {
FailingRunnable(ExecutionMode m) { super(m); }
final CFException ex;
FailingRunnable(ExecutionMode m) { super(m); ex = new CFException(); }
public void run() {
invoked();
throw new CFException();
throw ex;
}
}
@ -561,11 +568,21 @@ public class CompletableFutureTest extends JSR166TestCase {
class FailingCompletableFutureFunction extends CheckedIntegerAction
implements Function<Integer, CompletableFuture<Integer>>
{
FailingCompletableFutureFunction(ExecutionMode m) { super(m); }
final CFException ex;
FailingCompletableFutureFunction(ExecutionMode m) { super(m); ex = new CFException(); }
public CompletableFuture<Integer> apply(Integer x) {
invoked();
value = x;
throw new CFException();
throw ex;
}
}
static class CountingRejectingExecutor implements Executor {
final RejectedExecutionException ex = new RejectedExecutionException();
final AtomicInteger count = new AtomicInteger(0);
public void execute(Runnable r) {
count.getAndIncrement();
throw ex;
}
}
@ -1249,10 +1266,22 @@ public class CompletableFutureTest extends JSR166TestCase {
{
final FailingRunnable r = new FailingRunnable(m);
final CompletableFuture<Void> f = m.runAsync(r);
checkCompletedWithWrappedCFException(f);
checkCompletedWithWrappedException(f, r.ex);
r.assertInvoked();
}}
public void testRunAsync_rejectingExecutor() {
CountingRejectingExecutor e = new CountingRejectingExecutor();
try {
CompletableFuture.runAsync(() -> {}, e);
shouldThrow();
} catch (Throwable t) {
assertSame(e.ex, t);
}
assertEquals(1, e.count.get());
}
/**
* supplyAsync completes with result of supplier
*/
@ -1283,10 +1312,22 @@ public class CompletableFutureTest extends JSR166TestCase {
{
FailingSupplier r = new FailingSupplier(m);
CompletableFuture<Integer> f = m.supplyAsync(r);
checkCompletedWithWrappedCFException(f);
checkCompletedWithWrappedException(f, r.ex);
r.assertInvoked();
}}
public void testSupplyAsync_rejectingExecutor() {
CountingRejectingExecutor e = new CountingRejectingExecutor();
try {
CompletableFuture.supplyAsync(() -> null, e);
shouldThrow();
} catch (Throwable t) {
assertSame(e.ex, t);
}
assertEquals(1, e.count.get());
}
// seq completion methods
/**
@ -1405,12 +1446,12 @@ public class CompletableFutureTest extends JSR166TestCase {
final CompletableFuture<Void> h4 = m.runAfterBoth(f, f, rs[4]);
final CompletableFuture<Void> h5 = m.runAfterEither(f, f, rs[5]);
checkCompletedWithWrappedCFException(h0);
checkCompletedWithWrappedCFException(h1);
checkCompletedWithWrappedCFException(h2);
checkCompletedWithWrappedCFException(h3);
checkCompletedWithWrappedCFException(h4);
checkCompletedWithWrappedCFException(h5);
checkCompletedWithWrappedException(h0, rs[0].ex);
checkCompletedWithWrappedException(h1, rs[1].ex);
checkCompletedWithWrappedException(h2, rs[2].ex);
checkCompletedWithWrappedException(h3, rs[3].ex);
checkCompletedWithWrappedException(h4, rs[4].ex);
checkCompletedWithWrappedException(h5, rs[5].ex);
checkCompletedNormally(f, v1);
}}
@ -1509,10 +1550,10 @@ public class CompletableFutureTest extends JSR166TestCase {
final CompletableFuture<Integer> h2 = m.thenApply(f, rs[2]);
final CompletableFuture<Integer> h3 = m.applyToEither(f, f, rs[3]);
checkCompletedWithWrappedCFException(h0);
checkCompletedWithWrappedCFException(h1);
checkCompletedWithWrappedCFException(h2);
checkCompletedWithWrappedCFException(h3);
checkCompletedWithWrappedException(h0, rs[0].ex);
checkCompletedWithWrappedException(h1, rs[1].ex);
checkCompletedWithWrappedException(h2, rs[2].ex);
checkCompletedWithWrappedException(h3, rs[3].ex);
checkCompletedNormally(f, v1);
}}
@ -1611,10 +1652,10 @@ public class CompletableFutureTest extends JSR166TestCase {
final CompletableFuture<Void> h2 = m.thenAccept(f, rs[2]);
final CompletableFuture<Void> h3 = m.acceptEither(f, f, rs[3]);
checkCompletedWithWrappedCFException(h0);
checkCompletedWithWrappedCFException(h1);
checkCompletedWithWrappedCFException(h2);
checkCompletedWithWrappedCFException(h3);
checkCompletedWithWrappedException(h0, rs[0].ex);
checkCompletedWithWrappedException(h1, rs[1].ex);
checkCompletedWithWrappedException(h2, rs[2].ex);
checkCompletedWithWrappedException(h3, rs[3].ex);
checkCompletedNormally(f, v1);
}}
@ -1776,9 +1817,9 @@ public class CompletableFutureTest extends JSR166TestCase {
assertTrue(snd.complete(w2));
final CompletableFuture<Integer> h3 = m.thenCombine(f, g, r3);
checkCompletedWithWrappedCFException(h1);
checkCompletedWithWrappedCFException(h2);
checkCompletedWithWrappedCFException(h3);
checkCompletedWithWrappedException(h1, r1.ex);
checkCompletedWithWrappedException(h2, r2.ex);
checkCompletedWithWrappedException(h3, r3.ex);
r1.assertInvoked();
r2.assertInvoked();
r3.assertInvoked();
@ -1940,9 +1981,9 @@ public class CompletableFutureTest extends JSR166TestCase {
assertTrue(snd.complete(w2));
final CompletableFuture<Void> h3 = m.thenAcceptBoth(f, g, r3);
checkCompletedWithWrappedCFException(h1);
checkCompletedWithWrappedCFException(h2);
checkCompletedWithWrappedCFException(h3);
checkCompletedWithWrappedException(h1, r1.ex);
checkCompletedWithWrappedException(h2, r2.ex);
checkCompletedWithWrappedException(h3, r3.ex);
r1.assertInvoked();
r2.assertInvoked();
r3.assertInvoked();
@ -2104,9 +2145,9 @@ public class CompletableFutureTest extends JSR166TestCase {
assertTrue(snd.complete(w2));
final CompletableFuture<Void> h3 = m.runAfterBoth(f, g, r3);
checkCompletedWithWrappedCFException(h1);
checkCompletedWithWrappedCFException(h2);
checkCompletedWithWrappedCFException(h3);
checkCompletedWithWrappedException(h1, r1.ex);
checkCompletedWithWrappedException(h2, r2.ex);
checkCompletedWithWrappedException(h3, r3.ex);
r1.assertInvoked();
r2.assertInvoked();
r3.assertInvoked();
@ -2396,10 +2437,10 @@ public class CompletableFutureTest extends JSR166TestCase {
f.complete(v1);
final CompletableFuture<Integer> h2 = m.applyToEither(f, g, rs[2]);
final CompletableFuture<Integer> h3 = m.applyToEither(g, f, rs[3]);
checkCompletedWithWrappedCFException(h0);
checkCompletedWithWrappedCFException(h1);
checkCompletedWithWrappedCFException(h2);
checkCompletedWithWrappedCFException(h3);
checkCompletedWithWrappedException(h0, rs[0].ex);
checkCompletedWithWrappedException(h1, rs[1].ex);
checkCompletedWithWrappedException(h2, rs[2].ex);
checkCompletedWithWrappedException(h3, rs[3].ex);
for (int i = 0; i < 4; i++) rs[i].assertValue(v1);
g.complete(v2);
@ -2408,10 +2449,10 @@ public class CompletableFutureTest extends JSR166TestCase {
final CompletableFuture<Integer> h4 = m.applyToEither(f, g, rs[4]);
final CompletableFuture<Integer> h5 = m.applyToEither(g, f, rs[5]);
checkCompletedWithWrappedCFException(h4);
checkCompletedWithWrappedException(h4, rs[4].ex);
assertTrue(Objects.equals(v1, rs[4].value) ||
Objects.equals(v2, rs[4].value));
checkCompletedWithWrappedCFException(h5);
checkCompletedWithWrappedException(h5, rs[5].ex);
assertTrue(Objects.equals(v1, rs[5].value) ||
Objects.equals(v2, rs[5].value));
@ -2655,10 +2696,10 @@ public class CompletableFutureTest extends JSR166TestCase {
f.complete(v1);
final CompletableFuture<Void> h2 = m.acceptEither(f, g, rs[2]);
final CompletableFuture<Void> h3 = m.acceptEither(g, f, rs[3]);
checkCompletedWithWrappedCFException(h0);
checkCompletedWithWrappedCFException(h1);
checkCompletedWithWrappedCFException(h2);
checkCompletedWithWrappedCFException(h3);
checkCompletedWithWrappedException(h0, rs[0].ex);
checkCompletedWithWrappedException(h1, rs[1].ex);
checkCompletedWithWrappedException(h2, rs[2].ex);
checkCompletedWithWrappedException(h3, rs[3].ex);
for (int i = 0; i < 4; i++) rs[i].assertValue(v1);
g.complete(v2);
@ -2667,10 +2708,10 @@ public class CompletableFutureTest extends JSR166TestCase {
final CompletableFuture<Void> h4 = m.acceptEither(f, g, rs[4]);
final CompletableFuture<Void> h5 = m.acceptEither(g, f, rs[5]);
checkCompletedWithWrappedCFException(h4);
checkCompletedWithWrappedException(h4, rs[4].ex);
assertTrue(Objects.equals(v1, rs[4].value) ||
Objects.equals(v2, rs[4].value));
checkCompletedWithWrappedCFException(h5);
checkCompletedWithWrappedException(h5, rs[5].ex);
assertTrue(Objects.equals(v1, rs[5].value) ||
Objects.equals(v2, rs[5].value));
@ -2686,6 +2727,7 @@ public class CompletableFutureTest extends JSR166TestCase {
for (ExecutionMode m : ExecutionMode.values())
for (Integer v1 : new Integer[] { 1, null })
for (Integer v2 : new Integer[] { 2, null })
for (boolean pushNop : new boolean[] { true, false })
{
final CompletableFuture<Integer> f = new CompletableFuture<>();
final CompletableFuture<Integer> g = new CompletableFuture<>();
@ -2698,6 +2740,10 @@ public class CompletableFutureTest extends JSR166TestCase {
checkIncomplete(h1);
rs[0].assertNotInvoked();
rs[1].assertNotInvoked();
if (pushNop) { // ad hoc test of intra-completion interference
m.thenRun(f, () -> {});
m.thenRun(g, () -> {});
}
f.complete(v1);
checkCompletedNormally(h0, null);
checkCompletedNormally(h1, null);
@ -2910,16 +2956,16 @@ public class CompletableFutureTest extends JSR166TestCase {
assertTrue(f.complete(v1));
final CompletableFuture<Void> h2 = m.runAfterEither(f, g, rs[2]);
final CompletableFuture<Void> h3 = m.runAfterEither(g, f, rs[3]);
checkCompletedWithWrappedCFException(h0);
checkCompletedWithWrappedCFException(h1);
checkCompletedWithWrappedCFException(h2);
checkCompletedWithWrappedCFException(h3);
checkCompletedWithWrappedException(h0, rs[0].ex);
checkCompletedWithWrappedException(h1, rs[1].ex);
checkCompletedWithWrappedException(h2, rs[2].ex);
checkCompletedWithWrappedException(h3, rs[3].ex);
for (int i = 0; i < 4; i++) rs[i].assertInvoked();
assertTrue(g.complete(v2));
final CompletableFuture<Void> h4 = m.runAfterEither(f, g, rs[4]);
final CompletableFuture<Void> h5 = m.runAfterEither(g, f, rs[5]);
checkCompletedWithWrappedCFException(h4);
checkCompletedWithWrappedCFException(h5);
checkCompletedWithWrappedException(h4, rs[4].ex);
checkCompletedWithWrappedException(h5, rs[5].ex);
checkCompletedNormally(f, v1);
checkCompletedNormally(g, v2);
@ -2980,7 +3026,7 @@ public class CompletableFutureTest extends JSR166TestCase {
final CompletableFuture<Integer> g = m.thenCompose(f, r);
if (createIncomplete) assertTrue(f.complete(v1));
checkCompletedWithWrappedCFException(g);
checkCompletedWithWrappedException(g, r.ex);
checkCompletedNormally(f, v1);
}}
@ -3089,7 +3135,7 @@ public class CompletableFutureTest extends JSR166TestCase {
}
}
public void testAllOf_backwards() throws Exception {
public void testAllOf_normal_backwards() throws Exception {
for (int k = 1; k < 10; k++) {
CompletableFuture<Integer>[] fs
= (CompletableFuture<Integer>[]) new CompletableFuture[k];
@ -3336,6 +3382,151 @@ public class CompletableFutureTest extends JSR166TestCase {
assertEquals(0, exec.count.get());
}
/**
* Test submissions to an executor that rejects all tasks.
*/
public void testRejectingExecutor() {
for (Integer v : new Integer[] { 1, null })
{
final CountingRejectingExecutor e = new CountingRejectingExecutor();
final CompletableFuture<Integer> complete = CompletableFuture.completedFuture(v);
final CompletableFuture<Integer> incomplete = new CompletableFuture<>();
List<CompletableFuture<?>> futures = new ArrayList<>();
List<CompletableFuture<Integer>> srcs = new ArrayList<>();
srcs.add(complete);
srcs.add(incomplete);
for (CompletableFuture<Integer> src : srcs) {
List<CompletableFuture<?>> fs = new ArrayList<>();
fs.add(src.thenRunAsync(() -> {}, e));
fs.add(src.thenAcceptAsync((z) -> {}, e));
fs.add(src.thenApplyAsync((z) -> z, e));
fs.add(src.thenCombineAsync(src, (x, y) -> x, e));
fs.add(src.thenAcceptBothAsync(src, (x, y) -> {}, e));
fs.add(src.runAfterBothAsync(src, () -> {}, e));
fs.add(src.applyToEitherAsync(src, (z) -> z, e));
fs.add(src.acceptEitherAsync(src, (z) -> {}, e));
fs.add(src.runAfterEitherAsync(src, () -> {}, e));
fs.add(src.thenComposeAsync((z) -> null, e));
fs.add(src.whenCompleteAsync((z, t) -> {}, e));
fs.add(src.handleAsync((z, t) -> null, e));
for (CompletableFuture<?> future : fs) {
if (src.isDone())
checkCompletedWithWrappedException(future, e.ex);
else
checkIncomplete(future);
}
futures.addAll(fs);
}
{
List<CompletableFuture<?>> fs = new ArrayList<>();
fs.add(complete.thenCombineAsync(incomplete, (x, y) -> x, e));
fs.add(incomplete.thenCombineAsync(complete, (x, y) -> x, e));
fs.add(complete.thenAcceptBothAsync(incomplete, (x, y) -> {}, e));
fs.add(incomplete.thenAcceptBothAsync(complete, (x, y) -> {}, e));
fs.add(complete.runAfterBothAsync(incomplete, () -> {}, e));
fs.add(incomplete.runAfterBothAsync(complete, () -> {}, e));
for (CompletableFuture<?> future : fs)
checkIncomplete(future);
futures.addAll(fs);
}
{
List<CompletableFuture<?>> fs = new ArrayList<>();
fs.add(complete.applyToEitherAsync(incomplete, (z) -> z, e));
fs.add(incomplete.applyToEitherAsync(complete, (z) -> z, e));
fs.add(complete.acceptEitherAsync(incomplete, (z) -> {}, e));
fs.add(incomplete.acceptEitherAsync(complete, (z) -> {}, e));
fs.add(complete.runAfterEitherAsync(incomplete, () -> {}, e));
fs.add(incomplete.runAfterEitherAsync(complete, () -> {}, e));
for (CompletableFuture<?> future : fs)
checkCompletedWithWrappedException(future, e.ex);
futures.addAll(fs);
}
incomplete.complete(v);
for (CompletableFuture<?> future : futures)
checkCompletedWithWrappedException(future, e.ex);
assertEquals(futures.size(), e.count.get());
}}
/**
* Test submissions to an executor that rejects all tasks, but
* should never be invoked because the dependent future is
* explicitly completed.
*/
public void testRejectingExecutorNeverInvoked() {
for (Integer v : new Integer[] { 1, null })
{
final CountingRejectingExecutor e = new CountingRejectingExecutor();
final CompletableFuture<Integer> complete = CompletableFuture.completedFuture(v);
final CompletableFuture<Integer> incomplete = new CompletableFuture<>();
List<CompletableFuture<?>> futures = new ArrayList<>();
List<CompletableFuture<Integer>> srcs = new ArrayList<>();
srcs.add(complete);
srcs.add(incomplete);
List<CompletableFuture<?>> fs = new ArrayList<>();
fs.add(incomplete.thenRunAsync(() -> {}, e));
fs.add(incomplete.thenAcceptAsync((z) -> {}, e));
fs.add(incomplete.thenApplyAsync((z) -> z, e));
fs.add(incomplete.thenCombineAsync(incomplete, (x, y) -> x, e));
fs.add(incomplete.thenAcceptBothAsync(incomplete, (x, y) -> {}, e));
fs.add(incomplete.runAfterBothAsync(incomplete, () -> {}, e));
fs.add(incomplete.applyToEitherAsync(incomplete, (z) -> z, e));
fs.add(incomplete.acceptEitherAsync(incomplete, (z) -> {}, e));
fs.add(incomplete.runAfterEitherAsync(incomplete, () -> {}, e));
fs.add(incomplete.thenComposeAsync((z) -> null, e));
fs.add(incomplete.whenCompleteAsync((z, t) -> {}, e));
fs.add(incomplete.handleAsync((z, t) -> null, e));
fs.add(complete.thenCombineAsync(incomplete, (x, y) -> x, e));
fs.add(incomplete.thenCombineAsync(complete, (x, y) -> x, e));
fs.add(complete.thenAcceptBothAsync(incomplete, (x, y) -> {}, e));
fs.add(incomplete.thenAcceptBothAsync(complete, (x, y) -> {}, e));
fs.add(complete.runAfterBothAsync(incomplete, () -> {}, e));
fs.add(incomplete.runAfterBothAsync(complete, () -> {}, e));
for (CompletableFuture<?> future : fs)
checkIncomplete(future);
for (CompletableFuture<?> future : fs)
future.complete(null);
incomplete.complete(v);
for (CompletableFuture<?> future : fs)
checkCompletedNormally(future, null);
assertEquals(0, e.count.get());
}}
/**
* toCompletableFuture returns this CompletableFuture.
*/
@ -3659,12 +3850,25 @@ public class CompletableFutureTest extends JSR166TestCase {
//--- tests of implementation details; not part of official tck ---
Object resultOf(CompletableFuture<?> f) {
SecurityManager sm = System.getSecurityManager();
if (sm != null) {
try {
System.setSecurityManager(null);
} catch (SecurityException giveUp) {
return "Reflection not available";
}
}
try {
java.lang.reflect.Field resultField
= CompletableFuture.class.getDeclaredField("result");
resultField.setAccessible(true);
return resultField.get(f);
} catch (Throwable t) { throw new AssertionError(t); }
} catch (Throwable t) {
throw new AssertionError(t);
} finally {
if (sm != null) System.setSecurityManager(sm);
}
}
public void testExceptionPropagationReusesResultObject() {
@ -3675,33 +3879,44 @@ public class CompletableFutureTest extends JSR166TestCase {
final CompletableFuture<Integer> v42 = CompletableFuture.completedFuture(42);
final CompletableFuture<Integer> incomplete = new CompletableFuture<>();
final Runnable noopRunnable = new Noop(m);
final Consumer<Integer> noopConsumer = new NoopConsumer(m);
final Function<Integer, Integer> incFunction = new IncFunction(m);
List<Function<CompletableFuture<Integer>, CompletableFuture<?>>> funs
= new ArrayList<>();
funs.add((y) -> m.thenRun(y, new Noop(m)));
funs.add((y) -> m.thenAccept(y, new NoopConsumer(m)));
funs.add((y) -> m.thenApply(y, new IncFunction(m)));
funs.add((y) -> m.thenRun(y, noopRunnable));
funs.add((y) -> m.thenAccept(y, noopConsumer));
funs.add((y) -> m.thenApply(y, incFunction));
funs.add((y) -> m.runAfterEither(y, incomplete, new Noop(m)));
funs.add((y) -> m.acceptEither(y, incomplete, new NoopConsumer(m)));
funs.add((y) -> m.applyToEither(y, incomplete, new IncFunction(m)));
funs.add((y) -> m.runAfterEither(y, incomplete, noopRunnable));
funs.add((y) -> m.acceptEither(y, incomplete, noopConsumer));
funs.add((y) -> m.applyToEither(y, incomplete, incFunction));
funs.add((y) -> m.runAfterBoth(y, v42, new Noop(m)));
funs.add((y) -> m.runAfterBoth(y, v42, noopRunnable));
funs.add((y) -> m.runAfterBoth(v42, y, noopRunnable));
funs.add((y) -> m.thenAcceptBoth(y, v42, new SubtractAction(m)));
funs.add((y) -> m.thenAcceptBoth(v42, y, new SubtractAction(m)));
funs.add((y) -> m.thenCombine(y, v42, new SubtractFunction(m)));
funs.add((y) -> m.thenCombine(v42, y, new SubtractFunction(m)));
funs.add((y) -> m.whenComplete(y, (Integer r, Throwable t) -> {}));
funs.add((y) -> m.thenCompose(y, new CompletableFutureInc(m)));
funs.add((y) -> CompletableFuture.allOf(new CompletableFuture<?>[] {y, v42}));
funs.add((y) -> CompletableFuture.anyOf(new CompletableFuture<?>[] {y, incomplete}));
funs.add((y) -> CompletableFuture.allOf(y));
funs.add((y) -> CompletableFuture.allOf(y, v42));
funs.add((y) -> CompletableFuture.allOf(v42, y));
funs.add((y) -> CompletableFuture.anyOf(y));
funs.add((y) -> CompletableFuture.anyOf(y, incomplete));
funs.add((y) -> CompletableFuture.anyOf(incomplete, y));
for (Function<CompletableFuture<Integer>, CompletableFuture<?>>
fun : funs) {
CompletableFuture<Integer> f = new CompletableFuture<>();
f.completeExceptionally(ex);
CompletableFuture<Integer> src = m.thenApply(f, new IncFunction(m));
CompletableFuture<Integer> src = m.thenApply(f, incFunction);
checkCompletedWithWrappedException(src, ex);
CompletableFuture<?> dep = fun.apply(src);
checkCompletedWithWrappedException(dep, ex);
@ -3711,7 +3926,7 @@ public class CompletableFutureTest extends JSR166TestCase {
for (Function<CompletableFuture<Integer>, CompletableFuture<?>>
fun : funs) {
CompletableFuture<Integer> f = new CompletableFuture<>();
CompletableFuture<Integer> src = m.thenApply(f, new IncFunction(m));
CompletableFuture<Integer> src = m.thenApply(f, incFunction);
CompletableFuture<?> dep = fun.apply(src);
f.completeExceptionally(ex);
checkCompletedWithWrappedException(src, ex);
@ -3725,7 +3940,7 @@ public class CompletableFutureTest extends JSR166TestCase {
CompletableFuture<Integer> f = new CompletableFuture<>();
f.cancel(mayInterruptIfRunning);
checkCancelled(f);
CompletableFuture<Integer> src = m.thenApply(f, new IncFunction(m));
CompletableFuture<Integer> src = m.thenApply(f, incFunction);
checkCompletedWithWrappedCancellationException(src);
CompletableFuture<?> dep = fun.apply(src);
checkCompletedWithWrappedCancellationException(dep);
@ -3736,7 +3951,7 @@ public class CompletableFutureTest extends JSR166TestCase {
for (Function<CompletableFuture<Integer>, CompletableFuture<?>>
fun : funs) {
CompletableFuture<Integer> f = new CompletableFuture<>();
CompletableFuture<Integer> src = m.thenApply(f, new IncFunction(m));
CompletableFuture<Integer> src = m.thenApply(f, incFunction);
CompletableFuture<?> dep = fun.apply(src);
f.cancel(mayInterruptIfRunning);
checkCancelled(f);
@ -3747,7 +3962,7 @@ public class CompletableFutureTest extends JSR166TestCase {
}}
/**
* Minimal completion stages throw UOE for all non-CompletionStage methods
* Minimal completion stages throw UOE for most non-CompletionStage methods
*/
public void testMinimalCompletionStage_minimality() {
if (!testImplementationDetails) return;
@ -3776,8 +3991,10 @@ public class CompletableFutureTest extends JSR166TestCase {
.filter((method) -> !permittedMethodSignatures.contains(toSignature.apply(method)))
.collect(Collectors.toList());
CompletionStage<Integer> minimalStage =
new CompletableFuture<Integer>().minimalCompletionStage();
List<CompletionStage<Integer>> stages = new ArrayList<>();
stages.add(new CompletableFuture<Integer>().minimalCompletionStage());
stages.add(CompletableFuture.completedStage(1));
stages.add(CompletableFuture.failedStage(new CFException()));
List<Method> bugs = new ArrayList<>();
for (Method method : allMethods) {
@ -3793,20 +4010,22 @@ public class CompletableFutureTest extends JSR166TestCase {
else if (parameterTypes[i] == long.class)
args[i] = 0L;
}
try {
method.invoke(minimalStage, args);
bugs.add(method);
}
catch (java.lang.reflect.InvocationTargetException expected) {
if (! (expected.getCause() instanceof UnsupportedOperationException)) {
for (CompletionStage<Integer> stage : stages) {
try {
method.invoke(stage, args);
bugs.add(method);
// expected.getCause().printStackTrace();
}
catch (java.lang.reflect.InvocationTargetException expected) {
if (! (expected.getCause() instanceof UnsupportedOperationException)) {
bugs.add(method);
// expected.getCause().printStackTrace();
}
}
catch (ReflectiveOperationException bad) { throw new Error(bad); }
}
catch (ReflectiveOperationException bad) { throw new Error(bad); }
}
if (!bugs.isEmpty())
throw new Error("Methods did not throw UOE: " + bugs.toString());
throw new Error("Methods did not throw UOE: " + bugs);
}
static class Monad {
@ -3955,12 +4174,33 @@ public class CompletableFutureTest extends JSR166TestCase {
Monad.plus(godot, Monad.unit(5L)));
}
/** Test long recursive chains of CompletableFutures with cascading completions */
public void testRecursiveChains() throws Throwable {
for (ExecutionMode m : ExecutionMode.values())
for (boolean addDeadEnds : new boolean[] { true, false })
{
final int val = 42;
final int n = expensiveTests ? 1_000 : 2;
CompletableFuture<Integer> head = new CompletableFuture<>();
CompletableFuture<Integer> tail = head;
for (int i = 0; i < n; i++) {
if (addDeadEnds) m.thenApply(tail, v -> v + 1);
tail = m.thenApply(tail, v -> v + 1);
if (addDeadEnds) m.applyToEither(tail, tail, v -> v + 1);
tail = m.applyToEither(tail, tail, v -> v + 1);
if (addDeadEnds) m.thenCombine(tail, tail, (v, w) -> v + 1);
tail = m.thenCombine(tail, tail, (v, w) -> v + 1);
}
head.complete(val);
assertEquals(val + 3 * n, (int) tail.join());
}}
/**
* A single CompletableFuture with many dependents.
* A demo of scalability - runtime is O(n).
*/
public void testManyDependents() throws Throwable {
final int n = 1_000;
final int n = expensiveTests ? 1_000_000 : 10;
final CompletableFuture<Void> head = new CompletableFuture<>();
final CompletableFuture<Void> complete = CompletableFuture.completedFuture((Void)null);
final AtomicInteger count = new AtomicInteger(0);
@ -3987,6 +4227,78 @@ public class CompletableFutureTest extends JSR166TestCase {
assertEquals(5 * 3 * n, count.get());
}
/** ant -Dvmoptions=-Xmx8m -Djsr166.expensiveTests=true -Djsr166.tckTestClass=CompletableFutureTest tck */
public void testCoCompletionGarbageRetention() throws Throwable {
final int n = expensiveTests ? 1_000_000 : 10;
final CompletableFuture<Integer> incomplete = new CompletableFuture<>();
CompletableFuture<Integer> f;
for (int i = 0; i < n; i++) {
f = new CompletableFuture<>();
f.runAfterEither(incomplete, () -> {});
f.complete(null);
f = new CompletableFuture<>();
f.acceptEither(incomplete, (x) -> {});
f.complete(null);
f = new CompletableFuture<>();
f.applyToEither(incomplete, (x) -> x);
f.complete(null);
f = new CompletableFuture<>();
CompletableFuture.anyOf(new CompletableFuture<?>[] { f, incomplete });
f.complete(null);
}
for (int i = 0; i < n; i++) {
f = new CompletableFuture<>();
incomplete.runAfterEither(f, () -> {});
f.complete(null);
f = new CompletableFuture<>();
incomplete.acceptEither(f, (x) -> {});
f.complete(null);
f = new CompletableFuture<>();
incomplete.applyToEither(f, (x) -> x);
f.complete(null);
f = new CompletableFuture<>();
CompletableFuture.anyOf(new CompletableFuture<?>[] { incomplete, f });
f.complete(null);
}
}
/*
* Tests below currently fail in stress mode due to memory retention.
* ant -Dvmoptions=-Xmx8m -Djsr166.expensiveTests=true -Djsr166.tckTestClass=CompletableFutureTest tck
*/
/** Checks for garbage retention with anyOf. */
public void testAnyOfGarbageRetention() throws Throwable {
for (Integer v : new Integer[] { 1, null })
{
final int n = expensiveTests ? 100_000 : 10;
CompletableFuture<Integer>[] fs
= (CompletableFuture<Integer>[]) new CompletableFuture<?>[100];
for (int i = 0; i < fs.length; i++)
fs[i] = new CompletableFuture<>();
fs[fs.length - 1].complete(v);
for (int i = 0; i < n; i++)
checkCompletedNormally(CompletableFuture.anyOf(fs), v);
}}
/** Checks for garbage retention with allOf. */
public void testCancelledAllOfGarbageRetention() throws Throwable {
final int n = expensiveTests ? 100_000 : 10;
CompletableFuture<Integer>[] fs
= (CompletableFuture<Integer>[]) new CompletableFuture<?>[100];
for (int i = 0; i < fs.length; i++)
fs[i] = new CompletableFuture<>();
for (int i = 0; i < n; i++)
assertTrue(CompletableFuture.allOf(fs).cancel(false));
}
// static <U> U join(CompletionStage<U> stage) {
// CompletableFuture<U> f = new CompletableFuture<>();
// stage.whenComplete((v, ex) -> {