This is an automated email from the ASF dual-hosted git repository.
bossenti pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/streampipes.git
The following commit(s) were added to refs/heads/dev by this push:
new 90b62ba0d Migrate boolean inverter (#1388)
90b62ba0d is described below
commit 90b62ba0dc251f95110ffa44819ead5ec3c8b3d5
Author: Felix Qvarfordt <[email protected]>
AuthorDate: Sun Mar 5 21:47:48 2023 +0100
Migrate boolean inverter (#1388)
* migrate boolean inverter (#1387)
Removed the three classes (BooleanInverter, BooleanInverterController,
BooleanInverterParameters) used to represent the boolean inverter
processing element and replaced them with a single class,
BooleanInverterProcessor.
* Added tests for the Boolean Inverter (#1387)
Added a few tests to make sure the Boolean inverter works as expected.
---------
Co-authored-by: Felix Qvarfordt <[email protected]>
---
.../transformation/jvm/TransformationJvmInit.java | 4 +-
.../booloperator/inverter/BooleanInverter.java | 54 -------
.../inverter/BooleanInverterParameters.java | 35 -----
...ntroller.java => BooleanInverterProcessor.java} | 40 ++++--
.../inverter/TestBooleanInverterProcessor.java | 159 +++++++++++++++++++++
5 files changed, 188 insertions(+), 104 deletions(-)
diff --git
a/streampipes-extensions/streampipes-processors-transformation-jvm/src/main/java/org/apache/streampipes/processors/transformation/jvm/TransformationJvmInit.java
b/streampipes-extensions/streampipes-processors-transformation-jvm/src/main/java/org/apache/streampipes/processors/transformation/jvm/TransformationJvmInit.java
index 9a033517d..b060d0858 100644
---
a/streampipes-extensions/streampipes-processors-transformation-jvm/src/main/java/org/apache/streampipes/processors/transformation/jvm/TransformationJvmInit.java
+++
b/streampipes-extensions/streampipes-processors-transformation-jvm/src/main/java/org/apache/streampipes/processors/transformation/jvm/TransformationJvmInit.java
@@ -31,7 +31,7 @@ import
org.apache.streampipes.processors.transformation.jvm.processor.array.coun
import
org.apache.streampipes.processors.transformation.jvm.processor.array.split.SplitArrayController;
import
org.apache.streampipes.processors.transformation.jvm.processor.booloperator.counter.BooleanCounterProcessor;
import
org.apache.streampipes.processors.transformation.jvm.processor.booloperator.edge.SignalEdgeFilterController;
-import
org.apache.streampipes.processors.transformation.jvm.processor.booloperator.inverter.BooleanInverterController;
+import
org.apache.streampipes.processors.transformation.jvm.processor.booloperator.inverter.BooleanInverterProcessor;
import
org.apache.streampipes.processors.transformation.jvm.processor.booloperator.logical.BooleanOperatorProcessor;
import
org.apache.streampipes.processors.transformation.jvm.processor.booloperator.state.BooleanToStateController;
import
org.apache.streampipes.processors.transformation.jvm.processor.booloperator.timekeeping.BooleanTimekeepingController;
@@ -71,7 +71,7 @@ public class TransformationJvmInit extends
ExtensionsModelSubmitter {
new ChangedValueDetectionController(),
new TimestampExtractorController(),
new BooleanCounterProcessor(),
- new BooleanInverterController(),
+ new BooleanInverterProcessor(),
new BooleanTimekeepingController(),
new BooleanTimerController(),
new CsvMetadataEnrichmentController(),
diff --git
a/streampipes-extensions/streampipes-processors-transformation-jvm/src/main/java/org/apache/streampipes/processors/transformation/jvm/processor/booloperator/inverter/BooleanInverter.java
b/streampipes-extensions/streampipes-processors-transformation-jvm/src/main/java/org/apache/streampipes/processors/transformation/jvm/processor/booloperator/inverter/BooleanInverter.java
deleted file mode 100644
index a5137d5ec..000000000
---
a/streampipes-extensions/streampipes-processors-transformation-jvm/src/main/java/org/apache/streampipes/processors/transformation/jvm/processor/booloperator/inverter/BooleanInverter.java
+++ /dev/null
@@ -1,54 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements. See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- *
- */
-
-package
org.apache.streampipes.processors.transformation.jvm.processor.booloperator.inverter;
-
-import org.apache.streampipes.logging.api.Logger;
-import org.apache.streampipes.model.runtime.Event;
-import org.apache.streampipes.wrapper.context.EventProcessorRuntimeContext;
-import org.apache.streampipes.wrapper.routing.SpOutputCollector;
-import org.apache.streampipes.wrapper.runtime.EventProcessor;
-
-public class BooleanInverter implements
EventProcessor<BooleanInverterParameters> {
-
- private static Logger log;
-
- private String invertFieldName;
-
-
- @Override
- public void onInvocation(BooleanInverterParameters booleanInverterParameters,
- SpOutputCollector spOutputCollector,
- EventProcessorRuntimeContext runtimeContext) {
- log =
booleanInverterParameters.getGraph().getLogger(BooleanInverter.class);
- this.invertFieldName = booleanInverterParameters.getInvertFieldName();
- }
-
- @Override
- public void onEvent(Event inputEvent, SpOutputCollector out) {
-
- boolean field =
inputEvent.getFieldBySelector(invertFieldName).getAsPrimitive().getAsBoolean();
- inputEvent.updateFieldBySelector(invertFieldName, !field);
-
- out.collect(inputEvent);
- }
-
- @Override
- public void onDetach() {
- }
-}
diff --git
a/streampipes-extensions/streampipes-processors-transformation-jvm/src/main/java/org/apache/streampipes/processors/transformation/jvm/processor/booloperator/inverter/BooleanInverterParameters.java
b/streampipes-extensions/streampipes-processors-transformation-jvm/src/main/java/org/apache/streampipes/processors/transformation/jvm/processor/booloperator/inverter/BooleanInverterParameters.java
deleted file mode 100644
index 636bb7908..000000000
---
a/streampipes-extensions/streampipes-processors-transformation-jvm/src/main/java/org/apache/streampipes/processors/transformation/jvm/processor/booloperator/inverter/BooleanInverterParameters.java
+++ /dev/null
@@ -1,35 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements. See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- *
- */
-
-package
org.apache.streampipes.processors.transformation.jvm.processor.booloperator.inverter;
-
-import org.apache.streampipes.model.graph.DataProcessorInvocation;
-import
org.apache.streampipes.wrapper.params.binding.EventProcessorBindingParams;
-
-public class BooleanInverterParameters extends EventProcessorBindingParams {
- private String invertFieldName;
-
- public BooleanInverterParameters(DataProcessorInvocation graph, String
invertFieldName) {
- super(graph);
- this.invertFieldName = invertFieldName;
- }
-
- public String getInvertFieldName() {
- return invertFieldName;
- }
-}
diff --git
a/streampipes-extensions/streampipes-processors-transformation-jvm/src/main/java/org/apache/streampipes/processors/transformation/jvm/processor/booloperator/inverter/BooleanInverterController.java
b/streampipes-extensions/streampipes-processors-transformation-jvm/src/main/java/org/apache/streampipes/processors/transformation/jvm/processor/booloperator/inverter/BooleanInverterProcessor.java
similarity index 62%
rename from
streampipes-extensions/streampipes-processors-transformation-jvm/src/main/java/org/apache/streampipes/processors/transformation/jvm/processor/booloperator/inverter/BooleanInverterController.java
rename to
streampipes-extensions/streampipes-processors-transformation-jvm/src/main/java/org/apache/streampipes/processors/transformation/jvm/processor/booloperator/inverter/BooleanInverterProcessor.java
index 8ae674e07..7c75ed72b 100644
---
a/streampipes-extensions/streampipes-processors-transformation-jvm/src/main/java/org/apache/streampipes/processors/transformation/jvm/processor/booloperator/inverter/BooleanInverterController.java
+++
b/streampipes-extensions/streampipes-processors-transformation-jvm/src/main/java/org/apache/streampipes/processors/transformation/jvm/processor/booloperator/inverter/BooleanInverterProcessor.java
@@ -18,9 +18,11 @@
package
org.apache.streampipes.processors.transformation.jvm.processor.booloperator.inverter;
+import org.apache.streampipes.commons.exceptions.SpRuntimeException;
+import org.apache.streampipes.logging.api.Logger;
import org.apache.streampipes.model.DataProcessorType;
import org.apache.streampipes.model.graph.DataProcessorDescription;
-import org.apache.streampipes.model.graph.DataProcessorInvocation;
+import org.apache.streampipes.model.runtime.Event;
import org.apache.streampipes.model.schema.PropertyScope;
import org.apache.streampipes.sdk.builder.ProcessingElementBuilder;
import org.apache.streampipes.sdk.builder.StreamRequirementsBuilder;
@@ -30,13 +32,15 @@ import org.apache.streampipes.sdk.helpers.Labels;
import org.apache.streampipes.sdk.helpers.Locales;
import org.apache.streampipes.sdk.helpers.OutputStrategies;
import org.apache.streampipes.sdk.utils.Assets;
-import org.apache.streampipes.wrapper.standalone.ConfiguredEventProcessor;
-import
org.apache.streampipes.wrapper.standalone.declarer.StandaloneEventProcessingDeclarer;
-
-public class BooleanInverterController extends
StandaloneEventProcessingDeclarer<BooleanInverterParameters> {
+import org.apache.streampipes.wrapper.context.EventProcessorRuntimeContext;
+import org.apache.streampipes.wrapper.routing.SpOutputCollector;
+import org.apache.streampipes.wrapper.standalone.ProcessorParams;
+import org.apache.streampipes.wrapper.standalone.StreamPipesDataProcessor;
+public class BooleanInverterProcessor extends StreamPipesDataProcessor {
public static final String INVERT_FIELD_ID = "invert-field";
-
+ private static Logger log;
+ private String invertFieldName;
@Override
public DataProcessorDescription declareModel() {
return
ProcessingElementBuilder.create("org.apache.streampipes.processors.transformation.jvm.booloperator.inverter")
@@ -54,13 +58,23 @@ public class BooleanInverterController extends
StandaloneEventProcessingDeclarer
}
@Override
- public ConfiguredEventProcessor<BooleanInverterParameters> onInvocation(
- DataProcessorInvocation graph,
- ProcessingElementParameterExtractor extractor) {
+ public void onInvocation(ProcessorParams parameters,
+ SpOutputCollector spOutputCollector,
+ EventProcessorRuntimeContext runtimeContext) throws
SpRuntimeException {
+ ProcessingElementParameterExtractor extractor = parameters.extractor();
+ log = parameters.getGraph().getLogger(BooleanInverterProcessor.class);
+ this.invertFieldName = extractor.mappingPropertyValue(INVERT_FIELD_ID);
+ }
- String invertFieldName = extractor.mappingPropertyValue(INVERT_FIELD_ID);
- BooleanInverterParameters params = new BooleanInverterParameters(graph,
invertFieldName);
+ @Override
+ public void onEvent(Event inputEvent, SpOutputCollector collector) throws
SpRuntimeException {
+ boolean field =
inputEvent.getFieldBySelector(invertFieldName).getAsPrimitive().getAsBoolean();
+ inputEvent.updateFieldBySelector(invertFieldName, !field);
+ collector.collect(inputEvent);
+ }
+
+ @Override
+ public void onDetach() throws SpRuntimeException {
- return new ConfiguredEventProcessor<>(params, BooleanInverter::new);
}
-}
+}
\ No newline at end of file
diff --git
a/streampipes-extensions/streampipes-processors-transformation-jvm/src/test/java/org/apache/streampipes/processors/transformation/jvm/processors/booloperator/inverter/TestBooleanInverterProcessor.java
b/streampipes-extensions/streampipes-processors-transformation-jvm/src/test/java/org/apache/streampipes/processors/transformation/jvm/processors/booloperator/inverter/TestBooleanInverterProcessor.java
new file mode 100644
index 000000000..2a46a4308
--- /dev/null
+++
b/streampipes-extensions/streampipes-processors-transformation-jvm/src/test/java/org/apache/streampipes/processors/transformation/jvm/processors/booloperator/inverter/TestBooleanInverterProcessor.java
@@ -0,0 +1,159 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ *
+ */
+
+package
org.apache.streampipes.processors.transformation.jvm.processor.booloperator.inverter;
+
+import org.apache.streampipes.commons.exceptions.SpRuntimeException;
+import org.apache.streampipes.messaging.InternalEventProcessor;
+import org.apache.streampipes.model.graph.DataProcessorDescription;
+import org.apache.streampipes.model.graph.DataProcessorInvocation;
+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.model.staticproperty.MappingPropertyUnary;
+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.routing.SpOutputCollector;
+import org.apache.streampipes.wrapper.standalone.ProcessorParams;
+
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.junit.runners.Parameterized;
+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 static org.junit.Assert.assertEquals;
+
+@RunWith(Parameterized.class)
+public class TestBooleanInverterProcessor {
+
+ private static final Logger LOG =
LoggerFactory.getLogger(TestBooleanInverterProcessor.class);
+
+ @org.junit.runners.Parameterized.Parameters
+ public static Iterable<Object[]> data() {
+ return Arrays.asList(new Object[][]{
+ {"Test", Arrays.asList(false, true), false},
+ {"Test", Arrays.asList(false, true, false), true},
+ {"Test", Arrays.asList(false), true},
+ {"Test", Arrays.asList(false, true, false, false, true), false},
+ {"Test", Arrays.asList(true, false), true},
+ {"Test", Arrays.asList(true), false},
+ });
+ }
+
+ @org.junit.runners.Parameterized.Parameter
+ public String invertFieldName;
+
+ @org.junit.runners.Parameterized.Parameter(1)
+ public List<Boolean> eventBooleans;
+
+ @org.junit.runners.Parameterized.Parameter(2)
+ public Boolean expectedBooleanCount;
+
+ @Test
+ public void testBoolenInverter() {
+ BooleanInverterProcessor bip = new BooleanInverterProcessor();
+ DataProcessorDescription originalGraph = bip.declareModel();
+
originalGraph.setSupportedGrounding(EventGroundingGenerator.makeDummyGrounding());
+
+ DataProcessorInvocation graph =
+ InvocationGraphGenerator.makeEmptyInvocation(originalGraph);
+
+ graph.setInputStreams(Collections
+ .singletonList(EventStreamGenerator
+
.makeStreamWithProperties(Collections.singletonList(invertFieldName))));
+
+
graph.setOutputStream(EventStreamGenerator.makeStreamWithProperties(Collections.singletonList(invertFieldName)));
+
+
graph.getOutputStream().getEventGrounding().getTransportProtocol().getTopicDefinition()
+ .setActualTopicName("output-topic");
+
+ graph.getStaticProperties().stream()
+ .filter(p -> p instanceof MappingPropertyUnary)
+ .map((p -> (MappingPropertyUnary) p))
+ .filter(p ->
p.getInternalName().equals(BooleanInverterProcessor.INVERT_FIELD_ID))
+ .findFirst().get().setSelectedProperty("s0::" + invertFieldName);
+ ProcessorParams params = new ProcessorParams(graph);
+ SpOutputCollector spOut = new SpOutputCollector() {
+ @Override
+ public void collect(Event event) {}
+
+ @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 {}
+ };
+
+ bip.onInvocation(params, spOut, null);
+
+ boolean result = sendEvents(bip, spOut);
+
+ LOG.info("Expected boolean is {}", expectedBooleanCount);
+ LOG.info("Actual boolean is {}", result);
+ assertEquals(expectedBooleanCount, result);
+ }
+
+
+ private boolean sendEvents(BooleanInverterProcessor trend, SpOutputCollector
spOut) {
+ boolean result = false;
+ List<Event> events = makeEvents();
+ for (Event event : events) {
+ LOG.info("Sending event with value " + event.getFieldBySelector("s0::" +
invertFieldName));
+ trend.onEvent(event, spOut);
+ try {
+ Thread.sleep(100);
+ } catch (InterruptedException e) {
+ e.printStackTrace();
+ }
+ result = event.getFieldBySelector("s0::" +
invertFieldName).getAsPrimitive().getAsBoolean();
+ }
+
+ return result;
+ }
+
+ private List<Event> makeEvents() {
+ List<Event> events = new ArrayList<>();
+ for (Boolean eventSetting : eventBooleans) {
+ events.add(makeEvent(eventSetting));
+ }
+ return events;
+ }
+
+ private Event makeEvent(Boolean value) {
+ Map<String, Object> map = new HashMap<>();
+ map.put(invertFieldName, value);
+ return EventFactory.fromMap(map, new SourceInfo("test" + "-topic", "s0"),
+ new SchemaInfo(null, new ArrayList<>()));
+ }
+}
\ No newline at end of file