summaryrefslogtreecommitdiff
path: root/src/main
diff options
context:
space:
mode:
authorGravatar Matthias Jaros <jarlucmat@mailbox.org>2026-08-20 17:14:51 +0200
committerGravatar Matthias Jaros <jarlucmat@mailbox.org>2026-08-21 14:18:37 +0200
commit8abc5f6d209778b678fb882b1213f596dfb42f7b (patch)
treeba501285d2ebfa2b40f118a9cce86066b703a96c /src/main
parent38111c6204cbbdc1f784fdedb03170fe33d000af (diff)
Refactored consumer into custom method
Diffstat (limited to 'src/main')
-rw-r--r--src/main/java/com/it_jaros/jns/scan/Scanner.java50
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<>();
}