Skip to content
Open
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 @@ -112,10 +112,15 @@ public void waitAllComplete(List<Pair<PException, Result>> results, int timeoutM
if (fu.isSuccess()) {
results.add(Pair.of(null, fu.getNow()));
} else {
// fu.cause() is null when the future never completed (e.g. await timed out),
// in which case dereferencing it throws a NullPointerException. Guard against
// that so a failing-but-incomplete task is reported instead of crashing the
// caller. See https://github.com/apache/incubator-pegasus/issues/2152.
Throwable cause = fu.cause();
String causeMsg =
cause != null ? cause.getMessage() : "unknown cause (future not completed)";
results.add(
Pair.of(
new PException("async task #[" + i + "] await failed: " + fu.cause().getMessage()),
null));
Pair.of(new PException("async task #[" + i + "] await failed: " + causeMsg), null));
}
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,13 +18,23 @@
*/
package org.apache.pegasus.client;

import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.fail;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;

import io.netty.util.concurrent.Future;
import io.netty.util.concurrent.Promise;
import io.netty.util.concurrent.SingleThreadEventExecutor;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.Executors;
import java.util.concurrent.atomic.AtomicBoolean;
import org.apache.commons.lang3.tuple.Pair;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.TestInfo;

Expand Down Expand Up @@ -100,4 +110,28 @@ public void testFutureWaitTimeout(TestInfo testInfo) throws Exception {
}
fail();
}

@Test
public void testWaitAllCompleteHandlesNullCause(TestInfo testInfo) throws Exception {
// Reproduces #2152: a Future that is neither completed nor successful (e.g. its
// await timed out before it ever resolved) reports isSuccess()==false and
// cause()==null. waitAllComplete must not dereference the null cause.
@SuppressWarnings("unchecked")
Future<String> unfinished = mock(Future.class);
when(unfinished.await(anyLong())).thenReturn(false); // timed out, still not done
when(unfinished.isSuccess()).thenReturn(false);
when(unfinished.cause()).thenReturn(null);

FutureGroup<String> group = new FutureGroup<>(1);
group.add(unfinished);

List<Pair<PException, String>> results = new ArrayList<>();
// Before the fix this threw NullPointerException inside waitAllComplete.
group.waitAllComplete(results, 5000);

assertEquals(1, results.size());
assertNotNull(results.get(0).getLeft()); // a PException is reported
assertNull(results.get(0).getRight()); // no result
System.err.println(testInfo.getDisplayName() + ": " + results.get(0).getLeft().toString());
}
}
Loading