From 56644a0982100f4622dde71bc12bd3228517e466 Mon Sep 17 00:00:00 2001 From: Matthias Jaros Date: Sat, 22 Aug 2026 22:03:29 +0200 Subject: Renamed Class to better match what it is doing --- src/main/java/com/it_jaros/jns/scan/Scanner.java | 2 +- .../jns/scan/engine/CompletionConsumer.java | 60 ++++++++++++++++++++++ .../it_jaros/jns/scan/engine/ConsumerThread.java | 59 --------------------- .../com/it_jaros/jns/scan/engine/ScanHostTask.java | 2 +- 4 files changed, 62 insertions(+), 61 deletions(-) create mode 100644 src/main/java/com/it_jaros/jns/scan/engine/CompletionConsumer.java delete mode 100644 src/main/java/com/it_jaros/jns/scan/engine/ConsumerThread.java (limited to 'src/main/java/com/it_jaros') 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 12d882a..82707d6 100644 --- a/src/main/java/com/it_jaros/jns/scan/Scanner.java +++ b/src/main/java/com/it_jaros/jns/scan/Scanner.java @@ -46,7 +46,7 @@ public class Scanner implements AutoCloseable { ); try { - new ConsumerThread<>(cancelledToken, state).runConsumerLoop(poll -> { + new CompletionConsumer<>(cancelledToken, state).runConsumerLoop(poll -> { switch (poll) { case PollState.Success(ScanResult value) -> consumer.accept(value); case PollState.Failure(Throwable error) -> mapToScanResult(error).ifPresent(consumer); diff --git a/src/main/java/com/it_jaros/jns/scan/engine/CompletionConsumer.java b/src/main/java/com/it_jaros/jns/scan/engine/CompletionConsumer.java new file mode 100644 index 0000000..2dfedf5 --- /dev/null +++ b/src/main/java/com/it_jaros/jns/scan/engine/CompletionConsumer.java @@ -0,0 +1,60 @@ +package com.it_jaros.jns.scan.engine; + +import java.util.concurrent.CancellationException; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.function.Consumer; + +public class CompletionConsumer { + + private final CancelledToken cancelledToken; + private final ProducerState state; + + public CompletionConsumer(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> consumer) { + while (!cancelledToken.isCancelled() && (state.running().get() || state.inPipeline().get() > 0)) { + consumer.accept(getPollResult()); + } + } + + private PollState getPollResult() { + try { + return pollResult(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + + return new PollState.Unavailable<>(); + } + + private PollState pollResult() 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 | CancellationException 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/ConsumerThread.java b/src/main/java/com/it_jaros/jns/scan/engine/ConsumerThread.java deleted file mode 100644 index 8f93317..0000000 --- a/src/main/java/com/it_jaros/jns/scan/engine/ConsumerThread.java +++ /dev/null @@ -1,59 +0,0 @@ -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 { - - 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> consumer) { - while (!cancelledToken.isCancelled() && (state.running().get() || state.inPipeline().get() > 0)) { - consumer.accept(getPollResult()); - } - } - - private PollState getPollResult() { - try { - return pollResult(); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - } - - return new PollState.Unavailable<>(); - } - - private PollState pollResult() 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 f90d4ce..7fee311 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 @@ -71,7 +71,7 @@ public class ScanHostTask implements Callable { port -> new ScanPortTask(port, context, hostContext) ); - new ConsumerThread<>(cancelledToken, state).runConsumerLoop(poll -> { + new CompletionConsumer<>(cancelledToken, state).runConsumerLoop(poll -> { switch (poll) { case PollState.Success(PortResult value) -> accumulator.add(value); case PollState.Failure(Throwable error) -> accumulator.recordFailure(error, -1); -- cgit v1.3.1