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 | |
| parent | 540b09b1a019f95a9322a6d19e8928369bb20fcb (diff) | |
Major refactoring of package structure
Diffstat (limited to 'src/main')
29 files changed, 614 insertions, 505 deletions
diff --git a/src/main/java/com/it_jaros/jscanner/App.java b/src/main/java/com/it_jaros/jscanner/App.java index 67f8cdb..fdb9ce6 100644 --- a/src/main/java/com/it_jaros/jscanner/App.java +++ b/src/main/java/com/it_jaros/jscanner/App.java @@ -1,5 +1,9 @@ package com.it_jaros.jscanner; +import com.it_jaros.jscanner.scan.Scan; +import com.it_jaros.jscanner.scan.ScanOptions; +import com.it_jaros.jscanner.scan.Scanner; + import java.io.IOException; import java.util.concurrent.CountDownLatch; @@ -34,7 +38,7 @@ public class App { private static void runScan(Scanner scanner, ScanOptions options) throws IOException { addShutdownHook(scanner); Scan scan = Scan.create(options); - ProgressBar progressBar = new ProgressBar(scan, options.quiet()); + CliPrinter progressBar = new CliPrinter(scan, options.quiet()); progressBar.start(); scanner.runScan(scan, result -> progressBar.printResult(result, options.showFilteredPorts())); progressBar.stop(); diff --git a/src/main/java/com/it_jaros/jscanner/CliParser.java b/src/main/java/com/it_jaros/jscanner/CliParser.java index d7f5143..7e27064 100644 --- a/src/main/java/com/it_jaros/jscanner/CliParser.java +++ b/src/main/java/com/it_jaros/jscanner/CliParser.java @@ -1,5 +1,7 @@ package com.it_jaros.jscanner; +import com.it_jaros.jscanner.scan.ScanOptions; + import java.util.HashSet; import java.util.Set; diff --git a/src/main/java/com/it_jaros/jscanner/ProgressBar.java b/src/main/java/com/it_jaros/jscanner/CliPrinter.java index 2c323d7..709c3dd 100644 --- a/src/main/java/com/it_jaros/jscanner/ProgressBar.java +++ b/src/main/java/com/it_jaros/jscanner/CliPrinter.java @@ -1,5 +1,10 @@ package com.it_jaros.jscanner; +import com.it_jaros.jscanner.scan.Scan; +import com.it_jaros.jscanner.scan.domain.ScanFailure; +import com.it_jaros.jscanner.scan.domain.ScanResult; +import com.it_jaros.jscanner.scan.service.ServiceType; + import java.io.BufferedReader; import java.io.InputStreamReader; import java.util.List; @@ -11,14 +16,14 @@ import java.util.regex.Matcher; import java.util.regex.Pattern; import java.util.stream.Collectors; -public class ProgressBar { +public class CliPrinter { private final ScheduledExecutorService ui = Executors.newSingleThreadScheduledExecutor(); private final Scan scan; private final boolean quiet; private final Object outputLock = new Object(); - public ProgressBar(Scan scan, boolean quiet) { + public CliPrinter(Scan scan, boolean quiet) { this.scan = scan; this.quiet = quiet; } @@ -209,7 +214,7 @@ public class ProgressBar { stats.threadMaxConcurrent() ); - int availableWidth = ProgressBar.getColumnWidth(); + int availableWidth = CliPrinter.getColumnWidth(); String statusLine = appendIfEnoughSpace("", availableWidth, hostStat, portStat, durationStat, socketStat, workerStat); String barFormat = " [%s] %3d%%"; diff --git a/src/main/java/com/it_jaros/jscanner/ScanFailure.java b/src/main/java/com/it_jaros/jscanner/ScanFailure.java deleted file mode 100644 index 2c38790..0000000 --- a/src/main/java/com/it_jaros/jscanner/ScanFailure.java +++ /dev/null @@ -1,7 +0,0 @@ -package com.it_jaros.jscanner; - -public record ScanFailure( - Integer port, - ExceptionInfo exception -) { -} diff --git a/src/main/java/com/it_jaros/jscanner/Scanner.java b/src/main/java/com/it_jaros/jscanner/Scanner.java deleted file mode 100644 index 7f6b743..0000000 --- a/src/main/java/com/it_jaros/jscanner/Scanner.java +++ /dev/null @@ -1,476 +0,0 @@ -package com.it_jaros.jscanner; - -import java.io.ByteArrayOutputStream; -import java.io.IOException; -import java.net.*; -import java.nio.ByteBuffer; -import java.nio.channels.SocketChannel; -import java.time.Duration; -import java.util.*; -import java.util.concurrent.*; -import java.util.concurrent.atomic.AtomicBoolean; -import java.util.concurrent.atomic.AtomicInteger; -import java.util.concurrent.locks.LockSupport; -import java.util.function.Consumer; -import java.util.function.Function; - -public class Scanner implements AutoCloseable { - - private static final Duration pollInterval = Duration.ofSeconds(1); - private static final int READ_BUFFER_SIZE = 1024; - - private volatile boolean cancelled = false; - - private final ExecutorService executor = Executors.newVirtualThreadPerTaskExecutor(); - private final Semaphore socketLimit; - private final boolean bannerRecognition; - private final boolean disableOnlineCheck; - private final int maxHostsLimit; - private final int maxWorkersPerHost; - private final int timeoutInMillis; - private final long delayInNanos; - - public Scanner( - int socketLimit, - int timeoutInMillis, - int delayInMillis, - int maxWorkersPerHost, - int maxHostsLimit, - boolean disableOnlineCheck, - boolean bannerRecognition - ) { - this.delayInNanos = TimeUnit.MILLISECONDS.toNanos(Math.max(0, delayInMillis)); - this.bannerRecognition = bannerRecognition; - this.disableOnlineCheck = disableOnlineCheck; - this.maxHostsLimit = maxHostsLimit; - this.maxWorkersPerHost = maxWorkersPerHost; - this.socketLimit = new Semaphore(socketLimit); - this.timeoutInMillis = timeoutInMillis; - } - - public Scanner(ScanOptions options) { - this(options.socketLimit(), options.timeoutInMillis(), options.delayInMillis(), options.maxWorkersPerHost(), options.maxHostsLimit(), options.disableOnlineCheck(), options.bannerRecognition()); - } - - /** - * 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(); - - // start producer thread - scan.producerStart(); - final ProducerState<ScanResult> state = startProducer( - scan.getHosts().iterator(), - maxHostsLimit, - host -> new ScanHostTask(scan, host) - ); - - // 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(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(); - } - } - - /** - * This method helps to cleanup the code a bit and remove redundancy - * The producer for providing hosts and the one for providing ports - * are similar and the small differences can be handled using a function - * - * @param queue - * @param maxWorkers - * @param taskFactory - * @param <INPUT> - * @param <OUTPUT> - * @return - */ - private <INPUT, OUTPUT> ProducerState<OUTPUT> startProducer( - Iterator<INPUT> queue, - int maxWorkers, - Function<INPUT, Callable<OUTPUT>> taskFactory - ) { - final AtomicInteger inPipeline = new AtomicInteger(0); - final AtomicBoolean running = new AtomicBoolean(true); - final Semaphore activeWorkers = new Semaphore(maxWorkers); - CompletionService<OUTPUT> completionService = new ExecutorCompletionService<>(executor); - executor.submit(() -> { - try { - while (!cancelled && 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 (cancelled) { - 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(); - } - } - } - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - } finally { - running.set(false); - } - }); - return new ProducerState<>(running, inPipeline, activeWorkers, completionService); - } - - public boolean awaitTermination(Duration duration) throws InterruptedException { - return executor.awaitTermination(duration.toMillis(), TimeUnit.MILLISECONDS); - } - - public void cancel() { - if (!cancelled) { - cancelled = true; - } - executor.shutdown(); - } - - public void cancelNow() { - if (!cancelled) { - cancelled = true; - } - executor.shutdownNow(); - } - - @Override - public void close() throws Exception { - cancel(); - } - - - private final class ScanHostTask implements Callable<ScanResult> { - private final PortScanRateLimiter rateLimiter = new PortScanRateLimiter(); - - private final Scan scan; - private final String host; // input parameter - - ScanHostTask(Scan scan, String host) { - this.scan = scan; - this.host = host; - } - - @Override - public ScanResult call() { - try { - scan.hostStart(); - if (cancelled) { - return ScanResult.empty(host); - } - return scanHostPorts(); - } finally { - scan.hostFinish(); - } - } - - private ScanResult scanHostPorts() { - if (!disableOnlineCheck) { - boolean isHostOnline = checkHostOnline(); - if (!isHostOnline) { - // Unreachable host - return ScanResult.empty(host); - } - // online check also sends packets to the target system. - // in order not to violate set delay time - // we wait here too - LockSupport.parkNanos(delayInNanos); - } - - final PortRange portRange = new PortRange(scan.getPorts()); - // producer thread - scan.producerStart(); - ProducerState<PortResult> state = startProducer( - portRange.iterator(), - maxWorkersPerHost, - port -> new ScanPortTask(host, port, scan, rateLimiter) - ); - - // consumer is the main thread - // we run as long as the producer is running OR - // as long as things are in pipeline waiting to be processed - // ONLY exception is when cancelled is set - final PortResultAccumulator accumulator = new PortResultAccumulator(host); - while (!cancelled && (state.running().get() || state.inPipeline().get() > 0)) { - try { - PollState<PortResult> poll = getPortResult(state); - if (poll instanceof PollState.Success<PortResult>(PortResult value)) { - accumulator.add(value); - } else if (poll instanceof PollState.Failure(Throwable error)) { - System.err.printf("scanHostPorts(%host): ScanPortTask() failed for with error %s -> %s%n", host, error.getClass().getSimpleName(), error.getMessage()); - } - } catch (InterruptedException ignored) { - Thread.currentThread().interrupt(); - } - } - scan.producerStop(); - - return accumulator.build(); - } - - private boolean checkHostOnline() { - try { - return InetAddress.getByName(host).isReachable(timeoutInMillis); - } catch (IOException e) { - // we ignore this error because it means that the host is probably not online - } - - return false; - } - - private PollState<PortResult> getPortResult(ProducerState<PortResult> state) throws InterruptedException { - Future<PortResult> portResultFuture = state.completionService().poll(pollInterval.toMillis(), TimeUnit.MILLISECONDS); - if (portResultFuture == null) { - return new PollState.Unavailable<>(); - } - - PortResult portResult; - try { - portResult = portResultFuture.get(); - return new PollState.Success<>(portResult); - } catch (ExecutionException e) { - Throwable cause = e.getCause(); - return new PollState.Failure<>(cause); - } finally { - state.activeWorkers().release(); - state.inPipeline().decrementAndGet(); - } - } - } - - /** - * Per-host rate limiter. Each ScanHostTask creates its own instance and - * shares it with all its ScanPortTasks via constructor. - * - * Java allows one inner class to access another's private members, so this works. - */ - private final class PortScanRateLimiter { - private final Object lock = new Object(); - private volatile long nextAllowedTime; - - void apply() { - if (delayInNanos <= 0) { - return; - } - synchronized (lock) { - if (cancelled) { - return; - } - long now = System.nanoTime(); - if (nextAllowedTime > now) { - LockSupport.parkNanos(nextAllowedTime - now); - now = System.nanoTime(); // re-read after waking - } - nextAllowedTime = now + delayInNanos; - } - } - } - - private final class ScanPortTask implements Callable<PortResult> { - private final Scan scan; - private final String host; - private final int port; - private final PortScanRateLimiter portScanRateLimiter; // per-host shared limiter - - private ScanPortTask(String host, int port, Scan scan, PortScanRateLimiter portScanRateLimiter) { - this.host = host; - this.port = port; - this.scan = scan; - this.portScanRateLimiter = portScanRateLimiter; - } - - @Override - public PortResult call() throws Exception { - try { - socketLimit.acquire(); - scan.portStart(); - portScanRateLimiter.apply(); - if (cancelled) { - return PortResult.empty(); - } - return checkPort(); - } finally { - scan.portFinish(); - socketLimit.release(); - } - } - - private PortResult checkPort() { - PortResult result = new PortResult(); - result.setPort(port); - result.setState(PortState.UNKNOWN); - try(SocketChannel socketChannel = SocketChannel.open()) { - socketChannel.configureBlocking(false); - socketChannel.connect(new InetSocketAddress(host, port)); - final long deadlineNanos = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(timeoutInMillis); - final long waitInNanos = TimeUnit.MILLISECONDS.toNanos(1000); - boolean isConnected = socketChannel.finishConnect(); - while (!cancelled && !isConnected) { - long remainingNanos = deadlineNanos - System.nanoTime(); - if (remainingNanos <= 0) { - break; - } - LockSupport.parkNanos(Math.min(waitInNanos, remainingNanos)); - isConnected = socketChannel.finishConnect(); - } - - if (cancelled) { - return result; - } - - if (isConnected) { - result.setState(PortState.OPEN); - if (bannerRecognition) { - result.setBanner(getBanner(socketChannel)); - } - } else { - result.setState(PortState.FILTERED); - } - } catch (SocketTimeoutException ignored) { - result.setState(PortState.FILTERED); - } catch (ConnectException ignored) { - result.setState(PortState.CLOSED); - } catch (NoRouteToHostException ignored) { - // NoRouteToHostException: this can be safely ignored because the port is closed if a host is unreachable - // Will happen a lot when scanning for open ports, so not needed - } catch (IOException e) { - result.setException(e); - } - - return result; - } - - private byte[] getBanner(SocketChannel socketChannel) { - try { - return tryReadFrom(socketChannel); - } catch (IOException ignore) { - // ignore - } - - return null; - } - - private byte[] tryReadFrom(SocketChannel socketChannel) throws IOException { - byte[] buffer = new byte[READ_BUFFER_SIZE]; - ByteArrayOutputStream byteArrayOutputStream = new ByteArrayOutputStream(); - - boolean timedout = false; - long deadlineNanos = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(timeoutInMillis); - while (!timedout) { - int bytesRead = socketChannel.read(ByteBuffer.wrap(buffer)); - if (bytesRead == -1) { - break; - } - if (bytesRead > 0) { - byteArrayOutputStream.write(buffer, 0, bytesRead); - } - long remaining = deadlineNanos - System.nanoTime(); - if (remaining <= 0) { - timedout = true; - } else { - LockSupport.parkNanos(Math.min(remaining, TimeUnit.MILLISECONDS.toNanos(1000))); - } - } - return byteArrayOutputStream.toByteArray(); - } - } - - private final class PortResultAccumulator { - private final String host; - private final BitSet openPorts = new BitSet(PortRange.MAX_PORT); - private final BitSet filteredPorts = new BitSet(PortRange.MAX_PORT); - private final Map<Integer, ServiceType> serviceTypes = new HashMap<>(); - private final List<ScanFailure> scanFailures = new ArrayList<>(); - - private PortResultAccumulator(String host) { - this.host = host; - } - - void add(PortResult portResult) { - switch (portResult.getState()) { - case OPEN -> { - openPorts.set(portResult.getPort()); - serviceTypes.put(portResult.getPort(), ServiceDetector.detect(portResult.getBanner())); - } - case FILTERED -> { - filteredPorts.set(portResult.getPort()); - } - default -> { - // intentional no-op for uncovered port states - } - } - Exception e = portResult.getException(); - if (e != null) { - scanFailures.add(new ScanFailure(portResult.getPort(), ExceptionInfo.from(e))); - } - } - - ScanResult build() { - return new ScanResult( - host, - new PortList(openPorts), - new PortList(filteredPorts), - Collections.unmodifiableMap(serviceTypes), - Collections.unmodifiableList(scanFailures) - ); - } - } -} diff --git a/src/main/java/com/it_jaros/jscanner/Counter.java b/src/main/java/com/it_jaros/jscanner/scan/Counter.java index af7966a..0e1d5c9 100644 --- a/src/main/java/com/it_jaros/jscanner/Counter.java +++ b/src/main/java/com/it_jaros/jscanner/scan/Counter.java @@ -1,4 +1,4 @@ -package com.it_jaros.jscanner; +package com.it_jaros.jscanner.scan; import java.util.concurrent.atomic.AtomicInteger; diff --git a/src/main/java/com/it_jaros/jscanner/ExceptionInfo.java b/src/main/java/com/it_jaros/jscanner/scan/ExceptionInfo.java index ab92a5b..f62caf1 100644 --- a/src/main/java/com/it_jaros/jscanner/ExceptionInfo.java +++ b/src/main/java/com/it_jaros/jscanner/scan/ExceptionInfo.java @@ -1,4 +1,4 @@ -package com.it_jaros.jscanner; +package com.it_jaros.jscanner.scan; public record ExceptionInfo( String type, diff --git a/src/main/java/com/it_jaros/jscanner/Scan.java b/src/main/java/com/it_jaros/jscanner/scan/Scan.java index 3428acf..0820abc 100644 --- a/src/main/java/com/it_jaros/jscanner/Scan.java +++ b/src/main/java/com/it_jaros/jscanner/scan/Scan.java @@ -1,4 +1,4 @@ -package com.it_jaros.jscanner; +package com.it_jaros.jscanner.scan; import java.io.BufferedReader; import java.io.IOException; diff --git a/src/main/java/com/it_jaros/jscanner/ScanException.java b/src/main/java/com/it_jaros/jscanner/scan/ScanException.java index 3b2cf24..d667346 100644 --- a/src/main/java/com/it_jaros/jscanner/ScanException.java +++ b/src/main/java/com/it_jaros/jscanner/scan/ScanException.java @@ -1,4 +1,4 @@ -package com.it_jaros.jscanner; +package com.it_jaros.jscanner.scan; public class ScanException extends java.lang.RuntimeException { public ScanException(String message) { diff --git a/src/main/java/com/it_jaros/jscanner/ScanOptions.java b/src/main/java/com/it_jaros/jscanner/scan/ScanOptions.java index 5dcb822..ce0d562 100644 --- a/src/main/java/com/it_jaros/jscanner/ScanOptions.java +++ b/src/main/java/com/it_jaros/jscanner/scan/ScanOptions.java @@ -1,4 +1,4 @@ -package com.it_jaros.jscanner; +package com.it_jaros.jscanner.scan; import java.util.List; 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(); + } + +} diff --git a/src/main/java/com/it_jaros/jscanner/PortList.java b/src/main/java/com/it_jaros/jscanner/scan/domain/PortList.java index 574d3eb..f145700 100644 --- a/src/main/java/com/it_jaros/jscanner/PortList.java +++ b/src/main/java/com/it_jaros/jscanner/scan/domain/PortList.java @@ -1,4 +1,4 @@ -package com.it_jaros.jscanner; +package com.it_jaros.jscanner.scan.domain; import java.util.BitSet; import java.util.List; diff --git a/src/main/java/com/it_jaros/jscanner/PortRange.java b/src/main/java/com/it_jaros/jscanner/scan/domain/PortRange.java index 8969b8f..4434763 100644 --- a/src/main/java/com/it_jaros/jscanner/PortRange.java +++ b/src/main/java/com/it_jaros/jscanner/scan/domain/PortRange.java @@ -1,4 +1,6 @@ -package com.it_jaros.jscanner; +package com.it_jaros.jscanner.scan.domain; + +import com.it_jaros.jscanner.scan.engine.PortRangeIterator; import java.util.BitSet; import java.util.Iterator; diff --git a/src/main/java/com/it_jaros/jscanner/PortResult.java b/src/main/java/com/it_jaros/jscanner/scan/domain/PortResult.java index f246302..9e70cc5 100644 --- a/src/main/java/com/it_jaros/jscanner/PortResult.java +++ b/src/main/java/com/it_jaros/jscanner/scan/domain/PortResult.java @@ -1,4 +1,4 @@ -package com.it_jaros.jscanner; +package com.it_jaros.jscanner.scan.domain; import java.io.IOException; diff --git a/src/main/java/com/it_jaros/jscanner/PortState.java b/src/main/java/com/it_jaros/jscanner/scan/domain/PortState.java index 656d530..6db4b2e 100644 --- a/src/main/java/com/it_jaros/jscanner/PortState.java +++ b/src/main/java/com/it_jaros/jscanner/scan/domain/PortState.java @@ -1,4 +1,4 @@ -package com.it_jaros.jscanner; +package com.it_jaros.jscanner.scan.domain; public enum PortState { CLOSED, diff --git a/src/main/java/com/it_jaros/jscanner/scan/domain/ScanFailure.java b/src/main/java/com/it_jaros/jscanner/scan/domain/ScanFailure.java new file mode 100644 index 0000000..db33465 --- /dev/null +++ b/src/main/java/com/it_jaros/jscanner/scan/domain/ScanFailure.java @@ -0,0 +1,9 @@ +package com.it_jaros.jscanner.scan.domain; + +import com.it_jaros.jscanner.scan.ExceptionInfo; + +public record ScanFailure( + Integer port, + ExceptionInfo exception +) { +} diff --git a/src/main/java/com/it_jaros/jscanner/ScanResult.java b/src/main/java/com/it_jaros/jscanner/scan/domain/ScanResult.java index 3eb4d0d..e18e75d 100644 --- a/src/main/java/com/it_jaros/jscanner/ScanResult.java +++ b/src/main/java/com/it_jaros/jscanner/scan/domain/ScanResult.java @@ -1,4 +1,6 @@ -package com.it_jaros.jscanner; +package com.it_jaros.jscanner.scan.domain; + +import com.it_jaros.jscanner.scan.service.ServiceType; import java.util.BitSet; import java.util.Collections; diff --git a/src/main/java/com/it_jaros/jscanner/scan/engine/CancelledToken.java b/src/main/java/com/it_jaros/jscanner/scan/engine/CancelledToken.java new file mode 100644 index 0000000..1834a12 --- /dev/null +++ b/src/main/java/com/it_jaros/jscanner/scan/engine/CancelledToken.java @@ -0,0 +1,15 @@ +package com.it_jaros.jscanner.scan.engine; + +public class CancelledToken { + private volatile boolean cancelled = false; + + public void cancel() { + if (!cancelled) { + this.cancelled = true; + } + } + + public boolean isCancelled() { + return cancelled; + } +} diff --git a/src/main/java/com/it_jaros/jscanner/PollState.java b/src/main/java/com/it_jaros/jscanner/scan/engine/PollState.java index 474c055..54ec767 100644 --- a/src/main/java/com/it_jaros/jscanner/PollState.java +++ b/src/main/java/com/it_jaros/jscanner/scan/engine/PollState.java @@ -1,4 +1,4 @@ -package com.it_jaros.jscanner; +package com.it_jaros.jscanner.scan.engine; /** * A sealed interface representing the three possible outcomes of a CompletionService diff --git a/src/main/java/com/it_jaros/jscanner/PortRangeIterator.java b/src/main/java/com/it_jaros/jscanner/scan/engine/PortRangeIterator.java index 6690957..0673f0f 100644 --- a/src/main/java/com/it_jaros/jscanner/PortRangeIterator.java +++ b/src/main/java/com/it_jaros/jscanner/scan/engine/PortRangeIterator.java @@ -1,11 +1,11 @@ -package com.it_jaros.jscanner; +package com.it_jaros.jscanner.scan.engine; import java.util.BitSet; import java.util.Iterator; import java.util.NoSuchElementException; -import static com.it_jaros.jscanner.PortRange.MAX_PORT; -import static com.it_jaros.jscanner.PortRange.MIN_PORT; +import static com.it_jaros.jscanner.scan.domain.PortRange.MAX_PORT; +import static com.it_jaros.jscanner.scan.domain.PortRange.MIN_PORT; public class PortRangeIterator implements Iterator<Integer> { @@ -13,7 +13,7 @@ public class PortRangeIterator implements Iterator<Integer> { private int done = 0; private final BitSet availablePorts = new BitSet(MAX_PORT); - PortRangeIterator(BitSet specifiedPorts) { + public PortRangeIterator(BitSet specifiedPorts) { availablePorts.or(specifiedPorts); currentPortCursor = availablePorts.nextSetBit(MIN_PORT); } diff --git a/src/main/java/com/it_jaros/jscanner/scan/engine/PortResultAccumulator.java b/src/main/java/com/it_jaros/jscanner/scan/engine/PortResultAccumulator.java new file mode 100644 index 0000000..0cc26b0 --- /dev/null +++ b/src/main/java/com/it_jaros/jscanner/scan/engine/PortResultAccumulator.java @@ -0,0 +1,49 @@ +package com.it_jaros.jscanner.scan.engine; + +import com.it_jaros.jscanner.scan.ExceptionInfo; +import com.it_jaros.jscanner.scan.domain.*; +import com.it_jaros.jscanner.scan.service.ServiceDetector; +import com.it_jaros.jscanner.scan.service.ServiceType; + +import java.util.*; + +final class PortResultAccumulator { + private final String host; + private final BitSet openPorts = new BitSet(PortRange.MAX_PORT); + private final BitSet filteredPorts = new BitSet(PortRange.MAX_PORT); + private final Map<Integer, ServiceType> serviceTypes = new HashMap<>(); + private final List<ScanFailure> scanFailures = new ArrayList<>(); + + PortResultAccumulator(String host) { + this.host = host; + } + + void add(PortResult portResult) { + switch (portResult.getState()) { + case OPEN -> { + openPorts.set(portResult.getPort()); + serviceTypes.put(portResult.getPort(), ServiceDetector.detect(portResult.getBanner())); + } + case FILTERED -> { + filteredPorts.set(portResult.getPort()); + } + default -> { + // intentional no-op for uncovered port states + } + } + Exception e = portResult.getException(); + if (e != null) { + scanFailures.add(new ScanFailure(portResult.getPort(), ExceptionInfo.from(e))); + } + } + + ScanResult build() { + return new ScanResult( + host, + new PortList(openPorts), + new PortList(filteredPorts), + Collections.unmodifiableMap(serviceTypes), + Collections.unmodifiableList(scanFailures) + ); + } +} diff --git a/src/main/java/com/it_jaros/jscanner/scan/engine/PortScanRateLimiter.java b/src/main/java/com/it_jaros/jscanner/scan/engine/PortScanRateLimiter.java new file mode 100644 index 0000000..92b3bcb --- /dev/null +++ b/src/main/java/com/it_jaros/jscanner/scan/engine/PortScanRateLimiter.java @@ -0,0 +1,37 @@ +package com.it_jaros.jscanner.scan.engine; + +import java.util.concurrent.locks.LockSupport; + +/** + * Per-host rate limiter. Each ScanHostTask creates its own instance and + * shares it with all its ScanPortTasks via constructor. + * + * Java allows one inner class to access another's private members, so this works. + */ +public class PortScanRateLimiter { + private final Object lock = new Object(); + private volatile long nextAllowedTime; + private final long delayInNanos; + private final CancelledToken cancelledToken; + + public PortScanRateLimiter(CancelledToken cancelledToken, long delayInNanos) { + this.cancelledToken = cancelledToken; + this.delayInNanos = delayInNanos; + } + + void apply() { if (delayInNanos <= 0) { + return; + } + synchronized (lock) { + if (cancelledToken.isCancelled()) { + return; + } + long now = System.nanoTime(); + if (nextAllowedTime > now) { + LockSupport.parkNanos(nextAllowedTime - now); + now = System.nanoTime(); // re-read after waking + } + nextAllowedTime = now + delayInNanos; + } + } +} diff --git a/src/main/java/com/it_jaros/jscanner/ProducerState.java b/src/main/java/com/it_jaros/jscanner/scan/engine/ProducerState.java index 36a5841..c49ea3e 100644 --- a/src/main/java/com/it_jaros/jscanner/ProducerState.java +++ b/src/main/java/com/it_jaros/jscanner/scan/engine/ProducerState.java @@ -1,4 +1,4 @@ -package com.it_jaros.jscanner; +package com.it_jaros.jscanner.scan.engine; import java.util.concurrent.CompletionService; import java.util.concurrent.Semaphore; diff --git a/src/main/java/com/it_jaros/jscanner/scan/engine/ProducerThread.java b/src/main/java/com/it_jaros/jscanner/scan/engine/ProducerThread.java new file mode 100644 index 0000000..e7bcbcc --- /dev/null +++ b/src/main/java/com/it_jaros/jscanner/scan/engine/ProducerThread.java @@ -0,0 +1,83 @@ +package com.it_jaros.jscanner.scan.engine; + +import java.time.Duration; +import java.util.Iterator; +import java.util.concurrent.*; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.Function; + +public class ProducerThread { + + public static final Duration pollInterval = Duration.ofSeconds(1); + + private final ExecutorService executor; + private final CancelledToken cancelledToken; + + public ProducerThread(ScanExecutionContext context) { + this.executor = context.executorService(); + this.cancelledToken = context.cancelledToken(); + } + + /** + * This method helps to cleanup the code a bit and remove redundancy + * The producer for providing hosts and the one for providing ports + * are similar and the small differences can be handled using a function + * + * @param queue + * @param maxWorkers + * @param taskFactory + * @param <INPUT> + * @param <OUTPUT> + * @return + */ + public <INPUT, OUTPUT> ProducerState<OUTPUT> startProducer( + Iterator<INPUT> queue, + int maxWorkers, + Function<INPUT, Callable<OUTPUT>> taskFactory + ) { + final AtomicInteger inPipeline = new AtomicInteger(0); + final AtomicBoolean running = new AtomicBoolean(true); + final Semaphore activeWorkers = new Semaphore(maxWorkers); + CompletionService<OUTPUT> completionService = new ExecutorCompletionService<>(executor); + executor.submit(() -> { + 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(); + } + } + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } finally { + running.set(false); + } + }); + return new ProducerState<>(running, inPipeline, activeWorkers, completionService); + } +} diff --git a/src/main/java/com/it_jaros/jscanner/scan/engine/ScanExecutionContext.java b/src/main/java/com/it_jaros/jscanner/scan/engine/ScanExecutionContext.java new file mode 100644 index 0000000..38e9cdb --- /dev/null +++ b/src/main/java/com/it_jaros/jscanner/scan/engine/ScanExecutionContext.java @@ -0,0 +1,16 @@ +package com.it_jaros.jscanner.scan.engine; + +import com.it_jaros.jscanner.scan.Scan; +import com.it_jaros.jscanner.scan.ScanOptions; + +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Semaphore; + +public record ScanExecutionContext( + ScanOptions scanOptions, + ExecutorService executorService, + Semaphore socketLimit, + CancelledToken cancelledToken, + Scan scan, + PortScanRateLimiter portScanRateLimiter +) {} diff --git a/src/main/java/com/it_jaros/jscanner/scan/engine/ScanHostTask.java b/src/main/java/com/it_jaros/jscanner/scan/engine/ScanHostTask.java new file mode 100644 index 0000000..7a01d54 --- /dev/null +++ b/src/main/java/com/it_jaros/jscanner/scan/engine/ScanHostTask.java @@ -0,0 +1,123 @@ +package com.it_jaros.jscanner.scan.engine; + +import com.it_jaros.jscanner.scan.Scan; +import com.it_jaros.jscanner.scan.domain.PortRange; +import com.it_jaros.jscanner.scan.domain.PortResult; +import com.it_jaros.jscanner.scan.domain.ScanResult; + +import java.io.IOException; +import java.net.InetAddress; +import java.util.concurrent.Callable; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.locks.LockSupport; + +public class ScanHostTask implements Callable<ScanResult> { + + private final Scan scan; + private final String host; // input parameter + private final ScanExecutionContext context; + private final CancelledToken cancelledToken; + private final boolean disableOnlineCheck; + private final int maxWorkersPerHost; + private final int timeoutInMillis; + private final long delayInNanos; + + public ScanHostTask(Scan scan, String host, ScanExecutionContext context) { + this.scan = scan; + this.host = host; + this.context = context; + this.cancelledToken = context.cancelledToken(); + this.delayInNanos = TimeUnit.MILLISECONDS.toNanos(Math.max(0, context.scanOptions().delayInMillis())); + this.disableOnlineCheck = context.scanOptions().disableOnlineCheck(); + this.maxWorkersPerHost = context.scanOptions().maxWorkersPerHost(); + this.timeoutInMillis = context.scanOptions().timeoutInMillis(); + } + + @Override + public ScanResult call() { + try { + scan.hostStart(); + if (context.cancelledToken().isCancelled()) { + return ScanResult.empty(host); + } + return scanHostPorts(); + } finally { + scan.hostFinish(); + } + } + + private ScanResult scanHostPorts() { + if (!disableOnlineCheck) { + boolean isHostOnline = checkHostOnline(); + if (!isHostOnline) { + // Unreachable host + return ScanResult.empty(host); + } + // online check also sends packets to the target system. + // in order not to violate set delay time + // we wait here too + LockSupport.parkNanos(delayInNanos); + } + + final PortRange portRange = new PortRange(scan.getPorts()); + // producer thread + scan.producerStart(); + ProducerState<PortResult> state = new ProducerThread(context).startProducer( + portRange.iterator(), + maxWorkersPerHost, + port -> new ScanPortTask(host, port, context) + ); + + // consumer is the main thread + // we run as long as the producer is running OR + // as long as things are in pipeline waiting to be processed + // ONLY exception is when cancelled is set + final PortResultAccumulator accumulator = new PortResultAccumulator(host); + while (!cancelledToken.isCancelled() && (state.running().get() || state.inPipeline().get() > 0)) { + try { + PollState<PortResult> poll = getPortResult(state); + if (poll instanceof PollState.Success<PortResult>(PortResult value)) { + accumulator.add(value); + } else if (poll instanceof PollState.Failure(Throwable error)) { + System.err.printf("scanHostPorts(%s): ScanPortTask() failed for with error %s -> %s%n", host, error.getClass().getSimpleName(), error.getMessage()); + } + } catch (InterruptedException ignored) { + Thread.currentThread().interrupt(); + } + } + scan.producerStop(); + + return accumulator.build(); + } + + private boolean checkHostOnline() { + try { + return InetAddress.getByName(host).isReachable(timeoutInMillis); + } catch (IOException e) { + // we ignore this error because it means that the host is probably not online + } + + return false; + } + + private PollState<PortResult> getPortResult(ProducerState<PortResult> state) throws InterruptedException { + Future<PortResult> portResultFuture = state.completionService().poll(ProducerThread.pollInterval.toMillis(), TimeUnit.MILLISECONDS); + if (portResultFuture == null) { + return new PollState.Unavailable<>(); + } + + PortResult portResult; + try { + portResult = portResultFuture.get(); + return new PollState.Success<>(portResult); + } catch (ExecutionException e) { + Throwable cause = e.getCause(); + return new PollState.Failure<>(cause); + } finally { + state.activeWorkers().release(); + state.inPipeline().decrementAndGet(); + } + } +} diff --git a/src/main/java/com/it_jaros/jscanner/scan/engine/ScanPortTask.java b/src/main/java/com/it_jaros/jscanner/scan/engine/ScanPortTask.java new file mode 100644 index 0000000..873ec95 --- /dev/null +++ b/src/main/java/com/it_jaros/jscanner/scan/engine/ScanPortTask.java @@ -0,0 +1,135 @@ +package com.it_jaros.jscanner.scan.engine; + +import com.it_jaros.jscanner.scan.Scan; +import com.it_jaros.jscanner.scan.domain.PortResult; +import com.it_jaros.jscanner.scan.domain.PortState; + +import java.io.ByteArrayOutputStream; +import java.io.IOException; +import java.net.ConnectException; +import java.net.InetSocketAddress; +import java.net.NoRouteToHostException; +import java.nio.ByteBuffer; +import java.nio.channels.SocketChannel; +import java.util.concurrent.Callable; +import java.util.concurrent.Semaphore; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.locks.LockSupport; + +public class ScanPortTask implements Callable<PortResult> { + + private static final int READ_BUFFER_SIZE = 1024; + + private final CancelledToken cancelledToken; + private final PortScanRateLimiter portScanRateLimiter; // per-host shared limiter + private final Scan scan; + private final Semaphore socketLimit; + private final String host; + private final boolean bannerRecognition; + private final int port; + private final long timeoutInNanos; + + ScanPortTask(String host, int port, ScanExecutionContext context) { + this.host = host; + this.port = port; + this.bannerRecognition = context.scanOptions().bannerRecognition(); + this.cancelledToken = context.cancelledToken(); + this.portScanRateLimiter = context.portScanRateLimiter(); + this.scan = context.scan(); + this.socketLimit = context.socketLimit(); + this.timeoutInNanos = TimeUnit.MILLISECONDS.toNanos(context.scanOptions().timeoutInMillis()); + } + + @Override + public PortResult call() throws Exception { + try { + socketLimit.acquire(); + scan.portStart(); + portScanRateLimiter.apply(); + if (cancelledToken.isCancelled()) { + return PortResult.empty(); + } + return checkPort(); + } finally { + scan.portFinish(); + socketLimit.release(); + } + } + + private PortResult checkPort() { + PortResult result = new PortResult(); + result.setPort(port); + result.setState(PortState.UNKNOWN); + try(SocketChannel socketChannel = SocketChannel.open()) { + socketChannel.configureBlocking(false); + socketChannel.connect(new InetSocketAddress(host, port)); + final long deadlineNanos = System.nanoTime() + timeoutInNanos; + final long waitInNanos = TimeUnit.MILLISECONDS.toNanos(1000); + boolean isConnected = socketChannel.finishConnect(); + while (!cancelledToken.isCancelled() && !isConnected) { + long remainingNanos = deadlineNanos - System.nanoTime(); + if (remainingNanos <= 0) { + break; + } + LockSupport.parkNanos(Math.min(waitInNanos, remainingNanos)); + isConnected = socketChannel.finishConnect(); + } + + if (cancelledToken.isCancelled()) { + return result; + } + + if (isConnected) { + result.setState(PortState.OPEN); + if (bannerRecognition) { + result.setBanner(getBanner(socketChannel)); + } + } else { + result.setState(PortState.FILTERED); + } + } catch (ConnectException ignored) { + result.setState(PortState.CLOSED); + } catch (NoRouteToHostException ignored) { + // NoRouteToHostException: this can be safely ignored because the port is closed if a host is unreachable + // Will happen a lot when scanning for open ports, so not needed + } catch (IOException e) { + result.setException(e); + } + + return result; + } + + private byte[] getBanner(SocketChannel socketChannel) { + try { + return tryReadFrom(socketChannel); + } catch (IOException ignore) { + // ignore + } + + return null; + } + + private byte[] tryReadFrom(SocketChannel socketChannel) throws IOException { + byte[] buffer = new byte[READ_BUFFER_SIZE]; + ByteArrayOutputStream byteArrayOutputStream = new ByteArrayOutputStream(); + + boolean timedout = false; + long deadlineNanos = System.nanoTime() + timeoutInNanos; + while (!timedout) { + int bytesRead = socketChannel.read(ByteBuffer.wrap(buffer)); + if (bytesRead == -1) { + break; + } + if (bytesRead > 0) { + byteArrayOutputStream.write(buffer, 0, bytesRead); + } + long remaining = deadlineNanos - System.nanoTime(); + if (remaining <= 0) { + timedout = true; + } else { + LockSupport.parkNanos(Math.min(remaining, TimeUnit.MILLISECONDS.toNanos(1000))); + } + } + return byteArrayOutputStream.toByteArray(); + } +} diff --git a/src/main/java/com/it_jaros/jscanner/ServiceDetector.java b/src/main/java/com/it_jaros/jscanner/scan/service/ServiceDetector.java index 4d4f76e..93e05d4 100644 --- a/src/main/java/com/it_jaros/jscanner/ServiceDetector.java +++ b/src/main/java/com/it_jaros/jscanner/scan/service/ServiceDetector.java @@ -1,4 +1,4 @@ -package com.it_jaros.jscanner; +package com.it_jaros.jscanner.scan.service; import java.nio.charset.StandardCharsets; import java.util.List; diff --git a/src/main/java/com/it_jaros/jscanner/ServiceType.java b/src/main/java/com/it_jaros/jscanner/scan/service/ServiceType.java index 3677faf..86b282a 100644 --- a/src/main/java/com/it_jaros/jscanner/ServiceType.java +++ b/src/main/java/com/it_jaros/jscanner/scan/service/ServiceType.java @@ -1,4 +1,4 @@ -package com.it_jaros.jscanner; +package com.it_jaros.jscanner.scan.service; public enum ServiceType { HTTP, FTP, SMTP, IMAP, POP3, UNKNOWN, DNS, SSH |
