summaryrefslogtreecommitdiff
path: root/src
diff options
context:
space:
mode:
Diffstat (limited to 'src')
-rw-r--r--src/main/java/com/it_jaros/jns/scan/Scanner.java1
-rw-r--r--src/main/java/com/it_jaros/jns/scan/engine/ProducerState.java2
-rw-r--r--src/main/java/com/it_jaros/jns/scan/engine/ProducerThread.java76
-rw-r--r--src/main/java/com/it_jaros/jns/scan/engine/ScanHostTask.java1
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();
}
}