From 87a16079004303998100cd3e5a9d761c305e4ea3 Mon Sep 17 00:00:00 2001 From: Matthias Jaros Date: Sat, 22 Aug 2026 12:30:55 +0200 Subject: Also removed redundancy of poll result for logic by extracting methods from scanner and ScanHostTask into ConsumerThread --- src/main/java/com/it_jaros/jns/scan/Scanner.java | 31 +------------------ .../it_jaros/jns/scan/engine/ConsumerThread.java | 35 ++++++++++++++++++++-- .../com/it_jaros/jns/scan/engine/ScanHostTask.java | 33 +------------------- 3 files changed, 35 insertions(+), 64 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 e76d804..12d882a 100644 --- a/src/main/java/com/it_jaros/jns/scan/Scanner.java +++ b/src/main/java/com/it_jaros/jns/scan/Scanner.java @@ -46,8 +46,7 @@ public class Scanner implements AutoCloseable { ); try { - new ConsumerThread<>(cancelledToken, state).runConsumerLoop(ignored -> { - PollState poll = getPollStateOfScanResult(state); + new ConsumerThread<>(cancelledToken, state).runConsumerLoop(poll -> { switch (poll) { case PollState.Success(ScanResult value) -> consumer.accept(value); case PollState.Failure(Throwable error) -> mapToScanResult(error).ifPresent(consumer); @@ -82,33 +81,6 @@ public class Scanner implements AutoCloseable { return Optional.of(ScanResult.from(error)); } - private PollState getPollStateOfScanResult(ProducerState state) { - try { - return pollCompletedHost(state); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - } - - return new PollState.Unavailable<>(); - } - - private static PollState pollCompletedHost(ProducerState state) throws InterruptedException { - Future finishedHost = state.completionService().poll(ProducerThread.pollInterval.toMillis(), TimeUnit.MILLISECONDS); - - if (finishedHost == null) { - return new PollState.Unavailable<>(); - } - - try { - ScanResult result = finishedHost.get(); - return new PollState.Success<>(result); - } catch (ExecutionException e) { - return new PollState.Failure<>(e.getCause()); - } finally { - state.inPipeline().decrementAndGet(); - } - } - public boolean awaitTermination(Duration duration) throws InterruptedException { return executor.awaitTermination(duration.toMillis(), TimeUnit.MILLISECONDS); } @@ -127,5 +99,4 @@ public class Scanner implements AutoCloseable { public void close() throws Exception { cancel(); } - } diff --git a/src/main/java/com/it_jaros/jns/scan/engine/ConsumerThread.java b/src/main/java/com/it_jaros/jns/scan/engine/ConsumerThread.java index 8cbfcde..8bc7fdf 100644 --- a/src/main/java/com/it_jaros/jns/scan/engine/ConsumerThread.java +++ b/src/main/java/com/it_jaros/jns/scan/engine/ConsumerThread.java @@ -1,5 +1,8 @@ package com.it_jaros.jns.scan.engine; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; import java.util.function.Consumer; public class ConsumerThread { @@ -20,9 +23,37 @@ public class ConsumerThread { * all results (also partial) collected for the consumer * with whatever is there already */ - public void runConsumerLoop(Consumer callable) { + public void runConsumerLoop(Consumer> callable) { while (!cancelledToken.isCancelled() && (state.running().get() || state.inPipeline().get() > 0)) { - callable.accept(null); + callable.accept(getPollResult(state)); + } + } + + private PollState getPollResult(ProducerState state) { + try { + return pollResult(state); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + + return new PollState.Unavailable<>(); + } + + private PollState pollResult(ProducerState state) throws InterruptedException { + Future future = state.completionService().poll(ProducerThread.pollInterval.toMillis(), TimeUnit.MILLISECONDS); + + if (future == null) { + return new PollState.Unavailable<>(); + } + + try { + T result = future.get(); + return new PollState.Success<>(result); + } catch (ExecutionException e) { + Throwable cause = e.getCause(); + return new PollState.Failure<>(cause); + } finally { + state.inPipeline().decrementAndGet(); } } } 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 87b5f95..f90d4ce 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 @@ -8,8 +8,6 @@ import com.it_jaros.jns.scan.domain.ScanResult; import java.net.InetAddress; import java.util.concurrent.Callable; -import java.util.concurrent.ExecutionException; -import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; public class ScanHostTask implements Callable { @@ -73,8 +71,7 @@ public class ScanHostTask implements Callable { port -> new ScanPortTask(port, context, hostContext) ); - new ConsumerThread<>(cancelledToken, state).runConsumerLoop(ignored -> { - PollState poll = getPortResult(state); + new ConsumerThread<>(cancelledToken, state).runConsumerLoop(poll -> { switch (poll) { case PollState.Success(PortResult value) -> accumulator.add(value); case PollState.Failure(Throwable error) -> accumulator.recordFailure(error, -1); @@ -99,32 +96,4 @@ public class ScanHostTask implements Callable { return false; } - - private PollState getPortResult(ProducerState state) { - try { - return pollPortResult(state); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - } - - return new PollState.Unavailable<>(); - } - - private static PollState pollPortResult(ProducerState state) throws InterruptedException { - Future portResultFuture = state.completionService().poll(ProducerThread.pollInterval.toMillis(), TimeUnit.MILLISECONDS); - - if (portResultFuture == null) { - return new PollState.Unavailable<>(); - } - - try { - PortResult portResult = portResultFuture.get(); - return new PollState.Success<>(portResult); - } catch (ExecutionException e) { - Throwable cause = e.getCause(); - return new PollState.Failure<>(cause); - } finally { - state.inPipeline().decrementAndGet(); - } - } } -- cgit v1.3.1