This commit is contained in:
kirillius 2026-07-04 02:35:39 +03:00
parent e488f9d1d6
commit ad5298bc39
11 changed files with 344 additions and 22 deletions

View File

@ -14,8 +14,7 @@ public class Action implements Function<NetworkScope, NetworkScope> {
private PipelineFunction function; private PipelineFunction function;
public Action(PipelineFunction function, Class<? extends PipelineFunction> functionClass, Map<String, String> properties) { public Action(Class<? extends PipelineFunction> functionClass, Map<String, String> properties) {
this.function = function;
this.functionClass = functionClass; this.functionClass = functionClass;
this.properties = properties; this.properties = properties;
instantiateFunction(); instantiateFunction();

View File

@ -1,10 +1,8 @@
package ru.kirillius.pf.sdn.api.pipeline; package ru.kirillius.pf.sdn.api.pipeline;
import lombok.AllArgsConstructor; import lombok.*;
import lombok.Getter;
import lombok.NoArgsConstructor;
import lombok.Setter;
import ru.kirillius.pf.sdn.api.Networking.NetworkScope; import ru.kirillius.pf.sdn.api.Networking.NetworkScope;
import ru.kirillius.pf.sdn.entity.PipelineConfig;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.List; 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.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.atomic.AtomicReference;
@AllArgsConstructor @RequiredArgsConstructor()
@NoArgsConstructor
public class ProcessingPipeline { public class ProcessingPipeline {
@Getter @Getter
@Setter private final PipelineConfig config;
private List<StartCondition> startConditions = new ArrayList<>(); private final List<Action> actions;
private final AtomicBoolean running = new AtomicBoolean(false); private final AtomicBoolean running = new AtomicBoolean(false);
@Getter private final AtomicBoolean interrupted = new AtomicBoolean(false);
@Setter
private List<Action> actions = new ArrayList<>();
private AtomicInteger currentStep = new AtomicInteger(0); private AtomicInteger currentStep = new AtomicInteger(0);
public boolean isRunning() { public boolean isRunning() {
@ -37,18 +33,27 @@ public class ProcessingPipeline {
return currentStep.get(); return currentStep.get();
} }
public void interrupt() {
interrupted.set(false);
}
private List<ExecutionInfoEntry> executionInfo = new ArrayList<>(); private List<ExecutionInfoEntry> executionInfo = new ArrayList<>();
public void execute() { public void execute() {
interrupted.set(false);
running.set(true); running.set(true);
currentStep.set(0); currentStep.set(0);
try { try {
var scope = new AtomicReference<>(new NetworkScope()); var scope = new AtomicReference<>(new NetworkScope());
actions.forEach(action -> { actions.forEach(action -> {
if (interrupted.get()) {
throw new RuntimeException("Interrupted");
}
currentStep.incrementAndGet(); currentStep.incrementAndGet();
scope.set(action.apply(scope.get())); scope.set(action.apply(scope.get()));
}); });
} finally { } finally {
interrupted.set(false);
running.set(false); running.set(false);
} }
} }

View File

@ -1,13 +1,32 @@
package ru.kirillius.pf.sdn.api.pipeline; package ru.kirillius.pf.sdn.api.pipeline;
import jakarta.persistence.*;
import lombok.Getter; import lombok.Getter;
import lombok.Setter; import lombok.Setter;
import java.util.Properties; import java.util.Map;
@Getter @Getter
@Setter @Setter
@Entity
public class StartCondition { public class StartCondition {
@Id
@GeneratedValue(strategy = GenerationType.IDENTITY)
private Long id;
@Column
private boolean dontStartIfRunning;
@Column(length = 100)
private TriggerType trigger; 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<String, String> properties;
} }

View File

@ -15,5 +15,6 @@ public class DataInitializer {
@EventListener(ApplicationReadyEvent.class) @EventListener(ApplicationReadyEvent.class)
public void init() { public void init() {
authService.createDefaultUserIfAbsent(); authService.createDefaultUserIfAbsent();
} }
} }

View File

@ -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<String, String> properties = new HashMap<>();
@Column(nullable = false)
private String functionId;
}

View File

@ -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<StartCondition> startConditions = new ArrayList<>();
@ManyToMany(fetch = FetchType.EAGER)//TODO проверить связность
private List<ActionConfig> actions = new ArrayList<>();
}

View File

@ -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<PipelineConfig, Long> {
boolean existsByGuid(UUID guid);
}

View File

@ -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<String, String> properties) {
throw new UnsupportedOperationException("Not implemented");
}
@Override
public Map<String, PropertyDescriptor> getProperties() {
return Map.of();
}
}

View File

@ -2,16 +2,20 @@ package ru.kirillius.pf.sdn.service;
import jakarta.annotation.PostConstruct; import jakarta.annotation.PostConstruct;
import jakarta.annotation.PreDestroy; import jakarta.annotation.PreDestroy;
import lombok.Getter;
import lombok.RequiredArgsConstructor; import lombok.RequiredArgsConstructor;
import lombok.SneakyThrows; import lombok.SneakyThrows;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service; 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.pipeline.ProcessingPipeline;
import java.util.Queue; import java.util.Queue;
import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.ExecutorService; import java.util.concurrent.ExecutorService;
import java.util.concurrent.Future; import java.util.concurrent.Future;
import java.util.concurrent.atomic.AtomicReference;
@Slf4j @Slf4j
@Service @Service
@ -19,6 +23,9 @@ import java.util.concurrent.Future;
public class PipelineExecutorService { public class PipelineExecutorService {
private final ExecutorService executor; private final ExecutorService executor;
private final Queue<ProcessingPipeline> executionQueue = new ConcurrentLinkedQueue<>(); private final Queue<ProcessingPipeline> executionQueue = new ConcurrentLinkedQueue<>();
private final AtomicReference<ProcessingPipeline> currentPipeline = new AtomicReference<>();
@Getter
private final EventHandler<ProcessingPipeline> onExecuted = new ConcurrentEventHandler<>();
public void triggerExecute(ProcessingPipeline pipeline) { public void triggerExecute(ProcessingPipeline pipeline) {
if (executionQueue.contains(pipeline)) { if (executionQueue.contains(pipeline)) {
@ -27,6 +34,11 @@ public class PipelineExecutorService {
executionQueue.add(pipeline); executionQueue.add(pipeline);
} }
public void cancel(ProcessingPipeline pipeline) {
pipeline.interrupt();
executionQueue.remove(pipeline);
}
@PostConstruct @PostConstruct
private void initialize() { private void initialize() {
worker = executor.submit(new PipelineWorker()); worker = executor.submit(new PipelineWorker());
@ -49,11 +61,14 @@ public class PipelineExecutorService {
while (!Thread.currentThread().isInterrupted()) { while (!Thread.currentThread().isInterrupted()) {
while (!executionQueue.isEmpty()) { while (!executionQueue.isEmpty()) {
var pipeline = executionQueue.poll(); var pipeline = executionQueue.poll();
currentPipeline.set(pipeline);
try { try {
pipeline.execute(); pipeline.execute();
onExecuted.invoke(pipeline);
} catch (Exception e) { } catch (Exception e) {
log.error("Exception while executing pipeline", e); log.error("Exception while executing pipeline", e);
} }
currentPipeline.set(null);
} }
Thread.sleep(100L); Thread.sleep(100L);
Thread.yield(); Thread.yield();

View File

@ -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<PipelineConfig, ProcessingPipeline> pipelines = new ConcurrentHashMap<>();
private final Map<String, Class<? extends PipelineFunction>> functions = new ConcurrentHashMap<>();
private EventListener<ProcessingPipeline> executedListener;
private EventListener<Subscription> subscriptionUpdateListener;
public void registerFunction(Class<? extends PipelineFunction> 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;
}
}
}

View File

@ -40,10 +40,11 @@ public class SubscriptionService {
private final Map<String, SubscriptionProvider> providers = new ConcurrentHashMap<>(); private final Map<String, SubscriptionProvider> providers = new ConcurrentHashMap<>();
private final AtomicInteger updateCounter = new AtomicInteger(0); private final AtomicInteger updateCounter = new AtomicInteger(0);
private final List<Subscription> subscriptions = new CopyOnWriteArrayList<>(); private final List<Subscription> subscriptions = new CopyOnWriteArrayList<>();
private final ApplicationContext context;
@Getter
private final EventHandler<Subscription> updateEvent = new ConcurrentEventHandler<>();
private Future<?> worker; private Future<?> worker;
public List<Subscription> getSubscriptions() { public List<Subscription> getSubscriptions() {
return Collections.unmodifiableList(subscriptions); return Collections.unmodifiableList(subscriptions);
} }
@ -60,8 +61,6 @@ public class SubscriptionService {
} }
} }
private final ApplicationContext context;
public void reloadProviders() { public void reloadProviders() {
var beanFactory = context.getAutowireCapableBeanFactory(); var beanFactory = context.getAutowireCapableBeanFactory();
synchronized (providers) { synchronized (providers) {
@ -81,9 +80,6 @@ public class SubscriptionService {
} }
} }
@Getter
private final EventHandler<Subscription> updateEvent = new ConcurrentEventHandler<>();
private void reloadSubscriptions() { private void reloadSubscriptions() {
var retrieved = new ArrayList<Subscription>(); var retrieved = new ArrayList<Subscription>();
var updated = new ArrayList<Subscription>(); var updated = new ArrayList<Subscription>();