From ad5298bc393f4bba0ebabab3d58a5bdbc8461e22 Mon Sep 17 00:00:00 2001 From: kirillius Date: Sat, 4 Jul 2026 02:35:39 +0300 Subject: [PATCH] WIP --- .../kirillius/pf/sdn/api/pipeline/Action.java | 3 +- .../sdn/api/pipeline/ProcessingPipeline.java | 27 +-- .../pf/sdn/api/pipeline/StartCondition.java | 23 ++- .../pf/sdn/config/DataInitializer.java | 1 + .../kirillius/pf/sdn/entity/ActionConfig.java | 29 +++ .../pf/sdn/entity/PipelineConfig.java | 46 +++++ .../repository/PipelineConfigRepository.java | 11 ++ .../pf/sdn/service/DummyFunction.java | 20 ++ .../sdn/service/PipelineExecutorService.java | 15 ++ .../pf/sdn/service/PipelineService.java | 181 ++++++++++++++++++ .../pf/sdn/service/SubscriptionService.java | 10 +- 11 files changed, 344 insertions(+), 22 deletions(-) create mode 100644 src/main/java/ru/kirillius/pf/sdn/entity/ActionConfig.java create mode 100644 src/main/java/ru/kirillius/pf/sdn/entity/PipelineConfig.java create mode 100644 src/main/java/ru/kirillius/pf/sdn/repository/PipelineConfigRepository.java create mode 100644 src/main/java/ru/kirillius/pf/sdn/service/DummyFunction.java create mode 100644 src/main/java/ru/kirillius/pf/sdn/service/PipelineService.java diff --git a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/Action.java b/src/main/java/ru/kirillius/pf/sdn/api/pipeline/Action.java index cf6cdba..4686076 100644 --- a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/Action.java +++ b/src/main/java/ru/kirillius/pf/sdn/api/pipeline/Action.java @@ -14,8 +14,7 @@ public class Action implements Function { private PipelineFunction function; - public Action(PipelineFunction function, Class functionClass, Map properties) { - this.function = function; + public Action(Class functionClass, Map properties) { this.functionClass = functionClass; this.properties = properties; instantiateFunction(); diff --git a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/ProcessingPipeline.java b/src/main/java/ru/kirillius/pf/sdn/api/pipeline/ProcessingPipeline.java index d1b9473..27b4015 100644 --- a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/ProcessingPipeline.java +++ b/src/main/java/ru/kirillius/pf/sdn/api/pipeline/ProcessingPipeline.java @@ -1,10 +1,8 @@ package ru.kirillius.pf.sdn.api.pipeline; -import lombok.AllArgsConstructor; -import lombok.Getter; -import lombok.NoArgsConstructor; -import lombok.Setter; +import lombok.*; import ru.kirillius.pf.sdn.api.Networking.NetworkScope; +import ru.kirillius.pf.sdn.entity.PipelineConfig; import java.util.ArrayList; import java.util.List; @@ -12,17 +10,15 @@ import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; -@AllArgsConstructor -@NoArgsConstructor +@RequiredArgsConstructor() public class ProcessingPipeline { @Getter - @Setter - private List startConditions = new ArrayList<>(); + private final PipelineConfig config; + private final List actions; private final AtomicBoolean running = new AtomicBoolean(false); - @Getter - @Setter - private List actions = new ArrayList<>(); + private final AtomicBoolean interrupted = new AtomicBoolean(false); + private AtomicInteger currentStep = new AtomicInteger(0); public boolean isRunning() { @@ -37,18 +33,27 @@ public class ProcessingPipeline { return currentStep.get(); } + public void interrupt() { + interrupted.set(false); + } + private List executionInfo = new ArrayList<>(); public void execute() { + interrupted.set(false); running.set(true); currentStep.set(0); try { var scope = new AtomicReference<>(new NetworkScope()); actions.forEach(action -> { + if (interrupted.get()) { + throw new RuntimeException("Interrupted"); + } currentStep.incrementAndGet(); scope.set(action.apply(scope.get())); }); } finally { + interrupted.set(false); running.set(false); } } diff --git a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/StartCondition.java b/src/main/java/ru/kirillius/pf/sdn/api/pipeline/StartCondition.java index 17df692..6c539de 100644 --- a/src/main/java/ru/kirillius/pf/sdn/api/pipeline/StartCondition.java +++ b/src/main/java/ru/kirillius/pf/sdn/api/pipeline/StartCondition.java @@ -1,13 +1,32 @@ package ru.kirillius.pf.sdn.api.pipeline; +import jakarta.persistence.*; import lombok.Getter; import lombok.Setter; -import java.util.Properties; +import java.util.Map; @Getter @Setter +@Entity public class StartCondition { + @Id + @GeneratedValue(strategy = GenerationType.IDENTITY) + private Long id; + + @Column + private boolean dontStartIfRunning; + + @Column(length = 100) private TriggerType trigger; - private Properties properties; + + @ElementCollection(fetch = FetchType.EAGER) // Подгружать сразу + @CollectionTable( + name = "start_condition_properties", + joinColumns = @JoinColumn(name = "condition_id") + ) + @MapKeyColumn(name = "prop_key") + @Column(name = "prop_value") + private Map properties; + } diff --git a/src/main/java/ru/kirillius/pf/sdn/config/DataInitializer.java b/src/main/java/ru/kirillius/pf/sdn/config/DataInitializer.java index f09a040..af69748 100644 --- a/src/main/java/ru/kirillius/pf/sdn/config/DataInitializer.java +++ b/src/main/java/ru/kirillius/pf/sdn/config/DataInitializer.java @@ -15,5 +15,6 @@ public class DataInitializer { @EventListener(ApplicationReadyEvent.class) public void init() { authService.createDefaultUserIfAbsent(); + } } diff --git a/src/main/java/ru/kirillius/pf/sdn/entity/ActionConfig.java b/src/main/java/ru/kirillius/pf/sdn/entity/ActionConfig.java new file mode 100644 index 0000000..7817879 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/entity/ActionConfig.java @@ -0,0 +1,29 @@ +package ru.kirillius.pf.sdn.entity; + +import jakarta.persistence.*; +import lombok.Getter; +import lombok.Setter; + +import java.util.HashMap; +import java.util.Map; + +@Getter +@Setter +@Entity +public class ActionConfig { + @Id + @GeneratedValue(strategy = GenerationType.IDENTITY) + private Long id; + + @ElementCollection(fetch = FetchType.EAGER) + @CollectionTable( + name = "action_properties", + joinColumns = @JoinColumn(name = "action_id") + ) + @MapKeyColumn(name = "prop_key") + @Column(name = "prop_value") + private Map properties = new HashMap<>(); + + @Column(nullable = false) + private String functionId; +} \ No newline at end of file diff --git a/src/main/java/ru/kirillius/pf/sdn/entity/PipelineConfig.java b/src/main/java/ru/kirillius/pf/sdn/entity/PipelineConfig.java new file mode 100644 index 0000000..80b9a3b --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/entity/PipelineConfig.java @@ -0,0 +1,46 @@ +package ru.kirillius.pf.sdn.entity; + +import jakarta.persistence.*; +import lombok.*; +import ru.kirillius.pf.sdn.api.pipeline.StartCondition; + +import java.util.ArrayList; +import java.util.List; +import java.util.Objects; +import java.util.UUID; + +@Entity +@Table(name = "pipeline_config") +@Getter +@Setter +@NoArgsConstructor +@Builder +@AllArgsConstructor +public class PipelineConfig { + @Id + @GeneratedValue(strategy = GenerationType.IDENTITY) + private Long id; + + @GeneratedValue(strategy = GenerationType.UUID) + private UUID guid; + + @Column(nullable = false, unique = true) + private String name; + + @Override + public boolean equals(Object o) { + if (!(o instanceof PipelineConfig that)) return false; + return Objects.equals(id, that.id) && Objects.equals(name, that.name); + } + + @Override + public int hashCode() { + return Objects.hash(id, name); + } + + @ManyToMany(fetch = FetchType.EAGER) + private List startConditions = new ArrayList<>(); + + @ManyToMany(fetch = FetchType.EAGER)//TODO проверить связность + private List actions = new ArrayList<>(); +} diff --git a/src/main/java/ru/kirillius/pf/sdn/repository/PipelineConfigRepository.java b/src/main/java/ru/kirillius/pf/sdn/repository/PipelineConfigRepository.java new file mode 100644 index 0000000..d58bdc4 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/repository/PipelineConfigRepository.java @@ -0,0 +1,11 @@ +package ru.kirillius.pf.sdn.repository; + +import org.springframework.data.jpa.repository.JpaRepository; +import ru.kirillius.pf.sdn.entity.PipelineConfig; + +import java.util.UUID; + +public interface PipelineConfigRepository extends JpaRepository { + + boolean existsByGuid(UUID guid); +} diff --git a/src/main/java/ru/kirillius/pf/sdn/service/DummyFunction.java b/src/main/java/ru/kirillius/pf/sdn/service/DummyFunction.java new file mode 100644 index 0000000..523df82 --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/service/DummyFunction.java @@ -0,0 +1,20 @@ +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/PipelineExecutorService.java index 1a0ac20..a27f90d 100644 --- a/src/main/java/ru/kirillius/pf/sdn/service/PipelineExecutorService.java +++ b/src/main/java/ru/kirillius/pf/sdn/service/PipelineExecutorService.java @@ -2,16 +2,20 @@ package ru.kirillius.pf.sdn.service; import jakarta.annotation.PostConstruct; import jakarta.annotation.PreDestroy; +import lombok.Getter; import lombok.RequiredArgsConstructor; import lombok.SneakyThrows; 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 java.util.Queue; import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.ExecutorService; import java.util.concurrent.Future; +import java.util.concurrent.atomic.AtomicReference; @Slf4j @Service @@ -19,6 +23,9 @@ import java.util.concurrent.Future; public class PipelineExecutorService { private final ExecutorService executor; private final Queue executionQueue = new ConcurrentLinkedQueue<>(); + private final AtomicReference currentPipeline = new AtomicReference<>(); + @Getter + private final EventHandler onExecuted = new ConcurrentEventHandler<>(); public void triggerExecute(ProcessingPipeline pipeline) { if (executionQueue.contains(pipeline)) { @@ -27,6 +34,11 @@ public class PipelineExecutorService { executionQueue.add(pipeline); } + public void cancel(ProcessingPipeline pipeline) { + pipeline.interrupt(); + executionQueue.remove(pipeline); + } + @PostConstruct private void initialize() { worker = executor.submit(new PipelineWorker()); @@ -49,11 +61,14 @@ public class PipelineExecutorService { while (!Thread.currentThread().isInterrupted()) { while (!executionQueue.isEmpty()) { var pipeline = executionQueue.poll(); + currentPipeline.set(pipeline); try { pipeline.execute(); + onExecuted.invoke(pipeline); } catch (Exception e) { log.error("Exception while executing pipeline", e); } + currentPipeline.set(null); } Thread.sleep(100L); Thread.yield(); diff --git a/src/main/java/ru/kirillius/pf/sdn/service/PipelineService.java b/src/main/java/ru/kirillius/pf/sdn/service/PipelineService.java new file mode 100644 index 0000000..73a912d --- /dev/null +++ b/src/main/java/ru/kirillius/pf/sdn/service/PipelineService.java @@ -0,0 +1,181 @@ +package ru.kirillius.pf.sdn.service; + +import jakarta.annotation.PostConstruct; +import jakarta.annotation.PreDestroy; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; +import ru.kirillius.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.entity.ActionConfig; +import ru.kirillius.pf.sdn.entity.PipelineConfig; +import ru.kirillius.pf.sdn.repository.PipelineConfigRepository; + +import java.util.Arrays; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ExecutorService; +import java.util.regex.Pattern; + +@Slf4j +@Service +@RequiredArgsConstructor +public class PipelineService { + + private final PipelineConfigRepository configRepository; + private final PipelineExecutorService pipelineExecutorService; + private final SubscriptionService subscriptionService; + private final ExecutorService executorService; + + 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) { + if (functions.containsKey(id)) { + throw new IllegalStateException("Function with id '" + id + "' already exists"); + } + functions.put(id, functionClass); + } + + public void trigger(long id) { + pipelines.keySet().stream().filter(config -> config.getId().equals(id)).findFirst().ifPresent(this::trigger); + } + + public void trigger(PipelineConfig pipelineConfig) { + if (pipelines.containsKey(pipelineConfig)) { + trigger(pipelines.get(pipelineConfig)); + } + } + + public void trigger(ProcessingPipeline 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()); + return; + } + matchedConditions.forEach(c -> { + if (c.isDontStartIfRunning() && pipeline.isRunning()) { + return; + } + pipelineExecutorService.triggerExecute(pipeline); + }); + + } + + public void update(PipelineConfig pipelineConfig) { + if (configRepository.existsByGuid(pipelineConfig.getGuid())) { + load(pipelineConfig); + } else { + pipelineExecutorService.cancel(pipelines.get(pipelineConfig)); + pipelines.remove(pipelineConfig); + } + } + + private void load(PipelineConfig pipelineConfig) { + if (pipelines.containsKey(pipelineConfig)) { + pipelineExecutorService.cancel(pipelines.get(pipelineConfig)); + } + pipelines.put(pipelineConfig, new ProcessingPipeline(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 void checkForIntervalExecution(ProcessingPipeline pipeline) { + + var matchedConditions = pipeline.getConfig().getStartConditions().stream().filter(c -> c.getTrigger() == TriggerType.Interval).toList(); + + if (matchedConditions.isEmpty()) { + return; + } + + matchedConditions.forEach(c -> { + var properties = c.getProperties(); + var timer = properties.getOrDefault("interval", null); + if (timer == null) { + return; + } + var delay = Long.parseLong(timer); + if (delay <= 0) { + return; + } + + executorService.submit(() -> { + try { + Thread.sleep(delay * 60000L); + } catch (InterruptedException e) { + return; + } + if (c.isDontStartIfRunning() && pipeline.isRunning()) { + return; + } + pipelineExecutorService.triggerExecute(pipeline); + }); + }); + + } + + @PostConstruct + private void initialize() { + subscriptionUpdateListener = subscriptionService.getUpdateEvent().add(this::subscriptionUpdate); + executedListener = pipelineExecutorService.getOnExecuted().add(this::checkForIntervalExecution); + registerFunction(DummyFunction.class, "Error:fallback"); + configRepository.findAll().forEach(this::load); + pipelines.values() + .stream() + .filter(p -> p + .getConfig() + .getStartConditions() + .stream() + .anyMatch(c -> c.getTrigger() == TriggerType.OnStart)) + .forEach(pipelineExecutorService::triggerExecute); + } + + private void subscriptionUpdate(Subscription subscription) { + + pipelines.values().forEach(pipeline -> { + var matchedConditions = pipeline.getConfig().getStartConditions().stream().filter(c -> c.getTrigger() == TriggerType.OnSubscriptionUpdate).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()) { + pipelineExecutorService.triggerExecute(pipeline); + return; + } + + if (Arrays.stream(names.split(Pattern.quote(","))).anyMatch(s -> subscription.getName().equals(s))) { + pipelineExecutorService.triggerExecute(pipeline); + } + }); + }); + + } + + @PreDestroy + private void destroy() { + if (executedListener != null) { + pipelineExecutorService.getOnExecuted().remove(executedListener); + executedListener = null; + } + if (subscriptionUpdateListener != null) { + subscriptionService.getUpdateEvent().remove(subscriptionUpdateListener); + subscriptionUpdateListener = 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 67871d7..0020ef7 100644 --- a/src/main/java/ru/kirillius/pf/sdn/service/SubscriptionService.java +++ b/src/main/java/ru/kirillius/pf/sdn/service/SubscriptionService.java @@ -40,10 +40,11 @@ public class SubscriptionService { private final Map providers = new ConcurrentHashMap<>(); private final AtomicInteger updateCounter = new AtomicInteger(0); private final List subscriptions = new CopyOnWriteArrayList<>(); + private final ApplicationContext context; + @Getter + private final EventHandler updateEvent = new ConcurrentEventHandler<>(); private Future worker; - - public List getSubscriptions() { return Collections.unmodifiableList(subscriptions); } @@ -60,8 +61,6 @@ public class SubscriptionService { } } - private final ApplicationContext context; - public void reloadProviders() { var beanFactory = context.getAutowireCapableBeanFactory(); synchronized (providers) { @@ -81,9 +80,6 @@ public class SubscriptionService { } } - @Getter - private final EventHandler updateEvent = new ConcurrentEventHandler<>(); - private void reloadSubscriptions() { var retrieved = new ArrayList(); var updated = new ArrayList();