From ab24066f816241f9db109f2827ffb3b64198cc13 Mon Sep 17 00:00:00 2001 From: kirillius Date: Sat, 11 Jul 2026 11:43:14 +0300 Subject: [PATCH] WIP --- .gitignore | 1 + .../pf/sdn/api/flow/CallContext.java | 20 ++++++ .../ExecutionInfoEntry.java | 2 +- .../Flow.java} | 6 +- .../Action.java => flow/FlowAction.java} | 19 +++--- .../pf/sdn/api/flow/FlowFunction.java | 15 +++++ .../{pipeline => flow}/StartCondition.java | 2 +- .../api/{pipeline => flow}/TriggerType.java | 2 +- .../pf/sdn/api/pipeline/PipelineFunction.java | 12 ---- .../pf/sdn/entity/PipelineConfig.java | 2 +- .../{pipeline => flow}/AutoResolveASN.java | 45 ++++++++------ .../AutoResolveDomains.java | 61 +++++++++++-------- .../sdn/{pipeline => flow}/DebugOutput.java | 23 ++++--- .../kirillius/pf/sdn/flow/DummyFunction.java | 30 +++++++++ .../pf/sdn/flow/DummyPassthrough.java | 36 +++++++++++ .../{pipeline => flow}/DuplicatesFilter.java | 22 +++++-- .../MergeNeighbourSubnetFilter.java | 21 +++++-- .../MergeWithCoverageSubnetFilter.java | 52 ++++++++++------ .../OverlappedDomainFilter.java | 22 +++++-- .../OverlappedSubnetFilter.java | 22 +++++-- .../{pipeline => flow}/ResourceFilter.java | 49 +++++++++------ .../{pipeline => flow}/SubscriptionInput.java | 48 +++++++++------ .../pf/sdn/pipeline/DummyPassthrough.java | 25 -------- .../pf/sdn/service/DummyFunction.java | 20 ------ ...rService.java => FlowExecutorService.java} | 20 +++--- ...{PipelineService.java => FlowService.java} | 49 +++++++-------- 26 files changed, 396 insertions(+), 230 deletions(-) create mode 100644 src/main/java/ru/kirillius/pf/sdn/api/flow/CallContext.java rename src/main/java/ru/kirillius/pf/sdn/api/{pipeline => flow}/ExecutionInfoEntry.java (97%) rename src/main/java/ru/kirillius/pf/sdn/api/{pipeline/ProcessingPipeline.java => flow/Flow.java} (93%) rename src/main/java/ru/kirillius/pf/sdn/api/{pipeline/Action.java => flow/FlowAction.java} (56%) create mode 100644 src/main/java/ru/kirillius/pf/sdn/api/flow/FlowFunction.java rename src/main/java/ru/kirillius/pf/sdn/api/{pipeline => flow}/StartCondition.java (94%) rename src/main/java/ru/kirillius/pf/sdn/api/{pipeline => flow}/TriggerType.java (68%) delete mode 100644 src/main/java/ru/kirillius/pf/sdn/api/pipeline/PipelineFunction.java rename src/main/java/ru/kirillius/pf/sdn/{pipeline => flow}/AutoResolveASN.java (67%) rename src/main/java/ru/kirillius/pf/sdn/{pipeline => flow}/AutoResolveDomains.java (76%) rename src/main/java/ru/kirillius/pf/sdn/{pipeline => flow}/DebugOutput.java (50%) create mode 100644 src/main/java/ru/kirillius/pf/sdn/flow/DummyFunction.java create mode 100644 src/main/java/ru/kirillius/pf/sdn/flow/DummyPassthrough.java rename src/main/java/ru/kirillius/pf/sdn/{pipeline => flow}/DuplicatesFilter.java (63%) rename src/main/java/ru/kirillius/pf/sdn/{pipeline => flow}/MergeNeighbourSubnetFilter.java (62%) rename src/main/java/ru/kirillius/pf/sdn/{pipeline => flow}/MergeWithCoverageSubnetFilter.java (74%) rename src/main/java/ru/kirillius/pf/sdn/{pipeline => flow}/OverlappedDomainFilter.java (61%) rename src/main/java/ru/kirillius/pf/sdn/{pipeline => flow}/OverlappedSubnetFilter.java (61%) rename src/main/java/ru/kirillius/pf/sdn/{pipeline => flow}/ResourceFilter.java (84%) rename src/main/java/ru/kirillius/pf/sdn/{pipeline => flow}/SubscriptionInput.java (77%) delete mode 100644 src/main/java/ru/kirillius/pf/sdn/pipeline/DummyPassthrough.java delete mode 100644 src/main/java/ru/kirillius/pf/sdn/service/DummyFunction.java rename src/main/java/ru/kirillius/pf/sdn/service/{PipelineExecutorService.java => FlowExecutorService.java} (74%) rename src/main/java/ru/kirillius/pf/sdn/service/{PipelineService.java => FlowService.java} (74%) diff --git a/.gitignore b/.gitignore index d6adc9d..b6a5a44 100644 --- a/.gitignore +++ b/.gitignore @@ -49,3 +49,4 @@ cache/ /data/ test/ pfsdn.mv.db +webui/node_modules/ 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 new file mode 100644 index 0000000..f34e40c --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/flow/CallContext.java @@ -0,0 +1,20 @@ +package ru.kirillius.pf.sdn.api.flow; + +import lombok.Getter; +import lombok.RequiredArgsConstructor; +import lombok.Setter; +import ru.kirillius.pf.sdn.api.Networking.NetworkScope; + +import java.util.Map; + +@RequiredArgsConstructor +public final class CallContext { + @Getter + private final NetworkScope source; + @Getter + @Setter + private NetworkScope output; + + @Getter + private final Map properties; +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/ExecutionInfoEntry.java b/src/main/java/ru/kirillius/pf/sdn/api/flow/ExecutionInfoEntry.java similarity index 97% rename from src/main/java/ru/kirillius/pf/sdn/api/pipeline/ExecutionInfoEntry.java rename to src/main/java/ru/kirillius/pf/sdn/api/flow/ExecutionInfoEntry.java index 1123c10..366316e 100644 --- a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/ExecutionInfoEntry.java +++ b/src/main/java/ru/kirillius/pf/sdn/api/flow/ExecutionInfoEntry.java @@ -1,4 +1,4 @@ -package ru.kirillius.pf.sdn.api.pipeline; +package ru.kirillius.pf.sdn.api.flow; import lombok.Getter; import ru.kirillius.pf.sdn.api.Networking.NetworkScope; diff --git a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/ProcessingPipeline.java b/src/main/java/ru/kirillius/pf/sdn/api/flow/Flow.java similarity index 93% rename from src/main/java/ru/kirillius/pf/sdn/api/pipeline/ProcessingPipeline.java rename to src/main/java/ru/kirillius/pf/sdn/api/flow/Flow.java index 27b4015..57602f5 100644 --- a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/ProcessingPipeline.java +++ b/src/main/java/ru/kirillius/pf/sdn/api/flow/Flow.java @@ -1,4 +1,4 @@ -package ru.kirillius.pf.sdn.api.pipeline; +package ru.kirillius.pf.sdn.api.flow; import lombok.*; import ru.kirillius.pf.sdn.api.Networking.NetworkScope; @@ -11,10 +11,10 @@ import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; @RequiredArgsConstructor() -public class ProcessingPipeline { +public class Flow { @Getter private final PipelineConfig config; - private final List actions; + private final List actions; private final AtomicBoolean running = new AtomicBoolean(false); private final AtomicBoolean interrupted = new AtomicBoolean(false); diff --git a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/Action.java b/src/main/java/ru/kirillius/pf/sdn/api/flow/FlowAction.java similarity index 56% rename from src/main/java/ru/kirillius/pf/sdn/api/pipeline/Action.java rename to src/main/java/ru/kirillius/pf/sdn/api/flow/FlowAction.java index 4686076..87d8129 100644 --- a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/Action.java +++ b/src/main/java/ru/kirillius/pf/sdn/api/flow/FlowAction.java @@ -1,4 +1,4 @@ -package ru.kirillius.pf.sdn.api.pipeline; +package ru.kirillius.pf.sdn.api.flow; import lombok.Builder; import lombok.Getter; @@ -10,11 +10,11 @@ import java.util.Map; import java.util.function.Function; @Builder -public class Action implements Function { +public class FlowAction implements Function { - private PipelineFunction function; + private FlowFunction function; - public Action(Class functionClass, Map properties) { + public FlowAction(Class functionClass, Map properties) { this.functionClass = functionClass; this.properties = properties; instantiateFunction(); @@ -25,19 +25,22 @@ public class Action implements Function { function = functionClass.getConstructor().newInstance(); } - public void setFunctionClass(Class functionClass) { + public void setFunctionClass(Class functionClass) { this.functionClass = functionClass; instantiateFunction(); } @Getter - private Class functionClass; + private Class functionClass; @Getter @Setter private Map properties; @Override - public NetworkScope apply(NetworkScope source) { - return function.apply(source, properties); + public CallContext apply(NetworkScope networkScope) { + var context = new CallContext(networkScope, 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 new file mode 100644 index 0000000..8ff9a4a --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/api/flow/FlowFunction.java @@ -0,0 +1,15 @@ +package ru.kirillius.pf.sdn.api.flow; + +import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor; + +import java.util.Map; + +public interface FlowFunction { + void apply(CallContext context); + + Map getProperties(); + + boolean hasInputs(); + + boolean hasOutputs(); +} diff --git a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/StartCondition.java b/src/main/java/ru/kirillius/pf/sdn/api/flow/StartCondition.java similarity index 94% rename from src/main/java/ru/kirillius/pf/sdn/api/pipeline/StartCondition.java rename to src/main/java/ru/kirillius/pf/sdn/api/flow/StartCondition.java index 6c539de..bfb0fa5 100644 --- a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/StartCondition.java +++ b/src/main/java/ru/kirillius/pf/sdn/api/flow/StartCondition.java @@ -1,4 +1,4 @@ -package ru.kirillius.pf.sdn.api.pipeline; +package ru.kirillius.pf.sdn.api.flow; import jakarta.persistence.*; import lombok.Getter; diff --git a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/TriggerType.java b/src/main/java/ru/kirillius/pf/sdn/api/flow/TriggerType.java similarity index 68% rename from src/main/java/ru/kirillius/pf/sdn/api/pipeline/TriggerType.java rename to src/main/java/ru/kirillius/pf/sdn/api/flow/TriggerType.java index 419e1e5..1f8f04f 100644 --- a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/TriggerType.java +++ b/src/main/java/ru/kirillius/pf/sdn/api/flow/TriggerType.java @@ -1,4 +1,4 @@ -package ru.kirillius.pf.sdn.api.pipeline; +package ru.kirillius.pf.sdn.api.flow; public enum TriggerType { Manual, 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 deleted file mode 100644 index 27d0232..0000000 --- a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/PipelineFunction.java +++ /dev/null @@ -1,12 +0,0 @@ -package ru.kirillius.pf.sdn.api.pipeline; - -import ru.kirillius.pf.sdn.api.Networking.NetworkScope; -import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor; - -import java.util.Map; - -public interface PipelineFunction { - NetworkScope apply(NetworkScope source, Map properties); - - Map getProperties(); -} diff --git a/src/main/java/ru/kirillius/pf/sdn/entity/PipelineConfig.java b/src/main/java/ru/kirillius/pf/sdn/entity/PipelineConfig.java index 80b9a3b..132055a 100644 --- a/src/main/java/ru/kirillius/pf/sdn/entity/PipelineConfig.java +++ b/src/main/java/ru/kirillius/pf/sdn/entity/PipelineConfig.java @@ -2,7 +2,7 @@ package ru.kirillius.pf.sdn.entity; import jakarta.persistence.*; import lombok.*; -import ru.kirillius.pf.sdn.api.pipeline.StartCondition; +import ru.kirillius.pf.sdn.api.flow.StartCondition; import java.util.ArrayList; import java.util.List; diff --git a/src/main/java/ru/kirillius/pf/sdn/pipeline/AutoResolveASN.java b/src/main/java/ru/kirillius/pf/sdn/flow/AutoResolveASN.java similarity index 67% rename from src/main/java/ru/kirillius/pf/sdn/pipeline/AutoResolveASN.java rename to src/main/java/ru/kirillius/pf/sdn/flow/AutoResolveASN.java index 8d0cc73..1ad29a1 100644 --- a/src/main/java/ru/kirillius/pf/sdn/pipeline/AutoResolveASN.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/AutoResolveASN.java @@ -1,8 +1,9 @@ -package ru.kirillius.pf.sdn.pipeline; +package ru.kirillius.pf.sdn.flow; import lombok.RequiredArgsConstructor; import ru.kirillius.pf.sdn.api.Networking.NetworkScope; -import ru.kirillius.pf.sdn.api.pipeline.PipelineFunction; +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; @@ -12,23 +13,9 @@ import java.util.HashMap; import java.util.Map; @RequiredArgsConstructor -public class AutoResolveASN implements PipelineFunction { +public class AutoResolveASN implements FlowFunction { private final AutonomousSystemCacheService autonomousSystemCacheService; - - @Override - public NetworkScope apply(NetworkScope source, Map properties) { - var output = new NetworkScope(); - output.add(source); - - if (properties.containsKey(CLEAR) && properties.get(CLEAR).equals(Boolean.TRUE.toString())) { - output.getASN().clear(); - } - - source.getASN().forEach(asn -> output.getSubnets().addAll(autonomousSystemCacheService.load(asn))); - return output; - } - private final static String CLEAR = "clear"; private final static Map properties; @@ -44,8 +31,32 @@ public class AutoResolveASN implements PipelineFunction { ); } + @Override + public void apply(CallContext context) { + var output = new NetworkScope(); + output.add(context.getSource()); + 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); + } + @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/pipeline/AutoResolveDomains.java b/src/main/java/ru/kirillius/pf/sdn/flow/AutoResolveDomains.java similarity index 76% rename from src/main/java/ru/kirillius/pf/sdn/pipeline/AutoResolveDomains.java rename to src/main/java/ru/kirillius/pf/sdn/flow/AutoResolveDomains.java index 5ab025a..68d617a 100644 --- a/src/main/java/ru/kirillius/pf/sdn/pipeline/AutoResolveDomains.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/AutoResolveDomains.java @@ -1,9 +1,10 @@ -package ru.kirillius.pf.sdn.pipeline; +package ru.kirillius.pf.sdn.flow; 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.PipelineFunction; +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; @@ -13,31 +14,10 @@ import java.util.HashMap; import java.util.Map; @RequiredArgsConstructor -public class AutoResolveDomains implements PipelineFunction { +public class AutoResolveDomains implements FlowFunction { private final DomainUpdaterService domainCacheService; - @Override - public NetworkScope apply(NetworkScope source, Map properties) { - var output = new NetworkScope(); - output.add(source); - - if (properties.containsKey(CLEAR) && properties.get(CLEAR).equals("true")) { - output.getAutoResolvedDomains().clear(); - } - - source.getAutoResolvedDomains() - .forEach(domain -> domainCacheService.getActualAddresses(domain) - .forEach(address -> output - .getSubnets() - .add(new IPv4Subnet(address, 32) - ) - ) - ); - - return output; - } - private final static String CLEAR = "clear"; private final static Map properties; @@ -53,8 +33,41 @@ public class AutoResolveDomains implements PipelineFunction { ); } + @Override + public void apply(CallContext context) { + var source = context.getSource(); + var properties = context.getProperties(); + var output = new NetworkScope(); + output.add(source); + + if (properties.containsKey(CLEAR) && properties.get(CLEAR).equals("true")) { + output.getAutoResolvedDomains().clear(); + } + + source.getAutoResolvedDomains() + .forEach(domain -> domainCacheService.getActualAddresses(domain) + .forEach(address -> output + .getSubnets() + .add(new IPv4Subnet(address, 32) + ) + ) + ); + + context.setOutput(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/pipeline/DebugOutput.java b/src/main/java/ru/kirillius/pf/sdn/flow/DebugOutput.java similarity index 50% rename from src/main/java/ru/kirillius/pf/sdn/pipeline/DebugOutput.java rename to src/main/java/ru/kirillius/pf/sdn/flow/DebugOutput.java index c25731c..34b8151 100644 --- a/src/main/java/ru/kirillius/pf/sdn/pipeline/DebugOutput.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/DebugOutput.java @@ -1,10 +1,10 @@ -package ru.kirillius.pf.sdn.pipeline; +package ru.kirillius.pf.sdn.flow; import com.fasterxml.jackson.databind.ObjectMapper; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; -import ru.kirillius.pf.sdn.api.Networking.NetworkScope; -import ru.kirillius.pf.sdn.api.pipeline.PipelineFunction; +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; @@ -12,13 +12,12 @@ import java.util.Map; @Slf4j @RequiredArgsConstructor -public class DebugOutput implements PipelineFunction { +public class DebugOutput implements FlowFunction { private final ObjectMapper objectMapper; @Override - public NetworkScope apply(NetworkScope source, Map properties) { - log.info(objectMapper.valueToTree(source).toString()); - return source; + public void apply(CallContext context) { + log.info(objectMapper.valueToTree(context.getSource()).toString()); } @Override @@ -26,4 +25,14 @@ public class DebugOutput implements PipelineFunction { return Collections.emptyMap(); } + @Override + public boolean hasInputs() { + return true; + } + + @Override + public boolean hasOutputs() { + return false; + } + } diff --git a/src/main/java/ru/kirillius/pf/sdn/flow/DummyFunction.java b/src/main/java/ru/kirillius/pf/sdn/flow/DummyFunction.java new file mode 100644 index 0000000..e458961 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/flow/DummyFunction.java @@ -0,0 +1,30 @@ +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 { + + @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 new file mode 100644 index 0000000..7984372 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/flow/DummyPassthrough.java @@ -0,0 +1,36 @@ +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 java.util.Collections; +import java.util.Map; + +@RequiredArgsConstructor +public class DummyPassthrough implements FlowFunction { + + @Override + public void apply(CallContext context) { + var output = new NetworkScope(); + output.add(context.getSource()); + context.setOutput(output); + } + + @Override + public Map getProperties() { + 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/pipeline/DuplicatesFilter.java b/src/main/java/ru/kirillius/pf/sdn/flow/DuplicatesFilter.java similarity index 63% rename from src/main/java/ru/kirillius/pf/sdn/pipeline/DuplicatesFilter.java rename to src/main/java/ru/kirillius/pf/sdn/flow/DuplicatesFilter.java index b6c7dff..dddd0a7 100644 --- a/src/main/java/ru/kirillius/pf/sdn/pipeline/DuplicatesFilter.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/DuplicatesFilter.java @@ -1,9 +1,10 @@ -package ru.kirillius.pf.sdn.pipeline; +package ru.kirillius.pf.sdn.flow; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import ru.kirillius.pf.sdn.api.Networking.NetworkScope; -import ru.kirillius.pf.sdn.api.pipeline.PipelineFunction; +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; @@ -11,16 +12,17 @@ import java.util.Map; @Slf4j @RequiredArgsConstructor -public class DuplicatesFilter implements PipelineFunction { +public class DuplicatesFilter implements FlowFunction { @Override - public NetworkScope apply(NetworkScope source, Map properties) { + public void apply(CallContext context) { + var source = context.getSource(); 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()); - return output; + context.setOutput(output); } @Override @@ -28,4 +30,14 @@ public class DuplicatesFilter implements PipelineFunction { 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/pipeline/MergeNeighbourSubnetFilter.java b/src/main/java/ru/kirillius/pf/sdn/flow/MergeNeighbourSubnetFilter.java similarity index 62% rename from src/main/java/ru/kirillius/pf/sdn/pipeline/MergeNeighbourSubnetFilter.java rename to src/main/java/ru/kirillius/pf/sdn/flow/MergeNeighbourSubnetFilter.java index da46504..77945e8 100644 --- a/src/main/java/ru/kirillius/pf/sdn/pipeline/MergeNeighbourSubnetFilter.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/MergeNeighbourSubnetFilter.java @@ -1,10 +1,11 @@ -package ru.kirillius.pf.sdn.pipeline; +package ru.kirillius.pf.sdn.flow; import lombok.RequiredArgsConstructor; 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.pipeline.PipelineFunction; +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; @@ -12,16 +13,16 @@ import java.util.Map; @Slf4j @RequiredArgsConstructor -public class MergeNeighbourSubnetFilter implements PipelineFunction { +public class MergeNeighbourSubnetFilter implements FlowFunction { @Override - public NetworkScope apply(NetworkScope source, Map properties) { + public void apply(CallContext context) { + var source = context.getSource(); var output = new NetworkScope(); output.setSubnets(IPv4Util.mergeNeighbours(source.getSubnets())); output.setDomains(source.getDomains()); output.setASN(source.getASN()); output.setAutoResolvedDomains(source.getAutoResolvedDomains()); - return output; } @Override @@ -29,4 +30,14 @@ public class MergeNeighbourSubnetFilter implements PipelineFunction { 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/pipeline/MergeWithCoverageSubnetFilter.java b/src/main/java/ru/kirillius/pf/sdn/flow/MergeWithCoverageSubnetFilter.java similarity index 74% rename from src/main/java/ru/kirillius/pf/sdn/pipeline/MergeWithCoverageSubnetFilter.java rename to src/main/java/ru/kirillius/pf/sdn/flow/MergeWithCoverageSubnetFilter.java index 40f65e5..4390fcd 100644 --- a/src/main/java/ru/kirillius/pf/sdn/pipeline/MergeWithCoverageSubnetFilter.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/MergeWithCoverageSubnetFilter.java @@ -1,10 +1,11 @@ -package ru.kirillius.pf.sdn.pipeline; +package ru.kirillius.pf.sdn.flow; import lombok.RequiredArgsConstructor; 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.pipeline.PipelineFunction; +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; @@ -15,25 +16,8 @@ import java.util.Map; @Slf4j @RequiredArgsConstructor -public class MergeWithCoverageSubnetFilter implements PipelineFunction { +public class MergeWithCoverageSubnetFilter implements FlowFunction { private final static String USAGE = "usage"; - - @Override - public NetworkScope apply(NetworkScope source, Map properties) { - var usage = properties.containsKey(USAGE) ? Integer.parseInt(properties.get(USAGE)) : 75; - if (usage < 51) { - usage = 51; - } else if (usage > 100) { - usage = 100; - } - var output = new NetworkScope(); - output.setSubnets(IPv4Util.mergeToLargerIfCovered(source.getSubnets(), usage)); - output.setDomains(source.getDomains()); - output.setASN(source.getASN()); - output.setAutoResolvedDomains(source.getAutoResolvedDomains()); - return output; - } - private final static Map properties; static { @@ -44,9 +28,37 @@ public class MergeWithCoverageSubnetFilter implements PipelineFunction { ); } + @Override + public void apply(CallContext context) { + var properties = context.getProperties(); + var source = context.getSource(); + var usage = properties.containsKey(USAGE) ? Integer.parseInt(properties.get(USAGE)) : 75; + if (usage < 51) { + usage = 51; + } else if (usage > 100) { + usage = 100; + } + var output = new NetworkScope(); + output.setSubnets(IPv4Util.mergeToLargerIfCovered(source.getSubnets(), usage)); + output.setDomains(source.getDomains()); + output.setASN(source.getASN()); + output.setAutoResolvedDomains(source.getAutoResolvedDomains()); + context.setOutput(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/pipeline/OverlappedDomainFilter.java b/src/main/java/ru/kirillius/pf/sdn/flow/OverlappedDomainFilter.java similarity index 61% rename from src/main/java/ru/kirillius/pf/sdn/pipeline/OverlappedDomainFilter.java rename to src/main/java/ru/kirillius/pf/sdn/flow/OverlappedDomainFilter.java index dda5f7f..f1aaa4c 100644 --- a/src/main/java/ru/kirillius/pf/sdn/pipeline/OverlappedDomainFilter.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/OverlappedDomainFilter.java @@ -1,10 +1,11 @@ -package ru.kirillius.pf.sdn.pipeline; +package ru.kirillius.pf.sdn.flow; import lombok.RequiredArgsConstructor; 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.pipeline.PipelineFunction; +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; @@ -12,16 +13,17 @@ import java.util.Map; @Slf4j @RequiredArgsConstructor -public class OverlappedDomainFilter implements PipelineFunction { +public class OverlappedDomainFilter implements FlowFunction { @Override - public NetworkScope apply(NetworkScope source, Map properties) { + public void apply(CallContext context) { + var source = context.getSource(); var output = new NetworkScope(); output.setSubnets(source.getSubnets()); output.setDomains(DomainUtil.removeOverlapped(source.getDomains())); output.setASN(source.getASN()); output.setAutoResolvedDomains(source.getAutoResolvedDomains()); - return output; + context.setOutput(output); } @Override @@ -29,4 +31,14 @@ public class OverlappedDomainFilter implements PipelineFunction { 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/pipeline/OverlappedSubnetFilter.java b/src/main/java/ru/kirillius/pf/sdn/flow/OverlappedSubnetFilter.java similarity index 61% rename from src/main/java/ru/kirillius/pf/sdn/pipeline/OverlappedSubnetFilter.java rename to src/main/java/ru/kirillius/pf/sdn/flow/OverlappedSubnetFilter.java index 90e14de..719120c 100644 --- a/src/main/java/ru/kirillius/pf/sdn/pipeline/OverlappedSubnetFilter.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/OverlappedSubnetFilter.java @@ -1,10 +1,11 @@ -package ru.kirillius.pf.sdn.pipeline; +package ru.kirillius.pf.sdn.flow; import lombok.RequiredArgsConstructor; 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.pipeline.PipelineFunction; +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; @@ -12,16 +13,17 @@ import java.util.Map; @Slf4j @RequiredArgsConstructor -public class OverlappedSubnetFilter implements PipelineFunction { +public class OverlappedSubnetFilter implements FlowFunction { @Override - public NetworkScope apply(NetworkScope source, Map properties) { + public void apply(CallContext context) { + var source = context.getSource(); var output = new NetworkScope(); output.setSubnets(IPv4Util.removeOverlapped(source.getSubnets())); output.setDomains(source.getDomains()); output.setASN(source.getASN()); output.setAutoResolvedDomains(source.getAutoResolvedDomains()); - return output; + context.setOutput(output); } @Override @@ -29,4 +31,14 @@ public class OverlappedSubnetFilter implements PipelineFunction { 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/pipeline/ResourceFilter.java b/src/main/java/ru/kirillius/pf/sdn/flow/ResourceFilter.java similarity index 84% rename from src/main/java/ru/kirillius/pf/sdn/pipeline/ResourceFilter.java rename to src/main/java/ru/kirillius/pf/sdn/flow/ResourceFilter.java index ec5ccaf..7a465b4 100644 --- a/src/main/java/ru/kirillius/pf/sdn/pipeline/ResourceFilter.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/ResourceFilter.java @@ -1,10 +1,11 @@ -package ru.kirillius.pf.sdn.pipeline; +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.pipeline.PipelineFunction; +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; @@ -17,11 +18,26 @@ import java.util.regex.Pattern; @Slf4j @RequiredArgsConstructor -public class ResourceFilter implements PipelineFunction { +public class ResourceFilter 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 NetworkScope apply(NetworkScope source, Map properties) { + public void apply(CallContext context) { var output = new NetworkScope(); + var source = context.getSource(); + var properties = context.getProperties(); if (properties.containsKey(SUBNETS)) { var filtered = Arrays.stream(properties.get(SUBNETS).split(Pattern.quote(","))) @@ -52,20 +68,7 @@ public class ResourceFilter implements PipelineFunction { output.setASN(source.getASN()); } - return output; - } - - 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?); + context.setOutput(output); } @Override @@ -73,4 +76,14 @@ public class ResourceFilter implements PipelineFunction { 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/pipeline/SubscriptionInput.java b/src/main/java/ru/kirillius/pf/sdn/flow/SubscriptionInput.java similarity index 77% rename from src/main/java/ru/kirillius/pf/sdn/pipeline/SubscriptionInput.java rename to src/main/java/ru/kirillius/pf/sdn/flow/SubscriptionInput.java index 396d6f9..9519a0a 100644 --- a/src/main/java/ru/kirillius/pf/sdn/pipeline/SubscriptionInput.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/SubscriptionInput.java @@ -1,8 +1,9 @@ -package ru.kirillius.pf.sdn.pipeline; +package ru.kirillius.pf.sdn.flow; import lombok.RequiredArgsConstructor; import ru.kirillius.pf.sdn.api.Networking.NetworkScope; -import ru.kirillius.pf.sdn.api.pipeline.PipelineFunction; +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; @@ -14,25 +15,10 @@ import java.util.Map; import java.util.regex.Pattern; @RequiredArgsConstructor -public class SubscriptionInput implements PipelineFunction { +public class SubscriptionInput implements FlowFunction { private final SubscriptionService subscriptionService; - @Override - public NetworkScope apply(NetworkScope source, Map properties) { - var subscriptions = subscriptionService.getSubscriptions(); - if (properties.containsKey(NAMES)) { - var names = properties.get(NAMES); - var filter = Arrays.stream(names.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 -> bundle.add(subscription.getScope())); - return bundle; - } - private final static String NAMES = "names"; private final static Map properties; @@ -41,8 +27,34 @@ public class SubscriptionInput implements PipelineFunction { properties.put(NAMES, PropertyDescriptor.builder().type(PropertyType.SUBSCRIPTION).array(true).required(false).defaultValue(null).build()); } + @Override + public void apply(CallContext context) { + + var properties = context.getProperties(); + var subscriptions = subscriptionService.getSubscriptions(); + if (properties.containsKey(NAMES)) { + var names = properties.get(NAMES); + var filter = Arrays.stream(names.split(Pattern.quote(","))).filter(s -> !s.isBlank()).toList(); + subscriptions = subscriptions.stream().filter(s -> filter.contains(s.getName())).toList(); + } + + var bundle = new NetworkScope(); + + subscriptions.forEach(subscription -> bundle.add(subscription.getScope())); + } + @Override public Map getProperties() { return Collections.unmodifiableMap(properties); } + + @Override + public boolean hasInputs() { + return false; + } + + @Override + public boolean hasOutputs() { + return true; + } } diff --git a/src/main/java/ru/kirillius/pf/sdn/pipeline/DummyPassthrough.java b/src/main/java/ru/kirillius/pf/sdn/pipeline/DummyPassthrough.java deleted file mode 100644 index bcdb31d..0000000 --- a/src/main/java/ru/kirillius/pf/sdn/pipeline/DummyPassthrough.java +++ /dev/null @@ -1,25 +0,0 @@ -package ru.kirillius.pf.sdn.pipeline; - -import lombok.RequiredArgsConstructor; -import ru.kirillius.pf.sdn.api.Networking.NetworkScope; -import ru.kirillius.pf.sdn.api.pipeline.PipelineFunction; -import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor; - -import java.util.Collections; -import java.util.Map; - -@RequiredArgsConstructor -public class DummyPassthrough implements PipelineFunction { - - @Override - public NetworkScope apply(NetworkScope source, Map properties) { - var output = new NetworkScope(); - output.add(source); - return output; - } - - @Override - public Map getProperties() { - return Collections.emptyMap(); - } -} diff --git a/src/main/java/ru/kirillius/pf/sdn/service/DummyFunction.java b/src/main/java/ru/kirillius/pf/sdn/service/DummyFunction.java deleted file mode 100644 index 523df82..0000000 --- a/src/main/java/ru/kirillius/pf/sdn/service/DummyFunction.java +++ /dev/null @@ -1,20 +0,0 @@ -package ru.kirillius.pf.sdn.service; - -import ru.kirillius.pf.sdn.api.Networking.NetworkScope; -import ru.kirillius.pf.sdn.api.pipeline.PipelineFunction; -import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor; - -import java.util.Map; - -public class DummyFunction implements PipelineFunction { - - @Override - public NetworkScope apply(NetworkScope source, Map properties) { - throw new UnsupportedOperationException("Not implemented"); - } - - @Override - public Map getProperties() { - return Map.of(); - } -} diff --git a/src/main/java/ru/kirillius/pf/sdn/service/PipelineExecutorService.java b/src/main/java/ru/kirillius/pf/sdn/service/FlowExecutorService.java similarity index 74% rename from src/main/java/ru/kirillius/pf/sdn/service/PipelineExecutorService.java rename to src/main/java/ru/kirillius/pf/sdn/service/FlowExecutorService.java index a27f90d..978629f 100644 --- a/src/main/java/ru/kirillius/pf/sdn/service/PipelineExecutorService.java +++ b/src/main/java/ru/kirillius/pf/sdn/service/FlowExecutorService.java @@ -9,7 +9,7 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; import ru.kirillius.java.utils.events.ConcurrentEventHandler; import ru.kirillius.java.utils.events.EventHandler; -import ru.kirillius.pf.sdn.api.pipeline.ProcessingPipeline; +import ru.kirillius.pf.sdn.api.flow.Flow; import java.util.Queue; import java.util.concurrent.ConcurrentLinkedQueue; @@ -20,28 +20,28 @@ import java.util.concurrent.atomic.AtomicReference; @Slf4j @Service @RequiredArgsConstructor -public class PipelineExecutorService { +public class FlowExecutorService { private final ExecutorService executor; - private final Queue executionQueue = new ConcurrentLinkedQueue<>(); - private final AtomicReference currentPipeline = new AtomicReference<>(); + private final Queue executionQueue = new ConcurrentLinkedQueue<>(); + private final AtomicReference currentPipeline = new AtomicReference<>(); @Getter - private final EventHandler onExecuted = new ConcurrentEventHandler<>(); + private final EventHandler onExecuted = new ConcurrentEventHandler<>(); - public void triggerExecute(ProcessingPipeline pipeline) { + public void triggerExecute(Flow pipeline) { if (executionQueue.contains(pipeline)) { return; } executionQueue.add(pipeline); } - public void cancel(ProcessingPipeline pipeline) { + public void cancel(Flow pipeline) { pipeline.interrupt(); executionQueue.remove(pipeline); } @PostConstruct private void initialize() { - worker = executor.submit(new PipelineWorker()); + worker = executor.submit(new FlowWorker()); } @PreDestroy @@ -53,11 +53,11 @@ public class PipelineExecutorService { private Future worker; - private class PipelineWorker implements Runnable { + private class FlowWorker implements Runnable { @SuppressWarnings("BusyWait") @SneakyThrows @Override - public void run() { + public void run() {//TODO заменить на worker while (!Thread.currentThread().isInterrupted()) { while (!executionQueue.isEmpty()) { var pipeline = executionQueue.poll(); diff --git a/src/main/java/ru/kirillius/pf/sdn/service/PipelineService.java b/src/main/java/ru/kirillius/pf/sdn/service/FlowService.java similarity index 74% rename from src/main/java/ru/kirillius/pf/sdn/service/PipelineService.java rename to src/main/java/ru/kirillius/pf/sdn/service/FlowService.java index 73a912d..35e9957 100644 --- a/src/main/java/ru/kirillius/pf/sdn/service/PipelineService.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.pipeline.Action; -import ru.kirillius.pf.sdn.api.pipeline.PipelineFunction; -import ru.kirillius.pf.sdn.api.pipeline.ProcessingPipeline; -import ru.kirillius.pf.sdn.api.pipeline.TriggerType; +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.PipelineConfig; +import ru.kirillius.pf.sdn.flow.DummyFunction; import ru.kirillius.pf.sdn.repository.PipelineConfigRepository; import java.util.Arrays; @@ -24,19 +25,19 @@ import java.util.regex.Pattern; @Slf4j @Service @RequiredArgsConstructor -public class PipelineService { +public class FlowService { private final PipelineConfigRepository configRepository; - private final PipelineExecutorService pipelineExecutorService; + private final FlowExecutorService flowExecutorService; private final SubscriptionService subscriptionService; private final ExecutorService executorService; - private final Map pipelines = new ConcurrentHashMap<>(); - private final Map> functions = new ConcurrentHashMap<>(); - private EventListener executedListener; + private final Map pipelines = new ConcurrentHashMap<>(); + private final Map> functions = new ConcurrentHashMap<>(); + private EventListener executedListener; private EventListener subscriptionUpdateListener; - public void registerFunction(Class functionClass, String id) { + public void registerFunction(Class functionClass, String id) { if (functions.containsKey(id)) { throw new IllegalStateException("Function with id '" + id + "' already exists"); } @@ -53,7 +54,7 @@ public class PipelineService { } } - public void trigger(ProcessingPipeline pipeline) { + public void trigger(Flow pipeline) { var matchedConditions = pipeline.getConfig().getStartConditions().stream().filter(c -> c.getTrigger() == TriggerType.Manual).toList(); if (matchedConditions.isEmpty()) { log.error("Unable to manual start pipeline {} because it has no manual trigger", pipeline.getConfig().getName()); @@ -63,7 +64,7 @@ public class PipelineService { if (c.isDontStartIfRunning() && pipeline.isRunning()) { return; } - pipelineExecutorService.triggerExecute(pipeline); + flowExecutorService.triggerExecute(pipeline); }); } @@ -72,23 +73,23 @@ public class PipelineService { if (configRepository.existsByGuid(pipelineConfig.getGuid())) { load(pipelineConfig); } else { - pipelineExecutorService.cancel(pipelines.get(pipelineConfig)); + flowExecutorService.cancel(pipelines.get(pipelineConfig)); pipelines.remove(pipelineConfig); } } private void load(PipelineConfig pipelineConfig) { if (pipelines.containsKey(pipelineConfig)) { - pipelineExecutorService.cancel(pipelines.get(pipelineConfig)); + flowExecutorService.cancel(pipelines.get(pipelineConfig)); } - pipelines.put(pipelineConfig, new ProcessingPipeline(pipelineConfig, pipelineConfig.getActions().stream().map(this::buildAction).toList())); + pipelines.put(pipelineConfig, new Flow(pipelineConfig, pipelineConfig.getActions().stream().map(this::buildAction).toList())); } - private Action buildAction(ActionConfig config) { - return new Action(functions.getOrDefault(config.getFunctionId(), DummyFunction.class), config.getProperties()); + private FlowAction buildAction(ActionConfig config) { + return new FlowAction(functions.getOrDefault(config.getFunctionId(), DummyFunction.class), config.getProperties()); } - private void checkForIntervalExecution(ProcessingPipeline pipeline) { + private void checkForIntervalExecution(Flow pipeline) { var matchedConditions = pipeline.getConfig().getStartConditions().stream().filter(c -> c.getTrigger() == TriggerType.Interval).toList(); @@ -116,7 +117,7 @@ public class PipelineService { if (c.isDontStartIfRunning() && pipeline.isRunning()) { return; } - pipelineExecutorService.triggerExecute(pipeline); + flowExecutorService.triggerExecute(pipeline); }); }); @@ -125,7 +126,7 @@ public class PipelineService { @PostConstruct private void initialize() { subscriptionUpdateListener = subscriptionService.getUpdateEvent().add(this::subscriptionUpdate); - executedListener = pipelineExecutorService.getOnExecuted().add(this::checkForIntervalExecution); + executedListener = flowExecutorService.getOnExecuted().add(this::checkForIntervalExecution); registerFunction(DummyFunction.class, "Error:fallback"); configRepository.findAll().forEach(this::load); pipelines.values() @@ -135,7 +136,7 @@ public class PipelineService { .getStartConditions() .stream() .anyMatch(c -> c.getTrigger() == TriggerType.OnStart)) - .forEach(pipelineExecutorService::triggerExecute); + .forEach(flowExecutorService::triggerExecute); } private void subscriptionUpdate(Subscription subscription) { @@ -154,12 +155,12 @@ public class PipelineService { var properties = c.getProperties(); var names = properties.getOrDefault("names", null); if (names == null || names.isEmpty()) { - pipelineExecutorService.triggerExecute(pipeline); + flowExecutorService.triggerExecute(pipeline); return; } if (Arrays.stream(names.split(Pattern.quote(","))).anyMatch(s -> subscription.getName().equals(s))) { - pipelineExecutorService.triggerExecute(pipeline); + flowExecutorService.triggerExecute(pipeline); } }); }); @@ -169,7 +170,7 @@ public class PipelineService { @PreDestroy private void destroy() { if (executedListener != null) { - pipelineExecutorService.getOnExecuted().remove(executedListener); + flowExecutorService.getOnExecuted().remove(executedListener); executedListener = null; } if (subscriptionUpdateListener != null) {