diff options
Diffstat (limited to 'src')
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<ScanResult> poll = getPollStateOfScanResult(state); + new ConsumerThread<>(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); @@ -82,33 +81,6 @@ public class Scanner implements AutoCloseable { return Optional.of(ScanResult.from(error)); } - private PollState<ScanResult> getPollStateOfScanResult(ProducerState<ScanResult> state) { - try { - return pollCompletedHost(state); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - } - - return new PollState.Unavailable<>(); - } - - private static PollState<ScanResult> pollCompletedHost(ProducerState<ScanResult> state) throws InterruptedException { - Future<ScanResult> 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<T> { @@ -20,9 +23,37 @@ public class ConsumerThread<T> { * all results (also partial) collected for the consumer * with whatever is there already */ - public void runConsumerLoop(Consumer<Void> callable) { + public void runConsumerLoop(Consumer<PollState<T>> callable) { while (!cancelledToken.isCancelled() && (state.running().get() || state.inPipeline().get() > 0)) { - callable.accept(null); + callable.accept(getPollResult(state)); + } + } + + private PollState<T> getPollResult(ProducerState<T> state) { + try { + return pollResult(state); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + + return new PollState.Unavailable<>(); + } + + private PollState<T> pollResult(ProducerState<T> state) throws InterruptedException { + Future<T> 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<ScanResult> { @@ -73,8 +71,7 @@ public class ScanHostTask implements Callable<ScanResult> { port -> new ScanPortTask(port, context, hostContext) ); - new ConsumerThread<>(cancelledToken, state).runConsumerLoop(ignored -> { - PollState<PortResult> poll = getPortResult(state); + new ConsumerThread<>(cancelledToken, state).runConsumerLoop(poll -> { switch (poll) { case PollState.Success<PortResult>(PortResult value) -> accumulator.add(value); case PollState.Failure(Throwable error) -> accumulator.recordFailure(error, -1); @@ -99,32 +96,4 @@ public class ScanHostTask implements Callable<ScanResult> { return false; } - - private PollState<PortResult> getPortResult(ProducerState<PortResult> state) { - try { - return pollPortResult(state); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - } - - return new PollState.Unavailable<>(); - } - - private static PollState<PortResult> pollPortResult(ProducerState<PortResult> state) throws InterruptedException { - Future<PortResult> 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(); - } - } } |
