diff options
3 files changed, 18 insertions, 13 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 32f1182..6608ee7 100644 --- a/src/main/java/com/it_jaros/jns/scan/Scanner.java +++ b/src/main/java/com/it_jaros/jns/scan/Scanner.java @@ -63,10 +63,12 @@ public class Scanner implements AutoCloseable { private void runConsumerLoop(ProducerState<ScanResult> state, Consumer<ScanResult> consumer) { while (state.running().get() || state.inPipeline().get() > 0) { PollState<ScanResult> poll = getPollStateOfScanResult(state); - if (poll instanceof PollState.Success<ScanResult>(ScanResult value)) { - consumer.accept(value); - } else if (poll instanceof PollState.Failure<ScanResult>(Throwable error)) { - mapToScanResult(error).ifPresent(consumer); + switch (poll) { + case PollState.Success<ScanResult>(ScanResult value) -> consumer.accept(value); + case PollState.Failure<ScanResult>(Throwable error) -> mapToScanResult(error).ifPresent(consumer); + case PollState.Unavailable<ScanResult> ignore -> { + // we do nothing + } } } } diff --git a/src/main/java/com/it_jaros/jns/scan/engine/ProducerThread.java b/src/main/java/com/it_jaros/jns/scan/engine/ProducerThread.java index f02e379..059c77b 100644 --- a/src/main/java/com/it_jaros/jns/scan/engine/ProducerThread.java +++ b/src/main/java/com/it_jaros/jns/scan/engine/ProducerThread.java @@ -36,9 +36,9 @@ public class ProducerThread { * @return */ public <INPUT, OUTPUT> ProducerState<OUTPUT> startProducer( - Iterator<INPUT> queue, - int maxWorkers, - Function<INPUT, Callable<OUTPUT>> taskFactory + final Iterator<INPUT> queue, + final int maxWorkers, + final Function<INPUT, Callable<OUTPUT>> taskFactory ) { final AtomicInteger inPipeline = new AtomicInteger(0); final AtomicBoolean running = new AtomicBoolean(true); @@ -64,13 +64,13 @@ public class ProducerThread { return new ProducerState<>(running, inPipeline, completionService); } - private <INPUT, OUTPUT> void runProducerLoop(ProducerContext<INPUT, OUTPUT> context) throws InterruptedException { + private <INPUT, OUTPUT> void runProducerLoop(final ProducerContext<INPUT, OUTPUT> context) throws InterruptedException { while (!cancelledToken.isCancelled() && context.queue().hasNext()) { submitNext(context); } } - private <INPUT, OUTPUT> void submitNext(ProducerContext<INPUT, OUTPUT> context) throws InterruptedException { + private <INPUT, OUTPUT> void submitNext(final ProducerContext<INPUT, OUTPUT> context) throws InterruptedException { context.activeWorkers().acquire(); // just in case something 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 02834ff..f99a248 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 @@ -100,11 +100,14 @@ public class ScanHostTask implements Callable<ScanResult> { private void runConsumerLoop(ProducerState<PortResult> state) { while (state.running().get() || state.inPipeline().get() > 0) { PollState<PortResult> poll = getPortResult(state); - if (poll instanceof PollState.Success<PortResult>(PortResult value)) { - accumulator.add(value); - } else if (poll instanceof PollState.Failure(Throwable error)) { - accumulator.recordFailure(error, -1); + switch (poll) { + case PollState.Success<PortResult>(PortResult value) -> accumulator.add(value); + case PollState.Failure(Throwable error) -> accumulator.recordFailure(error, -1); + case PollState.Unavailable<PortResult> ignore -> { + // we do nothing + } } + } } |
