summaryrefslogtreecommitdiff
path: root/src/main/java
diff options
context:
space:
mode:
Diffstat (limited to 'src/main/java')
-rw-r--r--src/main/java/com/it_jaros/jns/scan/Scanner.java2
-rw-r--r--src/main/java/com/it_jaros/jns/scan/engine/CompletionConsumer.java (renamed from src/main/java/com/it_jaros/jns/scan/engine/ConsumerThread.java)7
-rw-r--r--src/main/java/com/it_jaros/jns/scan/engine/ScanHostTask.java2
3 files changed, 6 insertions, 5 deletions
diff --git a/src/main/java/com/it_jaros/jns/scan/Scanner.java b/src/main/java/com/it_jaros/jns/scan/Scanner.java
index 12d882a..82707d6 100644
--- a/src/main/java/com/it_jaros/jns/scan/Scanner.java
+++ b/src/main/java/com/it_jaros/jns/scan/Scanner.java
@@ -46,7 +46,7 @@ public class Scanner implements AutoCloseable {
);
try {
- new ConsumerThread<>(cancelledToken, state).runConsumerLoop(poll -> {
+ new CompletionConsumer<>(cancelledToken, state).runConsumerLoop(poll -> {
switch (poll) {
case PollState.Success<ScanResult>(ScanResult value) -> consumer.accept(value);
case PollState.Failure<ScanResult>(Throwable error) -> mapToScanResult(error).ifPresent(consumer);
diff --git a/src/main/java/com/it_jaros/jns/scan/engine/ConsumerThread.java b/src/main/java/com/it_jaros/jns/scan/engine/CompletionConsumer.java
index 8f93317..2dfedf5 100644
--- a/src/main/java/com/it_jaros/jns/scan/engine/ConsumerThread.java
+++ b/src/main/java/com/it_jaros/jns/scan/engine/CompletionConsumer.java
@@ -1,16 +1,17 @@
package com.it_jaros.jns.scan.engine;
+import java.util.concurrent.CancellationException;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.function.Consumer;
-public class ConsumerThread<T> {
+public class CompletionConsumer<T> {
private final CancelledToken cancelledToken;
private final ProducerState<T> state;
- public ConsumerThread(CancelledToken cancelledToken, ProducerState<T> state) {
+ public CompletionConsumer(CancelledToken cancelledToken, ProducerState<T> state) {
this.cancelledToken = cancelledToken;
this.state = state;
}
@@ -49,7 +50,7 @@ public class ConsumerThread<T> {
try {
T result = future.get();
return new PollState.Success<>(result);
- } catch (ExecutionException e) {
+ } catch (ExecutionException | CancellationException e) {
Throwable cause = e.getCause();
return new PollState.Failure<>(cause);
} finally {
diff --git a/src/main/java/com/it_jaros/jns/scan/engine/ScanHostTask.java b/src/main/java/com/it_jaros/jns/scan/engine/ScanHostTask.java
index f90d4ce..7fee311 100644
--- a/src/main/java/com/it_jaros/jns/scan/engine/ScanHostTask.java
+++ b/src/main/java/com/it_jaros/jns/scan/engine/ScanHostTask.java
@@ -71,7 +71,7 @@ public class ScanHostTask implements Callable<ScanResult> {
port -> new ScanPortTask(port, context, hostContext)
);
- new ConsumerThread<>(cancelledToken, state).runConsumerLoop(poll -> {
+ new CompletionConsumer<>(cancelledToken, state).runConsumerLoop(poll -> {
switch (poll) {
case PollState.Success<PortResult>(PortResult value) -> accumulator.add(value);
case PollState.Failure(Throwable error) -> accumulator.recordFailure(error, -1);