Compare commits

...

2 Commits

Author SHA1 Message Date
kirillius 77044a2c64 WIP 2026-07-26 01:58:55 +03:00
kirillius e4745806c2 WIP: ExecutionGraph & Multiple inputs & Subscription groups 2026-07-21 00:12:31 +03:00
34 changed files with 843 additions and 288 deletions

View File

@ -25,13 +25,12 @@ public class NetworkScope {
@Getter @Getter
@Setter @Setter
private List<String> domains = new ArrayList<>(); private List<String> domains = new ArrayList<>();
@Getter @Getter
@Setter @Setter
private List<String> autoResolvedDomains = new ArrayList<>(); private List<String> resolveDomains = new ArrayList<>();
public int getResourceCount() { public int getResourceCount() {
return subnets.size() + ASN.size() + domains.size() + autoResolvedDomains.size(); return subnets.size() + ASN.size() + domains.size() + resolveDomains.size();
} }
@Override @Override
@ -40,12 +39,12 @@ public class NetworkScope {
return Objects.equals(ASN, that.ASN) && return Objects.equals(ASN, that.ASN) &&
Objects.equals(subnets, that.subnets) && Objects.equals(subnets, that.subnets) &&
Objects.equals(domains, that.domains) && Objects.equals(domains, that.domains) &&
Objects.equals(autoResolvedDomains, that.autoResolvedDomains); Objects.equals(resolveDomains, that.resolveDomains);
} }
@Override @Override
public int hashCode() { public int hashCode() {
return Objects.hash(ASN, subnets, domains, autoResolvedDomains); return Objects.hash(ASN, subnets, domains, resolveDomains);
} }
/** /**
@ -55,7 +54,7 @@ public class NetworkScope {
ASN.clear(); ASN.clear();
subnets.clear(); subnets.clear();
domains.clear(); domains.clear();
autoResolvedDomains.clear(); resolveDomains.clear();
} }
/** /**
@ -65,7 +64,7 @@ public class NetworkScope {
ASN.addAll(networkScope.getASN()); ASN.addAll(networkScope.getASN());
subnets.addAll(networkScope.getSubnets()); subnets.addAll(networkScope.getSubnets());
domains.addAll(networkScope.getDomains()); domains.addAll(networkScope.getDomains());
domains.addAll(networkScope.getAutoResolvedDomains()); resolveDomains.addAll(networkScope.getResolveDomains());
} }
} }

View File

@ -5,15 +5,16 @@ import lombok.RequiredArgsConstructor;
import lombok.Setter; import lombok.Setter;
import ru.kirillius.pf.sdn.api.Networking.NetworkScope; import ru.kirillius.pf.sdn.api.Networking.NetworkScope;
import java.util.List;
import java.util.Map; import java.util.Map;
@RequiredArgsConstructor @RequiredArgsConstructor
public final class CallContext { public final class CallContext {
@Getter @Getter
private final NetworkScope source; private final List<NetworkScope> sources;
@Getter @Getter
@Setter @Setter
private NetworkScope output; private List<NetworkScope> outputs;
@Getter @Getter
private final Map<String, String> properties; private final Map<String, String> properties;

View File

@ -0,0 +1,134 @@
package ru.kirillius.pf.sdn.api.flow;
import lombok.Builder;
import lombok.Getter;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import ru.kirillius.java.utils.events.ConcurrentEventHandler;
import ru.kirillius.java.utils.events.EventHandler;
import ru.kirillius.pf.sdn.api.Networking.NetworkScope;
import ru.kirillius.pf.sdn.entity.ActionConfig;
import ru.kirillius.pf.sdn.entity.ActionConnectionConfig;
import java.util.*;
import java.util.function.Function;
@Slf4j
@RequiredArgsConstructor
public class ExecutionGraph {
public int size() {
return nodes.size();
}
@Getter
private final EventHandler<FlowAction> startExecutingEvent = new ConcurrentEventHandler<>();
private final Map<ActionConfig, Node> nodes = new HashMap<>();
private final Map<Node, CallContext> executed = new HashMap<>();
public ExecutionGraph(List<ActionConfig> actionConfigs, List<ActionConnectionConfig> connectionConfigs, Function<ActionConfig, FlowAction> actionFactory) {
var pendingActions = new ArrayList<>(actionConfigs);
//создаём ноды у которых нет зависимостей
actionConfigs
.stream()
.filter(config -> connectionConfigs
.stream()
.noneMatch(conn -> conn.getTo().equals(config))
).forEach(config -> {
nodes.put(config, new Node(actionFactory.apply(config), Collections.emptyList()));
pendingActions.remove(config);
});
//создаём оставшиеся ноды
while (!pendingActions.isEmpty()) {
//ищем конфиги, все соединения которых можно зарезолвить
var spawned = new ArrayList<ActionConfig>();
pendingActions.forEach((actionConfig) -> {
var selfConnectionConfigs = connectionConfigs.stream().filter(conn -> conn.getTo().equals(actionConfig)).toList();
if (selfConnectionConfigs.stream().allMatch(conn -> nodes.containsKey(conn.getFrom()))) {
spawned.add(actionConfig);
nodes.put(actionConfig,
new Node(
actionFactory.apply(actionConfig),
selfConnectionConfigs
.stream()
.map(connectionConfig -> new Node.Connection(
nodes.get(connectionConfig.getFrom()),
connectionConfig.getFromPort(),
connectionConfig.getToPort()
)).toList()
)
);
}
});
pendingActions.removeAll(spawned);
if (spawned.isEmpty()) {
throw new IllegalStateException("Connection loop detected and could not be resolved");
}
}
}
public Iterator<ExecutionResult> execute() {
executed.clear();
final var configs = new HashSet<>(nodes.keySet());
return new Iterator<>() {
@Override
public boolean hasNext() {
return executed.size() != nodes.size();
}
@Override
public ExecutionResult next() {
var nodeToExecute = findNodeToExecute(configs);
var action = nodeToExecute.action();
try {
startExecutingEvent.invoke(action);
} catch (Exception e) {
log.error("Error on {}:startExecuting event", ExecutionGraph.class.getSimpleName(), e);
}
var inputs = new ArrayList<NetworkScope>();
nodeToExecute.inputs().forEach((input) -> {
var fromNode = input.from();
inputs.add(executed.get(fromNode).getOutputs().get(input.fromPort));
});
var callContext = action.apply(inputs);
executed.put(nodeToExecute, callContext);
configs.remove(nodeToExecute.action().getConfig());
return new ExecutionResult(action, callContext);
}
};
}
private Node findNodeToExecute(Set<ActionConfig> actionConfigs) {
for (var config : actionConfigs) {
var node = nodes.get(config);
//проверяем что все на всех входах есть данные
if (node.inputs().stream().allMatch(connection -> executed.containsKey(connection.from()))) {
return node;
}
}
throw new NoSuchElementException("There is No node to execute");
}
@Builder
public record ExecutionResult(FlowAction action, CallContext context) {
}
private record Node(FlowAction action, List<Connection> inputs) {
private record Connection(Node from, int fromPort, int toPort) {
}
}
}

View File

@ -16,16 +16,16 @@ public class ExecutionInfoEntry {
var compareSubnets = diff(source.getSubnets(), result.getSubnets()); var compareSubnets = diff(source.getSubnets(), result.getSubnets());
var compareDomains = diff(source.getDomains(), result.getDomains()); var compareDomains = diff(source.getDomains(), result.getDomains());
var compareASN = diff(source.getASN(), result.getASN()); var compareASN = diff(source.getASN(), result.getASN());
var compareAutoResolved = diff(source.getAutoResolvedDomains(), result.getAutoResolvedDomains()); var compareAutoResolved = diff(source.getResolveDomains(), result.getResolveDomains());
added.getSubnets().addAll(compareSubnets.added); added.getSubnets().addAll(compareSubnets.added);
added.getDomains().addAll(compareDomains.added); added.getDomains().addAll(compareDomains.added);
added.getASN().addAll(compareASN.added); added.getASN().addAll(compareASN.added);
added.getAutoResolvedDomains().addAll(compareAutoResolved.added); added.getResolveDomains().addAll(compareAutoResolved.added);
removed.getSubnets().addAll(compareSubnets.removed); removed.getSubnets().addAll(compareSubnets.removed);
removed.getDomains().addAll(compareDomains.removed); removed.getDomains().addAll(compareDomains.removed);
removed.getASN().addAll(compareASN.removed); removed.getASN().addAll(compareASN.removed);
removed.getAutoResolvedDomains().addAll(compareAutoResolved.removed); removed.getResolveDomains().addAll(compareAutoResolved.removed);
} }

View File

@ -1,7 +1,7 @@
package ru.kirillius.pf.sdn.api.flow; package ru.kirillius.pf.sdn.api.flow;
import lombok.*; import lombok.Getter;
import ru.kirillius.pf.sdn.api.Networking.NetworkScope; import ru.kirillius.pf.sdn.entity.ActionConfig;
import ru.kirillius.pf.sdn.entity.FlowConfig; import ru.kirillius.pf.sdn.entity.FlowConfig;
import java.util.ArrayList; import java.util.ArrayList;
@ -9,49 +9,67 @@ import java.util.List;
import java.util.concurrent.atomic.AtomicBoolean; 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;
import java.util.function.Function;
@RequiredArgsConstructor()
public class Flow { public class Flow {
public Flow(FlowConfig config, Function<ActionConfig, FlowAction> actionFactory) {
this.config = config;
graph = new ExecutionGraph(config.getActions(), config.getConnections(), actionFactory);
}
@Getter @Getter
private final FlowConfig config; private final FlowConfig config;
private final List<FlowAction> actions; private final ExecutionGraph graph;
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);
private AtomicInteger currentStep = new AtomicInteger(0); private AtomicInteger currentStep = new AtomicInteger(0);
private AtomicReference<FlowAction> currentAction = new AtomicReference<>(null);
public boolean isRunning() { public boolean isRunning() {
return running.get(); return running.get();
} }
public int getStepCount() { public int getStepCount() {
return actions.size(); return graph.size();
} }
public int getCurrentStep() { public int getCurrentStep() {
return currentStep.get(); return currentStep.get();
} }
public FlowAction getCurrentAction() {
return currentAction.get();
}
@Getter
private final List<ExecutionGraph.ExecutionResult> results = new ArrayList<>();
public void interrupt() { public void interrupt() {
interrupted.set(false); interrupted.set(false);
} }
private List<ExecutionInfoEntry> executionInfo = new ArrayList<>(); public void execute() throws InterruptedException {
public void execute() {
interrupted.set(false); interrupted.set(false);
running.set(true); running.set(true);
currentStep.set(0); currentStep.set(0);
results.clear();
try { try {
var scope = new AtomicReference<>(new NetworkScope()); var listener = graph.getStartExecutingEvent().add(action -> currentAction.set(action));
actions.forEach(action -> {
var iterator = graph.execute();
while (iterator.hasNext()) {
if (interrupted.get()) { if (interrupted.get()) {
throw new RuntimeException("Interrupted"); throw new InterruptedException();
} }
currentStep.incrementAndGet(); currentStep.incrementAndGet();
scope.set(action.apply(scope.get())); results.add(iterator.next());
}); }
graph.getStartExecutingEvent().remove(listener);
} finally { } finally {
interrupted.set(false); interrupted.set(false);
running.set(false); running.set(false);

View File

@ -1,22 +1,28 @@
package ru.kirillius.pf.sdn.api.flow; package ru.kirillius.pf.sdn.api.flow;
import lombok.Builder;
import lombok.Getter; import lombok.Getter;
import lombok.Setter; import lombok.Setter;
import lombok.SneakyThrows; import lombok.SneakyThrows;
import ru.kirillius.pf.sdn.api.Networking.NetworkScope; import ru.kirillius.pf.sdn.api.Networking.NetworkScope;
import ru.kirillius.pf.sdn.entity.ActionConfig;
import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.function.Function; import java.util.function.Function;
@Builder /**
public class FlowAction implements Function<NetworkScope, CallContext> { * Представляет собой вызов функции с параметрами
*/
public class FlowAction implements Function<List<NetworkScope>, CallContext> {
private FlowFunction function; private FlowFunction function;
@Getter
private final ActionConfig config;
public FlowAction(Class<? extends FlowFunction> functionClass, Map<String, String> properties) { public FlowAction(Class<? extends FlowFunction> functionClass, ActionConfig config) {
this.functionClass = functionClass; this.functionClass = functionClass;
this.properties = properties; this.properties = config.getProperties();
this.config = config;
instantiateFunction(); instantiateFunction();
} }
@ -25,20 +31,14 @@ public class FlowAction implements Function<NetworkScope, CallContext> {
function = functionClass.getConstructor().newInstance(); function = functionClass.getConstructor().newInstance();
} }
public void setFunctionClass(Class<? extends FlowFunction> functionClass) { private final Class<? extends FlowFunction> functionClass;
this.functionClass = functionClass;
instantiateFunction();
}
@Getter
private Class<? extends FlowFunction> functionClass;
@Getter @Getter
@Setter @Setter
private Map<String, String> properties; private Map<String, String> properties;
@Override @Override
public CallContext apply(NetworkScope networkScope) { public CallContext apply(List<NetworkScope> inputs) {
var context = new CallContext(networkScope, properties); var context = new CallContext(inputs, properties);
function.apply(context); function.apply(context);
return context; return context;
} }

