diff options
| author | 2026-08-12 13:24:27 +0200 | |
|---|---|---|
| committer | 2026-08-12 13:50:46 +0200 | |
| commit | ff14e0fc4e3fe4f5b8a67640d28850064661d7ab (patch) | |
| tree | f141a2d34d902d344207cfbbc6ac614ec1b3840d /src/main/java/com/it_jaros/jscanner/scan/Scanner.java | |
| parent | 540b09b1a019f95a9322a6d19e8928369bb20fcb (diff) | |
Major refactoring of package structure
Diffstat (limited to 'src/main/java/com/it_jaros/jscanner/scan/Scanner.java')
| -rw-r--r-- | src/main/java/com/it_jaros/jscanner/scan/Scanner.java | 110 |
1 files changed, 110 insertions, 0 deletions
diff --git a/src/main/java/com/it_jaros/jscanner/scan/Scanner.java b/src/main/java/com/it_jaros/jscanner/scan/Scanner.java new file mode 100644 index 0000000..e0fbb3a --- /dev/null +++ b/src/main/java/com/it_jaros/jscanner/scan/Scanner.java @@ -0,0 +1,110 @@ +package com.it_jaros.jscanner.scan; + +import com.it_jaros.jscanner.scan.domain.ScanResult; +import com.it_jaros.jscanner.scan.engine.*; + +import java.time.Duration; +import java.util.concurrent.*; +import java.util.function.Consumer; + +public class Scanner implements AutoCloseable { + + private final CancelledToken cancelledToken = new CancelledToken(); + private final ExecutorService executor = Executors.newVirtualThreadPerTaskExecutor(); + private final ScanOptions scanOptions; + private final Semaphore socketLimit; + + public Scanner(ScanOptions options) { + this.scanOptions = options; + this.socketLimit = new Semaphore(scanOptions.socketLimit()); + } + + /** + * Starts a given scan. + * + * @param scan + */ + public void runScan(final Scan scan, final Consumer<ScanResult> consumer) { + if (scan == null) { + throw new IllegalArgumentException("Scan argument cannot be null"); + } + + scan.start(); + long delayInNanos = TimeUnit.MILLISECONDS.toNanos(Math.max(0, scanOptions.delayInMillis())); + ScanExecutionContext context = new ScanExecutionContext( + scanOptions, + executor, + socketLimit, + cancelledToken, + scan, + new PortScanRateLimiter(cancelledToken, delayInNanos) + ); + + // start producer thread + scan.producerStart(); + final ProducerState<ScanResult> state = new ProducerThread(context).startProducer( + scan.getHosts().iterator(), + scanOptions.maxHostsLimit(), + host -> new ScanHostTask(scan, 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 + 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)) { + System.err.printf("runScan(): ScanHostTask() failed with error %s -> %s%n", error.getClass().getSimpleName(), error.getMessage()); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } + scan.producerStop(); + + scan.stop(); + } + + private PollState<ScanResult> getHostResult(ProducerState<ScanResult> state) throws InterruptedException { + Future<ScanResult> finishedHost = state.completionService().poll(ProducerThread.pollInterval.toMillis(), TimeUnit.MILLISECONDS); + if (finishedHost == null) { + return new PollState.Unavailable<>(); + } + + try { + ScanResult result = finishedHost.get(); + return new PollState.Success<>(result); + } catch (ExecutionException e) { + return new PollState.Failure<>(e.getCause()); + } finally { + state.activeWorkers().release(); + state.inPipeline().decrementAndGet(); + } + } + + public boolean awaitTermination(Duration duration) throws InterruptedException { + return executor.awaitTermination(duration.toMillis(), TimeUnit.MILLISECONDS); + } + + public void cancel() { + this.cancelledToken.cancel(); + executor.shutdown(); + } + + public void cancelNow() { + this.cancelledToken.cancel(); + executor.shutdownNow(); + } + + @Override + public void close() throws Exception { + cancel(); + } + +} |
