This is an automated email from the ASF dual-hosted git repository.
zehnder pushed a commit to branch 2779-teststringtostateprocessor
in repository https://gitbox.apache.org/repos/asf/streampipes.git
The following commit(s) were added to
refs/heads/2779-teststringtostateprocessor by this push:
new 5d79145b2d refactor(#2779): Fix the test for TestStringToStateProcessor
5d79145b2d is described below
commit 5d79145b2dbf78b80144e43a9587d561540e923e
Author: Philipp Zehnder <[email protected]>
AuthorDate: Mon Jul 15 17:59:37 2024 +0200
refactor(#2779): Fix the test for TestStringToStateProcessor
---
.../pom.xml | 6 +
.../state/StringToStateProcessor.java | 7 +-
.../state/TestStringToStateProcessor.java | 213 ++++++---------------
.../executors/ProcessingElementTestExecutor.java | 13 +-
4 files changed, 76 insertions(+), 163 deletions(-)
diff --git
a/streampipes-extensions/streampipes-processors-transformation-jvm/pom.xml
b/streampipes-extensions/streampipes-processors-transformation-jvm/pom.xml
index 26d9b8f430..4f05f1112d 100644
--- a/streampipes-extensions/streampipes-processors-transformation-jvm/pom.xml
+++ b/streampipes-extensions/streampipes-processors-transformation-jvm/pom.xml
@@ -69,6 +69,12 @@
<groupId>org.junit.jupiter</groupId>
<artifactId>junit-jupiter-params</artifactId>
</dependency>
+ <dependency>
+ <groupId>org.apache.streampipes</groupId>
+ <artifactId>streampipes-test-utils-executors</artifactId>
+ <version>0.97.0-SNAPSHOT</version>
+ <scope>test</scope>
+ </dependency>
</dependencies>
<build>
diff --git
a/streampipes-extensions/streampipes-processors-transformation-jvm/src/main/java/org/apache/streampipes/processors/transformation/jvm/processor/stringoperator/state/StringToStateProcessor.java
b/streampipes-extensions/streampipes-processors-transformation-jvm/src/main/java/org/apache/streampipes/processors/transformation/jvm/processor/stringoperator/state/StringToStateProcessor.java
index 50ecbab838..08bcff5349 100644
---
a/streampipes-extensions/streampipes-processors-transformation-jvm/src/main/java/org/apache/streampipes/processors/transformation/jvm/processor/stringoperator/state/StringToStateProcessor.java
+++
b/streampipes-extensions/streampipes-processors-transformation-jvm/src/main/java/org/apache/streampipes/processors/transformation/jvm/processor/stringoperator/state/StringToStateProcessor.java
@@ -38,8 +38,7 @@ import org.apache.streampipes.vocabulary.SPSensor;
import org.apache.streampipes.wrapper.params.compat.ProcessorParams;
import org.apache.streampipes.wrapper.standalone.StreamPipesDataProcessor;
-import com.google.common.collect.Lists;
-
+import java.util.ArrayList;
import java.util.List;
public class StringToStateProcessor extends StreamPipesDataProcessor {
@@ -77,13 +76,13 @@ public class StringToStateProcessor extends
StreamPipesDataProcessor {
@Override
public void onEvent(Event event, SpOutputCollector collector) throws
SpRuntimeException {
- List<String> states = Lists.newArrayList();
+ List<String> states = new ArrayList<>();
for (String stateField : stateFields) {
states.add(event.getFieldBySelector(stateField).getAsPrimitive().getAsString());
}
- event.addField(CURRENT_STATE, states.toArray());
+ event.addField(CURRENT_STATE, states);
collector.collect(event);
}
diff --git
a/streampipes-extensions/streampipes-processors-transformation-jvm/src/test/java/org/apache/streampipes/processors/transformation/jvm/processor/stringoperator/state/TestStringToStateProcessor.java
b/streampipes-extensions/streampipes-processors-transformation-jvm/src/test/java/org/apache/streampipes/processors/transformation/jvm/processor/stringoperator/state/TestStringToStateProcessor.java
index 332b07d9e0..43854c91cf 100644
---
a/streampipes-extensions/streampipes-processors-transformation-jvm/src/test/java/org/apache/streampipes/processors/transformation/jvm/processor/stringoperator/state/TestStringToStateProcessor.java
+++
b/streampipes-extensions/streampipes-processors-transformation-jvm/src/test/java/org/apache/streampipes/processors/transformation/jvm/processor/stringoperator/state/TestStringToStateProcessor.java
@@ -18,159 +18,64 @@
package
org.apache.streampipes.processors.transformation.jvm.processor.stringoperator.state;
-//@RunWith(Parameterized.class)
+import org.apache.streampipes.test.executors.ProcessingElementTestExecutor;
+import org.apache.streampipes.test.executors.TestConfiguration;
+
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.Arguments;
+import org.junit.jupiter.params.provider.MethodSource;
+
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+import java.util.stream.Stream;
+
public class TestStringToStateProcessor {
-//
-// private static final Logger LOG =
LoggerFactory.getLogger(TestStringToStateProcessor.class);
-//
-// @org.junit.runners.Parameterized.Parameters
-// public static Iterable<Object[]> data() {
-// return Arrays.asList(new Object[][] {
-// {
-// List.of(),
-// List.of("c1", "c2", "c3"),
-// List.of(Arrays.asList("t1", "t2", "t3")),
-// List.of()
-// },
-// {
-// List.of("c1"),
-// List.of("c1", "c2", "c3"),
-// List.of(Arrays.asList("t1", "t2", "t3")),
-// List.of("t1")
-// },
-// {
-// List.of("c1", "c2"),
-// List.of("c1", "c2", "c3"),
-// List.of(Arrays.asList("t1", "t2", "t3")),
-// Arrays.asList("t1", "t2")
-// },
-// {
-// List.of("c1", "c2"),
-// List.of("c1", "c2", "c3"),
-// Arrays.asList(
-// Arrays.asList("t1-1", "t2-1", "t3-1"),
-// Arrays.asList("t1-2", "t2-2", "t3-2")
-// ),
-// Arrays.asList("t1-2", "t2-2")
-// },
-// {
-// List.of("c1", "c2", "c3"),
-// List.of("c1", "c2", "c3"),
-// Arrays.asList(
-// Arrays.asList("t1-1", "t2-1", "t3-1"),
-// Arrays.asList("t1-2", "t2-2", "t3-2"),
-// Arrays.asList("t1-3", "t2-3", "t3-3")
-// ),
-// Arrays.asList("t1-3", "t2-3", "t3-3")
-// }
-// });
-// }
-//
-// @org.junit.runners.Parameterized.Parameter
-// public List<String> selectedFieldNames;
-//
-// @org.junit.runners.Parameterized.Parameter(1)
-// public List<String> fieldNames;
-//
-// @org.junit.runners.Parameterized.Parameter(2)
-// public List<List<String>> eventStrings;
-//
-// @org.junit.runners.Parameterized.Parameter(3)
-// public List<String> expectedValue;
-//
-// private static final String DEFAULT_STREAM_NAME = "stream1";
-//
-// @Test
-// public void testStringToState() {
-// StringToStateProcessor stringToStateProcessor = new
StringToStateProcessor();
-// DataProcessorDescription originalGraph =
stringToStateProcessor.declareModel();
-//
originalGraph.setSupportedGrounding(EventGroundingGenerator.makeDummyGrounding());
-//
-// DataProcessorInvocation graph =
InvocationGraphGenerator.makeEmptyInvocation(originalGraph);
-// graph.setInputStreams(Collections
-// .singletonList(EventStreamGenerator
-//
.makeStreamWithProperties(Collections.singletonList("stream-in"))));
-//
graph.setOutputStream(EventStreamGenerator.makeStreamWithProperties(Collections.singletonList("stream-out")));
-//
graph.getOutputStream().getEventGrounding().getTransportProtocol().getTopicDefinition()
-// .setActualTopicName("output-topic");
-//
-// MappingPropertyNary mappingPropertyNary =
graph.getStaticProperties().stream()
-// .filter(p -> p instanceof MappingPropertyNary)
-// .map(p -> (MappingPropertyNary) p)
-// .filter(p ->
p.getInternalName().equals(StringToStateProcessor.STRING_STATE_FIELD))
-// .findFirst().orElse(null);
-//
-// assert mappingPropertyNary != null;
-// mappingPropertyNary.setSelectedProperties(
-// selectedFieldNames.stream().map(field -> DEFAULT_STREAM_NAME + "::"
+ field).toList());
-//
-// ProcessorParams params = new ProcessorParams(graph);
-//
-// SpOutputCollector spOutputCollector = new SpOutputCollector() {
-// @Override
-// public void registerConsumer(String routeId,
InternalEventProcessor<Map<String, Object>> consumer) {
-// }
-//
-// @Override
-// public void unregisterConsumer(String routeId) {
-// }
-//
-// @Override
-// public void connect() throws SpRuntimeException {
-// }
-//
-// @Override
-// public void disconnect() throws SpRuntimeException {
-// }
-//
-// @Override
-// public void collect(Event event) {
-// }
-// };
-//
-// stringToStateProcessor.onInvocation(params, spOutputCollector, null);
-// Object[] states = sendEvents(stringToStateProcessor, spOutputCollector);
-// LOG.info("Expected states is {}.", expectedValue);
-// LOG.info("Actual states is {}.", Arrays.toString(states));
-// assertArrayEquals(expectedValue.toArray(), states);
-// }
-//
-// private Object[] sendEvents(StringToStateProcessor stateProcessor,
SpOutputCollector spOut) {
-// List<Event> events = makeEvents();
-// Object[] states = null;
-// for (Event event : events) {
-// stateProcessor.onEvent(event, spOut);
-// try {
-// TimeUnit.MILLISECONDS.sleep(100);
-// } catch (InterruptedException e) {
-// throw new RuntimeException(e);
-// }
-// try {
-// states = (Object[])
event.getFieldBySelector(StringToStateProcessor.CURRENT_STATE)
-// .getAsPrimitive().getRawValue();
-// LOG.info("Current states: " + Arrays.toString(states));
-// } catch (IllegalArgumentException e) {
-//
-// }
-// }
-// return states;
-// }
-//
-// private List<Event> makeEvents() {
-// List<Event> events = Lists.newArrayList();
-// for (List<String> eventSetting : eventStrings) {
-// events.add(makeEvent(eventSetting));
-// }
-// return events;
-// }
-//
-// private Event makeEvent(List<String> value) {
-// Map<String, Object> map = Maps.newHashMap();
-// for (int i = 0; i < selectedFieldNames.size(); i++) {
-// map.put(selectedFieldNames.get(i), value.get(i));
-// }
-// return EventFactory.fromMap(map,
-// new SourceInfo("test-topic", DEFAULT_STREAM_NAME),
-// new SchemaInfo(null, Lists.newArrayList()));
-// }
+
+ private StringToStateProcessor processor;
+
+ @BeforeEach
+ public void setup() {
+ processor = new StringToStateProcessor();
+ }
+
+ static Stream<Arguments> arguments() {
+ return Stream.of(
+ Arguments.of(
+ Collections.emptyList(),
+ List.of(Map.of("k1", "v1")),
+ List.of(Map.of("k1", "v1", StringToStateProcessor.CURRENT_STATE,
Collections.emptyList()))
+ ),
+ Arguments.of(
+ List.of("::k1"),
+ List.of(Map.of("k1", "v1")),
+ List.of(Map.of("k1", "v1", StringToStateProcessor.CURRENT_STATE,
List.of("v1")))
+ ),
+ Arguments.of(
+ List.of("::k1", "::k2"),
+ List.of(Map.of("k1", "v1", "k2", "v2")),
+ List.of(Map.of("k1", "v1", "k2", "v2",
StringToStateProcessor.CURRENT_STATE, List.of("v1", "v2")))
+ )
+ );
+ }
+
+ @ParameterizedTest
+ @MethodSource("arguments")
+ public void testStringToState(
+ List<String> selectedFieldNames,
+ List<Map<String, Object>> intpuEvents,
+ List<Map<String, Object>> outputEvents
+ ) {
+
+ var configuration = TestConfiguration
+ .builder()
+ .config(StringToStateProcessor.STRING_STATE_FIELD, selectedFieldNames)
+ .build();
+
+ var testExecutor = new ProcessingElementTestExecutor(processor,
configuration);
+
+ testExecutor.run(intpuEvents, outputEvents);
+ }
+
}
diff --git
a/streampipes-test-utils-executors/src/main/java/org/apache/streampipes/test/executors/ProcessingElementTestExecutor.java
b/streampipes-test-utils-executors/src/main/java/org/apache/streampipes/test/executors/ProcessingElementTestExecutor.java
index f2d87197ad..b27f0e887a 100644
---
a/streampipes-test-utils-executors/src/main/java/org/apache/streampipes/test/executors/ProcessingElementTestExecutor.java
+++
b/streampipes-test-utils-executors/src/main/java/org/apache/streampipes/test/executors/ProcessingElementTestExecutor.java
@@ -44,6 +44,9 @@ import java.util.Map;
import java.util.function.Consumer;
import java.util.stream.IntStream;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
public class ProcessingElementTestExecutor {
@@ -93,17 +96,17 @@ public class ProcessingElementTestExecutor {
invocationConfig.accept(dataProcessorInvocation);
}
- var e = getProcessingElementParameterExtractor(dataProcessorInvocation);
- var mockParams = Mockito.mock(IDataProcessorParameters.class);
+ var extractor =
getProcessingElementParameterExtractor(dataProcessorInvocation);
+ var mockParams = mock(IDataProcessorParameters.class);
- Mockito.when(mockParams.getModel()).thenReturn(dataProcessorInvocation);
- Mockito.when(mockParams.extractor()).thenReturn(e);
+ when(mockParams.getModel()).thenReturn(dataProcessorInvocation);
+ when(mockParams.extractor()).thenReturn(extractor);
// calls the onPipelineStarted method of the processor to initialize it
processor.onPipelineStarted(mockParams, null, null);
// mock the output collector to capture the output events and validate the
results later
- var mockCollector = Mockito.mock(SpOutputCollector.class);
+ var mockCollector = mock(SpOutputCollector.class);
var spOutputCollectorCaptor = ArgumentCaptor.forClass(Event.class);