ListenableFuture provides CompletableFuture adaptation via completable()
Issue: SPR-15696
This commit is contained in:
@@ -73,6 +73,12 @@ public class CompletableToListenableFutureAdapter<T> implements ListenableFuture
|
||||
this.callbacks.addFailureCallback(failureCallback);
|
||||
}
|
||||
|
||||
@Override
|
||||
public CompletableFuture<T> completable() {
|
||||
return this.completableFuture;
|
||||
}
|
||||
|
||||
|
||||
@Override
|
||||
public boolean cancel(boolean mayInterruptIfRunning) {
|
||||
return this.completableFuture.cancel(mayInterruptIfRunning);
|
||||
|
||||
@@ -0,0 +1,44 @@
|
||||
/*
|
||||
* Copyright 2002-2017 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
* You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.springframework.util.concurrent;
|
||||
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.Future;
|
||||
|
||||
/**
|
||||
* Extension of {@link CompletableFuture} which allows for cancelling
|
||||
* a delegate along with the {@link CompletableFuture} itself.
|
||||
*
|
||||
* @author Juergen Hoeller
|
||||
* @since 5.0
|
||||
*/
|
||||
class DelegatingCompletableFuture<T> extends CompletableFuture<T> {
|
||||
|
||||
private final Future<T> delegate;
|
||||
|
||||
public DelegatingCompletableFuture(Future<T> delegate) {
|
||||
this.delegate = delegate;
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean cancel(boolean mayInterruptIfRunning) {
|
||||
boolean result = this.delegate.cancel(mayInterruptIfRunning);
|
||||
super.cancel(mayInterruptIfRunning);
|
||||
return result;
|
||||
}
|
||||
|
||||
}
|
||||
@@ -16,6 +16,7 @@
|
||||
|
||||
package org.springframework.util.concurrent;
|
||||
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.Future;
|
||||
|
||||
/**
|
||||
@@ -27,6 +28,7 @@ import java.util.concurrent.Future;
|
||||
*
|
||||
* @author Arjen Poutsma
|
||||
* @author Sebastien Deleuze
|
||||
* @author Juergen Hoeller
|
||||
* @since 4.0
|
||||
*/
|
||||
public interface ListenableFuture<T> extends Future<T> {
|
||||
@@ -45,4 +47,15 @@ public interface ListenableFuture<T> extends Future<T> {
|
||||
*/
|
||||
void addCallback(SuccessCallback<? super T> successCallback, FailureCallback failureCallback);
|
||||
|
||||
|
||||
/**
|
||||
* Expose this {@link ListenableFuture} as a JDK {@link CompletableFuture}.
|
||||
* @since 5.0
|
||||
*/
|
||||
default CompletableFuture<T> completable() {
|
||||
CompletableFuture<T> completable = new DelegatingCompletableFuture<>(this);
|
||||
addCallback(completable::complete, completable::completeExceptionally);
|
||||
return completable;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2016 the original author or authors.
|
||||
* Copyright 2002-2017 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -18,6 +18,8 @@ package org.springframework.util.concurrent;
|
||||
|
||||
import java.util.concurrent.ExecutionException;
|
||||
|
||||
import org.springframework.lang.Nullable;
|
||||
|
||||
/**
|
||||
* Abstract class that adapts a {@link ListenableFuture} parameterized over S into a
|
||||
* {@code ListenableFuture} parameterized over T. All methods are delegated to the
|
||||
@@ -51,19 +53,21 @@ public abstract class ListenableFutureAdapter<T, S> extends FutureAdapter<T, S>
|
||||
ListenableFuture<S> listenableAdaptee = (ListenableFuture<S>) getAdaptee();
|
||||
listenableAdaptee.addCallback(new ListenableFutureCallback<S>() {
|
||||
@Override
|
||||
public void onSuccess(S result) {
|
||||
T adapted;
|
||||
try {
|
||||
adapted = adaptInternal(result);
|
||||
}
|
||||
catch (ExecutionException ex) {
|
||||
Throwable cause = ex.getCause();
|
||||
onFailure(cause != null ? cause : ex);
|
||||
return;
|
||||
}
|
||||
catch (Throwable ex) {
|
||||
onFailure(ex);
|
||||
return;
|
||||
public void onSuccess(@Nullable S result) {
|
||||
T adapted = null;
|
||||
if (result != null) {
|
||||
try {
|
||||
adapted = adaptInternal(result);
|
||||
}
|
||||
catch (ExecutionException ex) {
|
||||
Throwable cause = ex.getCause();
|
||||
onFailure(cause != null ? cause : ex);
|
||||
return;
|
||||
}
|
||||
catch (Throwable ex) {
|
||||
onFailure(ex);
|
||||
return;
|
||||
}
|
||||
}
|
||||
successCallback.onSuccess(adapted);
|
||||
}
|
||||
|
||||
@@ -17,6 +17,7 @@
|
||||
package org.springframework.util.concurrent;
|
||||
|
||||
import java.util.concurrent.Callable;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.ExecutionException;
|
||||
import java.util.concurrent.FutureTask;
|
||||
|
||||
@@ -65,6 +66,14 @@ public class ListenableFutureTask<T> extends FutureTask<T> implements Listenable
|
||||
this.callbacks.addFailureCallback(failureCallback);
|
||||
}
|
||||
|
||||
@Override
|
||||
public CompletableFuture<T> completable() {
|
||||
CompletableFuture<T> completable = new DelegatingCompletableFuture<>(this);
|
||||
this.callbacks.addSuccessCallback(completable::complete);
|
||||
this.callbacks.addFailureCallback(completable::completeExceptionally);
|
||||
return completable;
|
||||
}
|
||||
|
||||
|
||||
@Override
|
||||
protected void done() {
|
||||
|
||||
@@ -17,6 +17,7 @@
|
||||
package org.springframework.util.concurrent;
|
||||
|
||||
import java.util.concurrent.Callable;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.ExecutionException;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.TimeoutException;
|
||||
@@ -68,6 +69,7 @@ public class SettableListenableFuture<T> implements ListenableFuture<T> {
|
||||
return this.settableTask.setExceptionResult(exception);
|
||||
}
|
||||
|
||||
|
||||
@Override
|
||||
public void addCallback(ListenableFutureCallback<? super T> callback) {
|
||||
this.settableTask.addCallback(callback);
|
||||
@@ -78,6 +80,12 @@ public class SettableListenableFuture<T> implements ListenableFuture<T> {
|
||||
this.settableTask.addCallback(successCallback, failureCallback);
|
||||
}
|
||||
|
||||
@Override
|
||||
public CompletableFuture<T> completable() {
|
||||
return this.settableTask.completable();
|
||||
}
|
||||
|
||||
|
||||
@Override
|
||||
public boolean cancel(boolean mayInterruptIfRunning) {
|
||||
boolean cancelled = this.settableTask.cancel(mayInterruptIfRunning);
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2016 the original author or authors.
|
||||
* Copyright 2002-2017 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -18,6 +18,7 @@ package org.springframework.util.concurrent;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.util.concurrent.Callable;
|
||||
import java.util.concurrent.ExecutionException;
|
||||
|
||||
import org.junit.Test;
|
||||
|
||||
@@ -34,12 +35,8 @@ public class ListenableFutureTaskTests {
|
||||
@Test
|
||||
public void success() throws Exception {
|
||||
final String s = "Hello World";
|
||||
Callable<String> callable = new Callable<String>() {
|
||||
@Override
|
||||
public String call() throws Exception {
|
||||
return s;
|
||||
}
|
||||
};
|
||||
Callable<String> callable = () -> s;
|
||||
|
||||
ListenableFutureTask<String> task = new ListenableFutureTask<>(callable);
|
||||
task.addCallback(new ListenableFutureCallback<String>() {
|
||||
@Override
|
||||
@@ -52,17 +49,19 @@ public class ListenableFutureTaskTests {
|
||||
}
|
||||
});
|
||||
task.run();
|
||||
|
||||
assertSame(s, task.get());
|
||||
assertSame(s, task.completable().get());
|
||||
task.completable().thenAccept(v -> assertSame(s, v));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void failure() throws Exception {
|
||||
final String s = "Hello World";
|
||||
Callable<String> callable = new Callable<String>() {
|
||||
@Override
|
||||
public String call() throws Exception {
|
||||
throw new IOException(s);
|
||||
}
|
||||
Callable<String> callable = () -> {
|
||||
throw new IOException(s);
|
||||
};
|
||||
|
||||
ListenableFutureTask<String> task = new ListenableFutureTask<>(callable);
|
||||
task.addCallback(new ListenableFutureCallback<String>() {
|
||||
@Override
|
||||
@@ -75,12 +74,28 @@ public class ListenableFutureTaskTests {
|
||||
}
|
||||
});
|
||||
task.run();
|
||||
|
||||
try {
|
||||
task.get();
|
||||
fail("Should have thrown ExecutionException");
|
||||
}
|
||||
catch (ExecutionException ex) {
|
||||
assertSame(s, ex.getCause().getMessage());
|
||||
}
|
||||
try {
|
||||
task.completable().get();
|
||||
fail("Should have thrown ExecutionException");
|
||||
}
|
||||
catch (ExecutionException ex) {
|
||||
assertSame(s, ex.getCause().getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void successWithLambdas() throws Exception {
|
||||
final String s = "Hello World";
|
||||
Callable<String> callable = () -> s;
|
||||
|
||||
SuccessCallback<String> successCallback = mock(SuccessCallback.class);
|
||||
FailureCallback failureCallback = mock(FailureCallback.class);
|
||||
ListenableFutureTask<String> task = new ListenableFutureTask<>(callable);
|
||||
@@ -88,6 +103,10 @@ public class ListenableFutureTaskTests {
|
||||
task.run();
|
||||
verify(successCallback).onSuccess(s);
|
||||
verifyZeroInteractions(failureCallback);
|
||||
|
||||
assertSame(s, task.get());
|
||||
assertSame(s, task.completable().get());
|
||||
task.completable().thenAccept(v -> assertSame(s, v));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -97,6 +116,7 @@ public class ListenableFutureTaskTests {
|
||||
Callable<String> callable = () -> {
|
||||
throw ex;
|
||||
};
|
||||
|
||||
SuccessCallback<String> successCallback = mock(SuccessCallback.class);
|
||||
FailureCallback failureCallback = mock(FailureCallback.class);
|
||||
ListenableFutureTask<String> task = new ListenableFutureTask<>(callable);
|
||||
@@ -104,6 +124,21 @@ public class ListenableFutureTaskTests {
|
||||
task.run();
|
||||
verify(failureCallback).onFailure(ex);
|
||||
verifyZeroInteractions(successCallback);
|
||||
|
||||
try {
|
||||
task.get();
|
||||
fail("Should have thrown ExecutionException");
|
||||
}
|
||||
catch (ExecutionException ex2) {
|
||||
assertSame(s, ex2.getCause().getMessage());
|
||||
}
|
||||
try {
|
||||
task.completable().get();
|
||||
fail("Should have thrown ExecutionException");
|
||||
}
|
||||
catch (ExecutionException ex2) {
|
||||
assertSame(s, ex2.getCause().getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -18,6 +18,7 @@ package org.springframework.util.concurrent;
|
||||
|
||||
import java.util.concurrent.CancellationException;
|
||||
import java.util.concurrent.ExecutionException;
|
||||
import java.util.concurrent.Future;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.TimeoutException;
|
||||
|
||||
@@ -52,6 +53,16 @@ public class SettableListenableFutureTests {
|
||||
assertTrue(settableListenableFuture.isDone());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void returnsSetValueFromCompletable() throws ExecutionException, InterruptedException {
|
||||
String string = "hello";
|
||||
assertTrue(settableListenableFuture.set(string));
|
||||
Future<String> completable = settableListenableFuture.completable();
|
||||
assertThat(completable.get(), equalTo(string));
|
||||
assertFalse(completable.isCancelled());
|
||||
assertTrue(completable.isDone());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void setValueUpdatesDoneStatus() {
|
||||
settableListenableFuture.set("hello");
|
||||
@@ -60,7 +71,7 @@ public class SettableListenableFutureTests {
|
||||
}
|
||||
|
||||
@Test
|
||||
public void throwsSetExceptionWrappedInExecutionException() throws ExecutionException, InterruptedException {
|
||||
public void throwsSetExceptionWrappedInExecutionException() throws Exception {
|
||||
Throwable exception = new RuntimeException();
|
||||
assertTrue(settableListenableFuture.setException(exception));
|
||||
|
||||
@@ -77,7 +88,25 @@ public class SettableListenableFutureTests {
|
||||
}
|
||||
|
||||
@Test
|
||||
public void throwsSetErrorWrappedInExecutionException() throws ExecutionException, InterruptedException {
|
||||
public void throwsSetExceptionWrappedInExecutionExceptionFromCompletable() throws Exception {
|
||||
Throwable exception = new RuntimeException();
|
||||
assertTrue(settableListenableFuture.setException(exception));
|
||||
Future<String> completable = settableListenableFuture.completable();
|
||||
|
||||
try {
|
||||
completable.get();
|
||||
fail("Expected ExecutionException");
|
||||
}
|
||||
catch (ExecutionException ex) {
|
||||
assertThat(ex.getCause(), equalTo(exception));
|
||||
}
|
||||
|
||||
assertFalse(completable.isCancelled());
|
||||
assertTrue(completable.isDone());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void throwsSetErrorWrappedInExecutionException() throws Exception {
|
||||
Throwable exception = new OutOfMemoryError();
|
||||
assertTrue(settableListenableFuture.setException(exception));
|
||||
|
||||
@@ -93,6 +122,24 @@ public class SettableListenableFutureTests {
|
||||
assertTrue(settableListenableFuture.isDone());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void throwsSetErrorWrappedInExecutionExceptionFromCompletable() throws Exception {
|
||||
Throwable exception = new OutOfMemoryError();
|
||||
assertTrue(settableListenableFuture.setException(exception));
|
||||
Future<String> completable = settableListenableFuture.completable();
|
||||
|
||||
try {
|
||||
completable.get();
|
||||
fail("Expected ExecutionException");
|
||||
}
|
||||
catch (ExecutionException ex) {
|
||||
assertThat(ex.getCause(), equalTo(exception));
|
||||
}
|
||||
|
||||
assertFalse(completable.isCancelled());
|
||||
assertTrue(completable.isDone());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void setValueTriggersCallback() {
|
||||
String string = "hello";
|
||||
|
||||
Reference in New Issue
Block a user