WIP
This commit is contained in:
parent
ab24066f81
commit
fe71aae7aa
|
|
@ -2,7 +2,7 @@ package ru.kirillius.pf.sdn.api.flow;
|
||||||
|
|
||||||
import lombok.*;
|
import lombok.*;
|
||||||
import ru.kirillius.pf.sdn.api.Networking.NetworkScope;
|
import ru.kirillius.pf.sdn.api.Networking.NetworkScope;
|
||||||
import ru.kirillius.pf.sdn.entity.PipelineConfig;
|
import ru.kirillius.pf.sdn.entity.FlowConfig;
|
||||||
|
|
||||||
import java.util.ArrayList;
|
import java.util.ArrayList;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
|
|
@ -13,7 +13,7 @@ import java.util.concurrent.atomic.AtomicReference;
|
||||||
@RequiredArgsConstructor()
|
@RequiredArgsConstructor()
|
||||||
public class Flow {
|
public class Flow {
|
||||||
@Getter
|
@Getter
|
||||||
private final PipelineConfig config;
|
private final FlowConfig config;
|
||||||
private final List<FlowAction> actions;
|
private final List<FlowAction> actions;
|
||||||
|
|
||||||
private final AtomicBoolean running = new AtomicBoolean(false);
|
private final AtomicBoolean running = new AtomicBoolean(false);
|
||||||
|
|
|
||||||
|
|
@ -16,7 +16,7 @@ import java.util.UUID;
|
||||||
@NoArgsConstructor
|
@NoArgsConstructor
|
||||||
@Builder
|
@Builder
|
||||||
@AllArgsConstructor
|
@AllArgsConstructor
|
||||||
public class PipelineConfig {
|
public class FlowConfig {
|
||||||
@Id
|
@Id
|
||||||
@GeneratedValue(strategy = GenerationType.IDENTITY)
|
@GeneratedValue(strategy = GenerationType.IDENTITY)
|
||||||
private Long id;
|
private Long id;
|
||||||
|
|
@ -29,7 +29,7 @@ public class PipelineConfig {
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public boolean equals(Object o) {
|
public boolean equals(Object o) {
|
||||||
if (!(o instanceof PipelineConfig that)) return false;
|
if (!(o instanceof FlowConfig that)) return false;
|
||||||
return Objects.equals(id, that.id) && Objects.equals(name, that.name);
|
return Objects.equals(id, that.id) && Objects.equals(name, that.name);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -1,11 +1,11 @@
|
||||||
package ru.kirillius.pf.sdn.repository;
|
package ru.kirillius.pf.sdn.repository;
|
||||||
|
|
||||||
import org.springframework.data.jpa.repository.JpaRepository;
|
import org.springframework.data.jpa.repository.JpaRepository;
|
||||||
import ru.kirillius.pf.sdn.entity.PipelineConfig;
|
import ru.kirillius.pf.sdn.entity.FlowConfig;
|
||||||
|
|
||||||
import java.util.UUID;
|
import java.util.UUID;
|
||||||
|
|
||||||
public interface PipelineConfigRepository extends JpaRepository<PipelineConfig, Long> {
|
public interface PipelineConfigRepository extends JpaRepository<FlowConfig, Long> {
|
||||||
|
|
||||||
boolean existsByGuid(UUID guid);
|
boolean existsByGuid(UUID guid);
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -12,7 +12,7 @@ import ru.kirillius.pf.sdn.api.flow.FlowFunction;
|
||||||
import ru.kirillius.pf.sdn.api.flow.Flow;
|
import ru.kirillius.pf.sdn.api.flow.Flow;
|
||||||
import ru.kirillius.pf.sdn.api.flow.TriggerType;
|
import ru.kirillius.pf.sdn.api.flow.TriggerType;
|
||||||
import ru.kirillius.pf.sdn.entity.ActionConfig;
|
import ru.kirillius.pf.sdn.entity.ActionConfig;
|
||||||
import ru.kirillius.pf.sdn.entity.PipelineConfig;
|
import ru.kirillius.pf.sdn.entity.FlowConfig;
|
||||||
import ru.kirillius.pf.sdn.flow.DummyFunction;
|
import ru.kirillius.pf.sdn.flow.DummyFunction;
|
||||||
import ru.kirillius.pf.sdn.repository.PipelineConfigRepository;
|
import ru.kirillius.pf.sdn.repository.PipelineConfigRepository;
|
||||||
|
|
||||||
|
|
@ -32,7 +32,7 @@ public class FlowService {
|
||||||
private final SubscriptionService subscriptionService;
|
private final SubscriptionService subscriptionService;
|
||||||
private final ExecutorService executorService;
|
private final ExecutorService executorService;
|
||||||
|
|
||||||
private final Map<PipelineConfig, Flow> pipelines = new ConcurrentHashMap<>();
|
private final Map<FlowConfig, Flow> pipelines = new ConcurrentHashMap<>();
|
||||||
private final Map<String, Class<? extends FlowFunction>> functions = new ConcurrentHashMap<>();
|
private final Map<String, Class<? extends FlowFunction>> functions = new ConcurrentHashMap<>();
|
||||||
private EventListener<Flow> executedListener;
|
private EventListener<Flow> executedListener;
|
||||||
private EventListener<Subscription> subscriptionUpdateListener;
|
private EventListener<Subscription> subscriptionUpdateListener;
|
||||||
|
|
@ -48,9 +48,9 @@ public class FlowService {
|
||||||
pipelines.keySet().stream().filter(config -> config.getId().equals(id)).findFirst().ifPresent(this::trigger);
|
pipelines.keySet().stream().filter(config -> config.getId().equals(id)).findFirst().ifPresent(this::trigger);
|
||||||
}
|
}
|
||||||
|
|
||||||
public void trigger(PipelineConfig pipelineConfig) {
|
public void trigger(FlowConfig flowConfig) {
|
||||||
if (pipelines.containsKey(pipelineConfig)) {
|
if (pipelines.containsKey(flowConfig)) {
|
||||||
trigger(pipelines.get(pipelineConfig));
|
trigger(pipelines.get(flowConfig));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -69,20 +69,20 @@ public class FlowService {
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
public void update(PipelineConfig pipelineConfig) {
|
public void update(FlowConfig flowConfig) {
|
||||||
if (configRepository.existsByGuid(pipelineConfig.getGuid())) {
|
if (configRepository.existsByGuid(flowConfig.getGuid())) {
|
||||||
load(pipelineConfig);
|
load(flowConfig);
|
||||||
} else {
|
} else {
|
||||||
flowExecutorService.cancel(pipelines.get(pipelineConfig));
|
flowExecutorService.cancel(pipelines.get(flowConfig));
|
||||||
pipelines.remove(pipelineConfig);
|
pipelines.remove(flowConfig);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private void load(PipelineConfig pipelineConfig) {
|
private void load(FlowConfig flowConfig) {
|
||||||
if (pipelines.containsKey(pipelineConfig)) {
|
if (pipelines.containsKey(flowConfig)) {
|
||||||
flowExecutorService.cancel(pipelines.get(pipelineConfig));
|
flowExecutorService.cancel(pipelines.get(flowConfig));
|
||||||
}
|
}
|
||||||
pipelines.put(pipelineConfig, new Flow(pipelineConfig, pipelineConfig.getActions().stream().map(this::buildAction).toList()));
|
pipelines.put(flowConfig, new Flow(flowConfig, flowConfig.getActions().stream().map(this::buildAction).toList()));
|
||||||
}
|
}
|
||||||
|
|
||||||
private FlowAction buildAction(ActionConfig config) {
|
private FlowAction buildAction(ActionConfig config) {
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue