diff options
| -rw-r--r-- | src/main/java/com/it_jaros/jns/scan/Scanner.java | 50 |
1 files changed, 33 insertions, 17 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 9a62d52..3e84a26 100644 --- a/src/main/java/com/it_jaros/jns/scan/Scanner.java +++ b/src/main/java/com/it_jaros/jns/scan/Scanner.java @@ -45,25 +45,30 @@ public class Scanner implements AutoCloseable { host -> new ScanHostTask(host, context) ); - // the main thread is the consumer - // Let the consumer run as long as the producer runs - // or if still tasks are pending in pipeline - // we do not listen to canceled here because we want - // all results (also partial) collected for the consumer - // with whatever is there already + try { + runConsumerLoop(state, consumer); + } finally { + scan.stop(); + } + } + + /** + * the main thread is the consumer + * Let the consumer run as long as the producer runs + * or if still tasks are pending in pipeline + * we do not listen to canceled here because we want + * all results (also partial) collected for the consumer + * with whatever is there already + */ + private void runConsumerLoop(ProducerState<ScanResult> state, Consumer<ScanResult> consumer) { while (state.running().get() || state.inPipeline().get() > 0) { - try { - PollState<ScanResult> poll = getHostResult(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); - } - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); + 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); } } - scan.stop(); } private Optional<ScanResult> mapToScanResult(Throwable error) { @@ -87,8 +92,19 @@ public class Scanner implements AutoCloseable { return Optional.of(ScanResult.from(error)); } - private PollState<ScanResult> getHostResult(ProducerState<ScanResult> state) throws InterruptedException { + 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<>(); } |
