WIP
This commit is contained in:
parent
ad5298bc39
commit
ab24066f81
|
|
@ -49,3 +49,4 @@ cache/
|
||||||
/data/
|
/data/
|
||||||
test/
|
test/
|
||||||
pfsdn.mv.db
|
pfsdn.mv.db
|
||||||
|
webui/node_modules/
|
||||||
|
|
|
||||||
|
|
@ -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<String, String> properties;
|
||||||
|
}
|
||||||
|
|
@ -1,4 +1,4 @@
|
||||||
package ru.kirillius.pf.sdn.api.pipeline;
|
package ru.kirillius.pf.sdn.api.flow;
|
||||||
|
|
||||||
import lombok.Getter;
|
import lombok.Getter;
|
||||||
import ru.kirillius.pf.sdn.api.Networking.NetworkScope;
|
import ru.kirillius.pf.sdn.api.Networking.NetworkScope;
|
||||||
|
|
@ -1,4 +1,4 @@
|
||||||
package ru.kirillius.pf.sdn.api.pipeline;
|
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;
|
||||||
|
|
@ -11,10 +11,10 @@ import java.util.concurrent.atomic.AtomicInteger;
|
||||||
import java.util.concurrent.atomic.AtomicReference;
|
import java.util.concurrent.atomic.AtomicReference;
|
||||||
|
|
||||||
@RequiredArgsConstructor()
|
@RequiredArgsConstructor()
|
||||||
public class ProcessingPipeline {
|
public class Flow {
|
||||||
@Getter
|
@Getter
|
||||||
private final PipelineConfig config;
|
private final PipelineConfig config;
|
||||||
private final List<Action> actions;
|
private final List<FlowAction> actions;
|
||||||
|
|
||||||
private final AtomicBoolean running = new AtomicBoolean(false);
|
private final AtomicBoolean running = new AtomicBoolean(false);
|
||||||
private final AtomicBoolean interrupted = new AtomicBoolean(false);
|
private final AtomicBoolean interrupted = new AtomicBoolean(false);
|
||||||
|
|
@ -1,4 +1,4 @@
|
||||||
package ru.kirillius.pf.sdn.api.pipeline;
|
package ru.kirillius.pf.sdn.api.flow;
|
||||||
|
|
||||||
import lombok.Builder;
|
import lombok.Builder;
|
||||||
import lombok.Getter;
|
import lombok.Getter;
|
||||||
|
|
@ -10,11 +10,11 @@ import java.util.Map;
|
||||||
import java.util.function.Function;
|
import java.util.function.Function;
|
||||||
|
|
||||||
@Builder
|
@Builder
|
||||||
public class Action implements Function<NetworkScope, NetworkScope> {
|
public class FlowAction implements Function<NetworkScope, CallContext> {
|
||||||
|
|
||||||
private PipelineFunction function;
|
private FlowFunction function;
|
||||||
|
|
||||||
public Action(Class<? extends PipelineFunction> functionClass, Map<String, String> properties) {
|
public FlowAction(Class<? extends FlowFunction> functionClass, Map<String, String> properties) {
|
||||||
this.functionClass = functionClass;
|
this.functionClass = functionClass;
|
||||||
this.properties = properties;
|
this.properties = properties;
|
||||||
instantiateFunction();
|
instantiateFunction();
|
||||||
|
|
@ -25,19 +25,22 @@ public class Action implements Function<NetworkScope, NetworkScope> {
|
||||||
function = functionClass.getConstructor().newInstance();
|
function = functionClass.getConstructor().newInstance();
|
||||||
}
|
}
|
||||||
|
|
||||||
public void setFunctionClass(Class<? extends PipelineFunction> functionClass) {
|
public void setFunctionClass(Class<? extends FlowFunction> functionClass) {
|
||||||
this.functionClass = functionClass;
|
this.functionClass = functionClass;
|
||||||
instantiateFunction();
|
instantiateFunction();
|
||||||
}
|
}
|
||||||
|
|
||||||
@Getter
|
@Getter
|
||||||
private Class<? extends PipelineFunction> functionClass;
|
private Class<? extends FlowFunction> functionClass;
|
||||||
@Getter
|
@Getter
|
||||||
@Setter
|
@Setter
|
||||||
private Map<String, String> properties;
|
private Map<String, String> properties;
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public NetworkScope apply(NetworkScope source) {
|
public CallContext apply(NetworkScope networkScope) {
|
||||||
return function.apply(source, properties);
|
var context = new CallContext(networkScope, properties);
|
||||||
|
function.apply(context);
|
||||||
|
return context;
|
||||||
}
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
@ -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<String, PropertyDescriptor> getProperties();
|
||||||
|
|
||||||
|
boolean hasInputs();
|
||||||
|
|
||||||
|
boolean hasOutputs();
|
||||||
|
}
|
||||||
|
|
@ -1,4 +1,4 @@
|
||||||
package ru.kirillius.pf.sdn.api.pipeline;
|
package ru.kirillius.pf.sdn.api.flow;
|
||||||
|
|
||||||
import jakarta.persistence.*;
|
import jakarta.persistence.*;
|
||||||
import lombok.Getter;
|
import lombok.Getter;
|
||||||
|
|
@ -1,4 +1,4 @@
|
||||||
package ru.kirillius.pf.sdn.api.pipeline;
|
package ru.kirillius.pf.sdn.api.flow;
|
||||||
|
|
||||||
public enum TriggerType {
|
public enum TriggerType {
|
||||||
Manual,
|
Manual,
|
||||||
|
|
@ -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<String, String> properties);
|
|
||||||
|
|
||||||
Map<String, PropertyDescriptor> getProperties();
|
|
||||||
}
|
|
||||||
|
|
@ -2,7 +2,7 @@ package ru.kirillius.pf.sdn.entity;
|
||||||
|
|
||||||
import jakarta.persistence.*;
|
import jakarta.persistence.*;
|
||||||
import lombok.*;
|
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.ArrayList;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
|
|
|
||||||
|
|
@ -1,8 +1,9 @@
|
||||||
package ru.kirillius.pf.sdn.pipeline;
|
package ru.kirillius.pf.sdn.flow;
|
||||||
|
|
||||||
import lombok.RequiredArgsConstructor;
|
import lombok.RequiredArgsConstructor;
|
||||||
import ru.kirillius.pf.sdn.api.Networking.NetworkScope;
|
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.PropertyDescriptor;
|
||||||
import ru.kirillius.pf.sdn.api.properties.PropertyType;
|
import ru.kirillius.pf.sdn.api.properties.PropertyType;
|
||||||
import ru.kirillius.pf.sdn.service.AutonomousSystemCacheService;
|
import ru.kirillius.pf.sdn.service.AutonomousSystemCacheService;
|
||||||
|
|
@ -12,23 +13,9 @@ import java.util.HashMap;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
|
|
||||||
@RequiredArgsConstructor
|
@RequiredArgsConstructor
|
||||||
public class AutoResolveASN implements PipelineFunction {
|
public class AutoResolveASN implements FlowFunction {
|
||||||
|
|
||||||
private final AutonomousSystemCacheService autonomousSystemCacheService;
|
private final AutonomousSystemCacheService autonomousSystemCacheService;
|
||||||
|
|
||||||
@Override
|
|
||||||
public NetworkScope apply(NetworkScope source, Map<String, String> 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 String CLEAR = "clear";
|
||||||
private final static Map<String, PropertyDescriptor> properties;
|
private final static Map<String, PropertyDescriptor> 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
|
@Override
|
||||||
public Map<String, PropertyDescriptor> getProperties() {
|
public Map<String, PropertyDescriptor> getProperties() {
|
||||||
return Collections.unmodifiableMap(properties);
|
return Collections.unmodifiableMap(properties);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasInputs() {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasOutputs() {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -1,9 +1,10 @@
|
||||||
package ru.kirillius.pf.sdn.pipeline;
|
package ru.kirillius.pf.sdn.flow;
|
||||||
|
|
||||||
import lombok.RequiredArgsConstructor;
|
import lombok.RequiredArgsConstructor;
|
||||||
import ru.kirillius.pf.sdn.api.Networking.IPv4Subnet;
|
import ru.kirillius.pf.sdn.api.Networking.IPv4Subnet;
|
||||||
import ru.kirillius.pf.sdn.api.Networking.NetworkScope;
|
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.PropertyDescriptor;
|
||||||
import ru.kirillius.pf.sdn.api.properties.PropertyType;
|
import ru.kirillius.pf.sdn.api.properties.PropertyType;
|
||||||
import ru.kirillius.pf.sdn.service.DomainUpdaterService;
|
import ru.kirillius.pf.sdn.service.DomainUpdaterService;
|
||||||
|
|
@ -13,31 +14,10 @@ import java.util.HashMap;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
|
|
||||||
@RequiredArgsConstructor
|
@RequiredArgsConstructor
|
||||||
public class AutoResolveDomains implements PipelineFunction {
|
public class AutoResolveDomains implements FlowFunction {
|
||||||
|
|
||||||
private final DomainUpdaterService domainCacheService;
|
private final DomainUpdaterService domainCacheService;
|
||||||
|
|
||||||
@Override
|
|
||||||
public NetworkScope apply(NetworkScope source, Map<String, String> 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 String CLEAR = "clear";
|
||||||
private final static Map<String, PropertyDescriptor> properties;
|
private final static Map<String, PropertyDescriptor> 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
|
@Override
|
||||||
public Map<String, PropertyDescriptor> getProperties() {
|
public Map<String, PropertyDescriptor> getProperties() {
|
||||||
return Collections.unmodifiableMap(properties);
|
return Collections.unmodifiableMap(properties);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasInputs() {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasOutputs() {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -1,10 +1,10 @@
|
||||||
package ru.kirillius.pf.sdn.pipeline;
|
package ru.kirillius.pf.sdn.flow;
|
||||||
|
|
||||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||||
import lombok.RequiredArgsConstructor;
|
import lombok.RequiredArgsConstructor;
|
||||||
import lombok.extern.slf4j.Slf4j;
|
import lombok.extern.slf4j.Slf4j;
|
||||||
import ru.kirillius.pf.sdn.api.Networking.NetworkScope;
|
import ru.kirillius.pf.sdn.api.flow.CallContext;
|
||||||
import ru.kirillius.pf.sdn.api.pipeline.PipelineFunction;
|
import ru.kirillius.pf.sdn.api.flow.FlowFunction;
|
||||||
import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor;
|
import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor;
|
||||||
|
|
||||||
import java.util.Collections;
|
import java.util.Collections;
|
||||||
|
|
@ -12,13 +12,12 @@ import java.util.Map;
|
||||||
|
|
||||||
@Slf4j
|
@Slf4j
|
||||||
@RequiredArgsConstructor
|
@RequiredArgsConstructor
|
||||||
public class DebugOutput implements PipelineFunction {
|
public class DebugOutput implements FlowFunction {
|
||||||
private final ObjectMapper objectMapper;
|
private final ObjectMapper objectMapper;
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public NetworkScope apply(NetworkScope source, Map<String, String> properties) {
|
public void apply(CallContext context) {
|
||||||
log.info(objectMapper.valueToTree(source).toString());
|
log.info(objectMapper.valueToTree(context.getSource()).toString());
|
||||||
return source;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
|
|
@ -26,4 +25,14 @@ public class DebugOutput implements PipelineFunction {
|
||||||
return Collections.emptyMap();
|
return Collections.emptyMap();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasInputs() {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasOutputs() {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
@ -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<String, PropertyDescriptor> getProperties() {
|
||||||
|
return Map.of();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasInputs() {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasOutputs() {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -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<String, PropertyDescriptor> getProperties() {
|
||||||
|
return Collections.emptyMap();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasInputs() {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasOutputs() {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -1,9 +1,10 @@
|
||||||
package ru.kirillius.pf.sdn.pipeline;
|
package ru.kirillius.pf.sdn.flow;
|
||||||
|
|
||||||
import lombok.RequiredArgsConstructor;
|
import lombok.RequiredArgsConstructor;
|
||||||
import lombok.extern.slf4j.Slf4j;
|
import lombok.extern.slf4j.Slf4j;
|
||||||
import ru.kirillius.pf.sdn.api.Networking.NetworkScope;
|
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.PropertyDescriptor;
|
||||||
|
|
||||||
import java.util.Collections;
|
import java.util.Collections;
|
||||||
|
|
@ -11,16 +12,17 @@ import java.util.Map;
|
||||||
|
|
||||||
@Slf4j
|
@Slf4j
|
||||||
@RequiredArgsConstructor
|
@RequiredArgsConstructor
|
||||||
public class DuplicatesFilter implements PipelineFunction {
|
public class DuplicatesFilter implements FlowFunction {
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public NetworkScope apply(NetworkScope source, Map<String, String> properties) {
|
public void apply(CallContext context) {
|
||||||
|
var source = context.getSource();
|
||||||
var output = new NetworkScope();
|
var output = new NetworkScope();
|
||||||
output.setSubnets(source.getSubnets().stream().distinct().toList());
|
output.setSubnets(source.getSubnets().stream().distinct().toList());
|
||||||
output.setDomains(source.getDomains().stream().distinct().toList());
|
output.setDomains(source.getDomains().stream().distinct().toList());
|
||||||
output.setASN(source.getASN().stream().distinct().toList());
|
output.setASN(source.getASN().stream().distinct().toList());
|
||||||
output.setAutoResolvedDomains(source.getAutoResolvedDomains().stream().distinct().toList());
|
output.setAutoResolvedDomains(source.getAutoResolvedDomains().stream().distinct().toList());
|
||||||
return output;
|
context.setOutput(output);
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
|
|
@ -28,4 +30,14 @@ public class DuplicatesFilter implements PipelineFunction {
|
||||||
return Collections.emptyMap();
|
return Collections.emptyMap();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasInputs() {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasOutputs() {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
@ -1,10 +1,11 @@
|
||||||
package ru.kirillius.pf.sdn.pipeline;
|
package ru.kirillius.pf.sdn.flow;
|
||||||
|
|
||||||
import lombok.RequiredArgsConstructor;
|
import lombok.RequiredArgsConstructor;
|
||||||
import lombok.extern.slf4j.Slf4j;
|
import lombok.extern.slf4j.Slf4j;
|
||||||
import ru.kirillius.pf.sdn.Util.IPv4Util;
|
import ru.kirillius.pf.sdn.Util.IPv4Util;
|
||||||
import ru.kirillius.pf.sdn.api.Networking.NetworkScope;
|
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.PropertyDescriptor;
|
||||||
|
|
||||||
import java.util.Collections;
|
import java.util.Collections;
|
||||||
|
|
@ -12,16 +13,16 @@ import java.util.Map;
|
||||||
|
|
||||||
@Slf4j
|
@Slf4j
|
||||||
@RequiredArgsConstructor
|
@RequiredArgsConstructor
|
||||||
public class MergeNeighbourSubnetFilter implements PipelineFunction {
|
public class MergeNeighbourSubnetFilter implements FlowFunction {
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public NetworkScope apply(NetworkScope source, Map<String, String> properties) {
|
public void apply(CallContext context) {
|
||||||
|
var source = context.getSource();
|
||||||
var output = new NetworkScope();
|
var output = new NetworkScope();
|
||||||
output.setSubnets(IPv4Util.mergeNeighbours(source.getSubnets()));
|
output.setSubnets(IPv4Util.mergeNeighbours(source.getSubnets()));
|
||||||
output.setDomains(source.getDomains());
|
output.setDomains(source.getDomains());
|
||||||
output.setASN(source.getASN());
|
output.setASN(source.getASN());
|
||||||
output.setAutoResolvedDomains(source.getAutoResolvedDomains());
|
output.setAutoResolvedDomains(source.getAutoResolvedDomains());
|
||||||
return output;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
|
|
@ -29,4 +30,14 @@ public class MergeNeighbourSubnetFilter implements PipelineFunction {
|
||||||
return Collections.emptyMap();
|
return Collections.emptyMap();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasInputs() {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasOutputs() {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
@ -1,10 +1,11 @@
|
||||||
package ru.kirillius.pf.sdn.pipeline;
|
package ru.kirillius.pf.sdn.flow;
|
||||||
|
|
||||||
import lombok.RequiredArgsConstructor;
|
import lombok.RequiredArgsConstructor;
|
||||||
import lombok.extern.slf4j.Slf4j;
|
import lombok.extern.slf4j.Slf4j;
|
||||||
import ru.kirillius.pf.sdn.Util.IPv4Util;
|
import ru.kirillius.pf.sdn.Util.IPv4Util;
|
||||||
import ru.kirillius.pf.sdn.api.Networking.NetworkScope;
|
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.IntegerConstraint;
|
||||||
import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor;
|
import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor;
|
||||||
import ru.kirillius.pf.sdn.api.properties.PropertyType;
|
import ru.kirillius.pf.sdn.api.properties.PropertyType;
|
||||||
|
|
@ -15,25 +16,8 @@ import java.util.Map;
|
||||||
|
|
||||||
@Slf4j
|
@Slf4j
|
||||||
@RequiredArgsConstructor
|
@RequiredArgsConstructor
|
||||||
public class MergeWithCoverageSubnetFilter implements PipelineFunction {
|
public class MergeWithCoverageSubnetFilter implements FlowFunction {
|
||||||
private final static String USAGE = "usage";
|
private final static String USAGE = "usage";
|
||||||
|
|
||||||
@Override
|
|
||||||
public NetworkScope apply(NetworkScope source, Map<String, String> 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<String, PropertyDescriptor> properties;
|
private final static Map<String, PropertyDescriptor> properties;
|
||||||
|
|
||||||
static {
|
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
|
@Override
|
||||||
public Map<String, PropertyDescriptor> getProperties() {
|
public Map<String, PropertyDescriptor> getProperties() {
|
||||||
return Collections.unmodifiableMap(properties);
|
return Collections.unmodifiableMap(properties);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasInputs() {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasOutputs() {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
@ -1,10 +1,11 @@
|
||||||
package ru.kirillius.pf.sdn.pipeline;
|
package ru.kirillius.pf.sdn.flow;
|
||||||
|
|
||||||
import lombok.RequiredArgsConstructor;
|
import lombok.RequiredArgsConstructor;
|
||||||
import lombok.extern.slf4j.Slf4j;
|
import lombok.extern.slf4j.Slf4j;
|
||||||
import ru.kirillius.pf.sdn.Util.DomainUtil;
|
import ru.kirillius.pf.sdn.Util.DomainUtil;
|
||||||
import ru.kirillius.pf.sdn.api.Networking.NetworkScope;
|
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.PropertyDescriptor;
|
||||||
|
|
||||||
import java.util.Collections;
|
import java.util.Collections;
|
||||||
|
|
@ -12,16 +13,17 @@ import java.util.Map;
|
||||||
|
|
||||||
@Slf4j
|
@Slf4j
|
||||||
@RequiredArgsConstructor
|
@RequiredArgsConstructor
|
||||||
public class OverlappedDomainFilter implements PipelineFunction {
|
public class OverlappedDomainFilter implements FlowFunction {
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public NetworkScope apply(NetworkScope source, Map<String, String> properties) {
|
public void apply(CallContext context) {
|
||||||
|
var source = context.getSource();
|
||||||
var output = new NetworkScope();
|
var output = new NetworkScope();
|
||||||
output.setSubnets(source.getSubnets());
|
output.setSubnets(source.getSubnets());
|
||||||
output.setDomains(DomainUtil.removeOverlapped(source.getDomains()));
|
output.setDomains(DomainUtil.removeOverlapped(source.getDomains()));
|
||||||
output.setASN(source.getASN());
|
output.setASN(source.getASN());
|
||||||
output.setAutoResolvedDomains(source.getAutoResolvedDomains());
|
output.setAutoResolvedDomains(source.getAutoResolvedDomains());
|
||||||
return output;
|
context.setOutput(output);
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
|
|
@ -29,4 +31,14 @@ public class OverlappedDomainFilter implements PipelineFunction {
|
||||||
return Collections.emptyMap();
|
return Collections.emptyMap();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasInputs() {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasOutputs() {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
@ -1,10 +1,11 @@
|
||||||
package ru.kirillius.pf.sdn.pipeline;
|
package ru.kirillius.pf.sdn.flow;
|
||||||
|
|
||||||
import lombok.RequiredArgsConstructor;
|
import lombok.RequiredArgsConstructor;
|
||||||
import lombok.extern.slf4j.Slf4j;
|
import lombok.extern.slf4j.Slf4j;
|
||||||
import ru.kirillius.pf.sdn.Util.IPv4Util;
|
import ru.kirillius.pf.sdn.Util.IPv4Util;
|
||||||
import ru.kirillius.pf.sdn.api.Networking.NetworkScope;
|
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.PropertyDescriptor;
|
||||||
|
|
||||||
import java.util.Collections;
|
import java.util.Collections;
|
||||||
|
|
@ -12,16 +13,17 @@ import java.util.Map;
|
||||||
|
|
||||||
@Slf4j
|
@Slf4j
|
||||||
@RequiredArgsConstructor
|
@RequiredArgsConstructor
|
||||||
public class OverlappedSubnetFilter implements PipelineFunction {
|
public class OverlappedSubnetFilter implements FlowFunction {
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public NetworkScope apply(NetworkScope source, Map<String, String> properties) {
|
public void apply(CallContext context) {
|
||||||
|
var source = context.getSource();
|
||||||
var output = new NetworkScope();
|
var output = new NetworkScope();
|
||||||
output.setSubnets(IPv4Util.removeOverlapped(source.getSubnets()));
|
output.setSubnets(IPv4Util.removeOverlapped(source.getSubnets()));
|
||||||
output.setDomains(source.getDomains());
|
output.setDomains(source.getDomains());
|
||||||
output.setASN(source.getASN());
|
output.setASN(source.getASN());
|
||||||
output.setAutoResolvedDomains(source.getAutoResolvedDomains());
|
output.setAutoResolvedDomains(source.getAutoResolvedDomains());
|
||||||
return output;
|
context.setOutput(output);
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
|
|
@ -29,4 +31,14 @@ public class OverlappedSubnetFilter implements PipelineFunction {
|
||||||
return Collections.emptyMap();
|
return Collections.emptyMap();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasInputs() {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasOutputs() {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
@ -1,10 +1,11 @@
|
||||||
package ru.kirillius.pf.sdn.pipeline;
|
package ru.kirillius.pf.sdn.flow;
|
||||||
|
|
||||||
import lombok.RequiredArgsConstructor;
|
import lombok.RequiredArgsConstructor;
|
||||||
import lombok.extern.slf4j.Slf4j;
|
import lombok.extern.slf4j.Slf4j;
|
||||||
import ru.kirillius.pf.sdn.api.Networking.IPv4Subnet;
|
import ru.kirillius.pf.sdn.api.Networking.IPv4Subnet;
|
||||||
import ru.kirillius.pf.sdn.api.Networking.NetworkScope;
|
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.IntegerConstraint;
|
||||||
import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor;
|
import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor;
|
||||||
import ru.kirillius.pf.sdn.api.properties.PropertyType;
|
import ru.kirillius.pf.sdn.api.properties.PropertyType;
|
||||||
|
|
@ -17,11 +18,26 @@ import java.util.regex.Pattern;
|
||||||
|
|
||||||
@Slf4j
|
@Slf4j
|
||||||
@RequiredArgsConstructor
|
@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<String, PropertyDescriptor> 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
|
@Override
|
||||||
public NetworkScope apply(NetworkScope source, Map<String, String> properties) {
|
public void apply(CallContext context) {
|
||||||
var output = new NetworkScope();
|
var output = new NetworkScope();
|
||||||
|
var source = context.getSource();
|
||||||
|
var properties = context.getProperties();
|
||||||
|
|
||||||
if (properties.containsKey(SUBNETS)) {
|
if (properties.containsKey(SUBNETS)) {
|
||||||
var filtered = Arrays.stream(properties.get(SUBNETS).split(Pattern.quote(",")))
|
var filtered = Arrays.stream(properties.get(SUBNETS).split(Pattern.quote(",")))
|
||||||
|
|
@ -52,20 +68,7 @@ public class ResourceFilter implements PipelineFunction {
|
||||||
output.setASN(source.getASN());
|
output.setASN(source.getASN());
|
||||||
}
|
}
|
||||||
|
|
||||||
return output;
|
context.setOutput(output);
|
||||||
}
|
|
||||||
|
|
||||||
private final static String SUBNETS = "subnets";
|
|
||||||
private final static String DOMAINS = "domains";
|
|
||||||
private final static String ASN = "asn";
|
|
||||||
|
|
||||||
private final static Map<String, PropertyDescriptor> 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
|
@Override
|
||||||
|
|
@ -73,4 +76,14 @@ public class ResourceFilter implements PipelineFunction {
|
||||||
return Collections.unmodifiableMap(properties);
|
return Collections.unmodifiableMap(properties);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasInputs() {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasOutputs() {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
@ -1,8 +1,9 @@
|
||||||
package ru.kirillius.pf.sdn.pipeline;
|
package ru.kirillius.pf.sdn.flow;
|
||||||
|
|
||||||
import lombok.RequiredArgsConstructor;
|
import lombok.RequiredArgsConstructor;
|
||||||
import ru.kirillius.pf.sdn.api.Networking.NetworkScope;
|
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.PropertyDescriptor;
|
||||||
import ru.kirillius.pf.sdn.api.properties.PropertyType;
|
import ru.kirillius.pf.sdn.api.properties.PropertyType;
|
||||||
import ru.kirillius.pf.sdn.service.SubscriptionService;
|
import ru.kirillius.pf.sdn.service.SubscriptionService;
|
||||||
|
|
@ -14,25 +15,10 @@ import java.util.Map;
|
||||||
import java.util.regex.Pattern;
|
import java.util.regex.Pattern;
|
||||||
|
|
||||||
@RequiredArgsConstructor
|
@RequiredArgsConstructor
|
||||||
public class SubscriptionInput implements PipelineFunction {
|
public class SubscriptionInput implements FlowFunction {
|
||||||
|
|
||||||
private final SubscriptionService subscriptionService;
|
private final SubscriptionService subscriptionService;
|
||||||
|
|
||||||
@Override
|
|
||||||
public NetworkScope apply(NetworkScope source, Map<String, String> 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 String NAMES = "names";
|
||||||
private final static Map<String, PropertyDescriptor> properties;
|
private final static Map<String, PropertyDescriptor> 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());
|
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
|
@Override
|
||||||
public Map<String, PropertyDescriptor> getProperties() {
|
public Map<String, PropertyDescriptor> getProperties() {
|
||||||
return Collections.unmodifiableMap(properties);
|
return Collections.unmodifiableMap(properties);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasInputs() {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean hasOutputs() {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -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<String, String> properties) {
|
|
||||||
var output = new NetworkScope();
|
|
||||||
output.add(source);
|
|
||||||
return output;
|
|
||||||
}
|
|
||||||
|
|
||||||
@Override
|
|
||||||
public Map<String, PropertyDescriptor> getProperties() {
|
|
||||||
return Collections.emptyMap();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
@ -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<String, String> properties) {
|
|
||||||
throw new UnsupportedOperationException("Not implemented");
|
|
||||||
}
|
|
||||||
|
|
||||||
@Override
|
|
||||||
public Map<String, PropertyDescriptor> getProperties() {
|
|
||||||
return Map.of();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
@ -9,7 +9,7 @@ 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.ConcurrentEventHandler;
|
||||||
import ru.kirillius.java.utils.events.EventHandler;
|
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.Queue;
|
||||||
import java.util.concurrent.ConcurrentLinkedQueue;
|
import java.util.concurrent.ConcurrentLinkedQueue;
|
||||||
|
|
@ -20,28 +20,28 @@ import java.util.concurrent.atomic.AtomicReference;
|
||||||
@Slf4j
|
@Slf4j
|
||||||
@Service
|
@Service
|
||||||
@RequiredArgsConstructor
|
@RequiredArgsConstructor
|
||||||
public class PipelineExecutorService {
|
public class FlowExecutorService {
|
||||||
private final ExecutorService executor;
|
private final ExecutorService executor;
|
||||||
private final Queue<ProcessingPipeline> executionQueue = new ConcurrentLinkedQueue<>();
|
private final Queue<Flow> executionQueue = new ConcurrentLinkedQueue<>();
|
||||||
private final AtomicReference<ProcessingPipeline> currentPipeline = new AtomicReference<>();
|
private final AtomicReference<Flow> currentPipeline = new AtomicReference<>();
|
||||||
@Getter
|
@Getter
|
||||||
private final EventHandler<ProcessingPipeline> onExecuted = new ConcurrentEventHandler<>();
|
private final EventHandler<Flow> onExecuted = new ConcurrentEventHandler<>();
|
||||||
|
|
||||||
public void triggerExecute(ProcessingPipeline pipeline) {
|
public void triggerExecute(Flow pipeline) {
|
||||||
if (executionQueue.contains(pipeline)) {
|
if (executionQueue.contains(pipeline)) {
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
executionQueue.add(pipeline);
|
executionQueue.add(pipeline);
|
||||||
}
|
}
|
||||||
|
|
||||||
public void cancel(ProcessingPipeline pipeline) {
|
public void cancel(Flow pipeline) {
|
||||||
pipeline.interrupt();
|
pipeline.interrupt();
|
||||||
executionQueue.remove(pipeline);
|
executionQueue.remove(pipeline);
|
||||||
}
|
}
|
||||||
|
|
||||||
@PostConstruct
|
@PostConstruct
|
||||||
private void initialize() {
|
private void initialize() {
|
||||||
worker = executor.submit(new PipelineWorker());
|
worker = executor.submit(new FlowWorker());
|
||||||
}
|
}
|
||||||
|
|
||||||
@PreDestroy
|
@PreDestroy
|
||||||
|
|
@ -53,11 +53,11 @@ public class PipelineExecutorService {
|
||||||
|
|
||||||
private Future<?> worker;
|
private Future<?> worker;
|
||||||
|
|
||||||
private class PipelineWorker implements Runnable {
|
private class FlowWorker implements Runnable {
|
||||||
@SuppressWarnings("BusyWait")
|
@SuppressWarnings("BusyWait")
|
||||||
@SneakyThrows
|
@SneakyThrows
|
||||||
@Override
|
@Override
|
||||||
public void run() {
|
public void run() {//TODO заменить на worker
|
||||||
while (!Thread.currentThread().isInterrupted()) {
|
while (!Thread.currentThread().isInterrupted()) {
|
||||||
while (!executionQueue.isEmpty()) {
|
while (!executionQueue.isEmpty()) {
|
||||||
var pipeline = executionQueue.poll();
|
var pipeline = executionQueue.poll();
|
||||||
|
|
@ -7,12 +7,13 @@ import lombok.extern.slf4j.Slf4j;
|
||||||
import org.springframework.stereotype.Service;
|
import org.springframework.stereotype.Service;
|
||||||
import ru.kirillius.java.utils.events.EventListener;
|
import ru.kirillius.java.utils.events.EventListener;
|
||||||
import ru.kirillius.pf.sdn.api.Networking.Subscriptions.Subscription;
|
import ru.kirillius.pf.sdn.api.Networking.Subscriptions.Subscription;
|
||||||
import ru.kirillius.pf.sdn.api.pipeline.Action;
|
import ru.kirillius.pf.sdn.api.flow.FlowAction;
|
||||||
import ru.kirillius.pf.sdn.api.pipeline.PipelineFunction;
|
import ru.kirillius.pf.sdn.api.flow.FlowFunction;
|
||||||
import ru.kirillius.pf.sdn.api.pipeline.ProcessingPipeline;
|
import ru.kirillius.pf.sdn.api.flow.Flow;
|
||||||
import ru.kirillius.pf.sdn.api.pipeline.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.PipelineConfig;
|
||||||
|
import ru.kirillius.pf.sdn.flow.DummyFunction;
|
||||||
import ru.kirillius.pf.sdn.repository.PipelineConfigRepository;
|
import ru.kirillius.pf.sdn.repository.PipelineConfigRepository;
|
||||||
|
|
||||||
import java.util.Arrays;
|
import java.util.Arrays;
|
||||||
|
|
@ -24,19 +25,19 @@ import java.util.regex.Pattern;
|
||||||
@Slf4j
|
@Slf4j
|
||||||
@Service
|
@Service
|
||||||
@RequiredArgsConstructor
|
@RequiredArgsConstructor
|
||||||
public class PipelineService {
|
public class FlowService {
|
||||||
|
|
||||||
private final PipelineConfigRepository configRepository;
|
private final PipelineConfigRepository configRepository;
|
||||||
private final PipelineExecutorService pipelineExecutorService;
|
private final FlowExecutorService flowExecutorService;
|
||||||
private final SubscriptionService subscriptionService;
|
private final SubscriptionService subscriptionService;
|
||||||
private final ExecutorService executorService;
|
private final ExecutorService executorService;
|
||||||
|
|
||||||
private final Map<PipelineConfig, ProcessingPipeline> pipelines = new ConcurrentHashMap<>();
|
private final Map<PipelineConfig, Flow> pipelines = new ConcurrentHashMap<>();
|
||||||
private final Map<String, Class<? extends PipelineFunction>> functions = new ConcurrentHashMap<>();
|
private final Map<String, Class<? extends FlowFunction>> functions = new ConcurrentHashMap<>();
|
||||||
private EventListener<ProcessingPipeline> executedListener;
|
private EventListener<Flow> executedListener;
|
||||||
private EventListener<Subscription> subscriptionUpdateListener;
|
private EventListener<Subscription> subscriptionUpdateListener;
|
||||||
|
|
||||||
public void registerFunction(Class<? extends PipelineFunction> functionClass, String id) {
|
public void registerFunction(Class<? extends FlowFunction> functionClass, String id) {
|
||||||
if (functions.containsKey(id)) {
|
if (functions.containsKey(id)) {
|
||||||
throw new IllegalStateException("Function with id '" + id + "' already exists");
|
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();
|
var matchedConditions = pipeline.getConfig().getStartConditions().stream().filter(c -> c.getTrigger() == TriggerType.Manual).toList();
|
||||||
if (matchedConditions.isEmpty()) {
|
if (matchedConditions.isEmpty()) {
|
||||||
log.error("Unable to manual start pipeline {} because it has no manual trigger", pipeline.getConfig().getName());
|
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()) {
|
if (c.isDontStartIfRunning() && pipeline.isRunning()) {
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
pipelineExecutorService.triggerExecute(pipeline);
|
flowExecutorService.triggerExecute(pipeline);
|
||||||
});
|
});
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
@ -72,23 +73,23 @@ public class PipelineService {
|
||||||
if (configRepository.existsByGuid(pipelineConfig.getGuid())) {
|
if (configRepository.existsByGuid(pipelineConfig.getGuid())) {
|
||||||
load(pipelineConfig);
|
load(pipelineConfig);
|
||||||
} else {
|
} else {
|
||||||
pipelineExecutorService.cancel(pipelines.get(pipelineConfig));
|
flowExecutorService.cancel(pipelines.get(pipelineConfig));
|
||||||
pipelines.remove(pipelineConfig);
|
pipelines.remove(pipelineConfig);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private void load(PipelineConfig pipelineConfig) {
|
private void load(PipelineConfig pipelineConfig) {
|
||||||
if (pipelines.containsKey(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) {
|
private FlowAction buildAction(ActionConfig config) {
|
||||||
return new Action(functions.getOrDefault(config.getFunctionId(), DummyFunction.class), config.getProperties());
|
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();
|
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()) {
|
if (c.isDontStartIfRunning() && pipeline.isRunning()) {
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
pipelineExecutorService.triggerExecute(pipeline);
|
flowExecutorService.triggerExecute(pipeline);
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|
@ -125,7 +126,7 @@ public class PipelineService {
|
||||||
@PostConstruct
|
@PostConstruct
|
||||||
private void initialize() {
|
private void initialize() {
|
||||||
subscriptionUpdateListener = subscriptionService.getUpdateEvent().add(this::subscriptionUpdate);
|
subscriptionUpdateListener = subscriptionService.getUpdateEvent().add(this::subscriptionUpdate);
|
||||||
executedListener = pipelineExecutorService.getOnExecuted().add(this::checkForIntervalExecution);
|
executedListener = flowExecutorService.getOnExecuted().add(this::checkForIntervalExecution);
|
||||||
registerFunction(DummyFunction.class, "Error:fallback");
|
registerFunction(DummyFunction.class, "Error:fallback");
|
||||||
configRepository.findAll().forEach(this::load);
|
configRepository.findAll().forEach(this::load);
|
||||||
pipelines.values()
|
pipelines.values()
|
||||||
|
|
@ -135,7 +136,7 @@ public class PipelineService {
|
||||||
.getStartConditions()
|
.getStartConditions()
|
||||||
.stream()
|
.stream()
|
||||||
.anyMatch(c -> c.getTrigger() == TriggerType.OnStart))
|
.anyMatch(c -> c.getTrigger() == TriggerType.OnStart))
|
||||||
.forEach(pipelineExecutorService::triggerExecute);
|
.forEach(flowExecutorService::triggerExecute);
|
||||||
}
|
}
|
||||||
|
|
||||||
private void subscriptionUpdate(Subscription subscription) {
|
private void subscriptionUpdate(Subscription subscription) {
|
||||||
|
|
@ -154,12 +155,12 @@ public class PipelineService {
|
||||||
var properties = c.getProperties();
|
var properties = c.getProperties();
|
||||||
var names = properties.getOrDefault("names", null);
|
var names = properties.getOrDefault("names", null);
|
||||||
if (names == null || names.isEmpty()) {
|
if (names == null || names.isEmpty()) {
|
||||||
pipelineExecutorService.triggerExecute(pipeline);
|
flowExecutorService.triggerExecute(pipeline);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
if (Arrays.stream(names.split(Pattern.quote(","))).anyMatch(s -> subscription.getName().equals(s))) {
|
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
|
@PreDestroy
|
||||||
private void destroy() {
|
private void destroy() {
|
||||||
if (executedListener != null) {
|
if (executedListener != null) {
|
||||||
pipelineExecutorService.getOnExecuted().remove(executedListener);
|
flowExecutorService.getOnExecuted().remove(executedListener);
|
||||||
executedListener = null;
|
executedListener = null;
|
||||||
}
|
}
|
||||||
if (subscriptionUpdateListener != null) {
|
if (subscriptionUpdateListener != null) {
|
||||||
Loading…
Reference in New Issue