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);
   }
 
 }

Reply via email to