diff options
| author | 2026-08-21 17:58:30 +0200 | |
|---|---|---|
| committer | 2026-08-21 20:36:07 +0200 | |
| commit | b0f8d9ae7a8e2a327fbc3683eef3e9efb892d43c (patch) | |
| tree | ee7e1f068a5f6599d149aa72c39e5c2417bedb64 /src/main/java/com | |
| parent | 774e4e463f3c2f31f0991c303e43ca1a84d3fc82 (diff) | |
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
Diffstat (limited to 'src/main/java/com')
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<OUTPUT>( AtomicBoolean running, AtomicInteger inPipeline, - Semaphore activeWorkers, CompletionService<OUTPUT> 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<OUTPUT> 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 <INPUT, OUTPUT> void runProducerLoop(ProducerContext<INPUT, OUTPUT> context) throws InterruptedException { + while (!cancelledToken.isCancelled() && context.queue().hasNext()) { + submitNext(context); + } } + + private <OUTPUT, INPUT> void submitNext(ProducerContext<INPUT, OUTPUT> 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<OUTPUT> task = context.taskFactory().apply(item); + context.completionService().submit(() -> { + try { + return task.call(); + } finally { + context.activeWorkers().release(); + } + }); + context.inPipeline().incrementAndGet(); + } + + private record ProducerContext<INPUT, OUTPUT>( + Iterator<INPUT> queue, + Function<INPUT, Callable<OUTPUT>> taskFactory, + Semaphore activeWorkers, + AtomicInteger inPipeline, + CompletionService<OUTPUT> 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<ScanResult> { Throwable cause = e.getCause(); return new PollState.Failure<>(cause); } finally { - state.activeWorkers().release(); state.inPipeline().decrementAndGet(); } } |
