From 1aab7dedc9859e80fc0dd0cb1c2057d208244ba3 Mon Sep 17 00:00:00 2001 From: Matthias Jaros Date: Sat, 22 Aug 2026 12:23:10 +0200 Subject: Extracted consumer logic into ConsumerThread --- src/main/java/com/it_jaros/jns/scan/Scanner.java | 32 +++++++--------------- .../it_jaros/jns/scan/engine/ConsumerThread.java | 28 +++++++++++++++++++ .../com/it_jaros/jns/scan/engine/ScanHostTask.java | 31 +++++++-------------- 3 files changed, 48 insertions(+), 43 deletions(-) create mode 100644 src/main/java/com/it_jaros/jns/scan/engine/ConsumerThread.java 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 d9dfe47..e76d804 100644 --- a/src/main/java/com/it_jaros/jns/scan/Scanner.java +++ b/src/main/java/com/it_jaros/jns/scan/Scanner.java @@ -46,33 +46,21 @@ public class Scanner implements AutoCloseable { ); try { - runConsumerLoop(state, consumer); + new ConsumerThread<>(cancelledToken, state).runConsumerLoop(ignored -> { + PollState poll = getPollStateOfScanResult(state); + switch (poll) { + case PollState.Success(ScanResult value) -> consumer.accept(value); + case PollState.Failure(Throwable error) -> mapToScanResult(error).ifPresent(consumer); + case PollState.Unavailable ignore -> { + // we do nothing + } + } + }); } 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 state, Consumer consumer) { - while (!cancelledToken.isCancelled() && (state.running().get() || state.inPipeline().get() > 0)) { - PollState poll = getPollStateOfScanResult(state); - switch (poll) { - case PollState.Success(ScanResult value) -> consumer.accept(value); - case PollState.Failure(Throwable error) -> mapToScanResult(error).ifPresent(consumer); - case PollState.Unavailable ignore -> { - // we do nothing - } - } - } - } - private Optional mapToScanResult(Throwable error) { if (error instanceof RejectedExecutionException) { // happens when Ctrl-C is pressed 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 new file mode 100644 index 0000000..8cbfcde --- /dev/null +++ b/src/main/java/com/it_jaros/jns/scan/engine/ConsumerThread.java @@ -0,0 +1,28 @@ +package com.it_jaros.jns.scan.engine; + +import java.util.function.Consumer; + +public class ConsumerThread { + + private final CancelledToken cancelledToken; + private final ProducerState state; + + public ConsumerThread(CancelledToken cancelledToken, ProducerState state) { + this.cancelledToken = cancelledToken; + this.state = state; + } + + /** + * 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 + */ + public void runConsumerLoop(Consumer callable) { + while (!cancelledToken.isCancelled() && (state.running().get() || state.inPipeline().get() > 0)) { + callable.accept(null); + } + } +} 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 086a864..87b5f95 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 @@ -73,7 +73,16 @@ public class ScanHostTask implements Callable { port -> new ScanPortTask(port, context, hostContext) ); - runConsumerLoop(state); + new ConsumerThread<>(cancelledToken, state).runConsumerLoop(ignored -> { + PollState poll = getPortResult(state); + switch (poll) { + case PollState.Success(PortResult value) -> accumulator.add(value); + case PollState.Failure(Throwable error) -> accumulator.recordFailure(error, -1); + case PollState.Unavailable ignore -> { + // we do nothing + } + } + }); return accumulator.build(); } @@ -91,26 +100,6 @@ public class ScanHostTask implements Callable { return false; } - /** - * consumer is the main thread - * we run as long as the producer is running OR - * as long as things are in pipeline waiting to be processed - * ONLY exception is when cancelled is set - */ - private void runConsumerLoop(ProducerState state) { - while (!context.cancelledToken().isCancelled() && (state.running().get() || state.inPipeline().get() > 0)) { - PollState poll = getPortResult(state); - switch (poll) { - case PollState.Success(PortResult value) -> accumulator.add(value); - case PollState.Failure(Throwable error) -> accumulator.recordFailure(error, -1); - case PollState.Unavailable ignore -> { - // we do nothing - } - } - - } - } - private PollState getPortResult(ProducerState state) { try { return pollPortResult(state); -- cgit v1.3.1