diff options
Diffstat (limited to 'src/main/java/com')
| -rw-r--r-- | src/main/java/com/it_jaros/jns/scan/Scanner.java | 2 | ||||
| -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.java | 2 |
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); |
