From e4745806c2cf005a01ba26a55d6ce6d6577e8c02 Mon Sep 17 00:00:00 2001 From: kirillius Date: Tue, 21 Jul 2026 00:12:31 +0300 Subject: [PATCH] WIP: ExecutionGraph & Multiple inputs & Subscription groups --- .../pf/sdn/api/flow/CallContext.java | 5 +- .../pf/sdn/api/flow/ExecutionGraph.java | 134 ++++++++++++++++++ .../ru/kirillius/pf/sdn/api/flow/Flow.java | 44 ++++-- .../kirillius/pf/sdn/api/flow/FlowAction.java | 15 +- .../pf/sdn/api/flow/FlowFunction.java | 9 ++ .../pf/sdn/api/flow/TriggerType.java | 1 + .../pf/sdn/entity/ActionConnectionConfig.java | 24 ++++ .../kirillius/pf/sdn/entity/FlowConfig.java | 5 + .../pf/sdn/entity/SubscriptionCacheEntry.java | 1 + .../pf/sdn/entity/SubscriptionSetEntry.java | 25 ++++ .../kirillius/pf/sdn/flow/AutoResolveASN.java | 18 +-- .../pf/sdn/flow/AutoResolveDomains.java | 16 +-- .../kirillius/pf/sdn/flow/CommonFunction.java | 46 ++++++ .../ru/kirillius/pf/sdn/flow/DebugOutput.java | 23 +-- .../kirillius/pf/sdn/flow/DummyFunction.java | 20 +-- .../pf/sdn/flow/DummyPassthrough.java | 26 +++- .../pf/sdn/flow/DuplicatesFilter.java | 25 +++- .../sdn/flow/MergeNeighbourSubnetFilter.java | 34 ++--- .../flow/MergeWithCoverageSubnetFilter.java | 28 ++-- .../pf/sdn/flow/OverlappedDomainFilter.java | 28 ++-- .../pf/sdn/flow/OverlappedSubnetFilter.java | 36 ++--- .../kirillius/pf/sdn/flow/ResourceFilter.java | 122 ++++++++-------- .../pf/sdn/flow/ResourceFilterObsolete.java | 107 ++++++++++++++ .../pf/sdn/flow/StaticResourceInput.java | 66 +++++++++ .../pf/sdn/flow/SubscriptionInput.java | 17 +-- .../repository/SubscriptionSetRepository.java | 8 ++ .../kirillius/pf/sdn/service/FlowService.java | 39 ++++- .../pf/sdn/service/SubscriptionService.java | 33 +++-- 28 files changed, 723 insertions(+), 232 deletions(-) create mode 100644 src/main/java/ru/kirillius/pf/sdn/api/flow/ExecutionGraph.java create mode 100644 src/main/java/ru/kirillius/pf/sdn/entity/ActionConnectionConfig.java create mode 100644 src/main/java/ru/kirillius/pf/sdn/entity/SubscriptionSetEntry.java create mode 100644 src/main/java/ru/kirillius/pf/sdn/flow/CommonFunction.java create mode 100644 src/main/java/ru/kirillius/pf/sdn/flow/ResourceFilterObsolete.java create mode 100644 src/main/java/ru/kirillius/pf/sdn/flow/StaticResourceInput.java create mode 100644 src/main/java/ru/kirillius/pf/sdn/repository/SubscriptionSetRepository.java diff --git a/src/main/java/ru/kirillius/pf/sdn/api/flow/CallContext.java b/src/main/java/ru/kirillius/pf/sdn/api/flow/CallContext.java index f34e40c..449e7e6 100644 --- a/src/main/java/ru/kirillius/pf/sdn/api/flow/CallContext.java +++ b/src/main/java/ru/kirillius/pf/sdn/api/flow/CallContext.java @@ -5,15 +5,16 @@ import lombok.RequiredArgsConstructor; import lombok.Setter; import ru.kirillius.pf.sdn.api.Networking.NetworkScope; +import java.util.List; import java.util.Map; @RequiredArgsConstructor public final class CallContext { @Getter - private final NetworkScope source; + private final List sources; @Getter @Setter - private NetworkScope output; + private List outputs; @Getter private final Map properties; diff --git a/src/main/java/ru/kirillius/pf/sdn/api/flow/ExecutionGraph.java b/src/main/java/ru/kirillius/pf/sdn/api/flow/ExecutionGraph.java new file mode 100644 index 0000000..28462dd --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/flow/ExecutionGraph.java @@ -0,0 +1,134 @@ +package ru.kirillius.pf.sdn.api.flow; + +import lombok.Builder; +import lombok.Getter; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import ru.kirillius.java.utils.events.ConcurrentEventHandler; +import ru.kirillius.java.utils.events.EventHandler; +import ru.kirillius.pf.sdn.api.Networking.NetworkScope; +import ru.kirillius.pf.sdn.entity.ActionConfig; +import ru.kirillius.pf.sdn.entity.ActionConnectionConfig; + +import java.util.*; +import java.util.function.Function; + +@Slf4j +@RequiredArgsConstructor +public class ExecutionGraph { + + public int size() { + return nodes.size(); + } + + @Getter + private final EventHandler startExecutingEvent = new ConcurrentEventHandler<>(); + private final Map nodes = new HashMap<>(); + + private final Map executed = new HashMap<>(); + + public ExecutionGraph(List actionConfigs, List connectionConfigs, Function actionFactory) { + var pendingActions = new ArrayList<>(actionConfigs); + + //создаём ноды у которых нет зависимостей + actionConfigs + .stream() + .filter(config -> connectionConfigs + .stream() + .noneMatch(conn -> conn.getTo().equals(config)) + ).forEach(config -> { + nodes.put(config, new Node(actionFactory.apply(config), Collections.emptyList())); + pendingActions.remove(config); + }); + //создаём оставшиеся ноды + + while (!pendingActions.isEmpty()) { + //ищем конфиги, все соединения которых можно зарезолвить + var spawned = new ArrayList(); + + pendingActions.forEach((actionConfig) -> { + var selfConnectionConfigs = connectionConfigs.stream().filter(conn -> conn.getTo().equals(actionConfig)).toList(); + if (selfConnectionConfigs.stream().allMatch(conn -> nodes.containsKey(conn.getFrom()))) { + spawned.add(actionConfig); + + nodes.put(actionConfig, + new Node( + actionFactory.apply(actionConfig), + selfConnectionConfigs + .stream() + .map(connectionConfig -> new Node.Connection( + nodes.get(connectionConfig.getFrom()), + connectionConfig.getFromPort(), + connectionConfig.getToPort() + )).toList() + ) + ); + } + }); + + pendingActions.removeAll(spawned); + + if (spawned.isEmpty()) { + throw new IllegalStateException("Connection loop detected and could not be resolved"); + } + + } + + } + + public Iterator execute() { + executed.clear(); + final var configs = new HashSet<>(nodes.keySet()); + + return new Iterator<>() { + @Override + public boolean hasNext() { + return executed.size() != nodes.size(); + } + + @Override + public ExecutionResult next() { + var nodeToExecute = findNodeToExecute(configs); + var action = nodeToExecute.action(); + try { + startExecutingEvent.invoke(action); + } catch (Exception e) { + log.error("Error on {}:startExecuting event", ExecutionGraph.class.getSimpleName(), e); + } + + var inputs = new ArrayList(); + nodeToExecute.inputs().forEach((input) -> { + var fromNode = input.from(); + inputs.add(executed.get(fromNode).getOutputs().get(input.fromPort)); + }); + + var callContext = action.apply(inputs); + executed.put(nodeToExecute, callContext); + configs.remove(nodeToExecute.action().getConfig()); + + return new ExecutionResult(action, callContext); + } + }; + } + + private Node findNodeToExecute(Set actionConfigs) { + for (var config : actionConfigs) { + var node = nodes.get(config); + //проверяем что все на всех входах есть данные + if (node.inputs().stream().allMatch(connection -> executed.containsKey(connection.from()))) { + return node; + } + } + + throw new NoSuchElementException("There is No node to execute"); + } + + @Builder + public record ExecutionResult(FlowAction action, CallContext context) { + } + + private record Node(FlowAction action, List inputs) { + private record Connection(Node from, int fromPort, int toPort) { + } + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/flow/Flow.java b/src/main/java/ru/kirillius/pf/sdn/api/flow/Flow.java index 6a1f2b0..0b53993 100644 --- a/src/main/java/ru/kirillius/pf/sdn/api/flow/Flow.java +++ b/src/main/java/ru/kirillius/pf/sdn/api/flow/Flow.java @@ -1,7 +1,7 @@ package ru.kirillius.pf.sdn.api.flow; -import lombok.*; -import ru.kirillius.pf.sdn.api.Networking.NetworkScope; +import lombok.Getter; +import ru.kirillius.pf.sdn.entity.ActionConfig; import ru.kirillius.pf.sdn.entity.FlowConfig; import java.util.ArrayList; @@ -9,49 +9,67 @@ import java.util.List; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Function; -@RequiredArgsConstructor() public class Flow { + public Flow(FlowConfig config, Function actionFactory) { + this.config = config; + graph = new ExecutionGraph(config.getActions(), config.getConnections(), actionFactory); + } + @Getter private final FlowConfig config; - private final List actions; + private final ExecutionGraph graph; private final AtomicBoolean running = new AtomicBoolean(false); private final AtomicBoolean interrupted = new AtomicBoolean(false); private AtomicInteger currentStep = new AtomicInteger(0); + private AtomicReference currentAction = new AtomicReference<>(null); public boolean isRunning() { return running.get(); } public int getStepCount() { - return actions.size(); + return graph.size(); } public int getCurrentStep() { return currentStep.get(); } + public FlowAction getCurrentAction() { + return currentAction.get(); + } + + @Getter + private final List results = new ArrayList<>(); + public void interrupt() { interrupted.set(false); } - private List executionInfo = new ArrayList<>(); - - public void execute() { + public void execute() throws InterruptedException { interrupted.set(false); running.set(true); currentStep.set(0); + results.clear(); try { - var scope = new AtomicReference<>(new NetworkScope()); - actions.forEach(action -> { + var listener = graph.getStartExecutingEvent().add(action -> currentAction.set(action)); + + var iterator = graph.execute(); + + while (iterator.hasNext()) { if (interrupted.get()) { - throw new RuntimeException("Interrupted"); + throw new InterruptedException(); } currentStep.incrementAndGet(); - scope.set(action.apply(scope.get())); - }); + results.add(iterator.next()); + } + + graph.getStartExecutingEvent().remove(listener); + } finally { interrupted.set(false); running.set(false); diff --git a/src/main/java/ru/kirillius/pf/sdn/api/flow/FlowAction.java b/src/main/java/ru/kirillius/pf/sdn/api/flow/FlowAction.java index 87d8129..70b19f0 100644 --- a/src/main/java/ru/kirillius/pf/sdn/api/flow/FlowAction.java +++ b/src/main/java/ru/kirillius/pf/sdn/api/flow/FlowAction.java @@ -5,18 +5,23 @@ import lombok.Getter; import lombok.Setter; import lombok.SneakyThrows; import ru.kirillius.pf.sdn.api.Networking.NetworkScope; +import ru.kirillius.pf.sdn.entity.ActionConfig; +import java.util.List; import java.util.Map; import java.util.function.Function; @Builder -public class FlowAction implements Function { +public class FlowAction implements Function, CallContext> { private FlowFunction function; + @Getter + private final ActionConfig config; - public FlowAction(Class functionClass, Map properties) { + public FlowAction(Class functionClass, ActionConfig config) { this.functionClass = functionClass; - this.properties = properties; + this.properties = config.getProperties(); + this.config = config; instantiateFunction(); } @@ -37,8 +42,8 @@ public class FlowAction implements Function { private Map properties; @Override - public CallContext apply(NetworkScope networkScope) { - var context = new CallContext(networkScope, properties); + public CallContext apply(List inputs) { + var context = new CallContext(inputs, properties); function.apply(context); return context; } diff --git a/src/main/java/ru/kirillius/pf/sdn/api/flow/FlowFunction.java b/src/main/java/ru/kirillius/pf/sdn/api/flow/FlowFunction.java index 8ff9a4a..d2aa19d 100644 --- a/src/main/java/ru/kirillius/pf/sdn/api/flow/FlowFunction.java +++ b/src/main/java/ru/kirillius/pf/sdn/api/flow/FlowFunction.java @@ -2,6 +2,7 @@ package ru.kirillius.pf.sdn.api.flow; import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor; +import java.util.List; import java.util.Map; public interface FlowFunction { @@ -12,4 +13,12 @@ public interface FlowFunction { boolean hasInputs(); boolean hasOutputs(); + + int getInputCount(); + + int getOutputCount(); + + List getInputNames(); + + List getOutputNames(); } diff --git a/src/main/java/ru/kirillius/pf/sdn/api/flow/TriggerType.java b/src/main/java/ru/kirillius/pf/sdn/api/flow/TriggerType.java index 1f8f04f..432637e 100644 --- a/src/main/java/ru/kirillius/pf/sdn/api/flow/TriggerType.java +++ b/src/main/java/ru/kirillius/pf/sdn/api/flow/TriggerType.java @@ -4,5 +4,6 @@ public enum TriggerType { Manual, Interval, OnSubscriptionUpdate, + OnSubscriptionSetUpdate, OnStart } diff --git a/src/main/java/ru/kirillius/pf/sdn/entity/ActionConnectionConfig.java b/src/main/java/ru/kirillius/pf/sdn/entity/ActionConnectionConfig.java new file mode 100644 index 0000000..c8373e1 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/entity/ActionConnectionConfig.java @@ -0,0 +1,24 @@ +package ru.kirillius.pf.sdn.entity; + +import jakarta.persistence.*; +import lombok.Getter; +import lombok.Setter; + +@Getter +@Setter +@Entity +public class ActionConnectionConfig { + @Id + @GeneratedValue(strategy = GenerationType.IDENTITY) + private Long id; + + @ManyToOne(optional = false) + private ActionConfig from; + + private int fromPort; + + @ManyToOne(optional = false) + private ActionConfig to; + + private int toPort; +} \ No newline at end of file diff --git a/src/main/java/ru/kirillius/pf/sdn/entity/FlowConfig.java b/src/main/java/ru/kirillius/pf/sdn/entity/FlowConfig.java index c97ff49..1554207 100644 --- a/src/main/java/ru/kirillius/pf/sdn/entity/FlowConfig.java +++ b/src/main/java/ru/kirillius/pf/sdn/entity/FlowConfig.java @@ -43,4 +43,9 @@ public class FlowConfig { @ManyToMany(fetch = FetchType.EAGER)//TODO проверить связность private List actions = new ArrayList<>(); + + @ManyToMany(fetch = FetchType.EAGER)//TODO проверить связность + private List connections = new ArrayList<>(); + + } diff --git a/src/main/java/ru/kirillius/pf/sdn/entity/SubscriptionCacheEntry.java b/src/main/java/ru/kirillius/pf/sdn/entity/SubscriptionCacheEntry.java index 5bb69ac..d126d64 100644 --- a/src/main/java/ru/kirillius/pf/sdn/entity/SubscriptionCacheEntry.java +++ b/src/main/java/ru/kirillius/pf/sdn/entity/SubscriptionCacheEntry.java @@ -5,6 +5,7 @@ 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 diff --git a/src/main/java/ru/kirillius/pf/sdn/entity/SubscriptionSetEntry.java b/src/main/java/ru/kirillius/pf/sdn/entity/SubscriptionSetEntry.java new file mode 100644 index 0000000..317940b --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/entity/SubscriptionSetEntry.java @@ -0,0 +1,25 @@ +package ru.kirillius.pf.sdn.entity; + +import jakarta.persistence.*; +import lombok.*; + +import java.util.Set; + +@Entity +@Table(name = "subscription_set") +@Getter +@Setter +@NoArgsConstructor +@Builder +@AllArgsConstructor +public class SubscriptionSetEntry { + @Id + @GeneratedValue(strategy = GenerationType.IDENTITY) + private Long id; + + @Column(nullable = false) + private String name; + + @OneToMany(fetch = FetchType.EAGER) + private Set subscriptions; +} diff --git a/src/main/java/ru/kirillius/pf/sdn/flow/AutoResolveASN.java b/src/main/java/ru/kirillius/pf/sdn/flow/AutoResolveASN.java index 1ad29a1..12776c4 100644 --- a/src/main/java/ru/kirillius/pf/sdn/flow/AutoResolveASN.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/AutoResolveASN.java @@ -3,17 +3,17 @@ package ru.kirillius.pf.sdn.flow; import lombok.RequiredArgsConstructor; import ru.kirillius.pf.sdn.api.Networking.NetworkScope; import ru.kirillius.pf.sdn.api.flow.CallContext; -import ru.kirillius.pf.sdn.api.flow.FlowFunction; import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor; import ru.kirillius.pf.sdn.api.properties.PropertyType; import ru.kirillius.pf.sdn.service.AutonomousSystemCacheService; import java.util.Collections; import java.util.HashMap; +import java.util.List; import java.util.Map; @RequiredArgsConstructor -public class AutoResolveASN implements FlowFunction { +public class AutoResolveASN extends CommonFunction { private final AutonomousSystemCacheService autonomousSystemCacheService; private final static String CLEAR = "clear"; @@ -34,15 +34,15 @@ public class AutoResolveASN implements FlowFunction { @Override public void apply(CallContext context) { var output = new NetworkScope(); - output.add(context.getSource()); + output.add(context.getSources().getFirst()); var properties = context.getProperties(); if (properties.containsKey(CLEAR) && properties.get(CLEAR).equals(Boolean.TRUE.toString())) { output.getASN().clear(); } - context.getSource().getASN().forEach(asn -> output.getSubnets().addAll(autonomousSystemCacheService.load(asn))); - context.setOutput(output); + context.getSources().getFirst().getASN().forEach(asn -> output.getSubnets().addAll(autonomousSystemCacheService.load(asn))); + context.setOutputs(List.of(output)); } @Override @@ -51,12 +51,12 @@ public class AutoResolveASN implements FlowFunction { } @Override - public boolean hasInputs() { - return true; + public List getInputNames() { + return List.of("input"); } @Override - public boolean hasOutputs() { - return true; + public List getOutputNames() { + return List.of("output"); } } diff --git a/src/main/java/ru/kirillius/pf/sdn/flow/AutoResolveDomains.java b/src/main/java/ru/kirillius/pf/sdn/flow/AutoResolveDomains.java index 68d617a..a7e8139 100644 --- a/src/main/java/ru/kirillius/pf/sdn/flow/AutoResolveDomains.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/AutoResolveDomains.java @@ -4,17 +4,17 @@ 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.flow.CallContext; -import ru.kirillius.pf.sdn.api.flow.FlowFunction; import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor; import ru.kirillius.pf.sdn.api.properties.PropertyType; import ru.kirillius.pf.sdn.service.DomainUpdaterService; import java.util.Collections; import java.util.HashMap; +import java.util.List; import java.util.Map; @RequiredArgsConstructor -public class AutoResolveDomains implements FlowFunction { +public class AutoResolveDomains extends CommonFunction { private final DomainUpdaterService domainCacheService; @@ -35,7 +35,7 @@ public class AutoResolveDomains implements FlowFunction { @Override public void apply(CallContext context) { - var source = context.getSource(); + var source = context.getSources().getFirst(); var properties = context.getProperties(); var output = new NetworkScope(); output.add(source); @@ -53,7 +53,7 @@ public class AutoResolveDomains implements FlowFunction { ) ); - context.setOutput(output); + context.setOutputs(List.of(output)); } @Override @@ -62,12 +62,12 @@ public class AutoResolveDomains implements FlowFunction { } @Override - public boolean hasInputs() { - return true; + public List getInputNames() { + return List.of("input"); } @Override - public boolean hasOutputs() { - return true; + public List getOutputNames() { + return List.of("output"); } } diff --git a/src/main/java/ru/kirillius/pf/sdn/flow/CommonFunction.java b/src/main/java/ru/kirillius/pf/sdn/flow/CommonFunction.java new file mode 100644 index 0000000..c9c7163 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/flow/CommonFunction.java @@ -0,0 +1,46 @@ +package ru.kirillius.pf.sdn.flow; + +import ru.kirillius.pf.sdn.api.flow.FlowFunction; +import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor; + +import java.util.Collections; +import java.util.List; +import java.util.Map; + +abstract class CommonFunction implements FlowFunction { + + @Override + public final int getInputCount() { + return getInputNames().size(); + } + + @Override + public final int getOutputCount() { + return getOutputNames().size(); + } + + @Override + public Map getProperties() { + return Collections.emptyMap(); + } + + @Override + public final boolean hasInputs() { + return !getInputNames().isEmpty(); + } + + @Override + public final boolean hasOutputs() { + return !getOutputNames().isEmpty(); + } + + @Override + public List getInputNames() { + return Collections.emptyList(); + } + + @Override + public List getOutputNames() { + return Collections.emptyList(); + } +} diff --git a/src/main/java/ru/kirillius/pf/sdn/flow/DebugOutput.java b/src/main/java/ru/kirillius/pf/sdn/flow/DebugOutput.java index 34b8151..cef3e76 100644 --- a/src/main/java/ru/kirillius/pf/sdn/flow/DebugOutput.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/DebugOutput.java @@ -4,35 +4,22 @@ import com.fasterxml.jackson.databind.ObjectMapper; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import ru.kirillius.pf.sdn.api.flow.CallContext; -import ru.kirillius.pf.sdn.api.flow.FlowFunction; -import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor; -import java.util.Collections; -import java.util.Map; +import java.util.List; @Slf4j @RequiredArgsConstructor -public class DebugOutput implements FlowFunction { +public class DebugOutput extends CommonFunction { private final ObjectMapper objectMapper; @Override public void apply(CallContext context) { - log.info(objectMapper.valueToTree(context.getSource()).toString()); + log.info(objectMapper.valueToTree(context.getSources()).toString()); } @Override - public Map getProperties() { - return Collections.emptyMap(); - } - - @Override - public boolean hasInputs() { - return true; - } - - @Override - public boolean hasOutputs() { - return false; + public List getInputNames() { + return List.of("input"); } } diff --git a/src/main/java/ru/kirillius/pf/sdn/flow/DummyFunction.java b/src/main/java/ru/kirillius/pf/sdn/flow/DummyFunction.java index e458961..58f7d2a 100644 --- a/src/main/java/ru/kirillius/pf/sdn/flow/DummyFunction.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/DummyFunction.java @@ -1,30 +1,12 @@ package ru.kirillius.pf.sdn.flow; import ru.kirillius.pf.sdn.api.flow.CallContext; -import ru.kirillius.pf.sdn.api.flow.FlowFunction; -import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor; -import java.util.Map; - -public class DummyFunction implements FlowFunction { +public class DummyFunction extends CommonFunction { @Override public void apply(CallContext context) { throw new UnsupportedOperationException("Not implemented"); } - @Override - public Map getProperties() { - return Map.of(); - } - - @Override - public boolean hasInputs() { - return true; - } - - @Override - public boolean hasOutputs() { - return true; - } } diff --git a/src/main/java/ru/kirillius/pf/sdn/flow/DummyPassthrough.java b/src/main/java/ru/kirillius/pf/sdn/flow/DummyPassthrough.java index 7984372..e3ba7fd 100644 --- a/src/main/java/ru/kirillius/pf/sdn/flow/DummyPassthrough.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/DummyPassthrough.java @@ -7,16 +7,17 @@ import ru.kirillius.pf.sdn.api.flow.FlowFunction; import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor; import java.util.Collections; +import java.util.List; import java.util.Map; @RequiredArgsConstructor -public class DummyPassthrough implements FlowFunction { +public class DummyPassthrough implements FlowFunction {//TODO а зачем это надо вообще? @Override public void apply(CallContext context) { var output = new NetworkScope(); - output.add(context.getSource()); - context.setOutput(output); + output.add(context.getSources().getFirst()); + context.setOutputs(List.of(output)); } @Override @@ -33,4 +34,23 @@ public class DummyPassthrough implements FlowFunction { public boolean hasOutputs() { return true; } + @Override + public int getInputCount() { + return 1; + } + + @Override + public int getOutputCount() { + return 1; + } + + @Override + public List getInputNames() { + return List.of("input"); + } + + @Override + public List getOutputNames() { + return List.of("output"); + } } diff --git a/src/main/java/ru/kirillius/pf/sdn/flow/DuplicatesFilter.java b/src/main/java/ru/kirillius/pf/sdn/flow/DuplicatesFilter.java index dddd0a7..e4e113f 100644 --- a/src/main/java/ru/kirillius/pf/sdn/flow/DuplicatesFilter.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/DuplicatesFilter.java @@ -8,6 +8,7 @@ import ru.kirillius.pf.sdn.api.flow.FlowFunction; import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor; import java.util.Collections; +import java.util.List; import java.util.Map; @Slf4j @@ -16,13 +17,13 @@ public class DuplicatesFilter implements FlowFunction { @Override public void apply(CallContext context) { - var source = context.getSource(); + var source = context.getSources().getFirst(); var output = new NetworkScope(); output.setSubnets(source.getSubnets().stream().distinct().toList()); output.setDomains(source.getDomains().stream().distinct().toList()); output.setASN(source.getASN().stream().distinct().toList()); output.setAutoResolvedDomains(source.getAutoResolvedDomains().stream().distinct().toList()); - context.setOutput(output); + context.setOutputs(List.of(output)); } @Override @@ -40,4 +41,24 @@ public class DuplicatesFilter implements FlowFunction { return true; } + @Override + public int getInputCount() { + return 1; + } + + @Override + public int getOutputCount() { + return 1; + } + + @Override + public List getInputNames() { + return List.of("input"); + } + + @Override + public List getOutputNames() { + return List.of("output"); + } + } diff --git a/src/main/java/ru/kirillius/pf/sdn/flow/MergeNeighbourSubnetFilter.java b/src/main/java/ru/kirillius/pf/sdn/flow/MergeNeighbourSubnetFilter.java index 77945e8..e1858e6 100644 --- a/src/main/java/ru/kirillius/pf/sdn/flow/MergeNeighbourSubnetFilter.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/MergeNeighbourSubnetFilter.java @@ -5,19 +5,26 @@ import lombok.extern.slf4j.Slf4j; import ru.kirillius.pf.sdn.Util.IPv4Util; import ru.kirillius.pf.sdn.api.Networking.NetworkScope; import ru.kirillius.pf.sdn.api.flow.CallContext; -import ru.kirillius.pf.sdn.api.flow.FlowFunction; -import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor; -import java.util.Collections; -import java.util.Map; +import java.util.List; @Slf4j @RequiredArgsConstructor -public class MergeNeighbourSubnetFilter implements FlowFunction { +public class MergeNeighbourSubnetFilter extends CommonFunction { + + @Override + public List getInputNames() { + return List.of("input"); + } + + @Override + public List getOutputNames() { + return List.of("output"); + } @Override public void apply(CallContext context) { - var source = context.getSource(); + var source = context.getSources().getFirst(); var output = new NetworkScope(); output.setSubnets(IPv4Util.mergeNeighbours(source.getSubnets())); output.setDomains(source.getDomains()); @@ -25,19 +32,4 @@ public class MergeNeighbourSubnetFilter implements FlowFunction { output.setAutoResolvedDomains(source.getAutoResolvedDomains()); } - @Override - public Map getProperties() { - return Collections.emptyMap(); - } - - @Override - public boolean hasInputs() { - return false; - } - - @Override - public boolean hasOutputs() { - return true; - } - } diff --git a/src/main/java/ru/kirillius/pf/sdn/flow/MergeWithCoverageSubnetFilter.java b/src/main/java/ru/kirillius/pf/sdn/flow/MergeWithCoverageSubnetFilter.java index 4390fcd..0a0505d 100644 --- a/src/main/java/ru/kirillius/pf/sdn/flow/MergeWithCoverageSubnetFilter.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/MergeWithCoverageSubnetFilter.java @@ -5,18 +5,18 @@ import lombok.extern.slf4j.Slf4j; import ru.kirillius.pf.sdn.Util.IPv4Util; import ru.kirillius.pf.sdn.api.Networking.NetworkScope; import ru.kirillius.pf.sdn.api.flow.CallContext; -import ru.kirillius.pf.sdn.api.flow.FlowFunction; import ru.kirillius.pf.sdn.api.properties.IntegerConstraint; import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor; import ru.kirillius.pf.sdn.api.properties.PropertyType; import java.util.Collections; import java.util.HashMap; +import java.util.List; import java.util.Map; @Slf4j @RequiredArgsConstructor -public class MergeWithCoverageSubnetFilter implements FlowFunction { +public class MergeWithCoverageSubnetFilter extends CommonFunction { private final static String USAGE = "usage"; private final static Map properties; @@ -28,10 +28,20 @@ public class MergeWithCoverageSubnetFilter implements FlowFunction { ); } + @Override + public List getInputNames() { + return List.of("input"); + } + + @Override + public List getOutputNames() { + return List.of("output"); + } + @Override public void apply(CallContext context) { var properties = context.getProperties(); - var source = context.getSource(); + var source = context.getSources().getFirst(); var usage = properties.containsKey(USAGE) ? Integer.parseInt(properties.get(USAGE)) : 75; if (usage < 51) { usage = 51; @@ -43,7 +53,7 @@ public class MergeWithCoverageSubnetFilter implements FlowFunction { output.setDomains(source.getDomains()); output.setASN(source.getASN()); output.setAutoResolvedDomains(source.getAutoResolvedDomains()); - context.setOutput(output); + context.setOutputs(List.of(output)); } @Override @@ -51,14 +61,4 @@ public class MergeWithCoverageSubnetFilter implements FlowFunction { return Collections.unmodifiableMap(properties); } - @Override - public boolean hasInputs() { - return true; - } - - @Override - public boolean hasOutputs() { - return true; - } - } diff --git a/src/main/java/ru/kirillius/pf/sdn/flow/OverlappedDomainFilter.java b/src/main/java/ru/kirillius/pf/sdn/flow/OverlappedDomainFilter.java index f1aaa4c..6dd8ca7 100644 --- a/src/main/java/ru/kirillius/pf/sdn/flow/OverlappedDomainFilter.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/OverlappedDomainFilter.java @@ -5,25 +5,35 @@ import lombok.extern.slf4j.Slf4j; import ru.kirillius.pf.sdn.Util.DomainUtil; import ru.kirillius.pf.sdn.api.Networking.NetworkScope; import ru.kirillius.pf.sdn.api.flow.CallContext; -import ru.kirillius.pf.sdn.api.flow.FlowFunction; import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor; import java.util.Collections; +import java.util.List; import java.util.Map; @Slf4j @RequiredArgsConstructor -public class OverlappedDomainFilter implements FlowFunction { +public class OverlappedDomainFilter extends CommonFunction { + + @Override + public List getInputNames() { + return List.of("input"); + } + + @Override + public List getOutputNames() { + return List.of("output"); + } @Override public void apply(CallContext context) { - var source = context.getSource(); + var source = context.getSources().getFirst(); var output = new NetworkScope(); output.setSubnets(source.getSubnets()); output.setDomains(DomainUtil.removeOverlapped(source.getDomains())); output.setASN(source.getASN()); output.setAutoResolvedDomains(source.getAutoResolvedDomains()); - context.setOutput(output); + context.setOutputs(List.of(output)); } @Override @@ -31,14 +41,4 @@ public class OverlappedDomainFilter implements FlowFunction { return Collections.emptyMap(); } - @Override - public boolean hasInputs() { - return true; - } - - @Override - public boolean hasOutputs() { - return true; - } - } diff --git a/src/main/java/ru/kirillius/pf/sdn/flow/OverlappedSubnetFilter.java b/src/main/java/ru/kirillius/pf/sdn/flow/OverlappedSubnetFilter.java index 719120c..c206bec 100644 --- a/src/main/java/ru/kirillius/pf/sdn/flow/OverlappedSubnetFilter.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/OverlappedSubnetFilter.java @@ -5,40 +5,32 @@ import lombok.extern.slf4j.Slf4j; import ru.kirillius.pf.sdn.Util.IPv4Util; import ru.kirillius.pf.sdn.api.Networking.NetworkScope; import ru.kirillius.pf.sdn.api.flow.CallContext; -import ru.kirillius.pf.sdn.api.flow.FlowFunction; -import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor; -import java.util.Collections; -import java.util.Map; +import java.util.List; @Slf4j @RequiredArgsConstructor -public class OverlappedSubnetFilter implements FlowFunction { +public class OverlappedSubnetFilter extends CommonFunction { + + @Override + public List getInputNames() { + return List.of("input"); + } + + @Override + public List getOutputNames() { + return List.of("output"); + } @Override public void apply(CallContext context) { - var source = context.getSource(); + var source = context.getSources().getFirst(); var output = new NetworkScope(); output.setSubnets(IPv4Util.removeOverlapped(source.getSubnets())); output.setDomains(source.getDomains()); output.setASN(source.getASN()); output.setAutoResolvedDomains(source.getAutoResolvedDomains()); - context.setOutput(output); - } - - @Override - public Map getProperties() { - return Collections.emptyMap(); - } - - @Override - public boolean hasInputs() { - return true; - } - - @Override - public boolean hasOutputs() { - return true; + context.setOutputs(List.of(output)); } } diff --git a/src/main/java/ru/kirillius/pf/sdn/flow/ResourceFilter.java b/src/main/java/ru/kirillius/pf/sdn/flow/ResourceFilter.java index 7a465b4..ab2f86d 100644 --- a/src/main/java/ru/kirillius/pf/sdn/flow/ResourceFilter.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/ResourceFilter.java @@ -5,85 +5,91 @@ import lombok.extern.slf4j.Slf4j; import ru.kirillius.pf.sdn.api.Networking.IPv4Subnet; import ru.kirillius.pf.sdn.api.Networking.NetworkScope; import ru.kirillius.pf.sdn.api.flow.CallContext; -import ru.kirillius.pf.sdn.api.flow.FlowFunction; -import ru.kirillius.pf.sdn.api.properties.IntegerConstraint; -import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor; -import ru.kirillius.pf.sdn.api.properties.PropertyType; -import java.util.Arrays; -import java.util.Collections; -import java.util.HashMap; -import java.util.Map; -import java.util.regex.Pattern; +import java.util.ArrayList; +import java.util.List; @Slf4j @RequiredArgsConstructor -public class ResourceFilter implements FlowFunction { - private final static String SUBNETS = "subnets"; - private final static String DOMAINS = "domains"; - private final static String ASN = "asn"; +public class ResourceFilter extends CommonFunction { - private final static Map properties; + @Override + public List getInputNames() { + return List.of("input", "filter"); + } - static { - properties = new HashMap<>(); - properties.put(SUBNETS, PropertyDescriptor.builder().type(PropertyType.SUBNET).array(true).required(false).defaultValue(null).build()); - properties.put(DOMAINS, PropertyDescriptor.builder().type(PropertyType.STRING).array(true).required(false).defaultValue(null).build()); - properties.put(ASN, PropertyDescriptor.builder().type(PropertyType.INTEGER).array(true).required(false).defaultValue(null).build().addConstraint(new IntegerConstraint(1, 99999))); //TODO какой max asn?); + @Override + public List getOutputNames() { + return List.of("output", "filtered"); } @Override public void apply(CallContext context) { var output = new NetworkScope(); - var source = context.getSource(); - var properties = context.getProperties(); + var filtered = new NetworkScope(); + var source = context.getSources().get(0); + var filter = context.getSources().get(1); - if (properties.containsKey(SUBNETS)) { - var filtered = Arrays.stream(properties.get(SUBNETS).split(Pattern.quote(","))) - .filter(s -> !s.isBlank()) - .map(IPv4Subnet::new) - .toList(); - output.setSubnets(source.getSubnets().stream().filter(subnet -> !filtered.contains(subnet)).toList()); - } else { - output.setSubnets(source.getSubnets()); + if (!filter.getSubnets().isEmpty()) { + filterSubnets(source, filter, output, filtered); } - - if (properties.containsKey(DOMAINS)) { - var filtered = Arrays.stream(properties.get(DOMAINS).split(Pattern.quote(","))) - .filter(s -> !s.isBlank()) - .toList(); - output.setDomains(source.getDomains().stream().filter(domain -> !filtered.contains(domain)).toList()); - } else { - output.setDomains(source.getDomains()); + if (!filter.getASN().isEmpty()) { + filterASN(source, filter, output, filtered); } - - if (properties.containsKey(ASN)) { - var filtered = Arrays.stream(properties.get(ASN).split(Pattern.quote(","))) - .filter(s -> !s.isBlank()) - .map(Integer::parseInt) - .toList(); - output.setASN(source.getASN().stream().filter(asn -> !filtered.contains(asn)).toList()); - } else { - output.setASN(source.getASN()); + if (!filter.getDomains().isEmpty()) { + filterDomains(source, filter, output, filtered); } - - context.setOutput(output); + context.setOutputs(List.of(output, filtered)); } - @Override - public Map getProperties() { - return Collections.unmodifiableMap(properties); + private void filterSubnets(NetworkScope source, NetworkScope filter, NetworkScope output, NetworkScope filtered) { + var sourceSubnets = source.getSubnets(); + var filterSubnets = filter.getSubnets(); + var outputSubnets = new ArrayList(); + var filteredSubnets = new ArrayList(); + + sourceSubnets.forEach(subnet -> { + if (!filterSubnets.contains(subnet)) { + outputSubnets.add(subnet); + } else { + filteredSubnets.add(subnet); + } + }); + + output.setSubnets(outputSubnets); + filtered.setSubnets(filteredSubnets); } - @Override - public boolean hasInputs() { - return true; + private void filterASN(NetworkScope source, NetworkScope filter, NetworkScope output, NetworkScope filtered) { + var sourceASN = source.getASN(); + var filterASN = filter.getASN(); + var outputASN = new ArrayList(); + var filteredASN = new ArrayList(); + sourceASN.forEach(asn -> { + if (!filterASN.contains(asn)) { + outputASN.add(asn); + } else { + filteredASN.add(asn); + } + }); + output.setASN(outputASN); + filtered.setASN(filteredASN); } - @Override - public boolean hasOutputs() { - return true; + private void filterDomains(NetworkScope source, NetworkScope filter, NetworkScope output, NetworkScope filtered) { + var sourceDomains = source.getDomains(); + var filterDomains = filter.getDomains(); + var outputDomains = new ArrayList(); + var filteredDomains = new ArrayList(); + sourceDomains.forEach(domain -> { + if (!filterDomains.contains(domain)) { + outputDomains.add(domain); + } else { + filteredDomains.add(domain); + } + }); + output.setDomains(outputDomains); + filtered.setDomains(filteredDomains); } - } diff --git a/src/main/java/ru/kirillius/pf/sdn/flow/ResourceFilterObsolete.java b/src/main/java/ru/kirillius/pf/sdn/flow/ResourceFilterObsolete.java new file mode 100644 index 0000000..795a40b --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/flow/ResourceFilterObsolete.java @@ -0,0 +1,107 @@ +package ru.kirillius.pf.sdn.flow; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import ru.kirillius.pf.sdn.api.Networking.IPv4Subnet; +import ru.kirillius.pf.sdn.api.Networking.NetworkScope; +import ru.kirillius.pf.sdn.api.flow.CallContext; +import ru.kirillius.pf.sdn.api.flow.FlowFunction; +import ru.kirillius.pf.sdn.api.properties.IntegerConstraint; +import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor; +import ru.kirillius.pf.sdn.api.properties.PropertyType; + +import java.util.*; +import java.util.regex.Pattern; + +@Slf4j +@RequiredArgsConstructor +@Deprecated +public class ResourceFilterObsolete implements FlowFunction { + + private final static String SUBNETS = "subnets"; + private final static String DOMAINS = "domains"; + private final static String ASN = "asn"; + + private final static Map properties; + + static { + properties = new HashMap<>(); + properties.put(SUBNETS, PropertyDescriptor.builder().type(PropertyType.SUBNET).array(true).required(false).defaultValue(null).build()); + properties.put(DOMAINS, PropertyDescriptor.builder().type(PropertyType.STRING).array(true).required(false).defaultValue(null).build()); + properties.put(ASN, PropertyDescriptor.builder().type(PropertyType.INTEGER).array(true).required(false).defaultValue(null).build().addConstraint(new IntegerConstraint(1, 99999))); //TODO какой max asn?); + } + + @Override + public int getInputCount() { + return 1; + } + + @Override + public int getOutputCount() { + return 1; + } + + @Override + public List getInputNames() { + return List.of("input"); + } + + @Override + public List getOutputNames() { + return List.of("output"); + } + + @Override + public void apply(CallContext context) { + var output = new NetworkScope(); + var source = context.getSources().getFirst(); + var properties = context.getProperties(); + + if (properties.containsKey(SUBNETS)) { + var filtered = Arrays.stream(properties.get(SUBNETS).split(Pattern.quote(","))) + .filter(s -> !s.isBlank()) + .map(IPv4Subnet::new) + .toList(); + output.setSubnets(source.getSubnets().stream().filter(subnet -> !filtered.contains(subnet)).toList()); + } else { + output.setSubnets(source.getSubnets()); + } + + if (properties.containsKey(DOMAINS)) { + var filtered = Arrays.stream(properties.get(DOMAINS).split(Pattern.quote(","))) + .filter(s -> !s.isBlank()) + .toList(); + output.setDomains(source.getDomains().stream().filter(domain -> !filtered.contains(domain)).toList()); + } else { + output.setDomains(source.getDomains()); + } + + if (properties.containsKey(ASN)) { + var filtered = Arrays.stream(properties.get(ASN).split(Pattern.quote(","))) + .filter(s -> !s.isBlank()) + .map(Integer::parseInt) + .toList(); + output.setASN(source.getASN().stream().filter(asn -> !filtered.contains(asn)).toList()); + } else { + output.setASN(source.getASN()); + } + + context.setOutputs(List.of(output)); + } + + @Override + public Map getProperties() { + return Collections.unmodifiableMap(properties); + } + + @Override + public boolean hasInputs() { + return true; + } + + @Override + public boolean hasOutputs() { + return true; + } + +} diff --git a/src/main/java/ru/kirillius/pf/sdn/flow/StaticResourceInput.java b/src/main/java/ru/kirillius/pf/sdn/flow/StaticResourceInput.java new file mode 100644 index 0000000..6143670 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/flow/StaticResourceInput.java @@ -0,0 +1,66 @@ +package ru.kirillius.pf.sdn.flow; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import ru.kirillius.pf.sdn.api.Networking.IPv4Subnet; +import ru.kirillius.pf.sdn.api.Networking.NetworkScope; +import ru.kirillius.pf.sdn.api.flow.CallContext; +import ru.kirillius.pf.sdn.api.properties.IntegerConstraint; +import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor; +import ru.kirillius.pf.sdn.api.properties.PropertyType; + +import java.util.*; +import java.util.regex.Pattern; + +@Slf4j +@RequiredArgsConstructor +public class StaticResourceInput extends CommonFunction { + + private final static String SUBNETS = "subnets"; + private final static String DOMAINS = "domains"; + private final static String ASN = "asn"; + + private final static Map properties; + + static { + properties = new HashMap<>(); + properties.put(SUBNETS, PropertyDescriptor.builder().type(PropertyType.SUBNET).array(true).required(false).defaultValue(null).build()); + properties.put(DOMAINS, PropertyDescriptor.builder().type(PropertyType.STRING).array(true).required(false).defaultValue(null).build()); + properties.put(ASN, PropertyDescriptor.builder().type(PropertyType.INTEGER).array(true).required(false).defaultValue(null).build().addConstraint(new IntegerConstraint(1, 99999))); //TODO какой max asn?); + } + + @Override + public List getOutputNames() { + return List.of("output"); + } + + @Override + public void apply(CallContext context) { + var output = new NetworkScope(); + + var properties = context.getProperties(); + + if (properties.containsKey(SUBNETS)) { + output.setSubnets(Arrays.stream(properties.get(SUBNETS).split(Pattern.quote(","))).filter(s -> !s.isBlank()).map(IPv4Subnet::new).toList()); + } + + if (properties.containsKey(DOMAINS)) { + output.setDomains(Arrays.stream(properties.get(DOMAINS).split(Pattern.quote(","))).filter(s -> !s.isBlank()).toList()); + } + + if (properties.containsKey(ASN)) { + + output.setASN(Arrays.stream(properties.get(ASN).split(Pattern.quote(","))).filter(s -> !s.isBlank()).map(Integer::parseInt).toList()); + } + + context.setOutputs(List.of(output)); + } + + @Override + public Map getProperties() { + return Collections.unmodifiableMap(properties); + } + + + +} diff --git a/src/main/java/ru/kirillius/pf/sdn/flow/SubscriptionInput.java b/src/main/java/ru/kirillius/pf/sdn/flow/SubscriptionInput.java index 9519a0a..f935eeb 100644 --- a/src/main/java/ru/kirillius/pf/sdn/flow/SubscriptionInput.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/SubscriptionInput.java @@ -3,19 +3,15 @@ package ru.kirillius.pf.sdn.flow; import lombok.RequiredArgsConstructor; import ru.kirillius.pf.sdn.api.Networking.NetworkScope; import ru.kirillius.pf.sdn.api.flow.CallContext; -import ru.kirillius.pf.sdn.api.flow.FlowFunction; import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor; import ru.kirillius.pf.sdn.api.properties.PropertyType; import ru.kirillius.pf.sdn.service.SubscriptionService; -import java.util.Arrays; -import java.util.Collections; -import java.util.HashMap; -import java.util.Map; +import java.util.*; import java.util.regex.Pattern; @RequiredArgsConstructor -public class SubscriptionInput implements FlowFunction { +public class SubscriptionInput extends CommonFunction { private final SubscriptionService subscriptionService; @@ -49,12 +45,7 @@ public class SubscriptionInput implements FlowFunction { } @Override - public boolean hasInputs() { - return false; - } - - @Override - public boolean hasOutputs() { - return true; + public List getOutputNames() { + return List.of("output"); } } diff --git a/src/main/java/ru/kirillius/pf/sdn/repository/SubscriptionSetRepository.java b/src/main/java/ru/kirillius/pf/sdn/repository/SubscriptionSetRepository.java new file mode 100644 index 0000000..5a5f336 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/repository/SubscriptionSetRepository.java @@ -0,0 +1,8 @@ +package ru.kirillius.pf.sdn.repository; + +import org.springframework.data.jpa.repository.JpaRepository; +import ru.kirillius.pf.sdn.entity.SubscriptionSetEntry; + +public interface SubscriptionSetRepository extends JpaRepository { + +} diff --git a/src/main/java/ru/kirillius/pf/sdn/service/FlowService.java b/src/main/java/ru/kirillius/pf/sdn/service/FlowService.java index ec1fb4a..f0b2add 100644 --- a/src/main/java/ru/kirillius/pf/sdn/service/FlowService.java +++ b/src/main/java/ru/kirillius/pf/sdn/service/FlowService.java @@ -7,12 +7,13 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; import ru.kirillius.java.utils.events.EventListener; import ru.kirillius.pf.sdn.api.Networking.Subscriptions.Subscription; +import ru.kirillius.pf.sdn.api.flow.Flow; import ru.kirillius.pf.sdn.api.flow.FlowAction; import ru.kirillius.pf.sdn.api.flow.FlowFunction; -import ru.kirillius.pf.sdn.api.flow.Flow; import ru.kirillius.pf.sdn.api.flow.TriggerType; import ru.kirillius.pf.sdn.entity.ActionConfig; import ru.kirillius.pf.sdn.entity.FlowConfig; +import ru.kirillius.pf.sdn.entity.SubscriptionSetEntry; import ru.kirillius.pf.sdn.flow.DummyFunction; import ru.kirillius.pf.sdn.repository.PipelineConfigRepository; @@ -36,6 +37,7 @@ public class FlowService { private final Map> functions = new ConcurrentHashMap<>(); private EventListener executedListener; private EventListener subscriptionUpdateListener; + private EventListener subscriptionSetUpdateListener; public void registerFunction(Class functionClass, String id) { if (functions.containsKey(id)) { @@ -82,11 +84,11 @@ public class FlowService { if (pipelines.containsKey(flowConfig)) { flowExecutorService.cancel(pipelines.get(flowConfig)); } - pipelines.put(flowConfig, new Flow(flowConfig, flowConfig.getActions().stream().map(this::buildAction).toList())); + pipelines.put(flowConfig, new Flow(flowConfig, this::buildAction)); } private FlowAction buildAction(ActionConfig config) { - return new FlowAction(functions.getOrDefault(config.getFunctionId(), DummyFunction.class), config.getProperties()); + return new FlowAction(functions.getOrDefault(config.getFunctionId(), DummyFunction.class), config); } private void checkForIntervalExecution(Flow pipeline) { @@ -126,6 +128,7 @@ public class FlowService { @PostConstruct private void initialize() { subscriptionUpdateListener = subscriptionService.getUpdateEvent().add(this::subscriptionUpdate); + subscriptionSetUpdateListener = subscriptionService.getSetUpdateEvent().add(this::subscriptionSetUpdate); executedListener = flowExecutorService.getOnExecuted().add(this::checkForIntervalExecution); registerFunction(DummyFunction.class, "Error:fallback"); configRepository.findAll().forEach(this::load); @@ -139,6 +142,32 @@ public class FlowService { .forEach(flowExecutorService::triggerExecute); } + private void subscriptionSetUpdate(SubscriptionSetEntry setEntry) { + pipelines.values().forEach(pipeline -> { + var matchedConditions = pipeline.getConfig().getStartConditions().stream().filter(c -> c.getTrigger() == TriggerType.OnSubscriptionSetUpdate).toList(); + if (matchedConditions.isEmpty()) { + return; + } + + matchedConditions.forEach(c -> { + if (c.isDontStartIfRunning() && pipeline.isRunning()) { + return; + } + + var properties = c.getProperties(); + var names = properties.getOrDefault("names", null); + if (names == null || names.isEmpty()) { + flowExecutorService.triggerExecute(pipeline); + return; + } + + if (Arrays.stream(names.split(Pattern.quote(","))).anyMatch(s -> setEntry.getName().equals(s))) { + flowExecutorService.triggerExecute(pipeline); + } + }); + }); + } + private void subscriptionUpdate(Subscription subscription) { pipelines.values().forEach(pipeline -> { @@ -177,6 +206,10 @@ public class FlowService { subscriptionService.getUpdateEvent().remove(subscriptionUpdateListener); subscriptionUpdateListener = null; } + if (subscriptionSetUpdateListener != null) { + subscriptionService.getSetUpdateEvent().remove(subscriptionSetUpdateListener); + subscriptionSetUpdateListener = null; + } } } diff --git a/src/main/java/ru/kirillius/pf/sdn/service/SubscriptionService.java b/src/main/java/ru/kirillius/pf/sdn/service/SubscriptionService.java index 0020ef7..ccb6eb8 100644 --- a/src/main/java/ru/kirillius/pf/sdn/service/SubscriptionService.java +++ b/src/main/java/ru/kirillius/pf/sdn/service/SubscriptionService.java @@ -1,6 +1,5 @@ package ru.kirillius.pf.sdn.service; -import jakarta.annotation.PostConstruct; import jakarta.annotation.PreDestroy; import lombok.Getter; import lombok.RequiredArgsConstructor; @@ -14,7 +13,9 @@ import ru.kirillius.pf.sdn.api.Networking.Subscriptions.CacheFallbackProviderPro 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.entity.SubscriptionSetEntry; import ru.kirillius.pf.sdn.repository.SubscriptionProviderConfigRepository; +import ru.kirillius.pf.sdn.repository.SubscriptionSetRepository; import java.time.Duration; import java.time.Instant; @@ -41,8 +42,13 @@ public class SubscriptionService { private final AtomicInteger updateCounter = new AtomicInteger(0); private final List subscriptions = new CopyOnWriteArrayList<>(); private final ApplicationContext context; + private final SubscriptionSetRepository subscriptionSetRepository; + @Getter private final EventHandler updateEvent = new ConcurrentEventHandler<>(); + + @Getter + private final EventHandler setUpdateEvent = new ConcurrentEventHandler<>(); private Future worker; public List getSubscriptions() { @@ -84,6 +90,8 @@ public class SubscriptionService { var retrieved = new ArrayList(); var updated = new ArrayList(); + var sets = subscriptionSetRepository.findAll(); + providers.values().forEach(provider -> { updated.addAll(provider.update()); retrieved.addAll(provider.getSubscriptions()); @@ -95,16 +103,25 @@ public class SubscriptionService { } 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()); - } + invokeEvent(updateEvent, s); + sets.forEach(set -> { + if (set.getSubscriptions().contains(s)) { //TODO suspicious call??? + invokeEvent(setUpdateEvent, set); + } + }); }); } - @PostConstruct - private void initialize() { + private void invokeEvent(EventHandler eventHandler, T data) { + try { + eventHandler.invoke(data); + } catch (Exception e) { + log.error("Failed to invoke subscription update event because of error {}:{}", e.getClass().getSimpleName(), e.getMessage()); + } + } + + + private void start() { reloadProviders(); reloadSubscriptions();