From 77044a2c644466541388a4d159840666611dceb6 Mon Sep 17 00:00:00 2001 From: kirillius Date: Sun, 26 Jul 2026 01:58:55 +0300 Subject: [PATCH] WIP --- .../pf/sdn/api/Networking/NetworkScope.java | 13 +++-- .../pf/sdn/api/flow/ExecutionInfoEntry.java | 6 +-- .../kirillius/pf/sdn/api/flow/FlowAction.java | 13 ++--- .../pf/sdn/api/flow/FlowFunction.java | 3 ++ .../pf/sdn/api/flow/TriggerType.java | 2 +- .../kirillius/pf/sdn/entity/FlowConfig.java | 4 -- ...onSetEntry.java => SubscriptionGroup.java} | 4 +- .../pf/sdn/flow/AutoResolveDomains.java | 6 ++- .../pf/sdn/flow/DuplicatesFilter.java | 2 +- .../sdn/flow/MergeNeighbourSubnetFilter.java | 3 +- .../flow/MergeWithCoverageSubnetFilter.java | 2 +- .../pf/sdn/flow/OverlappedDomainFilter.java | 10 +--- .../pf/sdn/flow/OverlappedSubnetFilter.java | 2 +- .../kirillius/pf/sdn/flow/ResourceFilter.java | 28 ++++++++++ .../pf/sdn/flow/StaticResourceInput.java | 9 +++- .../ru/kirillius/pf/sdn/flow/Summator.java | 29 +++++++++++ .../repository/PipelineConfigRepository.java | 4 +- .../repository/SubscriptionSetRepository.java | 4 +- .../pf/sdn/service/DomainUpdaterService.java | 2 +- .../pf/sdn/service/FlowExecutorService.java | 14 +++-- .../kirillius/pf/sdn/service/FlowService.java | 52 +++++++++++-------- .../pf/sdn/service/SubscriptionService.java | 4 +- 22 files changed, 140 insertions(+), 76 deletions(-) rename src/main/java/ru/kirillius/pf/sdn/entity/{SubscriptionSetEntry.java => SubscriptionGroup.java} (85%) create mode 100644 src/main/java/ru/kirillius/pf/sdn/flow/Summator.java diff --git a/src/main/java/ru/kirillius/pf/sdn/api/Networking/NetworkScope.java b/src/main/java/ru/kirillius/pf/sdn/api/Networking/NetworkScope.java index e235962..57bbb3f 100644 --- a/src/main/java/ru/kirillius/pf/sdn/api/Networking/NetworkScope.java +++ b/src/main/java/ru/kirillius/pf/sdn/api/Networking/NetworkScope.java @@ -25,13 +25,12 @@ public class NetworkScope { @Getter @Setter private List domains = new ArrayList<>(); - @Getter @Setter - private List autoResolvedDomains = new ArrayList<>(); + private List resolveDomains = new ArrayList<>(); public int getResourceCount() { - return subnets.size() + ASN.size() + domains.size() + autoResolvedDomains.size(); + return subnets.size() + ASN.size() + domains.size() + resolveDomains.size(); } @Override @@ -40,12 +39,12 @@ public class NetworkScope { return Objects.equals(ASN, that.ASN) && Objects.equals(subnets, that.subnets) && Objects.equals(domains, that.domains) && - Objects.equals(autoResolvedDomains, that.autoResolvedDomains); + Objects.equals(resolveDomains, that.resolveDomains); } @Override public int hashCode() { - return Objects.hash(ASN, subnets, domains, autoResolvedDomains); + return Objects.hash(ASN, subnets, domains, resolveDomains); } /** @@ -55,7 +54,7 @@ public class NetworkScope { ASN.clear(); subnets.clear(); domains.clear(); - autoResolvedDomains.clear(); + resolveDomains.clear(); } /** @@ -65,7 +64,7 @@ public class NetworkScope { ASN.addAll(networkScope.getASN()); subnets.addAll(networkScope.getSubnets()); domains.addAll(networkScope.getDomains()); - domains.addAll(networkScope.getAutoResolvedDomains()); + resolveDomains.addAll(networkScope.getResolveDomains()); } } diff --git a/src/main/java/ru/kirillius/pf/sdn/api/flow/ExecutionInfoEntry.java b/src/main/java/ru/kirillius/pf/sdn/api/flow/ExecutionInfoEntry.java index 366316e..a90d9db 100644 --- a/src/main/java/ru/kirillius/pf/sdn/api/flow/ExecutionInfoEntry.java +++ b/src/main/java/ru/kirillius/pf/sdn/api/flow/ExecutionInfoEntry.java @@ -16,16 +16,16 @@ public class ExecutionInfoEntry { var compareSubnets = diff(source.getSubnets(), result.getSubnets()); var compareDomains = diff(source.getDomains(), result.getDomains()); var compareASN = diff(source.getASN(), result.getASN()); - var compareAutoResolved = diff(source.getAutoResolvedDomains(), result.getAutoResolvedDomains()); + var compareAutoResolved = diff(source.getResolveDomains(), result.getResolveDomains()); added.getSubnets().addAll(compareSubnets.added); added.getDomains().addAll(compareDomains.added); added.getASN().addAll(compareASN.added); - added.getAutoResolvedDomains().addAll(compareAutoResolved.added); + added.getResolveDomains().addAll(compareAutoResolved.added); removed.getSubnets().addAll(compareSubnets.removed); removed.getDomains().addAll(compareDomains.removed); removed.getASN().addAll(compareASN.removed); - removed.getAutoResolvedDomains().addAll(compareAutoResolved.removed); + removed.getResolveDomains().addAll(compareAutoResolved.removed); } 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 70b19f0..9f39f45 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 @@ -1,6 +1,5 @@ package ru.kirillius.pf.sdn.api.flow; -import lombok.Builder; import lombok.Getter; import lombok.Setter; import lombok.SneakyThrows; @@ -11,7 +10,9 @@ import java.util.List; import java.util.Map; import java.util.function.Function; -@Builder +/** + * Представляет собой вызов функции с параметрами + */ public class FlowAction implements Function, CallContext> { private FlowFunction function; @@ -30,13 +31,7 @@ public class FlowAction implements Function, CallContext> { function = functionClass.getConstructor().newInstance(); } - public void setFunctionClass(Class functionClass) { - this.functionClass = functionClass; - instantiateFunction(); - } - - @Getter - private Class functionClass; + private final Class functionClass; @Getter @Setter private Map properties; 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 d2aa19d..b07d3f6 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 @@ -5,6 +5,9 @@ import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor; import java.util.List; import java.util.Map; +/** + * Интерфейс, описывающий исполняемую функцию для Flow + */ public interface FlowFunction { void apply(CallContext context); 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 432637e..0e56abd 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,6 +4,6 @@ public enum TriggerType { Manual, Interval, OnSubscriptionUpdate, - OnSubscriptionSetUpdate, + OnSubscriptionGroupUpdate, OnStart } 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 1554207..fc48aaa 100644 --- a/src/main/java/ru/kirillius/pf/sdn/entity/FlowConfig.java +++ b/src/main/java/ru/kirillius/pf/sdn/entity/FlowConfig.java @@ -7,7 +7,6 @@ import ru.kirillius.pf.sdn.api.flow.StartCondition; import java.util.ArrayList; import java.util.List; import java.util.Objects; -import java.util.UUID; @Entity @Table(name = "pipeline_config") @@ -21,9 +20,6 @@ public class FlowConfig { @GeneratedValue(strategy = GenerationType.IDENTITY) private Long id; - @GeneratedValue(strategy = GenerationType.UUID) - private UUID guid; - @Column(nullable = false, unique = true) private String name; diff --git a/src/main/java/ru/kirillius/pf/sdn/entity/SubscriptionSetEntry.java b/src/main/java/ru/kirillius/pf/sdn/entity/SubscriptionGroup.java similarity index 85% rename from src/main/java/ru/kirillius/pf/sdn/entity/SubscriptionSetEntry.java rename to src/main/java/ru/kirillius/pf/sdn/entity/SubscriptionGroup.java index 317940b..270b55e 100644 --- a/src/main/java/ru/kirillius/pf/sdn/entity/SubscriptionSetEntry.java +++ b/src/main/java/ru/kirillius/pf/sdn/entity/SubscriptionGroup.java @@ -6,13 +6,13 @@ import lombok.*; import java.util.Set; @Entity -@Table(name = "subscription_set") +@Table(name = "subscription_group") @Getter @Setter @NoArgsConstructor @Builder @AllArgsConstructor -public class SubscriptionSetEntry { +public class SubscriptionGroup { @Id @GeneratedValue(strategy = GenerationType.IDENTITY) private Long id; 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 a7e8139..375e2b8 100644 --- a/src/main/java/ru/kirillius/pf/sdn/flow/AutoResolveDomains.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/AutoResolveDomains.java @@ -6,6 +6,7 @@ import ru.kirillius.pf.sdn.api.Networking.NetworkScope; import ru.kirillius.pf.sdn.api.flow.CallContext; import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor; import ru.kirillius.pf.sdn.api.properties.PropertyType; +import ru.kirillius.pf.sdn.service.DomainResolverService; import ru.kirillius.pf.sdn.service.DomainUpdaterService; import java.util.Collections; @@ -17,6 +18,7 @@ import java.util.Map; public class AutoResolveDomains extends CommonFunction { private final DomainUpdaterService domainCacheService; + private final DomainResolverService domainResolverService; private final static String CLEAR = "clear"; private final static Map properties; @@ -41,10 +43,10 @@ public class AutoResolveDomains extends CommonFunction { output.add(source); if (properties.containsKey(CLEAR) && properties.get(CLEAR).equals("true")) { - output.getAutoResolvedDomains().clear(); + output.getResolveDomains().clear(); } - source.getAutoResolvedDomains() + source.getResolveDomains() .forEach(domain -> domainCacheService.getActualAddresses(domain) .forEach(address -> output .getSubnets() 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 e4e113f..0864d2c 100644 --- a/src/main/java/ru/kirillius/pf/sdn/flow/DuplicatesFilter.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/DuplicatesFilter.java @@ -22,7 +22,7 @@ public class DuplicatesFilter implements FlowFunction { 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()); + output.setResolveDomains(source.getResolveDomains().stream().distinct().toList()); context.setOutputs(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 e1858e6..abef195 100644 --- a/src/main/java/ru/kirillius/pf/sdn/flow/MergeNeighbourSubnetFilter.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/MergeNeighbourSubnetFilter.java @@ -29,7 +29,8 @@ public class MergeNeighbourSubnetFilter extends CommonFunction { output.setSubnets(IPv4Util.mergeNeighbours(source.getSubnets())); output.setDomains(source.getDomains()); output.setASN(source.getASN()); - output.setAutoResolvedDomains(source.getAutoResolvedDomains()); + output.setResolveDomains(source.getResolveDomains()); + context.setOutputs(List.of(output)); } } 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 0a0505d..26d82bd 100644 --- a/src/main/java/ru/kirillius/pf/sdn/flow/MergeWithCoverageSubnetFilter.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/MergeWithCoverageSubnetFilter.java @@ -52,7 +52,7 @@ public class MergeWithCoverageSubnetFilter extends CommonFunction { output.setSubnets(IPv4Util.mergeToLargerIfCovered(source.getSubnets(), usage)); output.setDomains(source.getDomains()); output.setASN(source.getASN()); - output.setAutoResolvedDomains(source.getAutoResolvedDomains()); + output.setResolveDomains(source.getResolveDomains()); context.setOutputs(List.of(output)); } 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 6dd8ca7..9b83ac7 100644 --- a/src/main/java/ru/kirillius/pf/sdn/flow/OverlappedDomainFilter.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/OverlappedDomainFilter.java @@ -5,11 +5,8 @@ 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.properties.PropertyDescriptor; -import java.util.Collections; import java.util.List; -import java.util.Map; @Slf4j @RequiredArgsConstructor @@ -32,13 +29,8 @@ public class OverlappedDomainFilter extends CommonFunction { output.setSubnets(source.getSubnets()); output.setDomains(DomainUtil.removeOverlapped(source.getDomains())); output.setASN(source.getASN()); - output.setAutoResolvedDomains(source.getAutoResolvedDomains()); + output.setResolveDomains(source.getResolveDomains()); context.setOutputs(List.of(output)); } - @Override - public Map getProperties() { - return Collections.emptyMap(); - } - } 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 c206bec..035c271 100644 --- a/src/main/java/ru/kirillius/pf/sdn/flow/OverlappedSubnetFilter.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/OverlappedSubnetFilter.java @@ -29,7 +29,7 @@ public class OverlappedSubnetFilter extends CommonFunction { output.setSubnets(IPv4Util.removeOverlapped(source.getSubnets())); output.setDomains(source.getDomains()); output.setASN(source.getASN()); - output.setAutoResolvedDomains(source.getAutoResolvedDomains()); + output.setResolveDomains(source.getResolveDomains()); 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 ab2f86d..6e74a94 100644 --- a/src/main/java/ru/kirillius/pf/sdn/flow/ResourceFilter.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/ResourceFilter.java @@ -33,13 +33,25 @@ public class ResourceFilter extends CommonFunction { if (!filter.getSubnets().isEmpty()) { filterSubnets(source, filter, output, filtered); + } else { + output.setSubnets(source.getSubnets()); } if (!filter.getASN().isEmpty()) { filterASN(source, filter, output, filtered); + } else { + output.setASN(source.getASN()); } if (!filter.getDomains().isEmpty()) { filterDomains(source, filter, output, filtered); + } else { + output.setDomains(source.getDomains()); } + if (!filter.getResolveDomains().isEmpty()) { + filterResolveDomains(source, filter, output, filtered); + } else { + output.setResolveDomains(source.getResolveDomains()); + } + context.setOutputs(List.of(output, filtered)); } @@ -92,4 +104,20 @@ public class ResourceFilter extends CommonFunction { output.setDomains(outputDomains); filtered.setDomains(filteredDomains); } + + private void filterResolveDomains(NetworkScope source, NetworkScope filter, NetworkScope output, NetworkScope filtered) { + var sourceDomains = source.getResolveDomains(); + var filterDomains = filter.getResolveDomains(); + var outputDomains = new ArrayList(); + var filteredDomains = new ArrayList(); + sourceDomains.forEach(domain -> { + if (!filterDomains.contains(domain)) { + outputDomains.add(domain); + } else { + filteredDomains.add(domain); + } + }); + output.setResolveDomains(outputDomains); + filtered.setResolveDomains(filteredDomains); + } } diff --git a/src/main/java/ru/kirillius/pf/sdn/flow/StaticResourceInput.java b/src/main/java/ru/kirillius/pf/sdn/flow/StaticResourceInput.java index 6143670..7e446bf 100644 --- a/src/main/java/ru/kirillius/pf/sdn/flow/StaticResourceInput.java +++ b/src/main/java/ru/kirillius/pf/sdn/flow/StaticResourceInput.java @@ -18,6 +18,8 @@ public class StaticResourceInput extends CommonFunction { private final static String SUBNETS = "subnets"; private final static String DOMAINS = "domains"; + private final static String RESOLVE = "resolve"; + private final static String ASN = "asn"; private final static Map properties; @@ -26,6 +28,8 @@ public class StaticResourceInput extends CommonFunction { 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(RESOLVE, 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?); } @@ -47,6 +51,9 @@ public class StaticResourceInput extends CommonFunction { if (properties.containsKey(DOMAINS)) { output.setDomains(Arrays.stream(properties.get(DOMAINS).split(Pattern.quote(","))).filter(s -> !s.isBlank()).toList()); } + if (properties.containsKey(RESOLVE)) { + output.setResolveDomains(Arrays.stream(properties.get(RESOLVE).split(Pattern.quote(","))).filter(s -> !s.isBlank()).toList()); + } if (properties.containsKey(ASN)) { @@ -61,6 +68,4 @@ public class StaticResourceInput extends CommonFunction { return Collections.unmodifiableMap(properties); } - - } diff --git a/src/main/java/ru/kirillius/pf/sdn/flow/Summator.java b/src/main/java/ru/kirillius/pf/sdn/flow/Summator.java new file mode 100644 index 0000000..20d5ce2 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/flow/Summator.java @@ -0,0 +1,29 @@ +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 java.util.List; + +@RequiredArgsConstructor +public class Summator extends CommonFunction { + + @Override + public List getInputNames() { + return List.of("A", "B"); + } + + @Override + public List getOutputNames() { + return List.of("output"); + } + + @Override + public void apply(CallContext context) { + var output = new NetworkScope(); + context.getSources().forEach(output::add); + context.setOutputs(List.of(output)); + } + +} diff --git a/src/main/java/ru/kirillius/pf/sdn/repository/PipelineConfigRepository.java b/src/main/java/ru/kirillius/pf/sdn/repository/PipelineConfigRepository.java index 5345536..2528a2d 100644 --- a/src/main/java/ru/kirillius/pf/sdn/repository/PipelineConfigRepository.java +++ b/src/main/java/ru/kirillius/pf/sdn/repository/PipelineConfigRepository.java @@ -3,9 +3,7 @@ package ru.kirillius.pf.sdn.repository; import org.springframework.data.jpa.repository.JpaRepository; import ru.kirillius.pf.sdn.entity.FlowConfig; -import java.util.UUID; - public interface PipelineConfigRepository extends JpaRepository { - boolean existsByGuid(UUID guid); + } diff --git a/src/main/java/ru/kirillius/pf/sdn/repository/SubscriptionSetRepository.java b/src/main/java/ru/kirillius/pf/sdn/repository/SubscriptionSetRepository.java index 5a5f336..c9ac635 100644 --- a/src/main/java/ru/kirillius/pf/sdn/repository/SubscriptionSetRepository.java +++ b/src/main/java/ru/kirillius/pf/sdn/repository/SubscriptionSetRepository.java @@ -1,8 +1,8 @@ package ru.kirillius.pf.sdn.repository; import org.springframework.data.jpa.repository.JpaRepository; -import ru.kirillius.pf.sdn.entity.SubscriptionSetEntry; +import ru.kirillius.pf.sdn.entity.SubscriptionGroup; -public interface SubscriptionSetRepository extends JpaRepository { +public interface SubscriptionSetRepository extends JpaRepository { } diff --git a/src/main/java/ru/kirillius/pf/sdn/service/DomainUpdaterService.java b/src/main/java/ru/kirillius/pf/sdn/service/DomainUpdaterService.java index d78fa5a..04697c0 100644 --- a/src/main/java/ru/kirillius/pf/sdn/service/DomainUpdaterService.java +++ b/src/main/java/ru/kirillius/pf/sdn/service/DomainUpdaterService.java @@ -53,7 +53,7 @@ public class DomainUpdaterService { protected void doWork() { subscriptionService.getSubscriptions().forEach(subscription -> { var scope = subscription.getScope(); - if (!scope.isResolveDomains()) { + if (scope.getResolveDomains().isEmpty()) { return; } scope.getDomains().forEach(domain -> { diff --git a/src/main/java/ru/kirillius/pf/sdn/service/FlowExecutorService.java b/src/main/java/ru/kirillius/pf/sdn/service/FlowExecutorService.java index 978629f..4d235c8 100644 --- a/src/main/java/ru/kirillius/pf/sdn/service/FlowExecutorService.java +++ b/src/main/java/ru/kirillius/pf/sdn/service/FlowExecutorService.java @@ -17,13 +17,21 @@ import java.util.concurrent.ExecutorService; import java.util.concurrent.Future; import java.util.concurrent.atomic.AtomicReference; +/** + * Запускает и выполняет Flow в порядке очереди + */ @Slf4j @Service @RequiredArgsConstructor public class FlowExecutorService { private final ExecutorService executor; private final Queue executionQueue = new ConcurrentLinkedQueue<>(); - private final AtomicReference currentPipeline = new AtomicReference<>(); + + public Flow getCurrentFlow() { + return currentFlow.get(); + } + + private final AtomicReference currentFlow = new AtomicReference<>(); @Getter private final EventHandler onExecuted = new ConcurrentEventHandler<>(); @@ -61,14 +69,14 @@ public class FlowExecutorService { while (!Thread.currentThread().isInterrupted()) { while (!executionQueue.isEmpty()) { var pipeline = executionQueue.poll(); - currentPipeline.set(pipeline); + currentFlow.set(pipeline); try { pipeline.execute(); onExecuted.invoke(pipeline); } catch (Exception e) { log.error("Exception while executing pipeline", e); } - currentPipeline.set(null); + currentFlow.set(null); } Thread.sleep(100L); Thread.yield(); 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 f0b2add..3aefaa7 100644 --- a/src/main/java/ru/kirillius/pf/sdn/service/FlowService.java +++ b/src/main/java/ru/kirillius/pf/sdn/service/FlowService.java @@ -13,12 +13,13 @@ import ru.kirillius.pf.sdn.api.flow.FlowFunction; 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.entity.SubscriptionGroup; import ru.kirillius.pf.sdn.flow.DummyFunction; import ru.kirillius.pf.sdn.repository.PipelineConfigRepository; import java.util.Arrays; import java.util.Map; +import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutorService; import java.util.regex.Pattern; @@ -33,11 +34,11 @@ public class FlowService { private final SubscriptionService subscriptionService; private final ExecutorService executorService; - private final Map pipelines = new ConcurrentHashMap<>(); + private final Map flows = new ConcurrentHashMap<>(); private final Map> functions = new ConcurrentHashMap<>(); private EventListener executedListener; private EventListener subscriptionUpdateListener; - private EventListener subscriptionSetUpdateListener; + private EventListener subscriptionSetUpdateListener; public void registerFunction(Class functionClass, String id) { if (functions.containsKey(id)) { @@ -46,13 +47,17 @@ public class FlowService { functions.put(id, functionClass); } + public Set getFunctionNames() { + return functions.keySet(); + } + public void trigger(long id) { - pipelines.keySet().stream().filter(config -> config.getId().equals(id)).findFirst().ifPresent(this::trigger); + flows.keySet().stream().filter(config -> config.getId().equals(id)).findFirst().ifPresent(this::trigger); } public void trigger(FlowConfig flowConfig) { - if (pipelines.containsKey(flowConfig)) { - trigger(pipelines.get(flowConfig)); + if (flows.containsKey(flowConfig)) { + trigger(flows.get(flowConfig)); } } @@ -68,23 +73,26 @@ public class FlowService { } flowExecutorService.triggerExecute(pipeline); }); - } - public void update(FlowConfig flowConfig) { - if (configRepository.existsByGuid(flowConfig.getGuid())) { - load(flowConfig); + public void loadFlow(Flow flow) { + loadFlow(flow.getConfig()); + } + + public void loadFlow(FlowConfig flowConfig) { + if (configRepository.existsById(flowConfig.getId())) { + reloadFlowInternal(flowConfig); } else { - flowExecutorService.cancel(pipelines.get(flowConfig)); - pipelines.remove(flowConfig); + flowExecutorService.cancel(flows.get(flowConfig)); + flows.remove(flowConfig); } } - private void load(FlowConfig flowConfig) { - if (pipelines.containsKey(flowConfig)) { - flowExecutorService.cancel(pipelines.get(flowConfig)); + private void reloadFlowInternal(FlowConfig flowConfig) { + if (flows.containsKey(flowConfig)) { + flowExecutorService.cancel(flows.get(flowConfig)); } - pipelines.put(flowConfig, new Flow(flowConfig, this::buildAction)); + flows.put(flowConfig, new Flow(flowConfig, this::buildAction)); } private FlowAction buildAction(ActionConfig config) { @@ -131,8 +139,8 @@ public class FlowService { subscriptionSetUpdateListener = subscriptionService.getSetUpdateEvent().add(this::subscriptionSetUpdate); executedListener = flowExecutorService.getOnExecuted().add(this::checkForIntervalExecution); registerFunction(DummyFunction.class, "Error:fallback"); - configRepository.findAll().forEach(this::load); - pipelines.values() + configRepository.findAll().forEach(this::reloadFlowInternal); + flows.values() .stream() .filter(p -> p .getConfig() @@ -142,9 +150,9 @@ 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(); + private void subscriptionSetUpdate(SubscriptionGroup setEntry) { + flows.values().forEach(pipeline -> { + var matchedConditions = pipeline.getConfig().getStartConditions().stream().filter(c -> c.getTrigger() == TriggerType.OnSubscriptionGroupUpdate).toList(); if (matchedConditions.isEmpty()) { return; } @@ -170,7 +178,7 @@ public class FlowService { private void subscriptionUpdate(Subscription subscription) { - pipelines.values().forEach(pipeline -> { + flows.values().forEach(pipeline -> { var matchedConditions = pipeline.getConfig().getStartConditions().stream().filter(c -> c.getTrigger() == TriggerType.OnSubscriptionUpdate).toList(); if (matchedConditions.isEmpty()) { return; 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 ccb6eb8..4e462b9 100644 --- a/src/main/java/ru/kirillius/pf/sdn/service/SubscriptionService.java +++ b/src/main/java/ru/kirillius/pf/sdn/service/SubscriptionService.java @@ -13,7 +13,7 @@ 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.entity.SubscriptionGroup; import ru.kirillius.pf.sdn.repository.SubscriptionProviderConfigRepository; import ru.kirillius.pf.sdn.repository.SubscriptionSetRepository; @@ -48,7 +48,7 @@ public class SubscriptionService { private final EventHandler updateEvent = new ConcurrentEventHandler<>(); @Getter - private final EventHandler setUpdateEvent = new ConcurrentEventHandler<>(); + private final EventHandler setUpdateEvent = new ConcurrentEventHandler<>(); private Future worker; public List getSubscriptions() {