From b0f8d9ae7a8e2a327fbc3683eef3e9efb892d43c Mon Sep 17 00:00:00 2001 From: Matthias Jaros Date: Fri, 21 Aug 2026 17:58:30 +0200 Subject: Refactored ProducerThread: - Stopped exposing of semaphore used to control amount of workers - No outside modification of crucial ProducerState necessary - Moved aquire up before getting things out of the queue --- src/main/java/com/it_jaros/jns/scan/Scanner.java | 1 - .../it_jaros/jns/scan/engine/ProducerState.java | 2 - .../it_jaros/jns/scan/engine/ProducerThread.java | 76 +++++++++++++--------- .../com/it_jaros/jns/scan/engine/ScanHostTask.java | 1 - 4 files changed, 45 insertions(+), 35 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 3e84a26..32f1182 100644 --- a/src/main/java/com/it_jaros/jns/scan/Scanner.java +++ b/src/main/java/com/it_jaros/jns/scan/Scanner.java @@ -115,7 +115,6 @@ public class Scanner implements AutoCloseable { } catch (ExecutionException e) { return new PollState.Failure<>(e.getCause()); } finally { - state.activeWorkers().release(); state.inPipeline().decrementAndGet(); } } diff --git a/src/main/java/com/it_jaros/jns/scan/engine/ProducerState.java b/src/main/java/com/it_jaros/jns/scan/engine/ProducerState.java index 598fcbd..71e9887 100644 --- a/src/main/java/com/it_jaros/jns/scan/engine/ProducerState.java +++ b/src/main/java/com/it_jaros/jns/scan/engine/ProducerState.java @@ -1,13 +1,11 @@ package com.it_jaros.jns.scan.engine; import java.util.concurrent.CompletionService; -import java.util.concurrent.Semaphore; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; public record ProducerState( AtomicBoolean running, AtomicInteger inPipeline, - Semaphore activeWorkers, CompletionService completionService ) {} 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 056b4d9..013f3cc 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 @@ -47,36 +47,13 @@ public class ProducerThread { executor.submit(() -> { scan.producerStart(); try { - while (!cancelledToken.isCancelled() && queue.hasNext()) { - // get semaphore and remember if task got submitted - // so in case we fail to submit we release the semaphore - activeWorkers.acquire(); - boolean isTaskSubmitted = false; - try { - // just in case something - // changed while waiting - if (cancelledToken.isCancelled()) { - break; - } - - // get next item and create callable - // using lambda expression - final INPUT item = queue.next(); - Callable task = taskFactory.apply(item); - inPipeline.incrementAndGet(); - try { - completionService.submit(task); - isTaskSubmitted = true; - } catch (Throwable e) { - inPipeline.decrementAndGet(); - throw e; - } - } finally { - if (!isTaskSubmitted) { - activeWorkers.release(); - } - } - } + runProducerLoop(new ProducerContext<>( + queue, + taskFactory, + activeWorkers, + inPipeline, + completionService + )); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } finally { @@ -84,6 +61,43 @@ public class ProducerThread { running.set(false); } }); - return new ProducerState<>(running, inPipeline, activeWorkers, completionService); + return new ProducerState<>(running, inPipeline, completionService); + } + + private void runProducerLoop(ProducerContext context) throws InterruptedException { + while (!cancelledToken.isCancelled() && context.queue().hasNext()) { + submitNext(context); + } } + + private void submitNext(ProducerContext context) throws InterruptedException { + context.activeWorkers().acquire(); + + // just in case something + // changed while waiting + if (cancelledToken.isCancelled()) { + return; + } + + // get next item and create callable + // using lambda expression + final INPUT item = context.queue().next(); + final Callable task = context.taskFactory().apply(item); + context.completionService().submit(() -> { + try { + return task.call(); + } finally { + context.activeWorkers().release(); + } + }); + context.inPipeline().incrementAndGet(); + } + + private record ProducerContext( + Iterator queue, + Function> taskFactory, + Semaphore activeWorkers, + AtomicInteger inPipeline, + CompletionService completionService + ) {} } 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 2ae00d8..4a77cb1 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 @@ -135,7 +135,6 @@ public class ScanHostTask implements Callable { Throwable cause = e.getCause(); return new PollState.Failure<>(cause); } finally { - state.activeWorkers().release(); state.inPipeline().decrementAndGet(); } } -- cgit v1.3.1