View File

@ -2,8 +2,12 @@ package ru.kirillius.pf.sdn.api.flow;
import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor; import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor;
import java.util.List;
import java.util.Map; import java.util.Map;
/**
* Интерфейс, описывающий исполняемую функцию для Flow
*/
public interface FlowFunction { public interface FlowFunction {
void apply(CallContext context); void apply(CallContext context);
@ -12,4 +16,12 @@ public interface FlowFunction {
boolean hasInputs(); boolean hasInputs();
boolean hasOutputs(); boolean hasOutputs();
int getInputCount();
int getOutputCount();
List<String> getInputNames();
List<String> getOutputNames();
} }

View File

@ -4,5 +4,6 @@ public enum TriggerType {
Manual, Manual,
Interval, Interval,
OnSubscriptionUpdate, OnSubscriptionUpdate,
OnSubscriptionGroupUpdate,
OnStart OnStart
} }

View File

@ -0,0 +1,24 @@
package ru.kirillius.pf.sdn.entity;
import jakarta.persistence.*;
import lombok.Getter;
import lombok.Setter;
@Getter
@Setter
@Entity
public class ActionConnectionConfig {
@Id
@GeneratedValue(strategy = GenerationType.IDENTITY)
private Long id;
@ManyToOne(optional = false)
private ActionConfig from;
private int fromPort;
@ManyToOne(optional = false)
private ActionConfig to;
private int toPort;
}

