diff --git a/src/main/java/ru/kirillius/pf/sdn/PfSdnApplication.java b/src/main/java/ru/kirillius/pf/sdn/PfSdnApplication.java new file mode 100644 index 0000000..b50cc95 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/PfSdnApplication.java @@ -0,0 +1,31 @@ +package ru.kirillius.pf.sdn; + +import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.boot.builder.SpringApplicationBuilder; + +@SpringBootApplication +public class PfSdnApplication { + private final static String DATABASE_FILE = "pfsdn"; + + public static void main(String[] args) { + + /* + ENV: + SERVER_PORT=8080 + APP_DATA=/var/lib/pfsdn + + ARGS: + -Dserver.port=8080 + -app.data=/var/lib/pfsdn + */ + + //System.setProperty("server.port", String.valueOf(configData.port())); + //System.setProperty("jwt.secret", configData.jwtSecret()); + var appData = System.getProperty("app.data", "."); + System.setProperty("app.datasource.url", "jdbc:h2:file:" + appData + "/" + DATABASE_FILE); + + new SpringApplicationBuilder(PfSdnApplication.class).run(args); + } + + +} diff --git a/src/main/java/ru/kirillius/pf/sdn/Util/CommandLineUtils.java b/src/main/java/ru/kirillius/pf/sdn/Util/CommandLineUtils.java new file mode 100644 index 0000000..315f7c1 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/Util/CommandLineUtils.java @@ -0,0 +1,22 @@ +package ru.kirillius.pf.sdn.Util; + +import java.util.Arrays; + +/** + * Utility methods for parsing command-line arguments used by the launcher. + */ +public final class CommandLineUtils { + private CommandLineUtils() { + } + + /** + * Returns the first command-line argument starting with the given flag name. + */ + public static String getArgument(String argname, String[] args) { + var first = Arrays.stream(args).filter(arg -> arg.startsWith("-" + argname)).findFirst(); + if (first.isEmpty()) { + throw new IllegalArgumentException("Missing required argument: -" + argname); + } + return first.get().substring(argname.length() + 2); + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/Util/IPv4Util.java b/src/main/java/ru/kirillius/pf/sdn/Util/IPv4Util.java new file mode 100644 index 0000000..eb28c49 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/Util/IPv4Util.java @@ -0,0 +1,354 @@ +package ru.kirillius.pf.sdn.Util; + +import lombok.Getter; +import lombok.SneakyThrows; +import ru.kirillius.pf.sdn.api.Networking.IPv4Subnet; + +import java.net.InetAddress; +import java.util.*; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicLong; +import java.util.regex.Pattern; +import java.util.stream.Collectors; + + +/** + * Helper methods for IPv4 address manipulation and subnet aggregation logic. + */ +public class IPv4Util { + + private IPv4Util() { + } + + private static final Pattern pattern = Pattern.compile("^((25[0-5]|2[0-4][0-9]|[01]?[0-9][0-9]?)\\.){3}(25[0-5]|2[0-4][0-9]|[01]?[0-9][0-9]?)$"); + + /** + * Ensures the supplied string is a valid IPv4 address. + */ + public static void validateAddress(String address) { + if (!pattern.matcher(address).matches()) { + throw new IllegalArgumentException("Invalid IPv4 address: " + address); + } + } + + /** + * Ensures the prefix length is within the allowed range 0..32. + */ + public static void validatePrefix(int prefix) { + if (prefix < 0 || prefix > 32) { + throw new IllegalArgumentException("Invalid IPv4 prefix: " + prefix); + } + } + + + /** + * Converts a dotted IPv4 address into its numeric representation. + */ + @SneakyThrows + public static long ipAddressToLong(String address) { + validateAddress(address); + var ip = InetAddress.getByName(address); + var bytes = ip.getAddress(); + var result = 0L; + for (var b : bytes) { + result = (result << 8) | (b & 0xFF); + } + return result; + } + + /** + * Returns a bit mask representing the provided prefix length. + */ + public static long calculateMask(int prefixLength) { + validatePrefix(prefixLength); + return 0xFFFFFFFFL << (32 - prefixLength); + } + + + /** + * Converts a numeric IPv4 address into dotted notation. + */ + public static String longToIpAddress(long ipLong) { + if (ipLong < 0 || ipLong > 0xFFFFFFFFL) { + throw new IllegalArgumentException("Address number should be in range 0 - 4294967295"); + } + return ((ipLong >> 24) & 0xFF) + "." + ((ipLong >> 16) & 0xFF) + "." + ((ipLong >> 8) & 0xFF) + "." + (ipLong & 0xFF); + } + + /** + * Formats a numeric mask into dotted notation. + */ + public static String maskToString(long maskLong) { + return String.format("%d.%d.%d.%d", + (maskLong >> 24) & 0xff, + (maskLong >> 16) & 0xff, + (maskLong >> 8) & 0xff, + maskLong & 0xff); + } + + /** + * Result contract for subnet summarisation operations. + */ + public interface SummarisationResult { + /** + * Returns the resulting set of subnets after summarisation. + */ + List getResult(); + + /** + * Returns the original subnets that were merged during summarisation. + */ + Set getMergedSubnets(); + } + + /** + * Performs subnet merging strategies based on overlap and utilisation heuristics. + */ + private static class SubnetSummaryUtility implements SummarisationResult { + + @Getter + private final List result; + private final Collection source; + @Getter + private final Set mergedSubnets = new HashSet<>(); + + + /** + * Builds the summarisation utility with the supplied subnets and usage threshold. + */ + public SubnetSummaryUtility(Collection subnets, int usePercentage) { + source = subnets; + result = new ArrayList<>(subnets); + summaryOverlapped(); + mergeNeighbours(); + summaryWithUsage(usePercentage > 50 ? usePercentage : 51); + summaryOverlapped(); + result.sort(Comparator.comparing(IPv4Subnet::getLongAddress)); + } + + /** + * Removes redundant subnets wholly covered by other subnets. + */ + private void summaryOverlapped() { + if (result.size() < 2) { + return; + } + //check subnets overlaps + var overlapped = new ArrayList(); + var orderedByPrefix = result.stream().sorted(Comparator.comparing(IPv4Subnet::getPrefixLength)).toList(); + orderedByPrefix.stream() + .filter(subnet -> subnet.getPrefixLength() > 32) + .forEach(parent -> orderedByPrefix.forEach(subnet -> { + if (subnet.equals(parent)) { + return; + } + if (parent.overlaps(subnet)) { + overlapped.add(subnet); + if (source.contains(subnet)) { + mergedSubnets.add(subnet); + } + } + })); + overlapped.forEach(result::remove); + } + + /** + * Attempts to merge adjacent subnets into larger ones while preserving coverage. + */ + private void mergeNeighbours() { + if (result.size() < 2) { + return; + } + var availableLengths = result.stream().map(IPv4Subnet::getPrefixLength).collect(Collectors.toSet()); + for (var length = 32; length > 0; length--) { + if (!availableLengths.contains(length)) { + continue; + } + var finalLength = length; + var largerPrefixLength = length - 1; + var largerPrefixMask = IPv4Util.calculateMask(largerPrefixLength); + + + var selectedSubnets = result.stream().filter(subnet -> subnet.getPrefixLength() == finalLength).sorted(Comparator.comparing(IPv4Subnet::getLongAddress)).toList(); + //проверяем является ли адрес подсети таким же как адрес подсети с перфиксом -1 + for (var i = 0; i < selectedSubnets.size() - 1; i++) { + var subnet = selectedSubnets.get(i); + var next = selectedSubnets.get(i + 1); + var firstAddress = subnet.getLongAddress(); + if (firstAddress != (firstAddress & largerPrefixMask)) { + continue; + } + var largerSubnet = new IPv4Subnet(subnet.getAddress(), largerPrefixLength); + if (largerSubnet.overlaps(next)) { + //если подсеть перекрывает соседнюю, то удаляем обе и добавляем новую + availableLengths.add(largerPrefixLength); + result.remove(subnet); + result.remove(next); + result.add(largerSubnet); + if (source.contains(subnet)) { + mergedSubnets.add(subnet); + } + if (source.contains(next)) { + mergedSubnets.add(next); + } + i++; + } + } + } + } + + /** + * Returns the smallest prefix length currently present in the result set. + */ + private int findMinPrefixLength() { + return result.stream().mapToInt(IPv4Subnet::getPrefixLength).min().getAsInt(); + } + + /** + * Returns the largest prefix length currently present in the result set. + */ + private int findMaxPrefixLength() { + return result.stream().mapToInt(IPv4Subnet::getPrefixLength).max().getAsInt(); + } + + /** + * Returns the smallest network address found in the result set. + */ + private long findMinAddress() { + return result.stream().mapToLong(IPv4Subnet::getLongAddress).min().getAsLong(); + } + + /** + * Returns the largest network address found in the result set. + */ + private long findMaxAddress() { + return result.stream().mapToLong(IPv4Subnet::getLongAddress).max().getAsLong(); + } + + /** + * Produces candidate subnets capable of covering the current set at the given prefix length. + */ + private List findMergeCandidatesForPrefixLength(int prefixLength) { + //создаём подсети-кандидаты, которые покроют наш список + var maxAddress = findMaxAddress(); + + var mask = IPv4Util.calculateMask(prefixLength); + var firstAddress = (findMinAddress() & mask); + var lastAddress = (maxAddress & mask); + + var candidates = new ArrayList(); + var candidateAddress = firstAddress; + do { + var candidate = new IPv4Subnet(candidateAddress, prefixLength); + candidates.add(candidate); + if (candidates.size() > result.size()) { + throw new IllegalStateException("Too many IPv4 addresses when trying to summary " + result.size() + " subnets"); + } + //поиск следующего адреса кандидата + var nextAddress = candidateAddress + candidate.count(); + var nextSubnet = result.stream().filter(subnet -> { + var address = subnet.getLongAddress(); + return subnet.getPrefixLength() > prefixLength && address >= nextAddress && address <= maxAddress; + }).min(Comparator.comparingLong(IPv4Subnet::getLongAddress)).stream().findFirst(); + if (nextSubnet.isEmpty()) { + break; + } + candidateAddress = IPv4Util.createSubnetOverlapping(nextSubnet.get().getLongAddress(), prefixLength).getLongAddress(); + } while (candidateAddress <= lastAddress); + return candidates; + } + + /** + * Tests whether the candidate subnet meets utilisation requirements and performs the merge. + */ + private boolean testCandidate(IPv4Subnet candidate, int usePercentage) { + if (result.contains(candidate)) { + return false; + } + var min = candidate.getLongAddress(); + var max = candidate.getLongAddress() + candidate.count() - 1; + var overlapped = new ArrayList(); + var used = new AtomicLong(0L); + result.forEach(child -> { + if (child.getLongAddress() < min || child.getLongAddress() > max) { + return; + } + + if (candidate.overlaps(child)) { + overlapped.add(child); + used.addAndGet(child.count()); + } + }); + + if (100.0 * used.get() / candidate.count() >= usePercentage) { + //подсеть подходит под критерий + overlapped.forEach(subnet -> { + if (source.contains(subnet)) { + mergedSubnets.add(subnet); + } + }); + result.removeAll(overlapped); + result.add(candidate); + return true; + } + return false; + } + + /** + * Iteratively merges subnets that satisfy the utilisation threshold. + */ + private void summaryWithUsage(int usePercentage) { + if (result.isEmpty() || usePercentage >= 100 || usePercentage <= 0) { + return; + } + + var found = new AtomicBoolean(); + do { + found.set(false); + if (result.size() < 2) { + break; + } + var prefixMin = findMinPrefixLength(); + var prefixMax = findMaxPrefixLength(); + for (var testPrefixLength = prefixMin - 1; testPrefixLength < prefixMax; testPrefixLength++) { + //создаём подсети-кандидаты, которые покроют наш список + var candidates = findMergeCandidatesForPrefixLength(testPrefixLength); + + candidates.forEach(candidate -> { + if (testCandidate(candidate, usePercentage)) { + found.set(true); + } + }); + + if (candidates.stream().anyMatch(result::contains)) { + //если был добавлен хотя бы 1 кандидат, то нужно пересчитать maxPrefix + prefixMax = result.stream().mapToInt(IPv4Subnet::getPrefixLength).max().getAsInt(); + } + } + } while (found.get()); + } + } + + /** + * Summaries the provided subnets using merging heuristics and utilisation thresholds. + */ + public static SummarisationResult summarySubnets(Collection subnets, int usePercentage) { + return new SubnetSummaryUtility(subnets, usePercentage); + } + + /** + * Creates a subnet covering the given address at the specified prefix length. + */ + private static IPv4Subnet createSubnetOverlapping(long address, int prefixLength) { + return new IPv4Subnet(address & IPv4Util.calculateMask(prefixLength), prefixLength); + } + + /** + * Returns the number of addresses represented by a prefix length. + */ + public static long calculateCountForPrefixLength(long prefixLength) { + return 1L << (32L - prefixLength); + } + + +} diff --git a/src/main/java/ru/kirillius/pf/sdn/Util/Wait.java b/src/main/java/ru/kirillius/pf/sdn/Util/Wait.java new file mode 100644 index 0000000..5461240 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/Util/Wait.java @@ -0,0 +1,41 @@ +package ru.kirillius.pf.sdn.Util; + +import java.time.Duration; +import java.util.function.Supplier; + +/** + * Provides blocking wait utilities for polling boolean conditions with interruption support. + */ +public final class Wait { + private Wait() { + throw new AssertionError(); + } + + /** + * Blocks until the supplied condition becomes true or the thread is interrupted. + */ + public static void until(Supplier condition) throws InterruptedException { + if (condition == null) { + return; + } + + while (!condition.get() && !Thread.currentThread().isInterrupted()) { + Thread.sleep(Duration.ofSeconds(1)); + } + } + + /** + * Blocks while the supplied condition remains true or until interrupted. + */ + public static void when(Supplier condition) throws InterruptedException { + if (condition == null) { + return; + } + + while (condition.get() && !Thread.currentThread().isInterrupted()) { + Thread.sleep(Duration.ofSeconds(1)); + } + } + + +} diff --git a/src/main/java/ru/kirillius/pf/sdn/Util/WorkerUtils.java b/src/main/java/ru/kirillius/pf/sdn/Util/WorkerUtils.java new file mode 100644 index 0000000..488f07a --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/Util/WorkerUtils.java @@ -0,0 +1,39 @@ +package ru.kirillius.pf.sdn.Util; + +import java.time.Duration; +import java.time.Instant; +import java.util.function.BooleanSupplier; +import java.util.function.Supplier; + +public class WorkerUtils { + @SuppressWarnings("BusyWait") + public static void waitForCondition(BooleanSupplier condition, long checkInterval) { + while (!Thread.interrupted() && !condition.getAsBoolean()) { + try { + Thread.sleep(checkInterval); + } catch (InterruptedException e) { + return; + } + } + } + + public static void waitForCondition(BooleanSupplier condition) { + waitForCondition(condition, 1000); + } + + public static void waitForTimer(Supplier duration) { + var start = Instant.now(); + waitForCondition(() -> Duration.between(start, Instant.now()).compareTo(duration.get()) > 0, 1000); + } + + public abstract static class LoopedWorker implements Runnable { + protected abstract void doWork(); + + @Override + public final void run() { + while (!Thread.interrupted()) { + doWork(); + } + } + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/Networking/AutonomousSystemResolver.java b/src/main/java/ru/kirillius/pf/sdn/api/Networking/AutonomousSystemResolver.java new file mode 100644 index 0000000..f011b59 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/Networking/AutonomousSystemResolver.java @@ -0,0 +1,16 @@ +package ru.kirillius.pf.sdn.api.Networking; + +import ru.kirillius.pf.sdn.api.dto.IPQueryInfo; + +import java.net.Inet4Address; +import java.util.List; +import java.util.concurrent.Future; + +public interface AutonomousSystemResolver { + + void registerApi(AutonomousSystemResolverApi api); + + Future> getPrefixes(int as); + + Future getAddressInfo(Inet4Address address); +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/Networking/AutonomousSystemResolverApi.java b/src/main/java/ru/kirillius/pf/sdn/api/Networking/AutonomousSystemResolverApi.java new file mode 100644 index 0000000..862c3f2 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/Networking/AutonomousSystemResolverApi.java @@ -0,0 +1,19 @@ +package ru.kirillius.pf.sdn.api.Networking; + +import ru.kirillius.pf.sdn.api.dto.IPQueryInfo; + +import java.net.Inet4Address; +import java.util.List; + +/** + * Abstraction for retrieving prefixes announced by a specific autonomous system. + */ +public interface AutonomousSystemResolverApi { + /** + * Returns IPv4 subnets originated by the provided autonomous system number. + */ + List getPrefixes(int as); + + IPQueryInfo queryAddress(Inet4Address address); + +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/Networking/DomainResolver.java b/src/main/java/ru/kirillius/pf/sdn/api/Networking/DomainResolver.java new file mode 100644 index 0000000..144ba5d --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/Networking/DomainResolver.java @@ -0,0 +1,30 @@ +package ru.kirillius.pf.sdn.api.Networking; + +import java.net.Inet4Address; +import java.util.List; +import java.util.concurrent.Future; + +public interface DomainResolver { + /** + * Resolve domain on specified server + * @param domain + * @param server + * @return + */ + Future> getAddresses(String domain, Inet4Address server); + + /** + * Resolve domain using default servers + * @param domain + * @return + */ + Future> getAddresses(String domain); + + /** + * Resolve domain using specified servers + * @param domain + * @param servers + * @return + */ + Future> getAddresses(String domain, List servers); +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/Networking/IPv4Subnet.java b/src/main/java/ru/kirillius/pf/sdn/api/Networking/IPv4Subnet.java new file mode 100644 index 0000000..4c4148b --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/Networking/IPv4Subnet.java @@ -0,0 +1,122 @@ +package ru.kirillius.pf.sdn.api.Networking; + +import com.fasterxml.jackson.annotation.JsonCreator; +import com.fasterxml.jackson.annotation.JsonValue; +import lombok.Getter; +import ru.kirillius.pf.sdn.Util.IPv4Util; + +import java.net.Inet4Address; +import java.util.Objects; +import java.util.regex.Pattern; + +/** + * Immutable representation of an IPv4 subnet serialized as a CIDR string. + */ +public class IPv4Subnet { + + @Getter + private final long longAddress; + @Getter + private final int prefixLength; + + /** + * Compares two subnets for equality based on address and prefix length. + */ + @Override + public boolean equals(Object o) { + if (o == null || getClass() != o.getClass()) return false; + IPv4Subnet that = (IPv4Subnet) o; + return longAddress == that.longAddress && prefixLength == that.prefixLength; + } + + /** + * Computes a hash code using address and prefix length. + */ + @Override + public int hashCode() { + return Objects.hash(longAddress, prefixLength); + } + + /** + * Parses a subnet from textual CIDR notation. + */ + @JsonCreator + public IPv4Subnet(String subnet) { + var split = subnet.split(Pattern.quote("/")); + if (split.length != 2) { + throw new IllegalArgumentException("Invalid subnet: " + subnet); + } + var prefix = Integer.parseInt(split[1]); + IPv4Util.validatePrefix(prefix); + + longAddress = IPv4Util.ipAddressToLong(split[0]); + prefixLength = prefix; + } + + public IPv4Subnet(Inet4Address address, int prefixLength) { + IPv4Util.validatePrefix(prefixLength); + longAddress = IPv4Util.ipAddressToLong(address.getHostAddress()); + this.prefixLength = prefixLength; + } + + /** + * Creates a subnet from a numeric address and prefix length. + */ + public IPv4Subnet(long longAddress, int prefixLength) { + this.longAddress = longAddress; + this.prefixLength = prefixLength; + } + + /** + * Creates a subnet from a dotted address string and prefix length. + */ + public IPv4Subnet(String address, int prefixLength) { + IPv4Util.validatePrefix(prefixLength); + + this.longAddress = IPv4Util.ipAddressToLong(address); + this.prefixLength = prefixLength; + } + + /** + * Calculates the number of addresses within the subnet. + */ + public long count() { + return IPv4Util.calculateCountForPrefixLength(prefixLength); + } + + /** + * Returns the dotted-decimal representation of the subnet's network address. + */ + public String getAddress() { + return IPv4Util.longToIpAddress(longAddress); + } + + /** + * Formats the subnet using CIDR notation. + */ + @JsonValue + @Override + public String toString() { + return getAddress() + '/' + prefixLength; + } + + /** + * Determines whether this subnet overlaps with another. + */ + public boolean overlaps(IPv4Subnet subnet) { + var minPrefixLength = Math.min(prefixLength, subnet.prefixLength); + var commonMask = IPv4Util.calculateMask(minPrefixLength); + if (minPrefixLength != prefixLength) { + return false; //can't overlap larger prefix + } + + return (longAddress & commonMask) == (subnet.longAddress & commonMask); + } + + public boolean contains(String address) { + var longAddress = IPv4Util.ipAddressToLong(address); + var commonMask = IPv4Util.calculateMask(prefixLength); + return (this.longAddress & commonMask) == (longAddress & commonMask); + } + +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/Networking/NetworkScope.java b/src/main/java/ru/kirillius/pf/sdn/api/Networking/NetworkScope.java new file mode 100644 index 0000000..92f69b8 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/Networking/NetworkScope.java @@ -0,0 +1,62 @@ +package ru.kirillius.pf.sdn.api.Networking; + +import lombok.*; + +import java.util.ArrayList; +import java.util.List; +import java.util.Objects; + +/** + * Container for grouped network identifiers used to configure filtering and subscriptions. + */ +@AllArgsConstructor +@NoArgsConstructor +@Builder +public class NetworkScope { + @Getter + @Setter + private String description = ""; + @Getter + @Setter + private List ASN = new ArrayList<>(); + @Getter + @Setter + private List subnets = new ArrayList<>(); + @Getter + @Setter + private List domains = new ArrayList<>(); + + @Getter + @Setter + private boolean resolveDomains = false; + + @Override + public boolean equals(Object o) { + if (!(o instanceof NetworkScope that)) return false; + return Objects.equals(ASN, that.ASN) && Objects.equals(subnets, that.subnets) && Objects.equals(domains, that.domains); + } + + @Override + public int hashCode() { + return Objects.hash(ASN, subnets, domains); + } + + /** + * Clears all stored network identifiers. + */ + public void clear() { + ASN.clear(); + subnets.clear(); + domains.clear(); + } + + /** + * Adds all resources from the provided bundle into this bundle. + */ + public void add(NetworkScope networkScope) { + ASN.addAll(networkScope.getASN()); + subnets.addAll(networkScope.getSubnets()); + domains.addAll(networkScope.getDomains()); + } + +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/CacheFallbackProviderProtocol.java b/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/CacheFallbackProviderProtocol.java new file mode 100644 index 0000000..f86daca --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/CacheFallbackProviderProtocol.java @@ -0,0 +1,28 @@ +package ru.kirillius.pf.sdn.api.Networking.Subscriptions; + +import lombok.RequiredArgsConstructor; +import ru.kirillius.pf.sdn.entity.SubscriptionProviderConfig; +import ru.kirillius.pf.sdn.repository.SubscriptionCacheRepository; + +import java.util.Collections; + +@RequiredArgsConstructor +public class CacheFallbackProviderProtocol implements SubscriptionProviderProtocol { + @Override + public String getIdentifier() { + return "Fallback:cache"; + } + + @Override + public PropertyCollection getProperties() { + return new PropertyCollection(); + } + + private final SubscriptionCacheRepository cacheRepository; + + @Override + public SubscriptionQueryResult query(SubscriptionProviderConfig config) { + var found = cacheRepository.getAllByNameStartsWith(config.getName()); + return new SubscriptionQueryResult(Collections.unmodifiableList(found), Collections.emptyList()); + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/Descriptor.java b/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/Descriptor.java new file mode 100644 index 0000000..e7117d7 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/Descriptor.java @@ -0,0 +1,4 @@ +package ru.kirillius.pf.sdn.api.Networking.Subscriptions; + +public record Descriptor(PropertyType type, boolean list, String description) { +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/PropertyCollection.java b/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/PropertyCollection.java new file mode 100644 index 0000000..c5240e7 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/PropertyCollection.java @@ -0,0 +1,7 @@ +package ru.kirillius.pf.sdn.api.Networking.Subscriptions; + +import java.util.HashMap; + +public class PropertyCollection extends HashMap { + +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/PropertyType.java b/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/PropertyType.java new file mode 100644 index 0000000..6c1c862 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/PropertyType.java @@ -0,0 +1,9 @@ +package ru.kirillius.pf.sdn.api.Networking.Subscriptions; + +public enum PropertyType { + Integer, + String, + Boolean, + Float, + Subscription +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/Subscription.java b/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/Subscription.java new file mode 100644 index 0000000..058c89c --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/Subscription.java @@ -0,0 +1,9 @@ +package ru.kirillius.pf.sdn.api.Networking.Subscriptions; + +import ru.kirillius.pf.sdn.api.Networking.NetworkScope; + +public interface Subscription { + String getName(); + + NetworkScope getScope(); +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/SubscriptionFactory.java b/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/SubscriptionFactory.java new file mode 100644 index 0000000..1bb9dbd --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/SubscriptionFactory.java @@ -0,0 +1,7 @@ +package ru.kirillius.pf.sdn.api.Networking.Subscriptions; + +import ru.kirillius.pf.sdn.api.Networking.NetworkScope; + +public interface SubscriptionFactory { + Subscription create(String name, NetworkScope scope); +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/SubscriptionProvider.java b/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/SubscriptionProvider.java new file mode 100644 index 0000000..4b2045b --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/SubscriptionProvider.java @@ -0,0 +1,32 @@ +package ru.kirillius.pf.sdn.api.Networking.Subscriptions; + +import lombok.Builder; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import ru.kirillius.pf.sdn.entity.SubscriptionProviderConfig; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; + +@Slf4j +@RequiredArgsConstructor +@Builder +public class SubscriptionProvider { + private final SubscriptionProviderProtocol protocol; + private final SubscriptionProviderConfig config; + + public List getSubscriptions() { + return Collections.unmodifiableList(subscriptions); + } + + private final List subscriptions = new ArrayList<>(); + + public List update() { + var result = protocol.query(config); + subscriptions.clear(); + subscriptions.addAll(result.subscriptions()); + return result.updated(); + } + +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/SubscriptionProviderProtocol.java b/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/SubscriptionProviderProtocol.java new file mode 100644 index 0000000..79ec238 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/SubscriptionProviderProtocol.java @@ -0,0 +1,11 @@ +package ru.kirillius.pf.sdn.api.Networking.Subscriptions; + +import ru.kirillius.pf.sdn.entity.SubscriptionProviderConfig; + +public interface SubscriptionProviderProtocol { + String getIdentifier(); + + PropertyCollection getProperties(); + + SubscriptionQueryResult query(SubscriptionProviderConfig config); +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/SubscriptionQueryResult.java b/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/SubscriptionQueryResult.java new file mode 100644 index 0000000..763d331 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/Networking/Subscriptions/SubscriptionQueryResult.java @@ -0,0 +1,6 @@ +package ru.kirillius.pf.sdn.api.Networking.Subscriptions; + +import java.util.List; + +public record SubscriptionQueryResult(List subscriptions, List updated) { +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/dto/IPQueryInfo.java b/src/main/java/ru/kirillius/pf/sdn/api/dto/IPQueryInfo.java new file mode 100644 index 0000000..1868a3a --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/dto/IPQueryInfo.java @@ -0,0 +1,15 @@ +package ru.kirillius.pf.sdn.api.dto; + +import lombok.Builder; +import lombok.Getter; +import ru.kirillius.pf.sdn.api.Networking.IPv4Subnet; + +import java.util.List; + +@Builder +public class IPQueryInfo { + @Getter + private List ASN; + @Getter + private List prefixes; +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/Action.java b/src/main/java/ru/kirillius/pf/sdn/api/pipeline/Action.java new file mode 100644 index 0000000..9344254 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/pipeline/Action.java @@ -0,0 +1,44 @@ +package ru.kirillius.pf.sdn.api.pipeline; + +import lombok.Builder; +import lombok.Getter; +import lombok.Setter; +import lombok.SneakyThrows; +import ru.kirillius.pf.sdn.api.Networking.NetworkScope; + +import java.util.Properties; +import java.util.function.Function; + +@Builder +public class Action implements Function { + + private PipelineFunction function; + + public Action(PipelineFunction function, Class functionClass, Properties properties) { + this.function = function; + this.functionClass = functionClass; + this.properties = properties; + instantiateFunction(); + } + + @SneakyThrows + private void instantiateFunction() { + function = functionClass.getConstructor().newInstance(); + } + + public void setFunctionClass(Class functionClass) { + this.functionClass = functionClass; + instantiateFunction(); + } + + @Getter + private Class functionClass; + @Getter + @Setter + private Properties properties; + + @Override + public NetworkScope apply(NetworkScope source) { + return function.apply(source, properties); + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/InputResource.java b/src/main/java/ru/kirillius/pf/sdn/api/pipeline/InputResource.java new file mode 100644 index 0000000..5ae223a --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/pipeline/InputResource.java @@ -0,0 +1,5 @@ +package ru.kirillius.pf.sdn.api.pipeline; + +public interface InputResource extends PipelineFunction { + +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/OutputResource.java b/src/main/java/ru/kirillius/pf/sdn/api/pipeline/OutputResource.java new file mode 100644 index 0000000..71583d7 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/pipeline/OutputResource.java @@ -0,0 +1,7 @@ +package ru.kirillius.pf.sdn.api.pipeline; + +import ru.kirillius.pf.sdn.api.Networking.NetworkScope; + +public interface OutputResource extends PipelineFunction { + void apply(NetworkScope config); +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/PipelineFunction.java b/src/main/java/ru/kirillius/pf/sdn/api/pipeline/PipelineFunction.java new file mode 100644 index 0000000..deb83c3 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/pipeline/PipelineFunction.java @@ -0,0 +1,9 @@ +package ru.kirillius.pf.sdn.api.pipeline; + +import ru.kirillius.pf.sdn.api.Networking.NetworkScope; + +import java.util.Properties; + +public interface PipelineFunction { + NetworkScope apply(NetworkScope source, Properties properties); +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/ProcessingPipeline.java b/src/main/java/ru/kirillius/pf/sdn/api/pipeline/ProcessingPipeline.java new file mode 100644 index 0000000..5a3ddbb --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/pipeline/ProcessingPipeline.java @@ -0,0 +1,52 @@ +package ru.kirillius.pf.sdn.api.pipeline; + +import lombok.AllArgsConstructor; +import lombok.Getter; +import lombok.NoArgsConstructor; +import lombok.Setter; +import ru.kirillius.pf.sdn.api.Networking.NetworkScope; + +import java.util.List; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; + +@AllArgsConstructor +@NoArgsConstructor +public class ProcessingPipeline { + @Getter + @Setter + private List startConditions; + + private final AtomicBoolean running = new AtomicBoolean(false); + @Getter + @Setter + private List actions; + private AtomicInteger currentStep = new AtomicInteger(0); + + public boolean isRunning() { + return running.get(); + } + + public int getStepCount() { + return actions.size(); + } + + public int getCurrentStep() { + return currentStep.get(); + } + + public void execute() { + running.set(true); + currentStep.set(0); + try { + var config = new AtomicReference<>(new NetworkScope()); + actions.forEach(action -> { + currentStep.incrementAndGet(); + config.set(action.apply(config.get())); + }); + } finally { + running.set(false); + } + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/ResourceProcessor.java b/src/main/java/ru/kirillius/pf/sdn/api/pipeline/ResourceProcessor.java new file mode 100644 index 0000000..8e64e07 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/pipeline/ResourceProcessor.java @@ -0,0 +1,8 @@ + +package ru.kirillius.pf.sdn.api.pipeline; + +import ru.kirillius.pf.sdn.api.Networking.NetworkScope; + +public interface ResourceProcessor extends PipelineFunction { + NetworkScope apply(NetworkScope config); +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/StartCondition.java b/src/main/java/ru/kirillius/pf/sdn/api/pipeline/StartCondition.java new file mode 100644 index 0000000..17df692 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/pipeline/StartCondition.java @@ -0,0 +1,13 @@ +package ru.kirillius.pf.sdn.api.pipeline; + +import lombok.Getter; +import lombok.Setter; + +import java.util.Properties; + +@Getter +@Setter +public class StartCondition { + private TriggerType trigger; + private Properties properties; +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/TriggerType.java b/src/main/java/ru/kirillius/pf/sdn/api/pipeline/TriggerType.java new file mode 100644 index 0000000..419e1e5 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/pipeline/TriggerType.java @@ -0,0 +1,8 @@ +package ru.kirillius.pf.sdn.api.pipeline; + +public enum TriggerType { + Manual, + Interval, + OnSubscriptionUpdate, + OnStart +} diff --git a/src/main/java/ru/kirillius/pf/sdn/config/BeanConfiguration.java b/src/main/java/ru/kirillius/pf/sdn/config/BeanConfiguration.java new file mode 100644 index 0000000..407b6b6 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/config/BeanConfiguration.java @@ -0,0 +1,23 @@ +package ru.kirillius.pf.sdn.config; + +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.security.crypto.bcrypt.BCryptPasswordEncoder; +import org.springframework.security.crypto.password.PasswordEncoder; + +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; + +@Configuration +public class BeanConfiguration { + + @Bean + public ExecutorService executor() { + return Executors.newFixedThreadPool(Math.max(Runtime.getRuntime().availableProcessors() * 2, 10)); + } + + @Bean + public PasswordEncoder passwordEncoder() { + return new BCryptPasswordEncoder(); + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/config/DataInitializer.java b/src/main/java/ru/kirillius/pf/sdn/config/DataInitializer.java new file mode 100644 index 0000000..f09a040 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/config/DataInitializer.java @@ -0,0 +1,19 @@ +package ru.kirillius.pf.sdn.config; + +import lombok.RequiredArgsConstructor; +import org.springframework.boot.context.event.ApplicationReadyEvent; +import org.springframework.context.event.EventListener; +import org.springframework.stereotype.Component; +import ru.kirillius.pf.sdn.service.AuthService; + +@Component +@RequiredArgsConstructor +public class DataInitializer { + + private final AuthService authService; + + @EventListener(ApplicationReadyEvent.class) + public void init() { + authService.createDefaultUserIfAbsent(); + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/config/JwtService.java b/src/main/java/ru/kirillius/pf/sdn/config/JwtService.java new file mode 100644 index 0000000..b9abbb1 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/config/JwtService.java @@ -0,0 +1,50 @@ +package ru.kirillius.pf.sdn.config; + +import io.jsonwebtoken.Claims; +import io.jsonwebtoken.JwtException; +import io.jsonwebtoken.Jwts; +import io.jsonwebtoken.security.Keys; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Component; + +import javax.crypto.SecretKey; +import java.nio.charset.StandardCharsets; +import java.util.Date; +import java.util.Optional; + +@Component +public class JwtService { + + private final SecretKey key; + private final long expirationMs; + + public JwtService( + @Value("${jwt.secret}") String secret, + @Value("${jwt.expiration-ms:2592000000}") long expirationMs) { + this.key = Keys.hmacShaKeyFor(secret.getBytes(StandardCharsets.UTF_8)); + this.expirationMs = expirationMs; + } + + public String generateToken(Long userId) { + Date now = new Date(); + return Jwts.builder() + .subject(String.valueOf(userId)) + .issuedAt(now) + .expiration(new Date(now.getTime() + expirationMs)) + .signWith(key) + .compact(); + } + + public Optional validateToken(String token) { + try { + Claims claims = Jwts.parser() + .verifyWith(key) + .build() + .parseSignedClaims(token) + .getPayload(); + return Optional.of(Long.parseLong(claims.getSubject())); + } catch (JwtException | NumberFormatException e) { + return Optional.empty(); + } + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/config/WebConfig.java b/src/main/java/ru/kirillius/pf/sdn/config/WebConfig.java new file mode 100644 index 0000000..846bbd3 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/config/WebConfig.java @@ -0,0 +1,16 @@ +package ru.kirillius.pf.sdn.config; + +import org.springframework.context.annotation.Configuration; +import org.springframework.format.FormatterRegistry; +import org.springframework.web.servlet.config.annotation.WebMvcConfigurer; +import ru.kirillius.pf.sdn.api.Networking.IPv4Subnet; + +@Configuration +public class WebConfig implements WebMvcConfigurer { + + @Override + public void addFormatters(FormatterRegistry registry) { + registry.addConverter(String.class, IPv4Subnet.class, IPv4Subnet::new); + registry.addConverter(IPv4Subnet.class, String.class, IPv4Subnet::toString); + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/controller/ApiTokenController.java b/src/main/java/ru/kirillius/pf/sdn/controller/ApiTokenController.java new file mode 100644 index 0000000..46365b7 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/controller/ApiTokenController.java @@ -0,0 +1,76 @@ +package ru.kirillius.pf.sdn.controller; + +import lombok.RequiredArgsConstructor; +import org.springframework.http.HttpStatus; +import org.springframework.http.ResponseEntity; +import org.springframework.web.bind.annotation.*; +import ru.kirillius.pf.sdn.dto.ApiTokenResponse; +import ru.kirillius.pf.sdn.dto.CreateApiTokenRequest; +import ru.kirillius.pf.sdn.dto.ErrorResponse; +import ru.kirillius.pf.sdn.service.AuthService; +import ru.kirillius.pf.sdn.service.SessionService; + +import java.util.List; + +@RestController +@RequestMapping("/api/tokens") +@RequiredArgsConstructor +public class ApiTokenController { + + private final AuthService authService; + private final SessionService sessionService; + + @PostMapping + public ResponseEntity create( + @RequestHeader("Authorization") String authHeader, + @RequestBody CreateApiTokenRequest request) { + + var maybeUserId = sessionService.resolveUserId(authHeader); + if (maybeUserId.isEmpty()) { + return ResponseEntity.status(HttpStatus.UNAUTHORIZED) + .body(new ErrorResponse("Invalid or missing token")); + } + + if (request.getName() == null || request.getName().isBlank()) { + return ResponseEntity.badRequest() + .body(new ErrorResponse("Name is required")); + } + + String rawToken = authService.generateApiToken(maybeUserId.get(), request.getName().trim()); + return ResponseEntity.ok(new ApiTokenResponse(null, request.getName().trim(), rawToken, true, null, null)); + } + + @GetMapping + public ResponseEntity list(@RequestHeader("Authorization") String authHeader) { + var maybeUserId = sessionService.resolveUserId(authHeader); + if (maybeUserId.isEmpty()) { + return ResponseEntity.status(HttpStatus.UNAUTHORIZED) + .body(new ErrorResponse("Invalid or missing token")); + } + + var tokens = authService.getUserApiTokens(maybeUserId.get()).stream() + .map(t -> new ApiTokenResponse(t.getId(), t.getName(), null, t.isEnabled(), t.getCreatedAt(), t.getLastUsedAt())) + .toList(); + return ResponseEntity.ok(tokens); + } + + @DeleteMapping("/{id}") + public ResponseEntity delete( + @RequestHeader("Authorization") String authHeader, + @PathVariable Long id) { + + var maybeUserId = sessionService.resolveUserId(authHeader); + if (maybeUserId.isEmpty()) { + return ResponseEntity.status(HttpStatus.UNAUTHORIZED) + .body(new ErrorResponse("Invalid or missing token")); + } + + boolean deleted = authService.deleteApiToken(id, maybeUserId.get()); + if (!deleted) { + return ResponseEntity.status(HttpStatus.NOT_FOUND) + .body(new ErrorResponse("Token not found")); + } + + return ResponseEntity.noContent().build(); + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/controller/AuthController.java b/src/main/java/ru/kirillius/pf/sdn/controller/AuthController.java new file mode 100644 index 0000000..3deb898 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/controller/AuthController.java @@ -0,0 +1,68 @@ +package ru.kirillius.pf.sdn.controller; + +import lombok.RequiredArgsConstructor; +import org.springframework.http.HttpStatus; +import org.springframework.http.ResponseEntity; +import org.springframework.web.bind.annotation.*; +import ru.kirillius.pf.sdn.dto.AuthResponse; +import ru.kirillius.pf.sdn.dto.ChangePasswordRequest; +import ru.kirillius.pf.sdn.dto.ErrorResponse; +import ru.kirillius.pf.sdn.dto.LoginRequest; +import ru.kirillius.pf.sdn.service.AuthService; +import ru.kirillius.pf.sdn.service.SessionService; + +@RestController +@RequestMapping("/api/auth") +@RequiredArgsConstructor +public class AuthController { + + private final AuthService authService; + private final SessionService sessionService; + + @PostMapping("/login") + public ResponseEntity login(@RequestBody LoginRequest request) { + return authService.login(request.getUsername(), request.getPassword()) + .>map(token -> ResponseEntity.ok(new AuthResponse(token, request.getUsername()))) + .orElse(ResponseEntity.status(HttpStatus.UNAUTHORIZED) + .body(new ErrorResponse("Invalid username or password"))); + } + + @PostMapping("/logout") + public ResponseEntity logout(@RequestHeader("Authorization") String authHeader) { + var maybeToken = sessionService.extractToken(authHeader); + if (maybeToken.isEmpty()) { + return ResponseEntity.status(HttpStatus.UNAUTHORIZED) + .body(new ErrorResponse("Missing or invalid Authorization header")); + } + sessionService.invalidate(maybeToken.get()); + return ResponseEntity.ok().build(); + } + + @PostMapping("/change-password") + public ResponseEntity changePassword( + @RequestHeader("Authorization") String authHeader, + @RequestBody ChangePasswordRequest request) { + + var maybeToken = sessionService.extractToken(authHeader); + if (maybeToken.isEmpty()) { + return ResponseEntity.status(HttpStatus.UNAUTHORIZED) + .body(new ErrorResponse("Missing or invalid Authorization header")); + } + + String token = maybeToken.get(); + var maybeUserId = authService.validateToken(token); + if (maybeUserId.isEmpty()) { + return ResponseEntity.status(HttpStatus.UNAUTHORIZED) + .body(new ErrorResponse("Invalid token")); + } + + boolean changed = authService.changePassword(maybeUserId.get(), request.getOldPassword(), request.getNewPassword()); + if (changed) { + sessionService.invalidate(token); + return ResponseEntity.ok(new AuthResponse(null, null)); + } + + return ResponseEntity.status(HttpStatus.BAD_REQUEST) + .body(new ErrorResponse("Old password is incorrect")); + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/dto/ApiTokenResponse.java b/src/main/java/ru/kirillius/pf/sdn/dto/ApiTokenResponse.java new file mode 100644 index 0000000..3621e33 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/dto/ApiTokenResponse.java @@ -0,0 +1,19 @@ +package ru.kirillius.pf.sdn.dto; + +import lombok.AllArgsConstructor; +import lombok.Getter; +import lombok.Setter; + +import java.time.LocalDateTime; + +@Getter +@Setter +@AllArgsConstructor +public class ApiTokenResponse { + private Long id; + private String name; + private String token; + private boolean enabled; + private LocalDateTime createdAt; + private LocalDateTime lastUsedAt; +} diff --git a/src/main/java/ru/kirillius/pf/sdn/dto/AuthResponse.java b/src/main/java/ru/kirillius/pf/sdn/dto/AuthResponse.java new file mode 100644 index 0000000..6646e57 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/dto/AuthResponse.java @@ -0,0 +1,13 @@ +package ru.kirillius.pf.sdn.dto; + +import lombok.AllArgsConstructor; +import lombok.Getter; +import lombok.Setter; + +@Getter +@Setter +@AllArgsConstructor +public class AuthResponse { + private String token; + private String username; +} diff --git a/src/main/java/ru/kirillius/pf/sdn/dto/ChangePasswordRequest.java b/src/main/java/ru/kirillius/pf/sdn/dto/ChangePasswordRequest.java new file mode 100644 index 0000000..610bb53 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/dto/ChangePasswordRequest.java @@ -0,0 +1,15 @@ +package ru.kirillius.pf.sdn.dto; + +import lombok.AllArgsConstructor; +import lombok.Getter; +import lombok.NoArgsConstructor; +import lombok.Setter; + +@Getter +@Setter +@NoArgsConstructor +@AllArgsConstructor +public class ChangePasswordRequest { + private String oldPassword; + private String newPassword; +} diff --git a/src/main/java/ru/kirillius/pf/sdn/dto/CreateApiTokenRequest.java b/src/main/java/ru/kirillius/pf/sdn/dto/CreateApiTokenRequest.java new file mode 100644 index 0000000..3667b09 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/dto/CreateApiTokenRequest.java @@ -0,0 +1,14 @@ +package ru.kirillius.pf.sdn.dto; + +import lombok.AllArgsConstructor; +import lombok.Getter; +import lombok.NoArgsConstructor; +import lombok.Setter; + +@Getter +@Setter +@NoArgsConstructor +@AllArgsConstructor +public class CreateApiTokenRequest { + private String name; +} diff --git a/src/main/java/ru/kirillius/pf/sdn/dto/ErrorResponse.java b/src/main/java/ru/kirillius/pf/sdn/dto/ErrorResponse.java new file mode 100644 index 0000000..34b7161 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/dto/ErrorResponse.java @@ -0,0 +1,12 @@ +package ru.kirillius.pf.sdn.dto; + +import lombok.AllArgsConstructor; +import lombok.Getter; +import lombok.Setter; + +@Getter +@Setter +@AllArgsConstructor +public class ErrorResponse { + private String error; +} diff --git a/src/main/java/ru/kirillius/pf/sdn/dto/LoginRequest.java b/src/main/java/ru/kirillius/pf/sdn/dto/LoginRequest.java new file mode 100644 index 0000000..0ac778c --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/dto/LoginRequest.java @@ -0,0 +1,15 @@ +package ru.kirillius.pf.sdn.dto; + +import lombok.AllArgsConstructor; +import lombok.Getter; +import lombok.NoArgsConstructor; +import lombok.Setter; + +@Getter +@Setter +@NoArgsConstructor +@AllArgsConstructor +public class LoginRequest { + private String username; + private String password; +} diff --git a/src/main/java/ru/kirillius/pf/sdn/dto/ResolverCacheEntry.java b/src/main/java/ru/kirillius/pf/sdn/dto/ResolverCacheEntry.java new file mode 100644 index 0000000..e5453b4 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/dto/ResolverCacheEntry.java @@ -0,0 +1,18 @@ +package ru.kirillius.pf.sdn.dto; + +import lombok.Getter; +import lombok.NoArgsConstructor; +import lombok.Setter; +import ru.kirillius.pf.sdn.api.Networking.IPv4Subnet; + +import java.time.Instant; +import java.util.HashMap; +import java.util.Map; + +@Getter +@Setter +@NoArgsConstructor +public class ResolverCacheEntry { + private Instant lastUpdate = Instant.now(); + private Map addresses = new HashMap<>(); +} diff --git a/src/main/java/ru/kirillius/pf/sdn/entity/ApiToken.java b/src/main/java/ru/kirillius/pf/sdn/entity/ApiToken.java new file mode 100644 index 0000000..4a736e7 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/entity/ApiToken.java @@ -0,0 +1,42 @@ +package ru.kirillius.pf.sdn.entity; + +import jakarta.persistence.*; +import lombok.Getter; +import lombok.NoArgsConstructor; +import lombok.Setter; + +import java.time.LocalDateTime; + +@Entity +@Table(name = "api_tokens") +@Getter +@Setter +@NoArgsConstructor +public class ApiToken { + + @Id + @GeneratedValue(strategy = GenerationType.IDENTITY) + private Long id; + + @Column(nullable = false, length = 100) + private String name; + + @Column(nullable = false) + private String tokenHash; + + @Column(nullable = false) + private Long userId; + + @Column(nullable = false) + private boolean enabled = true; + + @Column(nullable = false, updatable = false) + private LocalDateTime createdAt; + + private LocalDateTime lastUsedAt; + + @PrePersist + protected void onCreate() { + this.createdAt = LocalDateTime.now(); + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/entity/AutonomousSystemCacheEntry.java b/src/main/java/ru/kirillius/pf/sdn/entity/AutonomousSystemCacheEntry.java new file mode 100644 index 0000000..1e19a9e --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/entity/AutonomousSystemCacheEntry.java @@ -0,0 +1,30 @@ +package ru.kirillius.pf.sdn.entity; + +import jakarta.persistence.*; +import lombok.*; +import ru.kirillius.pf.sdn.api.Networking.IPv4Subnet; + +import java.util.List; + +@Entity +@Table(name = "as_cache_entry") +@Getter +@Setter +@NoArgsConstructor +@Builder +@AllArgsConstructor +public class AutonomousSystemCacheEntry { + @Id + @GeneratedValue(strategy = GenerationType.IDENTITY) + private Long id; + + @Column(nullable = false, unique = true) + private int number; + + @ElementCollection + @CollectionTable(name = "as_cache_entry_prefixes", joinColumns = @JoinColumn(name = "entry_id")) + @Column(name = "prefix", nullable = false) + @Convert(converter = IPv4SubnetConverter.class) + private List prefixes; + +} diff --git a/src/main/java/ru/kirillius/pf/sdn/entity/DomainCacheEntry.java b/src/main/java/ru/kirillius/pf/sdn/entity/DomainCacheEntry.java new file mode 100644 index 0000000..70cb738 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/entity/DomainCacheEntry.java @@ -0,0 +1,32 @@ +package ru.kirillius.pf.sdn.entity; + +import jakarta.persistence.*; +import lombok.*; + +import java.net.Inet4Address; +import java.time.Instant; + +@Entity +@Table(name = "domain_cache_entry") +@Getter +@Setter +@NoArgsConstructor +@Builder +@AllArgsConstructor +public class DomainCacheEntry { + @Id + @GeneratedValue(strategy = GenerationType.IDENTITY) + private Long id; + + @Column(nullable = false) + private String domain; + + @Column(nullable = false) + private Instant lastUpdate; + + @Column(nullable = false) + private Instant lastSeen; + + @Column(nullable = false) + private Inet4Address address; +} diff --git a/src/main/java/ru/kirillius/pf/sdn/entity/IPv4SubnetConverter.java b/src/main/java/ru/kirillius/pf/sdn/entity/IPv4SubnetConverter.java new file mode 100644 index 0000000..896ee1e --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/entity/IPv4SubnetConverter.java @@ -0,0 +1,19 @@ +package ru.kirillius.pf.sdn.entity; + +import jakarta.persistence.AttributeConverter; +import jakarta.persistence.Converter; +import ru.kirillius.pf.sdn.api.Networking.IPv4Subnet; + +@Converter +public class IPv4SubnetConverter implements AttributeConverter { + + @Override + public String convertToDatabaseColumn(IPv4Subnet attribute) { + return attribute == null ? null : attribute.toString(); + } + + @Override + public IPv4Subnet convertToEntityAttribute(String dbData) { + return dbData == null ? null : new IPv4Subnet(dbData); + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/entity/NetworkScopeConverter.java b/src/main/java/ru/kirillius/pf/sdn/entity/NetworkScopeConverter.java new file mode 100644 index 0000000..78d65ff --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/entity/NetworkScopeConverter.java @@ -0,0 +1,33 @@ +package ru.kirillius.pf.sdn.entity; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.ObjectMapper; +import jakarta.persistence.AttributeConverter; +import jakarta.persistence.Converter; +import ru.kirillius.pf.sdn.api.Networking.NetworkScope; + +@Converter +public class NetworkScopeConverter implements AttributeConverter { + + private static final ObjectMapper MAPPER = new ObjectMapper(); + + @Override + public String convertToDatabaseColumn(NetworkScope attribute) { + if (attribute == null) return null; + try { + return MAPPER.writeValueAsString(attribute); + } catch (JsonProcessingException e) { + throw new RuntimeException("Failed to serialize NetworkScope", e); + } + } + + @Override + public NetworkScope convertToEntityAttribute(String dbData) { + if (dbData == null) return null; + try { + return MAPPER.readValue(dbData, NetworkScope.class); + } catch (JsonProcessingException e) { + throw new RuntimeException("Failed to deserialize NetworkScope", e); + } + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/entity/SubscriptionCacheEntry.java b/src/main/java/ru/kirillius/pf/sdn/entity/SubscriptionCacheEntry.java new file mode 100644 index 0000000..5bb69ac --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/entity/SubscriptionCacheEntry.java @@ -0,0 +1,26 @@ +package ru.kirillius.pf.sdn.entity; + +import jakarta.persistence.*; +import lombok.*; +import ru.kirillius.pf.sdn.api.Networking.NetworkScope; +import ru.kirillius.pf.sdn.api.Networking.Subscriptions.Subscription; + +@Entity +@Table(name = "subscription_cache") +@Getter +@Setter +@NoArgsConstructor +@Builder +@AllArgsConstructor +public class SubscriptionCacheEntry implements Subscription { + @Id + @GeneratedValue(strategy = GenerationType.IDENTITY) + private Long id; + + @Column(nullable = false, unique = true) + private String name; + + @Column(nullable = false, columnDefinition = "TEXT") + @Convert(converter = NetworkScopeConverter.class) + private NetworkScope scope; +} diff --git a/src/main/java/ru/kirillius/pf/sdn/entity/SubscriptionProviderConfig.java b/src/main/java/ru/kirillius/pf/sdn/entity/SubscriptionProviderConfig.java new file mode 100644 index 0000000..fd3d184 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/entity/SubscriptionProviderConfig.java @@ -0,0 +1,40 @@ +package ru.kirillius.pf.sdn.entity; + +import jakarta.persistence.*; +import lombok.*; + +import java.util.Properties; + +@Entity +@Table(name = "subscription_provider_config") +@Getter +@Setter +@NoArgsConstructor +@Builder +@AllArgsConstructor +public class SubscriptionProviderConfig { + @Id + @GeneratedValue(strategy = GenerationType.IDENTITY) + private Long id; + + @Column(nullable = false, unique = true) + private String name; + + @Column(nullable = false) + private String protocolId; + + @ElementCollection + @CollectionTable( + name = "subscription_provider_config_properties", + joinColumns = @JoinColumn(name = "subscription_provider_config_entry_id") + ) + @MapKeyColumn(name = "prop_key") + @Column(name = "prop_value", nullable = false) + private Properties properties; + + @Column(nullable = false) + private String shellScriptBefore; + + @Column(nullable = false) + private String shellScriptAfter; +} diff --git a/src/main/java/ru/kirillius/pf/sdn/entity/User.java b/src/main/java/ru/kirillius/pf/sdn/entity/User.java new file mode 100644 index 0000000..7296010 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/entity/User.java @@ -0,0 +1,37 @@ +package ru.kirillius.pf.sdn.entity; + +import jakarta.persistence.*; +import lombok.Getter; +import lombok.NoArgsConstructor; +import lombok.Setter; + +import java.time.LocalDateTime; + +@Entity +@Table(name = "users") +@Getter +@Setter +@NoArgsConstructor +public class User { + + @Id + @GeneratedValue(strategy = GenerationType.IDENTITY) + private Long id; + + @Column(nullable = false, unique = true, length = 50) + private String username; + + @Column(nullable = false) + private String passwordHash; + + @Column(nullable = false) + private boolean enabled = true; + + @Column(nullable = false, updatable = false) + private LocalDateTime createdAt; + + @PrePersist + protected void onCreate() { + this.createdAt = LocalDateTime.now(); + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/pipeline/SubscriptionInput.java b/src/main/java/ru/kirillius/pf/sdn/pipeline/SubscriptionInput.java new file mode 100644 index 0000000..d1fa17f --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/pipeline/SubscriptionInput.java @@ -0,0 +1,47 @@ +package ru.kirillius.pf.sdn.pipeline; + +import lombok.RequiredArgsConstructor; +import ru.kirillius.pf.sdn.api.Networking.IPv4Subnet; +import ru.kirillius.pf.sdn.api.Networking.NetworkScope; +import ru.kirillius.pf.sdn.api.pipeline.InputResource; +import ru.kirillius.pf.sdn.service.DomainUpdaterService; +import ru.kirillius.pf.sdn.service.SubscriptionService; + +import java.net.Inet4Address; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Properties; +import java.util.regex.Pattern; + +@RequiredArgsConstructor +public class SubscriptionInput implements InputResource { + + private final SubscriptionService subscriptionService; + private final DomainUpdaterService domainCacheService; + + @Override + public NetworkScope apply(NetworkScope source, Properties properties) { + var property = properties.getProperty("property", null); + var subscriptions = subscriptionService.getSubscriptions(); + if (property != null) { + var filter = Arrays.stream(property.split(Pattern.quote(","))).filter(s -> !s.isBlank()).toList(); + subscriptions = subscriptions.stream().filter(s -> filter.contains(s.getName())).toList(); + } + + var bundle = new NetworkScope(); + bundle.add(source); + + subscriptions.forEach(subscription -> { + var scope = subscription.getScope(); + bundle.add(scope); + var resolved = new ArrayList(); + if (scope.isResolveDomains()) { + scope.getDomains().forEach(domain -> resolved.addAll(domainCacheService.getActualAddresses(domain))); + } + if (!resolved.isEmpty()) { + bundle.getSubnets().addAll(resolved.stream().map(ip -> new IPv4Subnet(ip, 32)).toList()); + } + }); + return bundle; + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/plugins/LocalSubscriptionProvider/LocalSubscriptionProvider.java b/src/main/java/ru/kirillius/pf/sdn/plugins/LocalSubscriptionProvider/LocalSubscriptionProvider.java new file mode 100644 index 0000000..7dd6dd2 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/plugins/LocalSubscriptionProvider/LocalSubscriptionProvider.java @@ -0,0 +1,30 @@ +package ru.kirillius.pf.sdn.plugins.LocalSubscriptionProvider; + +import lombok.Getter; +import ru.kirillius.pf.sdn.api.Networking.Subscriptions.*; +import ru.kirillius.pf.sdn.entity.SubscriptionProviderConfig; + +public class LocalSubscriptionProvider implements SubscriptionProviderProtocol { + @Override + public String getIdentifier() { + return "local-filesystem"; + } + + @Getter + private final PropertyCollection properties; + + public LocalSubscriptionProvider() { + properties = new PropertyCollection(); + properties.put("path", new Descriptor( + PropertyType.String, + false, + "Path to files is local filesystem" + )); + } + + @Override + public SubscriptionQueryResult query(SubscriptionProviderConfig config) { + + return null; + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/plugins/Plugin.java b/src/main/java/ru/kirillius/pf/sdn/plugins/Plugin.java new file mode 100644 index 0000000..7ffbfe4 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/plugins/Plugin.java @@ -0,0 +1,6 @@ +package ru.kirillius.pf.sdn.plugins; + +public interface Plugin { + void onLoad(); + void onUnload(); +} diff --git a/src/main/java/ru/kirillius/pf/sdn/plugins/PluginLoader.java b/src/main/java/ru/kirillius/pf/sdn/plugins/PluginLoader.java new file mode 100644 index 0000000..409603a --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/plugins/PluginLoader.java @@ -0,0 +1,4 @@ +package ru.kirillius.pf.sdn.plugins; + +public class PluginLoader { +} diff --git a/src/main/java/ru/kirillius/pf/sdn/repository/ApiTokenRepository.java b/src/main/java/ru/kirillius/pf/sdn/repository/ApiTokenRepository.java new file mode 100644 index 0000000..fbaa9a3 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/repository/ApiTokenRepository.java @@ -0,0 +1,16 @@ +package ru.kirillius.pf.sdn.repository; + +import org.springframework.data.jpa.repository.JpaRepository; +import ru.kirillius.pf.sdn.entity.ApiToken; + +import java.util.List; +import java.util.Optional; + +public interface ApiTokenRepository extends JpaRepository { + + List findByUserId(Long userId); + + Optional findByIdAndUserId(Long id, Long userId); + + List findByEnabledTrue(); +} diff --git a/src/main/java/ru/kirillius/pf/sdn/repository/AutonomousSystemCacheRepository.java b/src/main/java/ru/kirillius/pf/sdn/repository/AutonomousSystemCacheRepository.java new file mode 100644 index 0000000..c911008 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/repository/AutonomousSystemCacheRepository.java @@ -0,0 +1,12 @@ +package ru.kirillius.pf.sdn.repository; + +import org.springframework.data.jpa.repository.JpaRepository; +import ru.kirillius.pf.sdn.entity.AutonomousSystemCacheEntry; + +import java.util.Optional; + +public interface AutonomousSystemCacheRepository extends JpaRepository { + Optional findByNumber(int number); + + boolean existsByNumber(int number); +} diff --git a/src/main/java/ru/kirillius/pf/sdn/repository/DomainCacheRepository.java b/src/main/java/ru/kirillius/pf/sdn/repository/DomainCacheRepository.java new file mode 100644 index 0000000..90a275e --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/repository/DomainCacheRepository.java @@ -0,0 +1,14 @@ +package ru.kirillius.pf.sdn.repository; + +import org.springframework.data.jpa.repository.JpaRepository; +import ru.kirillius.pf.sdn.entity.DomainCacheEntry; + +import java.util.List; + +public interface DomainCacheRepository extends JpaRepository { + boolean existsByDomain(String domain); + + List findAllByDomain(String domain); + + List getAllByDomain(String domain); +} diff --git a/src/main/java/ru/kirillius/pf/sdn/repository/SubscriptionCacheRepository.java b/src/main/java/ru/kirillius/pf/sdn/repository/SubscriptionCacheRepository.java new file mode 100644 index 0000000..477883c --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/repository/SubscriptionCacheRepository.java @@ -0,0 +1,11 @@ +package ru.kirillius.pf.sdn.repository; + +import org.springframework.data.jpa.repository.JpaRepository; +import ru.kirillius.pf.sdn.entity.SubscriptionCacheEntry; + +import java.util.List; + +public interface SubscriptionCacheRepository extends JpaRepository { + + List getAllByNameStartsWith(String s); +} diff --git a/src/main/java/ru/kirillius/pf/sdn/repository/SubscriptionProviderConfigRepository.java b/src/main/java/ru/kirillius/pf/sdn/repository/SubscriptionProviderConfigRepository.java new file mode 100644 index 0000000..b93cc06 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/repository/SubscriptionProviderConfigRepository.java @@ -0,0 +1,8 @@ +package ru.kirillius.pf.sdn.repository; + +import org.springframework.data.jpa.repository.JpaRepository; +import ru.kirillius.pf.sdn.entity.SubscriptionProviderConfig; + +public interface SubscriptionProviderConfigRepository extends JpaRepository { + +} diff --git a/src/main/java/ru/kirillius/pf/sdn/repository/UserRepository.java b/src/main/java/ru/kirillius/pf/sdn/repository/UserRepository.java new file mode 100644 index 0000000..cd136a7 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/repository/UserRepository.java @@ -0,0 +1,13 @@ +package ru.kirillius.pf.sdn.repository; + +import org.springframework.data.jpa.repository.JpaRepository; +import ru.kirillius.pf.sdn.entity.User; + +import java.util.Optional; + +public interface UserRepository extends JpaRepository { + + Optional findByUsername(String username); + + boolean existsByUsername(String username); +} diff --git a/src/main/java/ru/kirillius/pf/sdn/service/AuthService.java b/src/main/java/ru/kirillius/pf/sdn/service/AuthService.java new file mode 100644 index 0000000..7d3b456 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/service/AuthService.java @@ -0,0 +1,115 @@ +package ru.kirillius.pf.sdn.service; + +import lombok.RequiredArgsConstructor; +import org.springframework.security.crypto.password.PasswordEncoder; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; +import ru.kirillius.pf.sdn.config.JwtService; +import ru.kirillius.pf.sdn.entity.ApiToken; +import ru.kirillius.pf.sdn.entity.User; +import ru.kirillius.pf.sdn.repository.ApiTokenRepository; +import ru.kirillius.pf.sdn.repository.UserRepository; + +import java.security.SecureRandom; +import java.util.HexFormat; +import java.util.List; +import java.util.Optional; + +@Service +@RequiredArgsConstructor +public class AuthService { + + private final UserRepository userRepository; + private final ApiTokenRepository apiTokenRepository; + private final PasswordEncoder passwordEncoder; + private final JwtService jwtService; + + private static final int API_TOKEN_BYTES = 32; + + @Transactional + public void createDefaultUserIfAbsent() { + if (!userRepository.existsByUsername("admin")) { + User user = new User(); + user.setUsername("admin"); + user.setPasswordHash(passwordEncoder.encode("admin")); + userRepository.save(user); + } + } + + public Optional login(String username, String rawPassword) { + return userRepository.findByUsername(username) + .filter(User::isEnabled) + .filter(user -> passwordEncoder.matches(rawPassword, user.getPasswordHash())) + .map(user -> jwtService.generateToken(user.getId())); + } + + @Transactional + public boolean changePassword(Long userId, String oldPassword, String newPassword) { + return userRepository.findById(userId) + .filter(user -> passwordEncoder.matches(oldPassword, user.getPasswordHash())) + .map(user -> { + user.setPasswordHash(passwordEncoder.encode(newPassword)); + userRepository.save(user); + return true; + }) + .orElse(false); + } + + public Optional validateToken(String token) { + return jwtService.validateToken(token); + } + + public void invalidateToken(String token) { + // JWT is stateless — client just discards the token + } + + @Transactional + public String generateApiToken(Long userId, String name) { + SecureRandom random = new SecureRandom(); + byte[] bytes = new byte[API_TOKEN_BYTES]; + random.nextBytes(bytes); + String rawToken = HexFormat.of().formatHex(bytes); + + ApiToken apiToken = new ApiToken(); + apiToken.setName(name); + apiToken.setTokenHash(passwordEncoder.encode(rawToken)); + apiToken.setUserId(userId); + apiTokenRepository.save(apiToken); + + return apiToken.getId() + "." + rawToken; + } + + public Optional validateApiToken(String rawToken) { + for (ApiToken stored : apiTokenRepository.findByEnabledTrue()) { + if (passwordEncoder.matches(rawToken, stored.getTokenHash())) { + updateLastUsedAt(stored); + return Optional.of(stored.getUserId()); + } + } + return Optional.empty(); + } + + public Optional resolveUserId(String token) { + return validateToken(token) + .or(() -> validateApiToken(token)); + } + + @Transactional + public boolean deleteApiToken(Long tokenId, Long userId) { + return apiTokenRepository.findByIdAndUserId(tokenId, userId) + .map(token -> { + apiTokenRepository.delete(token); + return true; + }) + .orElse(false); + } + + public List getUserApiTokens(Long userId) { + return apiTokenRepository.findByUserId(userId); + } + + private void updateLastUsedAt(ApiToken token) { + token.setLastUsedAt(java.time.LocalDateTime.now()); + apiTokenRepository.save(token); + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/service/AutonomousSystemCacheService.java b/src/main/java/ru/kirillius/pf/sdn/service/AutonomousSystemCacheService.java new file mode 100644 index 0000000..600365d --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/service/AutonomousSystemCacheService.java @@ -0,0 +1,37 @@ +package ru.kirillius.pf.sdn.service; + +import lombok.RequiredArgsConstructor; +import lombok.Setter; +import org.springframework.stereotype.Service; +import ru.kirillius.pf.sdn.api.Networking.IPv4Subnet; +import ru.kirillius.pf.sdn.entity.AutonomousSystemCacheEntry; +import ru.kirillius.pf.sdn.repository.AutonomousSystemCacheRepository; + +import java.util.List; +import java.util.Optional; + +@Setter +@RequiredArgsConstructor +@Service +public class AutonomousSystemCacheService { + private final AutonomousSystemCacheRepository repository; + + public void store(int number, List subnets) { + synchronized (repository) { + if (repository.existsByNumber(number)) { + repository.findByNumber(number).ifPresent(entry -> { + entry.setPrefixes(subnets); + repository.save(entry); + }); + } else { + repository.save(AutonomousSystemCacheEntry.builder().number(number).prefixes(subnets).build()); + } + } + } + + public Optional load(int number) { + synchronized (repository) { + return repository.findByNumber(number); + } + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/service/AutonomousSystemResolverService.java b/src/main/java/ru/kirillius/pf/sdn/service/AutonomousSystemResolverService.java new file mode 100644 index 0000000..3862ab1 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/service/AutonomousSystemResolverService.java @@ -0,0 +1,83 @@ +package ru.kirillius.pf.sdn.service; + +import lombok.RequiredArgsConstructor; +import lombok.Setter; +import ru.kirillius.pf.sdn.api.Networking.AutonomousSystemResolver; +import ru.kirillius.pf.sdn.api.Networking.AutonomousSystemResolverApi; +import ru.kirillius.pf.sdn.api.Networking.IPv4Subnet; +import ru.kirillius.pf.sdn.api.dto.IPQueryInfo; + +import java.io.IOException; +import java.net.Inet4Address; +import java.util.List; +import java.util.concurrent.Callable; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Future; +import java.util.concurrent.atomic.AtomicReference; + +@Setter +@RequiredArgsConstructor +public class AutonomousSystemResolverService implements AutonomousSystemResolver { + private final ExecutorService executor; + private final List apis = new CopyOnWriteArrayList<>(); + private final AtomicReference currentApi = new AtomicReference<>(); + + private synchronized void selectNextApi() { + if (apis.isEmpty()) { + return; + } + var index = apis.indexOf(currentApi.get()); + if (index == -1 || index == apis.size() - 1) { + currentApi.set(apis.getFirst()); + return; + } + currentApi.set(apis.get(index + 1)); + } + + @Override + public void registerApi(AutonomousSystemResolverApi api) { + apis.add(api); + if (currentApi.get() == null) { + currentApi.set(api); + } + } + + @Override + public Future> getPrefixes(int as) { + return executor.submit(() -> { + checkApis(); + return executeWithFallbackSelection(() -> currentApi.get().getPrefixes(as)); + }); + } + + private void checkApis() throws IOException { + if (apis.isEmpty()) { + throw new IOException("You're trying to resolve an AS without any resolvers"); + } + if (currentApi.get() == null) { + selectNextApi(); + } + } + + private T executeWithFallbackSelection(Callable callable) throws IOException { + var retries = 0; + while (retries < apis.size()) { + try { + return callable.call(); + } catch (Exception e) { + retries++; + selectNextApi(); + } + } + throw new IOException("Unable to perform operation on every AS resolver"); + } + + @Override + public Future getAddressInfo(Inet4Address address) { + return executor.submit(() -> { + checkApis(); + return executeWithFallbackSelection(() -> currentApi.get().queryAddress(address)); + }); + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/service/DomainCacheService.java b/src/main/java/ru/kirillius/pf/sdn/service/DomainCacheService.java new file mode 100644 index 0000000..70dacd7 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/service/DomainCacheService.java @@ -0,0 +1,63 @@ +package ru.kirillius.pf.sdn.service; + +import lombok.RequiredArgsConstructor; +import lombok.Setter; +import org.springframework.stereotype.Service; +import ru.kirillius.pf.sdn.entity.DomainCacheEntry; +import ru.kirillius.pf.sdn.repository.DomainCacheRepository; + +import java.net.Inet4Address; +import java.time.Duration; +import java.time.Instant; +import java.util.List; + +@Setter +@RequiredArgsConstructor +@Service +public class DomainCacheService { + private final DomainCacheRepository repository; + private static final int DOMAIN_TTL_HOURS = 12; + + public void update(String domain, List addresses) { + var minDate = Instant.now().minus(Duration.ofHours(DOMAIN_TTL_HOURS)); + synchronized (repository) { + var existing = repository.getAllByDomain(domain); + var existingAddresses = existing.stream().map(DomainCacheEntry::getAddress).toList(); + existing.forEach(entry -> { + if (addresses.contains(entry.getAddress())) { + entry.setLastUpdate(Instant.now()); + entry.setLastSeen(Instant.now()); + } else { + if (entry.getLastSeen().isBefore(minDate)) { + repository.delete(entry); + } else { + entry.setLastUpdate(Instant.now()); + } + } + }); + + addresses.stream() + .filter(address -> !existingAddresses.contains(address)) + .forEach(address -> repository.save( + DomainCacheEntry.builder() + .domain(domain) + .address(address) + .lastUpdate(Instant.now()) + .lastSeen(Instant.now()) + .build() + ) + ); + } + } + + public List load(String domain, boolean filterOutdated) { + var minDate = Instant.now().minus(Duration.ofHours(DOMAIN_TTL_HOURS)); + synchronized (repository) { + var found = repository.getAllByDomain(domain); + if (filterOutdated && !found.isEmpty()) { + found = found.stream().filter(entry -> entry.getLastSeen().isAfter(minDate)).toList(); + } + return found.stream().map(DomainCacheEntry::getAddress).toList(); + } + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/service/DomainResolverService.java b/src/main/java/ru/kirillius/pf/sdn/service/DomainResolverService.java new file mode 100644 index 0000000..8862361 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/service/DomainResolverService.java @@ -0,0 +1,139 @@ +package ru.kirillius.pf.sdn.service; + +import lombok.RequiredArgsConstructor; +import org.springframework.stereotype.Service; +import ru.kirillius.pf.sdn.Util.IPv4Util; +import ru.kirillius.pf.sdn.api.Networking.DomainResolver; +import ru.kirillius.pf.sdn.config.BeanConfiguration; + +import java.io.IOException; +import java.io.InputStream; +import java.net.Inet4Address; +import java.net.UnknownHostException; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.Collections; +import java.util.LinkedHashSet; +import java.util.List; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; + +@Service +@RequiredArgsConstructor +public class DomainResolverService implements DomainResolver { + private final static int SUCCESS_CODE = 0; + private final BeanConfiguration config; + private final ExecutorService executor; + private final RuntimeConfigService configService; + + /** + * Returns a list of IPv4 addresses resolved for the provided domain. + */ + private static List lookup(String domain, Inet4Address server) throws IOException, InterruptedException { + if (domain == null || domain.isBlank()) { + throw new IllegalArgumentException("Domain must not be null or blank"); + } + + var processBuilder = new ProcessBuilder("nslookup", domain.trim(), server.getHostAddress().trim()); + processBuilder.redirectErrorStream(true); + + var process = processBuilder.start(); + try { + if (!process.waitFor(10, TimeUnit.SECONDS)) { + process.destroyForcibly(); + throw new InterruptedException("Request timed out for domain: " + domain); + } + + var output = readStream(process.getInputStream()); + if (process.exitValue() != SUCCESS_CODE) { + throw new IOException("nslookup failed for domain " + domain + ": " + output); + } + return extractIPv4(output); + } finally { + process.destroy(); + } + + } + + private static List extractIPv4(String output) { + var addresses = new LinkedHashSet(); + var allowAddresses = false; + + for (var rawLine : output.split("\\R")) { + var line = rawLine.trim(); + if (line.isEmpty()) { + allowAddresses = false; + continue; + } + + if (line.startsWith("Name:")) { + allowAddresses = true; + continue; + } + + if (!allowAddresses) { + continue; + } + + if (line.startsWith("Address:") || line.startsWith("Addresses:")) { + var colonIndex = line.indexOf(':'); + if (colonIndex == -1 || colonIndex == line.length() - 1) { + continue; + } + + var tokens = line.substring(colonIndex + 1).trim().split("\\s+"); + for (var token : tokens) { + try { + IPv4Util.validateAddress(token); + addresses.add((Inet4Address) Inet4Address.getByName(token)); + } catch (IllegalArgumentException ignored) { + // skip invalid addresses such as IPv6 + } catch (UnknownHostException e) { + throw new RuntimeException(e); + } + } + } + } + + return new ArrayList<>(addresses); + } + + private static String readStream(InputStream stream) throws IOException { + try (stream) { + return new String(stream.readAllBytes(), StandardCharsets.UTF_8); + } + } + + @Override + public Future> getAddresses(String domain, Inet4Address server) { + return getAddresses(domain, Collections.singletonList(server)); + } + + @Override + public Future> getAddresses(String domain) { + return getAddresses(domain, configService.getConfig().getDomainResolvers()); + } + + @Override + public Future> getAddresses(String domain, List servers) { + return executor.submit(() -> { + var resolvers = configService.getConfig().getDomainResolvers(); + + if (resolvers.isEmpty()) { + throw new UnknownHostException("You're trying to resolve a domain without any resolvers"); + } + + for (var server : resolvers) { + try { + return lookup(domain, server); + } catch (IOException e) { + throw new UnknownHostException("Failed to resolve domain: " + domain); + } catch (InterruptedException e) { + continue; //network error + } + } + throw new UnknownHostException("Failed to resolve domain " + domain + " due to network error or something nasty"); + }); + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/service/DomainUpdaterService.java b/src/main/java/ru/kirillius/pf/sdn/service/DomainUpdaterService.java new file mode 100644 index 0000000..d78fa5a --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/service/DomainUpdaterService.java @@ -0,0 +1,89 @@ +package ru.kirillius.pf.sdn.service; + +import jakarta.annotation.PostConstruct; +import jakarta.annotation.PreDestroy; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; +import ru.kirillius.pf.sdn.Util.WorkerUtils; +import ru.kirillius.pf.sdn.api.Networking.DomainResolver; + +import java.net.Inet4Address; +import java.time.Duration; +import java.util.List; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Future; + +@Service +@RequiredArgsConstructor +@Slf4j +public class DomainUpdaterService { + private final DomainCacheService domainCacheService; + private final DomainResolver domainResolver; + private final ExecutorService executor; + private final SubscriptionService subscriptionService; + + public List getActualAddresses(String domain) { + var cached = domainCacheService.load(domain, true); + if (!cached.isEmpty()) { + var resolved = domainResolver.getAddresses(domain); + while (!resolved.isDone() && !resolved.isCancelled()) { + Thread.yield(); + } + if (resolved.isDone()) { + try { + domainCacheService.update(domain, resolved.get()); + cached.addAll(resolved.get()); + } catch (InterruptedException | ExecutionException e) { + throw new RuntimeException(e); + } + } + } + return cached; + } + + private Future worker; + private final RuntimeConfigService configService; + + @PostConstruct + private void initialize() { + worker = executor.submit(() -> new WorkerUtils.LoopedWorker() { + @Override + protected void doWork() { + subscriptionService.getSubscriptions().forEach(subscription -> { + var scope = subscription.getScope(); + if (!scope.isResolveDomains()) { + return; + } + scope.getDomains().forEach(domain -> { + var future = domainResolver.getAddresses(domain); + while (!future.isDone() && !future.isCancelled()) { + Thread.yield(); + } + + if (future.isDone()) { + try { + domainCacheService.update(domain, future.get()); + } catch (InterruptedException | ExecutionException e) { + log.error("Domain {} resolve failed", domain); + } + } + }); + }); + + WorkerUtils.waitForTimer(() -> Duration.ofMinutes(configService.getConfig().getDomainResolveInterval())); + + } + }); + + } + + @PreDestroy + private void destroy() { + if (worker != null) { + worker.cancel(true); + worker = null; + } + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/service/NetworkingService.java.txt b/src/main/java/ru/kirillius/pf/sdn/service/NetworkingService.java.txt new file mode 100644 index 0000000..2a0561b --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/service/NetworkingService.java.txt @@ -0,0 +1,341 @@ +package ru.kirillius.pf.sdn.service; + +import lombok.Getter; +import org.json.JSONObject; +import org.json.JSONTokener; +import org.springframework.stereotype.Service; +import ru.kirillius.java.utils.events.EventListener; +import ru.kirillius.json.JSONUtility; +import ru.kirillius.pf.sdn.api.Networking.IPv4Subnet; +import ru.kirillius.pf.sdn.api.Networking.NetworkScope; +import ru.kirillius.pf.sdn.dto.ResolverCacheEntry; +import ru.kirillius.pf.sdn.core.AppService; +import ru.kirillius.pf.sdn.core.Context; +import ru.kirillius.pf.sdn.core.ContextEventsHandler; +import ru.kirillius.pf.sdn.core.Subscription.SubscriptionService; +import ru.kirillius.pf.sdn.core.Util.IPv4Util; +import ru.kirillius.utils.logging.SystemLogger; + +import java.io.File; +import java.io.FileInputStream; +import java.io.FileOutputStream; +import java.io.IOException; +import java.time.Instant; +import java.time.temporal.ChronoUnit; +import java.util.*; +import java.util.concurrent.*; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; + +/** + * Builds the effective set of network resources by combining subscriptions, caches, and filters. + */ +@Service + +public class NetworkingService extends AppService { + private final static String CTX = NetworkingService.class.getSimpleName(); + private final ExecutorService executor = Executors.newSingleThreadExecutor(); + private final File domainCacheFile; + private final File asCacheFile; + private final EventListener resourceUpdateSubscription; + private final EventListener configChangeSubscription; + private final AtomicReference> updateProcess = new AtomicReference<>(); + @Getter + private final NetworkResourceConfig inputResources = new NetworkResourceConfig(); + @Getter + private final NetworkResourceConfig outputResources = new NetworkResourceConfig(); + + private final Map> prefixCache = new ConcurrentHashMap<>(); + private final Map domainCache = new ConcurrentHashMap<>(); + + /** + * Creates the networking service, wiring subscriptions and restoring cached state. + */ + public NetworkingService(Context context) { + super(context); + inputResources.clear(); + inputResources.add(context.getConfig().getCustomResources()); + resourceUpdateSubscription = context.getEventsHandler().getSubscriptionsUpdateEvent().add(bundle -> rebuildInputs()); + configChangeSubscription = context.getEventsHandler().getConfigChangeEvent().add(changeContext -> { + var filtersChanges = !changeContext.getCurrent().getFilteredResources().equals(changeContext.getInitial().getFilteredResources()); + var resChanges = !changeContext.getCurrent().getCustomResources().equals(changeContext.getInitial().getCustomResources()); + if (resChanges || filtersChanges) { + NetworkingService.this.rebuildInputs(); + } + }); + domainCacheFile = new File(context.getConfig().getCacheDirectory(), "domain-cache.json"); + asCacheFile = new File(context.getConfig().getCacheDirectory(), "as-cache.json"); + if (asCacheFile.exists() && context.getConfig().isCachingAS()) { + SystemLogger.message("Loading as cache file", CTX); + try (var is = new FileInputStream(asCacheFile)) { + var json = new JSONObject(new JSONTokener(is)); + json.keySet().forEach(key -> { + var as = Integer.parseInt(key); + prefixCache.put(as, JSONUtility.deserializeCollection(json.getJSONArray(key), IPv4Subnet.class, null).toList()); + }); + } catch (Exception e) { + SystemLogger.error("Failed to load as cache file " + asCacheFile.getPath(), CTX, e); + } + } + + if (domainCacheFile.exists() && context.getConfig().isCachingAS()) { + SystemLogger.message("Loading domain cache file", CTX); + try (var is = new FileInputStream(domainCacheFile)) { + var json = new JSONObject(new JSONTokener(is)); + json.keySet().forEach(host -> { + domainCache.put(host, JSONUtility.deserializeStructure(json.getJSONObject(host), ResolverCacheEntry.class)); + }); + } catch (Exception e) { + SystemLogger.error("Failed to load domain cache file " + asCacheFile.getPath(), CTX, e); + } + } + + + } + + public void performAutoresolve() { + var current = new HashSet(); + domainCache.forEach((host, entry) -> current.addAll(entry.getAddresses().keySet())); + + resolveDomains(List.copyOf(context.getServiceManager().getService(SubscriptionService.class).getAutoResolvingDomains())); + + var resolved = new HashSet(); + domainCache.forEach((host, entry) -> resolved.addAll(entry.getAddresses().keySet())); + + if (resolved.size() != current.size()) { + rebuildInputs(); + } else { + var updated = false; + for (var subnet : resolved) { + if (!current.contains(subnet)) { + updated = true; + break; + } + } + if (updated) { + rebuildInputs(); + } + } + + if(!context.getConfig().isCachingDomains()){ + domainCache.clear(); + } + } + + private void rebuildInputs() { + inputResources.clear(); + inputResources.add(context.getConfig().getCustomResources()); + inputResources.add(context.getServiceManager().getService(SubscriptionService.class).getOutputResources()); + triggerUpdate(false); + } + + /** + * Indicates whether an update job is currently executing. + */ + public boolean isUpdatingNow() { + var future = updateProcess.get(); + return future != null && !future.isDone() && !future.isCancelled(); + } + + /** + * Schedules an update of network resources, optionally ignoring cached prefixes. + */ + public void triggerUpdate(boolean ignoreCache) { + if (isUpdatingNow()) { + return; + } + SystemLogger.message("Updating network manager", CTX); + + updateProcess.set(executor.submit(() -> { + try { + SystemLogger.message("Update is started", CTX); + var config = context.getConfig(); + var filteredResources = config.getFilteredResources(); + + var domains = new HashSet<>(inputResources.getDomains()); + filteredResources.getDomains().forEach(domains::remove); + //check domain overlaps + + var domainsToRemove = new HashSet(); + for (var domainToMatch : domains) { + var pattern = "." + domainToMatch; + for (var domain : domains) { + if (domain.endsWith(pattern)) { + domainsToRemove.add(domain); + } + } + } + + domains.removeAll(domainsToRemove); + + var asn = new ArrayList<>(inputResources.getASN()); + asn.removeAll(filteredResources.getASN()); + + var asnToFetch = new ArrayList<>(asn); + if (!ignoreCache) { + asnToFetch.removeAll(prefixCache.keySet()); + } + + fetchPrefixes(asnToFetch); + + if (config.isCachingAS()) { + try (var os = new FileOutputStream(asCacheFile)) { + var json = new JSONObject(); + prefixCache.forEach((key, asnList) -> { + json.put(String.valueOf(key), JSONUtility.serializeCollection(asnList, IPv4Subnet.class, null)); + }); + os.write(json.toString().getBytes()); + } catch (IOException e) { + SystemLogger.error("Unable to write file " + asCacheFile.getPath(), CTX, e); + } + } + + resolveDomains(List.copyOf(context.getServiceManager().getService(SubscriptionService.class).getAutoResolvingDomains())); + + if (config.isCachingDomains()) { + try (var os = new FileOutputStream(domainCacheFile)) { + var json = new JSONObject(); + domainCache.forEach((key, entry) -> { + var serialized = JSONUtility.serializeStructure(entry); + json.put(String.valueOf(key), serialized); + }); + os.write(json.toString().getBytes()); + } catch (IOException e) { + SystemLogger.error("Unable to write file " + domainCacheFile.getPath(), CTX, e); + } + } + + var subnets = new HashSet<>(inputResources.getSubnets()); + asn.forEach(n -> { + var cached = prefixCache.get(n); + if (cached == null) { + return; + } + subnets.addAll(cached); + SystemLogger.message("Using " + cached.size() + " subnets from AS" + n, CTX); + }); + + //добавляем отрезолвенные домены + domains.forEach(domain -> { + var entry = domainCache.get(domain); + if (entry != null) { + subnets.addAll(entry.getAddresses().keySet()); + } + }); + + filteredResources.getSubnets().forEach(subnets::remove); + + SystemLogger.message("Trying to summary " + subnets.size() + " subnets...", CTX); + + var merged = IPv4Util.summarySubnets(subnets, config.getMergeSubnetsWithUsage()); + var unmerged = new AtomicInteger(); + subnets.forEach(subnet -> { + if (!merged.getMergedSubnets().contains(subnet)) { + unmerged.getAndIncrement(); + } + }); + + SystemLogger.message(subnets.size() + " subnets has been summarized and merged to " + merged.getResult().size() + " new subnets. Unmerged: " + unmerged.get(), CTX); + + outputResources.setASN(Collections.unmodifiableList(asn)); + outputResources.setSubnets(merged.getResult()); + outputResources.setDomains(domains.stream().toList()); + + SystemLogger.message("Update is complete", CTX); + + try { + context.getEventsHandler().getNetworkManagerUpdateEvent().invoke(outputResources); + } catch (Exception e) { + SystemLogger.error("Unable to invoke update event", CTX, e); + } + } catch (Exception e) { + SystemLogger.error("Something went wrong on update", CTX, e); + } + })); + } + + private void resolveDomains(List domains) { + var resolvedSubnets = new ArrayList(); + var resolver = context.getServiceManager().getService(DomainResolverService.class); + for (var domain : domains) { + var task = resolver.resolve(domain); + while (!task.isDone() && !task.isCancelled()) { + Thread.yield(); + } + try { + var subnets = task.get(); + var entry = domainCache.get(domain); + + if(entry == null) { + entry = new ResolverCacheEntry(); + domainCache.put(domain, entry); + } + + var addresses = entry.getAddresses(); + entry.setLastUpdate(Instant.now()); + subnets.forEach(subnet -> addresses.put(subnet, Instant.now())); + resolvedSubnets.addAll(domainCache.get(domain).getAddresses().keySet()); + } catch (InterruptedException | ExecutionException e) { + SystemLogger.error("Error happened while resolving domain " + domain, CTX, e); + } + } + + //remove old entries + for (var domain : domainCache.keySet()) { + var entry = domainCache.get(domain); + var addresses = entry.getAddresses(); + for (var subnet : addresses.keySet()) { + var time = addresses.get(subnet); + if (time.isBefore(Instant.now().minus(context.getConfig().getDomainsTimeToLive(), ChronoUnit.HOURS))) { + addresses.remove(subnet); + } + } + if (addresses.isEmpty()) { + domainCache.remove(domain); + } + } + } + + /** + * Fetches prefixes for the given autonomous systems and stores them in the cache. + */ + private void fetchPrefixes(List systems) { + var service = context.getServiceManager().getService(BGPInfoService.class); + systems.forEach(as -> { + var currentProvider = service.getProvider(); + var done = false; + do { + SystemLogger.message("Fetching AS" + as + " prefixes...", CTX); + var future = service.getPrefixes(as); + + while (!future.isDone() && !future.isCancelled()) { + Thread.yield(); + } + + try { + var iPv4Subnets = future.get(); + prefixCache.put(as, iPv4Subnets); + done = true; + break; + } catch (InterruptedException | ExecutionException e) { + service.fallbackNextProvider(); + SystemLogger.error("Error happened while fetching AS" + as + " prefixes. Trying to use fallback BGP info provider:" + service.getProvider().getClass().getSimpleName(), CTX, e); + } + } while (service.getProvider() != currentProvider); + if (!done) { + SystemLogger.error("Unable to fetch AS" + as + " prefixes from all providers. Trying to use cache...", CTX); + } + }); + } + + /** + * Removes event subscriptions and shuts down the executor. + */ + @Override + public void close() throws IOException { + context.getEventsHandler().getSubscriptionsUpdateEvent().remove(resourceUpdateSubscription); + context.getEventsHandler().getConfigChangeEvent().remove(configChangeSubscription); + + executor.shutdown(); + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/service/PipelineExecutorService.java b/src/main/java/ru/kirillius/pf/sdn/service/PipelineExecutorService.java new file mode 100644 index 0000000..eff6339 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/service/PipelineExecutorService.java @@ -0,0 +1,61 @@ +package ru.kirillius.pf.sdn.service; + +import jakarta.annotation.PostConstruct; +import jakarta.annotation.PreDestroy; +import lombok.RequiredArgsConstructor; +import lombok.SneakyThrows; +import org.springframework.stereotype.Service; +import ru.kirillius.pf.sdn.api.pipeline.ProcessingPipeline; + +import java.util.Queue; +import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Future; + +@Service +@RequiredArgsConstructor +public class PipelineExecutorService { + private final ExecutorService executor; + private final Queue executionQueue = new ConcurrentLinkedQueue<>(); + + public void triggerExecute(ProcessingPipeline pipeline) { + if (executionQueue.contains(pipeline)) { + return; + } + executionQueue.add(pipeline); + } + + @PostConstruct + private void initialize() { + worker = executor.submit(new PipelineWorker()); + } + + @PreDestroy + private void destroy() { + if (worker != null) { + worker.cancel(true); + } + } + + private Future worker; + + private class PipelineWorker implements Runnable { + @SuppressWarnings("BusyWait") + @SneakyThrows + @Override + public void run() { + while (!Thread.currentThread().isInterrupted()) { + while (!executionQueue.isEmpty()) { + var pipeline = executionQueue.poll(); + try { + pipeline.execute(); + } catch (Exception e) { + //TODO write log + } + } + Thread.sleep(100L); + Thread.yield(); + } + } + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/service/RuntimeConfigService.java b/src/main/java/ru/kirillius/pf/sdn/service/RuntimeConfigService.java new file mode 100644 index 0000000..da7e9eb --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/service/RuntimeConfigService.java @@ -0,0 +1,23 @@ +package ru.kirillius.pf.sdn.service; + +import lombok.Getter; +import lombok.Setter; +import org.springframework.stereotype.Service; + +import java.net.Inet4Address; +import java.util.List; + +@Service +public class RuntimeConfigService { + @Getter + @Setter + public class Config { + private int domainResolveInterval = 1; + private int subscriptionUpdateInterval = 1; + + private List domainResolvers; + } + + @Getter + private final Config config = new Config(); +} diff --git a/src/main/java/ru/kirillius/pf/sdn/service/SessionService.java b/src/main/java/ru/kirillius/pf/sdn/service/SessionService.java new file mode 100644 index 0000000..a063f02 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/service/SessionService.java @@ -0,0 +1,30 @@ +package ru.kirillius.pf.sdn.service; + +import lombok.RequiredArgsConstructor; +import org.springframework.stereotype.Service; +import org.springframework.util.StringUtils; + +import java.util.Optional; + +@Service +@RequiredArgsConstructor +public class SessionService { + + private final AuthService authService; + + public Optional resolveUserId(String authHeader) { + return extractToken(authHeader) + .flatMap(authService::resolveUserId); + } + + public Optional extractToken(String authHeader) { + if (StringUtils.hasText(authHeader) && authHeader.startsWith("Bearer ")) { + return Optional.of(authHeader.substring(7)); + } + return Optional.empty(); + } + + public void invalidate(String token) { + authService.invalidateToken(token); + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/service/SubscriptionService.java b/src/main/java/ru/kirillius/pf/sdn/service/SubscriptionService.java new file mode 100644 index 0000000..67871d7 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/service/SubscriptionService.java @@ -0,0 +1,149 @@ +package ru.kirillius.pf.sdn.service; + +import jakarta.annotation.PostConstruct; +import jakarta.annotation.PreDestroy; +import lombok.Getter; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.context.ApplicationContext; +import org.springframework.stereotype.Service; +import ru.kirillius.java.utils.events.ConcurrentEventHandler; +import ru.kirillius.java.utils.events.EventHandler; +import ru.kirillius.pf.sdn.Util.WorkerUtils; +import ru.kirillius.pf.sdn.api.Networking.Subscriptions.CacheFallbackProviderProtocol; +import ru.kirillius.pf.sdn.api.Networking.Subscriptions.Subscription; +import ru.kirillius.pf.sdn.api.Networking.Subscriptions.SubscriptionProvider; +import ru.kirillius.pf.sdn.api.Networking.Subscriptions.SubscriptionProviderProtocol; +import ru.kirillius.pf.sdn.repository.SubscriptionProviderConfigRepository; + +import java.time.Duration; +import java.time.Instant; +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Future; +import java.util.concurrent.atomic.AtomicInteger; + +@Service +@RequiredArgsConstructor +@Slf4j +public class SubscriptionService { + + private final SubscriptionProviderConfigRepository configRepository; + private final ExecutorService executorService; + private final RuntimeConfigService runtimeConfigService; + private final Map protocols = new ConcurrentHashMap<>(); + private final Map providers = new ConcurrentHashMap<>(); + private final AtomicInteger updateCounter = new AtomicInteger(0); + private final List subscriptions = new CopyOnWriteArrayList<>(); + private Future worker; + + + + public List getSubscriptions() { + return Collections.unmodifiableList(subscriptions); + } + + public void update() { + updateCounter.incrementAndGet(); + } + + @PreDestroy + private void destroy() { + if (worker != null) { + worker.cancel(true); + worker = null; + } + } + + private final ApplicationContext context; + + public void reloadProviders() { + var beanFactory = context.getAutowireCapableBeanFactory(); + synchronized (providers) { + providers.clear(); + configRepository.findAll().forEach(config -> { + var protocol = protocols.get(config.getProtocolId()); + if (protocol == null) { + log.error("Provider protocol with id {} not found. Using fallback from cache.", config.getProtocolId()); + protocol = beanFactory.createBean(CacheFallbackProviderProtocol.class); + } + providers.put(config.getName(), SubscriptionProvider.builder() + .config(config) + .protocol(protocol) + .build() + ); + }); + } + } + + @Getter + private final EventHandler updateEvent = new ConcurrentEventHandler<>(); + + private void reloadSubscriptions() { + var retrieved = new ArrayList(); + var updated = new ArrayList(); + + providers.values().forEach(provider -> { + updated.addAll(provider.update()); + retrieved.addAll(provider.getSubscriptions()); + }); + + synchronized (subscriptions) { + subscriptions.clear(); + subscriptions.addAll(retrieved); + } + + updated.forEach(s -> { + try { + updateEvent.invoke(s); + } catch (Exception e) { + log.error("Failed to invoke subscription update event because of error {}:{}", e.getClass().getSimpleName(), e.getMessage()); + } + }); + } + + @PostConstruct + private void initialize() { + reloadProviders(); + reloadSubscriptions(); + + worker = executorService.submit(new WorkerUtils.LoopedWorker() { + private Instant lastUpdate = Instant.now(); + + @Override + protected void doWork() { + if (updateCounter.get() > 0) { + updateCounter.set(0); + } + reloadSubscriptions(); + lastUpdate = Instant.now(); + var config = runtimeConfigService.getConfig(); + + WorkerUtils.waitForCondition(() -> Duration.between(lastUpdate, Instant.now()) + .compareTo(Duration.ofMinutes(config.getSubscriptionUpdateInterval())) > 0 + || updateCounter.get() > 0, 1000L); + } + }); + + } + + public void registerProvider(SubscriptionProviderProtocol subscriptionProviderProtocol) { + if (subscriptionProviderProtocol == null) { + throw new IllegalArgumentException(); + } + var identifier = subscriptionProviderProtocol.getIdentifier(); + if (identifier == null || identifier.isBlank()) { + throw new IllegalArgumentException("Identifier is null or blank"); + } + if (protocols.containsKey(identifier)) { + return; + } + protocols.put(identifier, subscriptionProviderProtocol); + } + +} diff --git a/src/main/resources/application.properties b/src/main/resources/application.properties new file mode 100644 index 0000000..4e462de --- /dev/null +++ b/src/main/resources/application.properties @@ -0,0 +1 @@ +jwt.secret=very-strong1234568ds7f7df6g7fd5df67g5ds78g6ds75gdf678g78g623r432r3213few \ No newline at end of file diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml new file mode 100644 index 0000000..b7514f6 --- /dev/null +++ b/src/main/resources/application.yml @@ -0,0 +1,24 @@ +spring: + application: + name: pf-sdn + + datasource: + url: ${app.datasource.url} + driver-class-name: org.h2.Driver + username: sa + password: + + h2: + console: + enabled: true + path: /h2-console + + jpa: + hibernate: + ddl-auto: update + show-sql: true + properties: + hibernate: + dialect: org.hibernate.dialect.H2Dialect + + diff --git a/src/main/resources/static/index.html b/src/main/resources/static/index.html new file mode 100644 index 0000000..9e2c6a2 --- /dev/null +++ b/src/main/resources/static/index.html @@ -0,0 +1,121 @@ + + + + + pf-sdn Auth Test + + + +

pf-sdn Auth Test

+ +
+

Login

+ + + + + +

+    
+ + + + + +