From 8f94145316a998867a7bb077cd869b896dde48c5 Mon Sep 17 00:00:00 2001 From: Matthias Jaros Date: Thu, 13 Aug 2026 16:49:00 +0200 Subject: Moved ProducerThread counter to ProducerThread --- src/main/java/com/it_jaros/jns/scan/Scanner.java | 4 +--- src/main/java/com/it_jaros/jns/scan/engine/ProducerThread.java | 8 +++++++- src/main/java/com/it_jaros/jns/scan/engine/ScanHostTask.java | 6 ++---- 3 files changed, 10 insertions(+), 8 deletions(-) (limited to 'src/main') 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 4dad6a8..f081f29 100644 --- a/src/main/java/com/it_jaros/jns/scan/Scanner.java +++ b/src/main/java/com/it_jaros/jns/scan/Scanner.java @@ -40,11 +40,10 @@ public class Scanner implements AutoCloseable { ); scan.start(); - scan.producerStart(); final ProducerState state = new ProducerThread(context).startProducer( scan.getHosts().iterator(), scanOptions.maxHostsLimit(), - host -> new ScanHostTask(scan, host, context) + host -> new ScanHostTask(host, context) ); // the main thread is the consumer @@ -65,7 +64,6 @@ public class Scanner implements AutoCloseable { Thread.currentThread().interrupt(); } } - scan.producerStop(); scan.stop(); } diff --git a/src/main/java/com/it_jaros/jns/scan/engine/ProducerThread.java b/src/main/java/com/it_jaros/jns/scan/engine/ProducerThread.java index 0cae158..056b4d9 100644 --- a/src/main/java/com/it_jaros/jns/scan/engine/ProducerThread.java +++ b/src/main/java/com/it_jaros/jns/scan/engine/ProducerThread.java @@ -1,5 +1,7 @@ package com.it_jaros.jns.scan.engine; +import com.it_jaros.jns.scan.Scan; + import java.time.Duration; import java.util.Iterator; import java.util.concurrent.*; @@ -13,10 +15,12 @@ public class ProducerThread { private final ExecutorService executor; private final CancelledToken cancelledToken; + private final Scan scan; public ProducerThread(ScanExecutionContext context) { - this.executor = context.executorService(); this.cancelledToken = context.cancelledToken(); + this.executor = context.executorService(); + this.scan = context.scan(); } /** @@ -41,6 +45,7 @@ public class ProducerThread { final Semaphore activeWorkers = new Semaphore(maxWorkers); CompletionService completionService = new ExecutorCompletionService<>(executor); executor.submit(() -> { + scan.producerStart(); try { while (!cancelledToken.isCancelled() && queue.hasNext()) { // get semaphore and remember if task got submitted @@ -75,6 +80,7 @@ public class ProducerThread { } catch (InterruptedException e) { Thread.currentThread().interrupt(); } finally { + scan.producerStop(); running.set(false); } }); 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 ce06287..24928d9 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 @@ -21,14 +21,14 @@ public class ScanHostTask implements Callable { private final int maxWorkersPerHost; private final PortResultAccumulator accumulator; - public ScanHostTask(Scan scan, String host, ScanExecutionContext context) { - this.scan = scan; + public ScanHostTask(String host, ScanExecutionContext context) { this.host = host; this.context = context; this.accumulator = new PortResultAccumulator(host); this.cancelledToken = context.cancelledToken(); this.disableOnlineCheck = context.scanOptions().disableOnlineCheck(); this.maxWorkersPerHost = context.scanOptions().maxWorkersPerHost(); + this.scan = context.scan(); } @Override @@ -55,7 +55,6 @@ public class ScanHostTask implements Callable { final PortRange portRange = new PortRange(scan.getPorts()); // producer thread - scan.producerStart(); ProducerState state = new ProducerThread(context).startProducer( portRange.iterator(), maxWorkersPerHost, @@ -78,7 +77,6 @@ public class ScanHostTask implements Callable { Thread.currentThread().interrupt(); } } - scan.producerStop(); return accumulator.build(); } -- cgit v1.3.1