This is an automated email from the ASF dual-hosted git repository.
riemer pushed a commit to branch
2836-cannot-deploy-a-pipeline-with-trend-data-processor
in repository https://gitbox.apache.org/repos/asf/streampipes.git
The following commit(s) were added to
refs/heads/2836-cannot-deploy-a-pipeline-with-trend-data-processor by this push:
new 692eb1b625 fix(#2836): Sanitize blank space in Siddhi queries
692eb1b625 is described below
commit 692eb1b62597499176bc32b6afde7bf971651803
Author: Dominik Riemer <[email protected]>
AuthorDate: Tue May 7 23:36:45 2024 +0200
fix(#2836): Sanitize blank space in Siddhi queries
---
.../streampipes-processors-filters-siddhi/pom.xml | 5 +
.../siddhi/trend/TestTrendProcessor.java | 292 +++++++++++----------
.../engine/generator/SiddhiAppGenerator.java | 2 +-
.../wrapper/siddhi/model/EventPropertyDef.java | 8 +-
.../wrapper/siddhi/utils/SiddhiUtils.java | 10 +-
5 files changed, 176 insertions(+), 141 deletions(-)
diff --git
a/streampipes-extensions/streampipes-processors-filters-siddhi/pom.xml
b/streampipes-extensions/streampipes-processors-filters-siddhi/pom.xml
index 884581a1cd..5e7493d707 100644
--- a/streampipes-extensions/streampipes-processors-filters-siddhi/pom.xml
+++ b/streampipes-extensions/streampipes-processors-filters-siddhi/pom.xml
@@ -52,6 +52,11 @@
<artifactId>junit-jupiter-api</artifactId>
<scope>test</scope>
</dependency>
+ <dependency>
+ <groupId>org.junit.jupiter</groupId>
+ <artifactId>junit-jupiter-params</artifactId>
+ <scope>test</scope>
+ </dependency>
<dependency>
<groupId>org.apache.streampipes</groupId>
<artifactId>streampipes-test-utils</artifactId>
diff --git
a/streampipes-extensions/streampipes-processors-filters-siddhi/src/test/java/org/apache/streampipes/processors/siddhi/trend/TestTrendProcessor.java
b/streampipes-extensions/streampipes-processors-filters-siddhi/src/test/java/org/apache/streampipes/processors/siddhi/trend/TestTrendProcessor.java
index a77552bfbb..9b96316a35 100644
---
a/streampipes-extensions/streampipes-processors-filters-siddhi/src/test/java/org/apache/streampipes/processors/siddhi/trend/TestTrendProcessor.java
+++
b/streampipes-extensions/streampipes-processors-filters-siddhi/src/test/java/org/apache/streampipes/processors/siddhi/trend/TestTrendProcessor.java
@@ -17,141 +17,161 @@
*/
package org.apache.streampipes.processors.siddhi.trend;
-//@RunWith(Parameterized.class)
+
+import org.apache.streampipes.model.graph.DataProcessorDescription;
+import org.apache.streampipes.model.graph.DataProcessorInvocation;
+import org.apache.streampipes.model.output.CustomOutputStrategy;
+import org.apache.streampipes.model.runtime.Event;
+import org.apache.streampipes.model.runtime.EventFactory;
+import org.apache.streampipes.model.runtime.SchemaInfo;
+import org.apache.streampipes.model.runtime.SourceInfo;
+import org.apache.streampipes.sdk.helpers.Tuple2;
+import org.apache.streampipes.test.generator.EventStreamGenerator;
+import org.apache.streampipes.test.generator.InvocationGraphGenerator;
+import org.apache.streampipes.test.generator.grounding.EventGroundingGenerator;
+import
org.apache.streampipes.wrapper.params.generator.DataProcessorParameterGenerator;
+import
org.apache.streampipes.wrapper.siddhi.engine.callback.SiddhiDebugCallback;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.Arguments;
+import org.junit.jupiter.params.provider.MethodSource;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.stream.Collectors;
+import java.util.stream.Stream;
+
+
public class TestTrendProcessor {
-//
-// private static final Logger LOG =
LoggerFactory.getLogger(TrendProcessor.class);
-//
-// @org.junit.runners.Parameterized.Parameters
-// public static Iterable<Object[]> data() {
-// return Arrays.asList(new Object[][]{
-// {1, 100, TrendOperator.INCREASE, Arrays.asList(new Tuple2<>(100, 1),
-// new Tuple2<>(100, 2)), 1},
-// {1, 100, TrendOperator.INCREASE, Arrays.asList(
-// new Tuple2<>(100, 1),
-// new Tuple2<>(100, 1),
-// new Tuple2<>(100, 2))
-// , 1},
-// {1, 100, TrendOperator.DECREASE, Arrays.asList(
-// new Tuple2<>(100, 1),
-// new Tuple2<>(100, 1),
-// new Tuple2<>(100, 0))
-// , 1},
-// {1, 100, TrendOperator.INCREASE, Arrays.asList(
-// new Tuple2<>(100, 1),
-// new Tuple2<>(100, 1),
-// new Tuple2<>(100, 2),
-// new Tuple2<>(100, 1),
-// new Tuple2<>(100, 2))
-// , 2},
-// {1, 200, TrendOperator.INCREASE, Arrays.asList(
-// new Tuple2<>(100, 1),
-// new Tuple2<>(100, 1),
-// new Tuple2<>(100, 2),
-// new Tuple2<>(100, 1),
-// new Tuple2<>(100, 2))
-// , 0},
-//
-// });
-// }
-//
-// @org.junit.runners.Parameterized.Parameter
-// public Integer timeWindow;
-//
-// @org.junit.runners.Parameterized.Parameter(1)
-// public Integer increase;
-//
-// @org.junit.runners.Parameterized.Parameter(2)
-// public TrendOperator trendOperator;
-//
-// @org.junit.runners.Parameterized.Parameter(3)
-// public List<Tuple2<Integer, Integer>> eventSettings;
-//
-// @org.junit.runners.Parameterized.Parameter(4)
-// public Integer expectedMatchCount;
-//
-// @Test
-// public void testTrend() {
-// final Integer[] actualMatchCount = {0};
-// DataProcessorDescription originalGraph = new
TrendProcessor().declareModel();
-//
originalGraph.setSupportedGrounding(EventGroundingGenerator.makeDummyGrounding());
-//
-// DataProcessorInvocation graph =
-// InvocationGraphGenerator.makeEmptyInvocation(originalGraph);
-//
-// graph.setInputStreams(Collections
-// .singletonList(EventStreamGenerator
-//
.makeStreamWithProperties(Collections.singletonList("randomValue"))));
-//
-// graph.setOutputStrategies(
-// graph.getOutputStrategies()
-// .stream()
-// .filter(o -> o instanceof CustomOutputStrategy)
-// .peek(o -> ((CustomOutputStrategy)
o).setSelectedPropertyKeys(Arrays.asList("s0::randomValue")))
-// .collect(Collectors.toList())
-// );
-//
-//
graph.setOutputStream(EventStreamGenerator.makeStreamWithProperties(Collections.singletonList("randomValue")));
-//
-//
graph.getOutputStream().getEventGrounding().getTransportProtocol().getTopicDefinition()
-// .setActualTopicName("output-topic");
-//
-// var visitor = new TrendConfigurationVisitor("s0::randomValue",
-// trendOperator,
-// increase,
-// timeWindow);
-// graph.getStaticProperties().forEach(sp -> sp.accept(visitor));
-//
-// var processorParams = new
DataProcessorParameterGenerator().makeParameters(graph);
-//
-// SiddhiDebugCallback callback = new SiddhiDebugCallback() {
-// @Override
-// public void onEvent(io.siddhi.core.event.Event event) {
-// actualMatchCount[0]++;
-// }
-//
-// @Override
-// public void onEvent(List<io.siddhi.core.event.Event> events) {
-//
-// }
-// };
-//
-// TrendProcessor trend = new TrendProcessor(callback);
-// trend.onPipelineStarted(processorParams, null, null);
-//
-// sendEvents(trend);
-// LOG.info("Expected match count is {}", expectedMatchCount);
-// LOG.info("Actual match count is {}", actualMatchCount[0]);
-// assertEquals(expectedMatchCount, actualMatchCount[0]);
-// }
-//
-// private void sendEvents(TrendProcessor trend) {
-// List<Tuple2<Integer, Event>> events = makeEvents();
-// for (Tuple2<Integer, Event> event : events) {
-// LOG.info("Sending event with value " +
event.v.getFieldBySelector("s0::randomValue"));
-// trend.onEvent(event.v, null);
-// try {
-// Thread.sleep(event.k);
-// } catch (InterruptedException e) {
-// e.printStackTrace();
-// }
-// }
-// }
-//
-// private List<Tuple2<Integer, Event>> makeEvents() {
-// List<Tuple2<Integer, Event>> events = new ArrayList<>();
-// for (Tuple2<Integer, Integer> eventSetting : eventSettings) {
-// events.add(makeEvent(eventSetting.k, eventSetting.v));
-// }
-// return events;
-// }
-//
-// private Tuple2<Integer, Event> makeEvent(Integer timeout, Integer value) {
-// Map<String, Object> map = new HashMap<>();
-// map.put("randomValue", value);
-// return new Tuple2<>(timeout, EventFactory.fromMap(map, new
SourceInfo("test"
-// + "-topic", "s0"),
-// new SchemaInfo(null,
-// new ArrayList<>())));
-// }
+
+ private static final Logger LOG =
LoggerFactory.getLogger(TrendProcessor.class);
+
+ static Stream<Arguments> data() {
+ return Stream.of(
+ Arguments.of(1, 100, TrendOperator.INCREASE, Arrays.asList(new
Tuple2<>(100, 1),
+ new Tuple2<>(100, 2)), 1),
+ Arguments.of(1, 100, TrendOperator.INCREASE, Arrays.asList(
+ new Tuple2<>(100, 1),
+ new Tuple2<>(100, 1),
+ new Tuple2<>(100, 2))
+ , 1),
+ Arguments.of(1, 100, TrendOperator.DECREASE, Arrays.asList(
+ new Tuple2<>(100, 1),
+ new Tuple2<>(100, 1),
+ new Tuple2<>(100, 0))
+ , 1),
+ Arguments.of(1, 100, TrendOperator.INCREASE, Arrays.asList(
+ new Tuple2<>(100, 1),
+ new Tuple2<>(100, 1),
+ new Tuple2<>(100, 2),
+ new Tuple2<>(100, 1),
+ new Tuple2<>(100, 2))
+ , 2),
+ Arguments.of(1, 200, TrendOperator.INCREASE, Arrays.asList(
+ new Tuple2<>(100, 1),
+ new Tuple2<>(100, 1),
+ new Tuple2<>(100, 2),
+ new Tuple2<>(100, 1),
+ new Tuple2<>(100, 2))
+ , 0)
+ );
+ }
+
+ @ParameterizedTest
+ @MethodSource("data")
+ public void testTrend(int timeWindow,
+ int increase,
+ TrendOperator trendOperator,
+ List<Tuple2<Integer, Integer>> eventSettings,
+ int expectedMatchCount) {
+ final Integer[] actualMatchCount = {0};
+ DataProcessorDescription originalGraph = new
TrendProcessor().declareModel();
+
originalGraph.setSupportedGrounding(EventGroundingGenerator.makeDummyGrounding());
+
+ DataProcessorInvocation graph =
+ InvocationGraphGenerator.makeEmptyInvocation(originalGraph);
+
+ graph.setInputStreams(Collections
+ .singletonList(EventStreamGenerator
+
.makeStreamWithProperties(Collections.singletonList("randomValue"))));
+
+ graph.setOutputStrategies(
+ graph.getOutputStrategies()
+ .stream()
+ .filter(o -> o instanceof CustomOutputStrategy)
+ .peek(o -> ((CustomOutputStrategy)
o).setSelectedPropertyKeys(List.of("s0::randomValue")))
+ .collect(Collectors.toList())
+ );
+
+
graph.setOutputStream(EventStreamGenerator.makeStreamWithProperties(Collections.singletonList("randomValue")));
+
+
graph.getOutputStream().getEventGrounding().getTransportProtocol().getTopicDefinition()
+ .setActualTopicName("output-topic");
+
+ var visitor = new TrendConfigurationVisitor("s0::randomValue",
+ trendOperator,
+ increase,
+ timeWindow);
+ graph.getStaticProperties().forEach(sp -> sp.accept(visitor));
+
+ var processorParams = new
DataProcessorParameterGenerator().makeParameters(graph);
+
+ SiddhiDebugCallback callback = new SiddhiDebugCallback() {
+ @Override
+ public void onEvent(io.siddhi.core.event.Event event) {
+ actualMatchCount[0]++;
+ }
+
+ @Override
+ public void onEvent(List<io.siddhi.core.event.Event> events) {
+
+ }
+ };
+
+ TrendProcessor trend = new TrendProcessor(callback);
+ trend.onPipelineStarted(processorParams, null, null);
+
+ sendEvents(trend, eventSettings);
+ LOG.info("Expected match count is {}", expectedMatchCount);
+ LOG.info("Actual match count is {}", actualMatchCount[0]);
+ Assertions.assertEquals(expectedMatchCount, actualMatchCount[0]);
+ }
+
+ private void sendEvents(TrendProcessor trend,
+ List<Tuple2<Integer, Integer>> eventSettings) {
+ List<Tuple2<Integer, Event>> events = makeEvents(eventSettings);
+ for (Tuple2<Integer, Event> event : events) {
+ LOG.info("Sending event with value " +
event.v.getFieldBySelector("s0::randomValue"));
+ trend.onEvent(event.v, null);
+ try {
+ Thread.sleep(event.k);
+ } catch (InterruptedException e) {
+ e.printStackTrace();
+ }
+ }
+ }
+
+ private List<Tuple2<Integer, Event>> makeEvents(List<Tuple2<Integer,
Integer>> eventSettings) {
+ List<Tuple2<Integer, Event>> events = new ArrayList<>();
+ for (Tuple2<Integer, Integer> eventSetting : eventSettings) {
+ events.add(makeEvent(eventSetting.k, eventSetting.v));
+ }
+ return events;
+ }
+
+ private Tuple2<Integer, Event> makeEvent(Integer timeout, Integer value) {
+ Map<String, Object> map = new HashMap<>();
+ map.put("randomValue", value);
+ return new Tuple2<>(timeout, EventFactory.fromMap(map, new
SourceInfo("test"
+ + "-topic", "s0"),
+ new SchemaInfo(null,
+ new ArrayList<>())));
+ }
}
diff --git
a/streampipes-wrapper-siddhi/src/main/java/org/apache/streampipes/wrapper/siddhi/engine/generator/SiddhiAppGenerator.java
b/streampipes-wrapper-siddhi/src/main/java/org/apache/streampipes/wrapper/siddhi/engine/generator/SiddhiAppGenerator.java
index 4838f16a56..9b70a74a49 100644
---
a/streampipes-wrapper-siddhi/src/main/java/org/apache/streampipes/wrapper/siddhi/engine/generator/SiddhiAppGenerator.java
+++
b/streampipes-wrapper-siddhi/src/main/java/org/apache/streampipes/wrapper/siddhi/engine/generator/SiddhiAppGenerator.java
@@ -78,7 +78,7 @@ public class SiddhiAppGenerator {
.getQueries()
.forEach(query -> this.siddhiAppString.append(query).append("\n"));
- LOG.info("Registering statement: \n" + this.siddhiAppString);
+ LOG.info("Registering statement: \n{}", this.siddhiAppString);
}
}
diff --git
a/streampipes-wrapper-siddhi/src/main/java/org/apache/streampipes/wrapper/siddhi/model/EventPropertyDef.java
b/streampipes-wrapper-siddhi/src/main/java/org/apache/streampipes/wrapper/siddhi/model/EventPropertyDef.java
index 175f53e2c1..4d59c1fe8e 100644
---
a/streampipes-wrapper-siddhi/src/main/java/org/apache/streampipes/wrapper/siddhi/model/EventPropertyDef.java
+++
b/streampipes-wrapper-siddhi/src/main/java/org/apache/streampipes/wrapper/siddhi/model/EventPropertyDef.java
@@ -19,10 +19,16 @@ package org.apache.streampipes.wrapper.siddhi.model;
public class EventPropertyDef {
+ public static final String WHITESPACE_REPLACEMENT = "__w__";
+
private String selectorPrefix;
private final String fieldName;
private final String fieldType;
+ public static String toOriginalFieldName(String sanitizedFieldName) {
+ return sanitizedFieldName.replaceAll(WHITESPACE_REPLACEMENT, " ");
+ }
+
public EventPropertyDef(String fieldName,
String fieldType) {
this.fieldName = fieldName;
@@ -40,7 +46,7 @@ public class EventPropertyDef {
}
public String getFieldName() {
- return fieldName;
+ return fieldName.replaceAll(" ", WHITESPACE_REPLACEMENT);
}
public String getFieldType() {
diff --git
a/streampipes-wrapper-siddhi/src/main/java/org/apache/streampipes/wrapper/siddhi/utils/SiddhiUtils.java
b/streampipes-wrapper-siddhi/src/main/java/org/apache/streampipes/wrapper/siddhi/utils/SiddhiUtils.java
index 041bbce894..e2e589dc41 100644
---
a/streampipes-wrapper-siddhi/src/main/java/org/apache/streampipes/wrapper/siddhi/utils/SiddhiUtils.java
+++
b/streampipes-wrapper-siddhi/src/main/java/org/apache/streampipes/wrapper/siddhi/utils/SiddhiUtils.java
@@ -19,10 +19,12 @@ package org.apache.streampipes.wrapper.siddhi.utils;
import org.apache.streampipes.extensions.api.pe.param.IDataProcessorParameters;
+import org.apache.streampipes.model.constants.PropertySelectorConstants;
import org.apache.streampipes.model.runtime.EventFactory;
import org.apache.streampipes.model.runtime.SchemaInfo;
import org.apache.streampipes.model.runtime.SourceInfo;
import org.apache.streampipes.wrapper.siddhi.constants.SiddhiConstants;
+import org.apache.streampipes.wrapper.siddhi.model.EventPropertyDef;
import io.siddhi.core.event.Event;
import io.siddhi.query.api.definition.Attribute;
@@ -68,7 +70,7 @@ public class SiddhiUtils {
outputKey = outputKey.substring(2);
}
Object data = event.getData(i);
- outMap.put(outputKey, data);
+ outMap.put(EventPropertyDef.toOriginalFieldName(outputKey), data);
}
return outMap;
@@ -101,11 +103,13 @@ public class SiddhiUtils {
return eventName
.replaceAll("\\.", "")
.replaceAll("-", "")
- .replaceAll("::", "");
+ .replaceAll(PropertySelectorConstants.PROPERTY_DELIMITER, "");
}
public static String prepareProperty(String propertyName) {
- return propertyName.replaceAll("::", "");
+ return propertyName
+ .replaceAll(PropertySelectorConstants.PROPERTY_DELIMITER, "")
+ .replaceAll(" ", EventPropertyDef.WHITESPACE_REPLACEMENT);
}
}