View File

@ -7,7 +7,6 @@ import ru.kirillius.pf.sdn.api.flow.StartCondition;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.List; import java.util.List;
import java.util.Objects; import java.util.Objects;
import java.util.UUID;
@Entity @Entity
@Table(name = "pipeline_config") @Table(name = "pipeline_config")
@ -21,9 +20,6 @@ public class FlowConfig {
@GeneratedValue(strategy = GenerationType.IDENTITY) @GeneratedValue(strategy = GenerationType.IDENTITY)
private Long id; private Long id;
@GeneratedValue(strategy = GenerationType.UUID)
private UUID guid;
@Column(nullable = false, unique = true) @Column(nullable = false, unique = true)
private String name; private String name;
@ -43,4 +39,9 @@ public class FlowConfig {
@ManyToMany(fetch = FetchType.EAGER)//TODO проверить связность @ManyToMany(fetch = FetchType.EAGER)//TODO проверить связность
private List<ActionConfig> actions = new ArrayList<>(); private List<ActionConfig> actions = new ArrayList<>();
@ManyToMany(fetch = FetchType.EAGER)//TODO проверить связность
private List<ActionConnectionConfig> connections = new ArrayList<>();
} }

View File

@ -5,6 +5,7 @@ import lombok.*;
import ru.kirillius.pf.sdn.api.Networking.NetworkScope; import ru.kirillius.pf.sdn.api.Networking.NetworkScope;
import ru.kirillius.pf.sdn.api.Networking.Subscriptions.Subscription; import ru.kirillius.pf.sdn.api.Networking.Subscriptions.Subscription;
@Entity @Entity
@Table(name = "subscription_cache") @Table(name = "subscription_cache")
@Getter @Getter

View File

@ -0,0 +1,25 @@
package ru.kirillius.pf.sdn.entity;
import jakarta.persistence.*;
import lombok.*;
import java.util.Set;
@Entity
@Table(name = "subscription_group")
@Getter
@Setter
@NoArgsConstructor
@Builder
@AllArgsConstructor
public class SubscriptionGroup {
@Id
@GeneratedValue(strategy = GenerationType.IDENTITY)
private Long id;
@Column(nullable = false)
private String name;
@OneToMany(fetch = FetchType.EAGER)
private Set<SubscriptionCacheEntry> subscriptions;
}

View File

@ -3,17 +3,17 @@ 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.flow.CallContext; 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;
import java.util.Collections; import java.util.Collections;
import java.util.HashMap; import java.util.HashMap;
import java.util.List;
import java.util.Map; import java.util.Map;
@RequiredArgsConstructor @RequiredArgsConstructor
public class AutoResolveASN implements FlowFunction { public class AutoResolveASN extends CommonFunction {
private final AutonomousSystemCacheService autonomousSystemCacheService; private final AutonomousSystemCacheService autonomousSystemCacheService;
private final static String CLEAR = "clear"; private final static String CLEAR = "clear";
@ -34,15 +34,15 @@ public class AutoResolveASN implements FlowFunction {
@Override @Override
public void apply(CallContext context) { public void apply(CallContext context) {
var output = new NetworkScope(); var output = new NetworkScope();
output.add(context.getSource()); output.add(context.getSources().getFirst());
var properties = context.getProperties(); var properties = context.getProperties();
if (properties.containsKey(CLEAR) && properties.get(CLEAR).equals(Boolean.TRUE.toString())) { if (properties.containsKey(CLEAR) && properties.get(CLEAR).equals(Boolean.TRUE.toString())) {
output.getASN().clear(); output.getASN().clear();
} }
context.getSource().getASN().forEach(asn -> output.getSubnets().addAll(autonomousSystemCacheService.load(asn))); context.getSources().getFirst().getASN().forEach(asn -> output.getSubnets().addAll(autonomousSystemCacheService.load(asn)));
context.setOutput(output); context.setOutputs(List.of(output));
} }
@Override @Override
@ -51,12 +51,12 @@ public class AutoResolveASN implements FlowFunction {
} }
@Override @Override
public boolean hasInputs() { public List<String> getInputNames() {
return true; return List.of("input");
} }
@Override @Override
public boolean hasOutputs() { public List<String> getOutputNames() {
return true; return List.of("output");
} }
} }

View File

@ -4,19 +4,21 @@ 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.flow.CallContext; 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.DomainResolverService;
import ru.kirillius.pf.sdn.service.DomainUpdaterService; import ru.kirillius.pf.sdn.service.DomainUpdaterService;
import java.util.Collections; import java.util.Collections;
import java.util.HashMap; import java.util.HashMap;
import java.util.List;
import java.util.Map; import java.util.Map;
@RequiredArgsConstructor @RequiredArgsConstructor
public class AutoResolveDomains implements FlowFunction { public class AutoResolveDomains extends CommonFunction {
private final DomainUpdaterService domainCacheService; private final DomainUpdaterService domainCacheService;
private final DomainResolverService domainResolverService;
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;
@ -35,16 +37,16 @@ public class AutoResolveDomains implements FlowFunction {
@Override @Override
public void apply(CallContext context) { public void apply(CallContext context) {
var source = context.getSource(); var source = context.getSources().getFirst();
var properties = context.getProperties(); var properties = context.getProperties();
var output = new NetworkScope(); var output = new NetworkScope();
output.add(source); output.add(source);
if (properties.containsKey(CLEAR) && properties.get(CLEAR).equals("true")) { if (properties.containsKey(CLEAR) && properties.get(CLEAR).equals("true")) {
output.getAutoResolvedDomains().clear(); output.getResolveDomains().clear();
} }
source.getAutoResolvedDomains() source.getResolveDomains()
.forEach(domain -> domainCacheService.getActualAddresses(domain) .forEach(domain -> domainCacheService.getActualAddresses(domain)
.forEach(address -> output .forEach(address -> output
.getSubnets() .getSubnets()
@ -53,7 +55,7 @@ public class AutoResolveDomains implements FlowFunction {
) )
); );
context.setOutput(output); context.setOutputs(List.of(output));
} }
@Override @Override
@ -62,12 +64,12 @@ public class AutoResolveDomains implements FlowFunction {
} }
@Override @Override
public boolean hasInputs() { public List<String> getInputNames() {
return true; return List.of("input");
} }
@Override @Override
public boolean hasOutputs() { public List<String> getOutputNames() {
return true; return List.of("output");
} }
} }

View File

@ -0,0 +1,46 @@
package ru.kirillius.pf.sdn.flow;
import ru.kirillius.pf.sdn.api.flow.FlowFunction;
import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor;
import java.util.Collections;
import java.util.List;
import java.util.Map;
abstract class CommonFunction implements FlowFunction {
@Override
public final int getInputCount() {
return getInputNames().size();
}
@Override
public final int getOutputCount() {
return getOutputNames().size();
}
@Override
public Map<String, PropertyDescriptor> getProperties() {
return Collections.emptyMap();
}
@Override
public final boolean hasInputs() {
return !getInputNames().isEmpty();
}
@Override
public final boolean hasOutputs() {
return !getOutputNames().isEmpty();
}
@Override
public List<String> getInputNames() {
return Collections.emptyList();
}
@Override
public List<String> getOutputNames() {
return Collections.emptyList();
}
}

View File

@ -4,35 +4,22 @@ 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.flow.CallContext; 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.List;
import java.util.Map;
@Slf4j @Slf4j
@RequiredArgsConstructor @RequiredArgsConstructor
public class DebugOutput implements FlowFunction { public class DebugOutput extends CommonFunction {
private final ObjectMapper objectMapper; private final ObjectMapper objectMapper;
@Override @Override
public void apply(CallContext context) { public void apply(CallContext context) {
log.info(objectMapper.valueToTree(context.getSource()).toString()); log.info(objectMapper.valueToTree(context.getSources()).toString());
} }
@Override @Override
public Map<String, PropertyDescriptor> getProperties() { public List<String> getInputNames() {
return Collections.emptyMap(); return List.of("input");
}
@Override
public boolean hasInputs() {
return true;
}
@Override
public boolean hasOutputs() {
return false;
} }
} }

View File

@ -1,30 +1,12 @@
package ru.kirillius.pf.sdn.flow; package ru.kirillius.pf.sdn.flow;
import ru.kirillius.pf.sdn.api.flow.CallContext; 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 extends CommonFunction {
public class DummyFunction implements FlowFunction {
@Override @Override
public void apply(CallContext context) { public void apply(CallContext context) {
throw new UnsupportedOperationException("Not implemented"); 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;
}
} }

View File

@ -7,16 +7,17 @@ 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;
import java.util.List;
import java.util.Map; import java.util.Map;
@RequiredArgsConstructor @RequiredArgsConstructor
public class DummyPassthrough implements FlowFunction { public class DummyPassthrough implements FlowFunction {//TODO а зачем это надо вообще?
@Override @Override
public void apply(CallContext context) { public void apply(CallContext context) {
var output = new NetworkScope(); var output = new NetworkScope();
output.add(context.getSource()); output.add(context.getSources().getFirst());
context.setOutput(output); context.setOutputs(List.of(output));
} }
@Override @Override
@ -33,4 +34,23 @@ public class DummyPassthrough implements FlowFunction {
public boolean hasOutputs() { public boolean hasOutputs() {
return true; return true;
} }
@Override
public int getInputCount() {
return 1;
}
@Override
public int getOutputCount() {
return 1;
}
@Override
public List<String> getInputNames() {
return List.of("input");
}
@Override
public List<String> getOutputNames() {
return List.of("output");
}
} }

View File

@ -8,6 +8,7 @@ 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;
import java.util.List;
import java.util.Map; import java.util.Map;
@Slf4j @Slf4j
@ -16,13 +17,13 @@ public class DuplicatesFilter implements FlowFunction {
@Override @Override
public void apply(CallContext context) { public void apply(CallContext context) {
var source = context.getSource(); var source = context.getSources().getFirst();
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.setResolveDomains(source.getResolveDomains().stream().distinct().toList());
context.setOutput(output); context.setOutputs(List.of(output));
} }
@Override @Override
@ -40,4 +41,24 @@ public class DuplicatesFilter implements FlowFunction {
return true; return true;
} }
@Override
public int getInputCount() {
return 1;
}
@Override
public int getOutputCount() {
return 1;
}
@Override
public List<String> getInputNames() {
return List.of("input");
}
@Override
public List<String> getOutputNames() {
return List.of("output");
}
} }

View File

@ -5,39 +5,32 @@ 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.flow.CallContext; 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.List;
import java.util.Map;
@Slf4j @Slf4j
@RequiredArgsConstructor @RequiredArgsConstructor
public class MergeNeighbourSubnetFilter implements FlowFunction { public class MergeNeighbourSubnetFilter extends CommonFunction {
@Override
public List<String> getInputNames() {
return List.of("input");
}
@Override
public List<String> getOutputNames() {
return List.of("output");
}
@Override @Override
public void apply(CallContext context) { public void apply(CallContext context) {
var source = context.getSource(); var source = context.getSources().getFirst();
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.setResolveDomains(source.getResolveDomains());
} context.setOutputs(List.of(output));
@Override
public Map<String, PropertyDescriptor> getProperties() {
return Collections.emptyMap();
}
@Override
public boolean hasInputs() {
return false;
}
@Override
public boolean hasOutputs() {
return true;
} }
} }

View File

@ -5,18 +5,18 @@ 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.flow.CallContext; 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;
import java.util.Collections; import java.util.Collections;
import java.util.HashMap; import java.util.HashMap;
import java.util.List;
import java.util.Map; import java.util.Map;
@Slf4j @Slf4j
@RequiredArgsConstructor @RequiredArgsConstructor
public class MergeWithCoverageSubnetFilter implements FlowFunction { public class MergeWithCoverageSubnetFilter extends CommonFunction {
private final static String USAGE = "usage"; private final static String USAGE = "usage";
private final static Map<String, PropertyDescriptor> properties; private final static Map<String, PropertyDescriptor> properties;
@ -28,10 +28,20 @@ public class MergeWithCoverageSubnetFilter implements FlowFunction {
); );
} }
@Override
public List<String> getInputNames() {
return List.of("input");
}
@Override
public List<String> getOutputNames() {
return List.of("output");
}
@Override @Override
public void apply(CallContext context) { public void apply(CallContext context) {
var properties = context.getProperties(); var properties = context.getProperties();
var source = context.getSource(); var source = context.getSources().getFirst();
var usage = properties.containsKey(USAGE) ? Integer.parseInt(properties.get(USAGE)) : 75; var usage = properties.containsKey(USAGE) ? Integer.parseInt(properties.get(USAGE)) : 75;
if (usage < 51) { if (usage < 51) {
usage = 51; usage = 51;
@ -42,8 +52,8 @@ public class MergeWithCoverageSubnetFilter implements FlowFunction {
output.setSubnets(IPv4Util.mergeToLargerIfCovered(source.getSubnets(), usage)); output.setSubnets(IPv4Util.mergeToLargerIfCovered(source.getSubnets(), usage));
output.setDomains(source.getDomains()); output.setDomains(source.getDomains());
output.setASN(source.getASN()); output.setASN(source.getASN());
output.setAutoResolvedDomains(source.getAutoResolvedDomains()); output.setResolveDomains(source.getResolveDomains());
context.setOutput(output); context.setOutputs(List.of(output));
} }
@Override @Override
@ -51,14 +61,4 @@ public class MergeWithCoverageSubnetFilter implements FlowFunction {
return Collections.unmodifiableMap(properties); return Collections.unmodifiableMap(properties);
} }
@Override
public boolean hasInputs() {
return true;
}
@Override
public boolean hasOutputs() {
return true;
}
} }

View File

@ -5,40 +5,32 @@ 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.flow.CallContext; 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.List;
import java.util.Map;
@Slf4j @Slf4j
@RequiredArgsConstructor @RequiredArgsConstructor
public class OverlappedDomainFilter implements FlowFunction { public class OverlappedDomainFilter extends CommonFunction {
@Override
public List<String> getInputNames() {
return List.of("input");
}
@Override
public List<String> getOutputNames() {
return List.of("output");
}
@Override @Override
public void apply(CallContext context) { public void apply(CallContext context) {
var source = context.getSource(); var source = context.getSources().getFirst();
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.setResolveDomains(source.getResolveDomains());
context.setOutput(output); context.setOutputs(List.of(output));
}
@Override
public Map<String, PropertyDescriptor> getProperties() {
return Collections.emptyMap();
}
@Override
public boolean hasInputs() {
return true;
}
@Override
public boolean hasOutputs() {
return true;
} }
} }

View File

@ -5,40 +5,32 @@ 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.flow.CallContext; 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.List;
import java.util.Map;
@Slf4j @Slf4j
@RequiredArgsConstructor @RequiredArgsConstructor
public class OverlappedSubnetFilter implements FlowFunction { public class OverlappedSubnetFilter extends CommonFunction {
@Override
public List<String> getInputNames() {
return List.of("input");
}
@Override
public List<String> getOutputNames() {
return List.of("output");
}
@Override @Override
public void apply(CallContext context) { public void apply(CallContext context) {
var source = context.getSource(); var source = context.getSources().getFirst();
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.setResolveDomains(source.getResolveDomains());
context.setOutput(output); context.setOutputs(List.of(output));
}
@Override
public Map<String, PropertyDescriptor> getProperties() {
return Collections.emptyMap();
}
@Override
public boolean hasInputs() {
return true;
}
@Override
public boolean hasOutputs() {
return true;
} }
} }

View File

@ -5,85 +5,119 @@ 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.flow.CallContext; 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.PropertyDescriptor;
import ru.kirillius.pf.sdn.api.properties.PropertyType;
import java.util.Arrays; import java.util.ArrayList;
import java.util.Collections; import java.util.List;
import java.util.HashMap;
import java.util.Map;
import java.util.regex.Pattern;
@Slf4j @Slf4j
@RequiredArgsConstructor @RequiredArgsConstructor
public class ResourceFilter implements FlowFunction {
private final static String SUBNETS = "subnets"; public class ResourceFilter extends CommonFunction {
private final static String DOMAINS = "domains";
private final static String ASN = "asn";
private final static Map<String, PropertyDescriptor> properties; @Override
public List<String> getInputNames() {
return List.of("input", "filter");
}
static { @Override
properties = new HashMap<>(); public List<String> getOutputNames() {
properties.put(SUBNETS, PropertyDescriptor.builder().type(PropertyType.SUBNET).array(true).required(false).defaultValue(null).build()); return List.of("output", "filtered");
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 void apply(CallContext context) { public void apply(CallContext context) {
var output = new NetworkScope(); var output = new NetworkScope();
var source = context.getSource(); var filtered = new NetworkScope();
var properties = context.getProperties(); var source = context.getSources().get(0);
var filter = context.getSources().get(1);
if (properties.containsKey(SUBNETS)) { if (!filter.getSubnets().isEmpty()) {
var filtered = Arrays.stream(properties.get(SUBNETS).split(Pattern.quote(","))) filterSubnets(source, filter, output, filtered);
.filter(s -> !s.isBlank())
.map(IPv4Subnet::new)
.toList();
output.setSubnets(source.getSubnets().stream().filter(subnet -> !filtered.contains(subnet)).toList());
} else { } else {
output.setSubnets(source.getSubnets()); output.setSubnets(source.getSubnets());
} }
if (!filter.getASN().isEmpty()) {
if (properties.containsKey(DOMAINS)) { filterASN(source, filter, output, filtered);
var filtered = Arrays.stream(properties.get(DOMAINS).split(Pattern.quote(",")))
.filter(s -> !s.isBlank())
.toList();
output.setDomains(source.getDomains().stream().filter(domain -> !filtered.contains(domain)).toList());
} else {
output.setDomains(source.getDomains());
}
if (properties.containsKey(ASN)) {
var filtered = Arrays.stream(properties.get(ASN).split(Pattern.quote(",")))
.filter(s -> !s.isBlank())
.map(Integer::parseInt)
.toList();
output.setASN(source.getASN().stream().filter(asn -> !filtered.contains(asn)).toList());
} else { } else {
output.setASN(source.getASN()); output.setASN(source.getASN());
} }
if (!filter.getDomains().isEmpty()) {
filterDomains(source, filter, output, filtered);
} else {
output.setDomains(source.getDomains());
}
if (!filter.getResolveDomains().isEmpty()) {
filterResolveDomains(source, filter, output, filtered);
} else {
output.setResolveDomains(source.getResolveDomains());
}
context.setOutput(output); context.setOutputs(List.of(output, filtered));
} }
@Override private void filterSubnets(NetworkScope source, NetworkScope filter, NetworkScope output, NetworkScope filtered) {
public Map<String, PropertyDescriptor> getProperties() { var sourceSubnets = source.getSubnets();
return Collections.unmodifiableMap(properties); var filterSubnets = filter.getSubnets();
var outputSubnets = new ArrayList<IPv4Subnet>();
var filteredSubnets = new ArrayList<IPv4Subnet>();
sourceSubnets.forEach(subnet -> {
if (!filterSubnets.contains(subnet)) {
outputSubnets.add(subnet);
} else {
filteredSubnets.add(subnet);
}
});
output.setSubnets(outputSubnets);
filtered.setSubnets(filteredSubnets);
} }
@Override private void filterASN(NetworkScope source, NetworkScope filter, NetworkScope output, NetworkScope filtered) {
public boolean hasInputs() { var sourceASN = source.getASN();
return true; var filterASN = filter.getASN();
var outputASN = new ArrayList<Integer>();
var filteredASN = new ArrayList<Integer>();
sourceASN.forEach(asn -> {
if (!filterASN.contains(asn)) {
outputASN.add(asn);
} else {
filteredASN.add(asn);
}
});
output.setASN(outputASN);
filtered.setASN(filteredASN);
} }
@Override private void filterDomains(NetworkScope source, NetworkScope filter, NetworkScope output, NetworkScope filtered) {
public boolean hasOutputs() { var sourceDomains = source.getDomains();
return true; var filterDomains = filter.getDomains();
var outputDomains = new ArrayList<String>();
var filteredDomains = new ArrayList<String>();
sourceDomains.forEach(domain -> {
if (!filterDomains.contains(domain)) {
outputDomains.add(domain);
} else {
filteredDomains.add(domain);
}
});
output.setDomains(outputDomains);
filtered.setDomains(filteredDomains);
} }
private void filterResolveDomains(NetworkScope source, NetworkScope filter, NetworkScope output, NetworkScope filtered) {
var sourceDomains = source.getResolveDomains();
var filterDomains = filter.getResolveDomains();
var outputDomains = new ArrayList<String>();
var filteredDomains = new ArrayList<String>();
sourceDomains.forEach(domain -> {
if (!filterDomains.contains(domain)) {
outputDomains.add(domain);
} else {
filteredDomains.add(domain);
}
});
output.setResolveDomains(outputDomains);
filtered.setResolveDomains(filteredDomains);
}
} }

View File

@ -0,0 +1,107 @@
package ru.kirillius.pf.sdn.flow;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import ru.kirillius.pf.sdn.api.Networking.IPv4Subnet;
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.IntegerConstraint;
import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor;
import ru.kirillius.pf.sdn.api.properties.PropertyType;
import java.util.*;
import java.util.regex.Pattern;
@Slf4j
@RequiredArgsConstructor
@Deprecated
public class ResourceFilterObsolete 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
public int getInputCount() {
return 1;
}
@Override
public int getOutputCount() {
return 1;
}
@Override
public List<String> getInputNames() {
return List.of("input");
}
@Override
public List<String> getOutputNames() {
return List.of("output");
}
@Override
public void apply(CallContext context) {
var output = new NetworkScope();
var source = context.getSources().getFirst();
var properties = context.getProperties();
if (properties.containsKey(SUBNETS)) {
var filtered = Arrays.stream(properties.get(SUBNETS).split(Pattern.quote(",")))
.filter(s -> !s.isBlank())
.map(IPv4Subnet::new)
.toList();
output.setSubnets(source.getSubnets().stream().filter(subnet -> !filtered.contains(subnet)).toList());
} else {
output.setSubnets(source.getSubnets());
}
if (properties.containsKey(DOMAINS)) {
var filtered = Arrays.stream(properties.get(DOMAINS).split(Pattern.quote(",")))
.filter(s -> !s.isBlank())
.toList();
output.setDomains(source.getDomains().stream().filter(domain -> !filtered.contains(domain)).toList());
} else {
output.setDomains(source.getDomains());
}
if (properties.containsKey(ASN)) {
var filtered = Arrays.stream(properties.get(ASN).split(Pattern.quote(",")))
.filter(s -> !s.isBlank())
.map(Integer::parseInt)
.toList();
output.setASN(source.getASN().stream().filter(asn -> !filtered.contains(asn)).toList());
} else {
output.setASN(source.getASN());
}
context.setOutputs(List.of(output));
}
@Override
public Map<String, PropertyDescriptor> getProperties() {
return Collections.unmodifiableMap(properties);
}
@Override
public boolean hasInputs() {
return true;
}
@Override
public boolean hasOutputs() {
return true;
}
}

View File

@ -0,0 +1,71 @@
package ru.kirillius.pf.sdn.flow;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import ru.kirillius.pf.sdn.api.Networking.IPv4Subnet;
import ru.kirillius.pf.sdn.api.Networking.NetworkScope;
import ru.kirillius.pf.sdn.api.flow.CallContext;
import ru.kirillius.pf.sdn.api.properties.IntegerConstraint;
import ru.kirillius.pf.sdn.api.properties.PropertyDescriptor;
import ru.kirillius.pf.sdn.api.properties.PropertyType;
import java.util.*;
import java.util.regex.Pattern;
@Slf4j
@RequiredArgsConstructor
public class StaticResourceInput extends CommonFunction {
private final static String SUBNETS = "subnets";
private final static String DOMAINS = "domains";
private final static String RESOLVE = "resolve";
private final static String ASN = "asn";
private final static Map<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(RESOLVE, PropertyDescriptor.builder().type(PropertyType.STRING).array(true).required(false).defaultValue(null).build());
properties.put(ASN, PropertyDescriptor.builder().type(PropertyType.INTEGER).array(true).required(false).defaultValue(null).build().addConstraint(new IntegerConstraint(1, 99999))); //TODO какой max asn?);
}
@Override
public List<String> getOutputNames() {
return List.of("output");
}
@Override
public void apply(CallContext context) {
var output = new NetworkScope();
var properties = context.getProperties();
if (properties.containsKey(SUBNETS)) {
output.setSubnets(Arrays.stream(properties.get(SUBNETS).split(Pattern.quote(","))).filter(s -> !s.isBlank()).map(IPv4Subnet::new).toList());
}
if (properties.containsKey(DOMAINS)) {
output.setDomains(Arrays.stream(properties.get(DOMAINS).split(Pattern.quote(","))).filter(s -> !s.isBlank()).toList());
}
if (properties.containsKey(RESOLVE)) {
output.setResolveDomains(Arrays.stream(properties.get(RESOLVE).split(Pattern.quote(","))).filter(s -> !s.isBlank()).toList());
}
if (properties.containsKey(ASN)) {
output.setASN(Arrays.stream(properties.get(ASN).split(Pattern.quote(","))).filter(s -> !s.isBlank()).map(Integer::parseInt).toList());
}
context.setOutputs(List.of(output));
}
@Override
public Map<String, PropertyDescriptor> getProperties() {
return Collections.unmodifiableMap(properties);
}
}

View File

@ -3,19 +3,15 @@ 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.flow.CallContext; 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;
import java.util.Arrays; import java.util.*;
import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
import java.util.regex.Pattern; import java.util.regex.Pattern;
@RequiredArgsConstructor @RequiredArgsConstructor
public class SubscriptionInput implements FlowFunction { public class SubscriptionInput extends CommonFunction {
private final SubscriptionService subscriptionService; private final SubscriptionService subscriptionService;
@ -49,12 +45,7 @@ public class SubscriptionInput implements FlowFunction {
} }
@Override @Override
public boolean hasInputs() { public List<String> getOutputNames() {
return false; return List.of("output");
}
@Override
public boolean hasOutputs() {
return true;
} }
} }

View File

@ -0,0 +1,29 @@
package ru.kirillius.pf.sdn.flow;
import lombok.RequiredArgsConstructor;
import ru.kirillius.pf.sdn.api.Networking.NetworkScope;
import ru.kirillius.pf.sdn.api.flow.CallContext;
import java.util.List;
@RequiredArgsConstructor
public class Summator extends CommonFunction {
@Override
public List<String> getInputNames() {
return List.of("A", "B");
}
@Override
public List<String> getOutputNames() {
return List.of("output");
}
@Override
public void apply(CallContext context) {
var output = new NetworkScope();
context.getSources().forEach(output::add);
context.setOutputs(List.of(output));
}
}

View File

@ -3,9 +3,7 @@ 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.FlowConfig; import ru.kirillius.pf.sdn.entity.FlowConfig;
import java.util.UUID;
public interface PipelineConfigRepository extends JpaRepository<FlowConfig, Long> { public interface PipelineConfigRepository extends JpaRepository<FlowConfig, Long> {
boolean existsByGuid(UUID guid);
} }

View File

@ -0,0 +1,8 @@
package ru.kirillius.pf.sdn.repository;
import org.springframework.data.jpa.repository.JpaRepository;
import ru.kirillius.pf.sdn.entity.SubscriptionGroup;
public interface SubscriptionSetRepository extends JpaRepository<SubscriptionGroup, Long> {
}

View File

@ -53,7 +53,7 @@ public class DomainUpdaterService {
protected void doWork() { protected void doWork() {
subscriptionService.getSubscriptions().forEach(subscription -> { subscriptionService.getSubscriptions().forEach(subscription -> {
var scope = subscription.getScope(); var scope = subscription.getScope();
if (!scope.isResolveDomains()) { if (scope.getResolveDomains().isEmpty()) {
return; return;
} }
scope.getDomains().forEach(domain -> { scope.getDomains().forEach(domain -> {

View File

@ -17,13 +17,21 @@ import java.util.concurrent.ExecutorService;
import java.util.concurrent.Future; import java.util.concurrent.Future;
import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.atomic.AtomicReference;
/**
* Запускает и выполняет Flow в порядке очереди
*/
@Slf4j @Slf4j
@Service @Service
@RequiredArgsConstructor @RequiredArgsConstructor
public class FlowExecutorService { public class FlowExecutorService {
private final ExecutorService executor; private final ExecutorService executor;
private final Queue<Flow> executionQueue = new ConcurrentLinkedQueue<>(); private final Queue<Flow> executionQueue = new ConcurrentLinkedQueue<>();
private final AtomicReference<Flow> currentPipeline = new AtomicReference<>();
public Flow getCurrentFlow() {
return currentFlow.get();
}
private final AtomicReference<Flow> currentFlow = new AtomicReference<>();
@Getter @Getter
private final EventHandler<Flow> onExecuted = new ConcurrentEventHandler<>(); private final EventHandler<Flow> onExecuted = new ConcurrentEventHandler<>();
@ -61,14 +69,14 @@ public class FlowExecutorService {
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); currentFlow.set(pipeline);
try { try {
pipeline.execute(); pipeline.execute();
onExecuted.invoke(pipeline); 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); currentFlow.set(null);
} }
Thread.sleep(100L); Thread.sleep(100L);
Thread.yield(); Thread.yield();

View File

@ -7,17 +7,19 @@ 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.flow.Flow;
import ru.kirillius.pf.sdn.api.flow.FlowAction; import ru.kirillius.pf.sdn.api.flow.FlowAction;
import ru.kirillius.pf.sdn.api.flow.FlowFunction; import ru.kirillius.pf.sdn.api.flow.FlowFunction;
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.FlowConfig; import ru.kirillius.pf.sdn.entity.FlowConfig;
import ru.kirillius.pf.sdn.entity.SubscriptionGroup;
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;
import java.util.Arrays; import java.util.Arrays;
import java.util.Map; import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutorService; import java.util.concurrent.ExecutorService;
import java.util.regex.Pattern; import java.util.regex.Pattern;
@ -32,10 +34,11 @@ public class FlowService {
private final SubscriptionService subscriptionService; private final SubscriptionService subscriptionService;
private final ExecutorService executorService; private final ExecutorService executorService;
private final Map<FlowConfig, Flow> pipelines = new ConcurrentHashMap<>(); private final Map<FlowConfig, Flow> flows = 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;
private EventListener<SubscriptionGroup> subscriptionSetUpdateListener;
public void registerFunction(Class<? extends FlowFunction> functionClass, String id) { public void registerFunction(Class<? extends FlowFunction> functionClass, String id) {
if (functions.containsKey(id)) { if (functions.containsKey(id)) {
@ -44,13 +47,17 @@ public class FlowService {
functions.put(id, functionClass); functions.put(id, functionClass);
} }
public Set<String> getFunctionNames() {
return functions.keySet();
}
public void trigger(long id) { public void trigger(long id) {
pipelines.keySet().stream().filter(config -> config.getId().equals(id)).findFirst().ifPresent(this::trigger); flows.keySet().stream().filter(config -> config.getId().equals(id)).findFirst().ifPresent(this::trigger);
} }
public void trigger(FlowConfig flowConfig) { public void trigger(FlowConfig flowConfig) {
if (pipelines.containsKey(flowConfig)) { if (flows.containsKey(flowConfig)) {
trigger(pipelines.get(flowConfig)); trigger(flows.get(flowConfig));
} }
} }
@ -66,27 +73,30 @@ public class FlowService {
} }
flowExecutorService.triggerExecute(pipeline); flowExecutorService.triggerExecute(pipeline);
}); });
} }
public void update(FlowConfig flowConfig) { public void loadFlow(Flow flow) {
if (configRepository.existsByGuid(flowConfig.getGuid())) { loadFlow(flow.getConfig());
load(flowConfig); }
public void loadFlow(FlowConfig flowConfig) {
if (configRepository.existsById(flowConfig.getId())) {
reloadFlowInternal(flowConfig);
} else { } else {
flowExecutorService.cancel(pipelines.get(flowConfig)); flowExecutorService.cancel(flows.get(flowConfig));
pipelines.remove(flowConfig); flows.remove(flowConfig);
} }
} }
private void load(FlowConfig flowConfig) { private void reloadFlowInternal(FlowConfig flowConfig) {
if (pipelines.containsKey(flowConfig)) { if (flows.containsKey(flowConfig)) {
flowExecutorService.cancel(pipelines.get(flowConfig)); flowExecutorService.cancel(flows.get(flowConfig));
} }
pipelines.put(flowConfig, new Flow(flowConfig, flowConfig.getActions().stream().map(this::buildAction).toList())); flows.put(flowConfig, new Flow(flowConfig, this::buildAction));
} }
private FlowAction buildAction(ActionConfig config) { private FlowAction buildAction(ActionConfig config) {
return new FlowAction(functions.getOrDefault(config.getFunctionId(), DummyFunction.class), config.getProperties()); return new FlowAction(functions.getOrDefault(config.getFunctionId(), DummyFunction.class), config);
} }
private void checkForIntervalExecution(Flow pipeline) { private void checkForIntervalExecution(Flow pipeline) {
@ -126,10 +136,11 @@ public class FlowService {
@PostConstruct @PostConstruct
private void initialize() { private void initialize() {
subscriptionUpdateListener = subscriptionService.getUpdateEvent().add(this::subscriptionUpdate); subscriptionUpdateListener = subscriptionService.getUpdateEvent().add(this::subscriptionUpdate);
subscriptionSetUpdateListener = subscriptionService.getSetUpdateEvent().add(this::subscriptionSetUpdate);
executedListener = flowExecutorService.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::reloadFlowInternal);
pipelines.values() flows.values()
.stream() .stream()
.filter(p -> p .filter(p -> p
.getConfig() .getConfig()
@ -139,9 +150,35 @@ public class FlowService {
.forEach(flowExecutorService::triggerExecute); .forEach(flowExecutorService::triggerExecute);
} }
private void subscriptionSetUpdate(SubscriptionGroup setEntry) {
flows.values().forEach(pipeline -> {
var matchedConditions = pipeline.getConfig().getStartConditions().stream().filter(c -> c.getTrigger() == TriggerType.OnSubscriptionGroupUpdate).toList();
if (matchedConditions.isEmpty()) {
return;
}
matchedConditions.forEach(c -> {
if (c.isDontStartIfRunning() && pipeline.isRunning()) {
return;
}
var properties = c.getProperties();
var names = properties.getOrDefault("names", null);
if (names == null || names.isEmpty()) {
flowExecutorService.triggerExecute(pipeline);
return;
}
if (Arrays.stream(names.split(Pattern.quote(","))).anyMatch(s -> setEntry.getName().equals(s))) {
flowExecutorService.triggerExecute(pipeline);
}
});
});
}
private void subscriptionUpdate(Subscription subscription) { private void subscriptionUpdate(Subscription subscription) {
pipelines.values().forEach(pipeline -> { flows.values().forEach(pipeline -> {
var matchedConditions = pipeline.getConfig().getStartConditions().stream().filter(c -> c.getTrigger() == TriggerType.OnSubscriptionUpdate).toList(); var matchedConditions = pipeline.getConfig().getStartConditions().stream().filter(c -> c.getTrigger() == TriggerType.OnSubscriptionUpdate).toList();
if (matchedConditions.isEmpty()) { if (matchedConditions.isEmpty()) {
return; return;
@ -177,6 +214,10 @@ public class FlowService {
subscriptionService.getUpdateEvent().remove(subscriptionUpdateListener); subscriptionService.getUpdateEvent().remove(subscriptionUpdateListener);
subscriptionUpdateListener = null; subscriptionUpdateListener = null;
} }
if (subscriptionSetUpdateListener != null) {
subscriptionService.getSetUpdateEvent().remove(subscriptionSetUpdateListener);
subscriptionSetUpdateListener = null;
}
} }
} }

View File

@ -1,6 +1,5 @@
package ru.kirillius.pf.sdn.service; package ru.kirillius.pf.sdn.service;
import jakarta.annotation.PostConstruct;
import jakarta.annotation.PreDestroy; import jakarta.annotation.PreDestroy;
import lombok.Getter; import lombok.Getter;
import lombok.RequiredArgsConstructor; import lombok.RequiredArgsConstructor;
@ -14,7 +13,9 @@ import ru.kirillius.pf.sdn.api.Networking.Subscriptions.CacheFallbackProviderPro
import ru.kirillius.pf.sdn.api.Networking.Subscriptions.Subscription; import ru.kirillius.pf.sdn.api.Networking.Subscriptions.Subscription;
import ru.kirillius.pf.sdn.api.Networking.Subscriptions.SubscriptionProvider; import ru.kirillius.pf.sdn.api.Networking.Subscriptions.SubscriptionProvider;
import ru.kirillius.pf.sdn.api.Networking.Subscriptions.SubscriptionProviderProtocol; import ru.kirillius.pf.sdn.api.Networking.Subscriptions.SubscriptionProviderProtocol;
import ru.kirillius.pf.sdn.entity.SubscriptionGroup;
import ru.kirillius.pf.sdn.repository.SubscriptionProviderConfigRepository; import ru.kirillius.pf.sdn.repository.SubscriptionProviderConfigRepository;
import ru.kirillius.pf.sdn.repository.SubscriptionSetRepository;
import java.time.Duration; import java.time.Duration;
import java.time.Instant; import java.time.Instant;
@ -41,8 +42,13 @@ public class SubscriptionService {
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; private final ApplicationContext context;
private final SubscriptionSetRepository subscriptionSetRepository;
@Getter @Getter
private final EventHandler<Subscription> updateEvent = new ConcurrentEventHandler<>(); private final EventHandler<Subscription> updateEvent = new ConcurrentEventHandler<>();
@Getter
private final EventHandler<SubscriptionGroup> setUpdateEvent = new ConcurrentEventHandler<>();
private Future<?> worker; private Future<?> worker;
public List<Subscription> getSubscriptions() { public List<Subscription> getSubscriptions() {
@ -84,6 +90,8 @@ public class SubscriptionService {
var retrieved = new ArrayList<Subscription>(); var retrieved = new ArrayList<Subscription>();
var updated = new ArrayList<Subscription>(); var updated = new ArrayList<Subscription>();
var sets = subscriptionSetRepository.findAll();
providers.values().forEach(provider -> { providers.values().forEach(provider -> {
updated.addAll(provider.update()); updated.addAll(provider.update());
retrieved.addAll(provider.getSubscriptions()); retrieved.addAll(provider.getSubscriptions());
@ -95,16 +103,25 @@ public class SubscriptionService {
} }
updated.forEach(s -> { updated.forEach(s -> {
try { invokeEvent(updateEvent, s);
updateEvent.invoke(s); sets.forEach(set -> {
} catch (Exception e) { if (set.getSubscriptions().contains(s)) { //TODO suspicious call???
log.error("Failed to invoke subscription update event because of error {}:{}", e.getClass().getSimpleName(), e.getMessage()); invokeEvent(setUpdateEvent, set);
} }
});
}); });
} }
@PostConstruct private <T> void invokeEvent(EventHandler<T> eventHandler, T data) {
private void initialize() { try {
eventHandler.invoke(data);
} catch (Exception e) {
log.error("Failed to invoke subscription update event because of error {}:{}", e.getClass().getSimpleName(), e.getMessage());
}
}
private void start() {
reloadProviders(); reloadProviders();
reloadSubscriptions(); reloadSubscriptions();