diff options
| author | 2026-08-22 12:23:10 +0200 | |
|---|---|---|
| committer | 2026-08-22 12:23:10 +0200 | |
| commit | 1aab7dedc9859e80fc0dd0cb1c2057d208244ba3 (patch) | |
| tree | e482706b2f634488c55819add3265ebe86b282f3 /src/main/java | |
| parent | c2853e0a0f02c863336c672db92c1b6ba41a1e3d (diff) | |
Extracted consumer logic into ConsumerThread
Diffstat (limited to 'src/main/java')
3 files changed, 48 insertions, 43 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 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<ScanResult> poll = getPollStateOfScanResult(state); + 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 + } + } + }); } 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 (!cancelledToken.isCancelled() && (state.running().get() || state.inPipeline().get() > 0)) { - PollState<ScanResult> poll = getPollStateOfScanResult(state); - 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 - } - } - } - } - private Optional<ScanResult> 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<T> { + + private final CancelledToken cancelledToken; + private final ProducerState<T> state; + + public ConsumerThread(CancelledToken cancelledToken, ProducerState<T> 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<Void> 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<ScanResult> { port -> new ScanPortTask(port, context, hostContext) ); - runConsumerLoop(state); + new ConsumerThread<>(cancelledToken, state).runConsumerLoop(ignored -> { + PollState<PortResult> poll = getPortResult(state); + 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 + } + } + }); return accumulator.build(); } @@ -91,26 +100,6 @@ public class ScanHostTask implements Callable<ScanResult> { 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<PortResult> state) { - while (!context.cancelledToken().isCancelled() && (state.running().get() || state.inPipeline().get() > 0)) { - PollState<PortResult> poll = getPortResult(state); - 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 - } - } - - } - } - private PollState<PortResult> getPortResult(ProducerState<PortResult> state) { try { return pollPortResult(state); |
