This is an automated email from the ASF dual-hosted git repository.
zehnder pushed a commit to branch
2964-timestamp-conversion-is-broken-in-adapters
in repository https://gitbox.apache.org/repos/asf/streampipes.git
The following commit(s) were added to
refs/heads/2964-timestamp-conversion-is-broken-in-adapters by this push:
new 7e21430694 refactor(#2964): Fix timestamp transformations within
FileReplayAdapter
7e21430694 is described below
commit 7e21430694866db5497ec7707798dc83f12e1828
Author: Philipp Zehnder <[email protected]>
AuthorDate: Mon Jul 1 15:13:48 2024 +0200
refactor(#2964): Fix timestamp transformations within FileReplayAdapter
---
.../SupportsNestedTransformationRule.java | 2 +-
.../AdapterTransformationPipelineElement.java | 2 +-
.../TransformationRuleGeneratorVisitor.java | 2 +-
.../schema/AddValueTransformationRule.java | 2 +-
.../transform/schema/MoveTransformationRule.java | 2 +-
.../stream/DuplicateFilterPipelineElement.java | 2 +-
.../stream/EventRateTransformationRule.java | 2 +-
.../value/AddTimestampTransformationRule.java | 2 +-
.../value/DatatypeTransformationRule.java | 2 +-
.../schema/SchemaEventTransformerTest.java | 2 +-
.../transform/value/ValueEventTransformerTest.java | 2 +-
.../extensions/api/connect/StreamPipesAdapter.java | 15 +++++
.../api/connect}/TransformationRule.java | 2 +-
.../connect/AdapterWorkerManagement.java | 5 ++
.../iiot/protocol/stream/FileReplayAdapter.java | 65 +++++++++++++++++++---
.../protocol/stream/FileReplayAdapterTest.java | 65 ++++++++++++++++++++++
ui/cypress/tests/adapter/rules/valueRules.spec.ts | 12 ++--
17 files changed, 160 insertions(+), 26 deletions(-)
diff --git
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/SupportsNestedTransformationRule.java
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/SupportsNestedTransformationRule.java
index 86a36b7d6c..d4ee006249 100644
---
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/SupportsNestedTransformationRule.java
+++
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/SupportsNestedTransformationRule.java
@@ -18,7 +18,7 @@
package org.apache.streampipes.connect.shared.preprocessing;
-import
org.apache.streampipes.connect.shared.preprocessing.transform.TransformationRule;
+import org.apache.streampipes.extensions.api.connect.TransformationRule;
import java.util.List;
import java.util.Map;
diff --git
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/elements/AdapterTransformationPipelineElement.java
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/elements/AdapterTransformationPipelineElement.java
index 2b6d5070b2..c3d6bb2445 100644
---
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/elements/AdapterTransformationPipelineElement.java
+++
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/elements/AdapterTransformationPipelineElement.java
@@ -19,9 +19,9 @@
package org.apache.streampipes.connect.shared.preprocessing.elements;
import
org.apache.streampipes.connect.shared.preprocessing.generator.TransformationRuleGeneratorVisitor;
-import
org.apache.streampipes.connect.shared.preprocessing.transform.TransformationRule;
import org.apache.streampipes.connect.shared.preprocessing.utils.Utils;
import org.apache.streampipes.extensions.api.connect.IAdapterPipelineElement;
+import org.apache.streampipes.extensions.api.connect.TransformationRule;
import
org.apache.streampipes.model.connect.rules.TransformationRuleDescription;
import java.util.List;
diff --git
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/generator/TransformationRuleGeneratorVisitor.java
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/generator/TransformationRuleGeneratorVisitor.java
index 2b6a4edaff..87c9bb6881 100644
---
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/generator/TransformationRuleGeneratorVisitor.java
+++
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/generator/TransformationRuleGeneratorVisitor.java
@@ -18,7 +18,7 @@
package org.apache.streampipes.connect.shared.preprocessing.generator;
-import
org.apache.streampipes.connect.shared.preprocessing.transform.TransformationRule;
+import org.apache.streampipes.extensions.api.connect.TransformationRule;
import org.apache.streampipes.model.connect.rules.ITransformationRuleVisitor;
import java.util.ArrayList;
diff --git
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/schema/AddValueTransformationRule.java
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/schema/AddValueTransformationRule.java
index 77b8b12f20..8a31382ebe 100644
---
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/schema/AddValueTransformationRule.java
+++
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/schema/AddValueTransformationRule.java
@@ -19,7 +19,7 @@
package org.apache.streampipes.connect.shared.preprocessing.transform.schema;
import org.apache.streampipes.connect.shared.DatatypeUtils;
-import
org.apache.streampipes.connect.shared.preprocessing.transform.TransformationRule;
+import org.apache.streampipes.extensions.api.connect.TransformationRule;
import java.util.Map;
diff --git
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/schema/MoveTransformationRule.java
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/schema/MoveTransformationRule.java
index f6506f376c..b584b9eff2 100644
---
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/schema/MoveTransformationRule.java
+++
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/schema/MoveTransformationRule.java
@@ -18,7 +18,7 @@
package org.apache.streampipes.connect.shared.preprocessing.transform.schema;
-import
org.apache.streampipes.connect.shared.preprocessing.transform.TransformationRule;
+import org.apache.streampipes.extensions.api.connect.TransformationRule;
import java.util.HashMap;
import java.util.List;
diff --git
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/stream/DuplicateFilterPipelineElement.java
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/stream/DuplicateFilterPipelineElement.java
index bbca98fb1e..fedf106a24 100644
---
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/stream/DuplicateFilterPipelineElement.java
+++
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/stream/DuplicateFilterPipelineElement.java
@@ -18,7 +18,7 @@
package org.apache.streampipes.connect.shared.preprocessing.transform.stream;
-import
org.apache.streampipes.connect.shared.preprocessing.transform.TransformationRule;
+import org.apache.streampipes.extensions.api.connect.TransformationRule;
import java.util.HashMap;
import java.util.Map;
diff --git
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/stream/EventRateTransformationRule.java
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/stream/EventRateTransformationRule.java
index 229952c8b3..b5b586f65e 100644
---
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/stream/EventRateTransformationRule.java
+++
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/stream/EventRateTransformationRule.java
@@ -18,7 +18,7 @@
package org.apache.streampipes.connect.shared.preprocessing.transform.stream;
-import
org.apache.streampipes.connect.shared.preprocessing.transform.TransformationRule;
+import org.apache.streampipes.extensions.api.connect.TransformationRule;
import java.util.Map;
diff --git
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/value/AddTimestampTransformationRule.java
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/value/AddTimestampTransformationRule.java
index c8267ef34a..f075933c33 100644
---
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/value/AddTimestampTransformationRule.java
+++
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/value/AddTimestampTransformationRule.java
@@ -18,7 +18,7 @@
package org.apache.streampipes.connect.shared.preprocessing.transform.value;
-import
org.apache.streampipes.connect.shared.preprocessing.transform.TransformationRule;
+import org.apache.streampipes.extensions.api.connect.TransformationRule;
import java.util.Map;
diff --git
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/value/DatatypeTransformationRule.java
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/value/DatatypeTransformationRule.java
index ce0f0509a5..df61f2e131 100644
---
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/value/DatatypeTransformationRule.java
+++
b/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/value/DatatypeTransformationRule.java
@@ -20,7 +20,7 @@ package
org.apache.streampipes.connect.shared.preprocessing.transform.value;
import org.apache.streampipes.connect.shared.DatatypeUtils;
-import
org.apache.streampipes.connect.shared.preprocessing.transform.TransformationRule;
+import org.apache.streampipes.extensions.api.connect.TransformationRule;
import java.util.Map;
diff --git
a/streampipes-connect-shared/src/test/java/org/apache/streampipes/connect/shared/preprocessing/transform/schema/SchemaEventTransformerTest.java
b/streampipes-connect-shared/src/test/java/org/apache/streampipes/connect/shared/preprocessing/transform/schema/SchemaEventTransformerTest.java
index ec87937de1..56400a8066 100644
---
a/streampipes-connect-shared/src/test/java/org/apache/streampipes/connect/shared/preprocessing/transform/schema/SchemaEventTransformerTest.java
+++
b/streampipes-connect-shared/src/test/java/org/apache/streampipes/connect/shared/preprocessing/transform/schema/SchemaEventTransformerTest.java
@@ -18,7 +18,7 @@
package org.apache.streampipes.connect.shared.preprocessing.transform.schema;
-import
org.apache.streampipes.connect.shared.preprocessing.transform.TransformationRule;
+import org.apache.streampipes.extensions.api.connect.TransformationRule;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
diff --git
a/streampipes-connect-shared/src/test/java/org/apache/streampipes/connect/shared/preprocessing/transform/value/ValueEventTransformerTest.java
b/streampipes-connect-shared/src/test/java/org/apache/streampipes/connect/shared/preprocessing/transform/value/ValueEventTransformerTest.java
index 89cd25be8a..286be99a00 100644
---
a/streampipes-connect-shared/src/test/java/org/apache/streampipes/connect/shared/preprocessing/transform/value/ValueEventTransformerTest.java
+++
b/streampipes-connect-shared/src/test/java/org/apache/streampipes/connect/shared/preprocessing/transform/value/ValueEventTransformerTest.java
@@ -18,7 +18,7 @@
package org.apache.streampipes.connect.shared.preprocessing.transform.value;
-import
org.apache.streampipes.connect.shared.preprocessing.transform.TransformationRule;
+import org.apache.streampipes.extensions.api.connect.TransformationRule;
import org.apache.streampipes.model.schema.EventProperty;
import org.apache.streampipes.model.schema.EventPropertyPrimitive;
import org.apache.streampipes.model.schema.EventSchema;
diff --git
a/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/connect/StreamPipesAdapter.java
b/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/connect/StreamPipesAdapter.java
index 4a5fb1d362..8651b05ee5 100644
---
a/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/connect/StreamPipesAdapter.java
+++
b/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/connect/StreamPipesAdapter.java
@@ -22,11 +22,26 @@ import
org.apache.streampipes.commons.exceptions.connect.AdapterException;
import
org.apache.streampipes.extensions.api.connect.context.IAdapterGuessSchemaContext;
import
org.apache.streampipes.extensions.api.connect.context.IAdapterRuntimeContext;
import
org.apache.streampipes.extensions.api.extractor.IAdapterParameterExtractor;
+import org.apache.streampipes.model.connect.adapter.AdapterDescription;
import org.apache.streampipes.model.connect.guess.GuessSchema;
public interface StreamPipesAdapter {
IAdapterConfiguration declareConfig();
+ /**
+ * Preprocesses the adapter description before the adapter is invoked.
+ *
+ * <p>This method is designed to allow adapters to modify the adapter
description prior to invocation.
+ * It is particularly useful for adapters that need to manipulate certain
values internally,
+ * e.g. bypassing the adapter preprocessing pipeline. An example of such an
adapter is the FileReplayAdapter,
+ * which manipulates timestamp values.</p>
+ *
+ * <p>This is a default method and does not need to be overridden unless
specific preprocessing is required.</p>
+ *
+ * @param adapterDescription The adapter description to be preprocessed.
+ */
+ default void preprocessAdapterDescription(AdapterDescription
adapterDescription) {};
+
void onAdapterStarted(IAdapterParameterExtractor extractor,
IEventCollector collector,
IAdapterRuntimeContext adapterRuntimeContext) throws
AdapterException;
diff --git
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/TransformationRule.java
b/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/connect/TransformationRule.java
similarity index 92%
rename from
streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/TransformationRule.java
rename to
streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/connect/TransformationRule.java
index f6443afc62..91b877825d 100644
---
a/streampipes-connect-shared/src/main/java/org/apache/streampipes/connect/shared/preprocessing/transform/TransformationRule.java
+++
b/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/connect/TransformationRule.java
@@ -16,7 +16,7 @@
*
*/
-package org.apache.streampipes.connect.shared.preprocessing.transform;
+package org.apache.streampipes.extensions.api.connect;
import java.util.Map;
diff --git
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagement.java
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagement.java
index 958d413f1f..1b17ab6bd3 100644
---
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagement.java
+++
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagement.java
@@ -63,6 +63,11 @@ public class AdapterWorkerManagement {
newAdapterInstance,
adapterDescription);
+ // This method allows adapters to modify the adapter description prior
to invocation.
+ // It is particularly useful for adapters like FileReplayAdapter that
need to manipulate timestamp values internally,
+ // bypassing the adapter preprocessing pipeline.
+ newAdapterInstance.preprocessAdapterDescription(adapterDescription);
+
var registeredParsers =
newAdapterInstance.declareConfig().getSupportedParsers();
var extractor = AdapterParameterExtractor.from(adapterDescription,
registeredParsers);
var eventCollector = EventCollector.from(adapterDescription);
diff --git
a/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/FileReplayAdapter.java
b/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/FileReplayAdapter.java
index 8ed3697a3c..6a8b14650f 100644
---
a/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/FileReplayAdapter.java
+++
b/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/FileReplayAdapter.java
@@ -20,6 +20,7 @@ package org.apache.streampipes.connect.iiot.protocol.stream;
import org.apache.streampipes.commons.exceptions.connect.AdapterException;
import org.apache.streampipes.connect.iiot.utils.FileProtocolUtils;
+import
org.apache.streampipes.connect.shared.preprocessing.generator.StatelessTransformationRuleGeneratorVisitor;
import org.apache.streampipes.extensions.api.connect.IAdapterConfiguration;
import org.apache.streampipes.extensions.api.connect.IEventCollector;
import org.apache.streampipes.extensions.api.connect.IParser;
@@ -32,8 +33,10 @@ import
org.apache.streampipes.extensions.management.connect.adapter.parser.Image
import
org.apache.streampipes.extensions.management.connect.adapter.parser.JsonParsers;
import
org.apache.streampipes.extensions.management.connect.adapter.parser.xml.XmlParser;
import org.apache.streampipes.model.AdapterType;
+import org.apache.streampipes.model.connect.adapter.AdapterDescription;
import org.apache.streampipes.model.connect.guess.GuessSchema;
import org.apache.streampipes.model.connect.rules.schema.RenameRuleDescription;
+import
org.apache.streampipes.model.connect.rules.value.TimestampTranfsformationRuleDescription;
import org.apache.streampipes.model.extensions.ExtensionAssetType;
import org.apache.streampipes.sdk.StaticProperties;
import org.apache.streampipes.sdk.builder.adapter.AdapterConfigurationBuilder;
@@ -78,6 +81,7 @@ public class FileReplayAdapter implements StreamPipesAdapter {
private boolean replaceTimestamp;
private String timestampRuntimeName;
+ private TimestampTranfsformationRuleDescription
timestampTranfsformationRuleDescription;
private String timestampSourceFieldName;
@@ -293,29 +297,46 @@ public class FileReplayAdapter implements
StreamPipesAdapter {
collector.collect(event);
}
- // TODO Add description that explains that the timestamp is transformed to a
unix timestamp
- private long getTimestampFromEvent(Map<String, Object> event) throws
AdapterException {
+ protected long getTimestampFromEvent(Map<String, Object> event) throws
AdapterException {
long actualEventTimestamp = -1;
var timestampFieldValue = event.get(timestampSourceFieldName);
+
if (timestampFieldValue instanceof Long) {
actualEventTimestamp = (Long) timestampFieldValue;
} else if (timestampFieldValue instanceof Integer) {
actualEventTimestamp = (Integer) timestampFieldValue;
- } else if (!(timestampFieldValue == null && replaceTimestamp)) {
- throw new AdapterException(
- "Timestamp field is not a unix timestamp in ms, skipping event. "
- + "Value: %s".formatted(event.get(timestampSourceFieldName)));
}
- // TODO should be replaced
+ // transform timestamp if transformation rule is present
+ actualEventTimestamp =
transformTimestampIfTransformationRuleIsPresent(event, actualEventTimestamp);
+
+
if (actualEventTimestamp == -1 && !replaceTimestamp) {
- throw new AdapterException("TODO");
+ throw new AdapterException("Timestamp field could not be parsed,
skipping event. "
+ + "Value:
%s".formatted(event.get(timestampSourceFieldName)));
}
return actualEventTimestamp;
}
+ private long transformTimestampIfTransformationRuleIsPresent(Map<String,
Object> event, long actualEventTimestamp) {
+ if (timestampTranfsformationRuleDescription != null) {
+ var transformationRuleDescription =
timestampTranfsformationRuleDescription;
+
+ var transformationRuleVisitor = new
StatelessTransformationRuleGeneratorVisitor();
+ transformationRuleVisitor.visit(transformationRuleDescription);
+ var timestampTransformationRule =
transformationRuleVisitor.getTransformationRules()
+ .get(0);
+
+ actualEventTimestamp = (Long) (
+ timestampTransformationRule.apply(event)
+ .get(timestampSourceFieldName)
+ );
+ }
+ return actualEventTimestamp;
+ }
+
private void reduceReplaySpeedIfRequired(long actualEventTimestamp) {
long sleepTime;
if (timestampLastEvent != -1 && actualEventTimestamp != -1) {
@@ -374,4 +395,32 @@ public class FileReplayAdapter implements
StreamPipesAdapter {
protected void setReplaceTimestamp(boolean replaceTimestamp) {
this.replaceTimestamp = replaceTimestamp;
}
+
+ protected void
setTimestampTranfsformationRuleDescription(TimestampTranfsformationRuleDescription
timestampTranfsformationRuleDescription) {
+ this.timestampTranfsformationRuleDescription =
timestampTranfsformationRuleDescription;
+ }
+
+ /**
+ * Removes the timestamp transformation rules from the adapter description.
+ *
+ * <p>The FileReplay adapter manages timestamp transformations internally to
accurately simulate the replay frequency.
+ * This is necessary as the timestamp field values are crucial for this
simulation. As a result, the timestamp rule
+ * description is stored locally within the FileReplay adapter and is
applied when the onAdapterStarted method is invoked.</p>
+ */
+ @Override
+ public void preprocessAdapterDescription(AdapterDescription
adapterDescription) {
+
+ this.timestampTranfsformationRuleDescription = adapterDescription
+ .getRules()
+ .stream()
+ .filter(rule -> rule instanceof
TimestampTranfsformationRuleDescription)
+ .map(rule -> (TimestampTranfsformationRuleDescription) rule)
+ .findFirst()
+ .get();
+
+ // remove timestamp preprocessing rule
+ adapterDescription
+ .getRules()
+ .removeIf(rule -> rule instanceof
TimestampTranfsformationRuleDescription);
+ }
}
diff --git
a/streampipes-extensions/streampipes-connect-adapters-iiot/src/test/java/org/apache/streampipes/connect/iiot/protocol/stream/FileReplayAdapterTest.java
b/streampipes-extensions/streampipes-connect-adapters-iiot/src/test/java/org/apache/streampipes/connect/iiot/protocol/stream/FileReplayAdapterTest.java
index 0efbd947af..e9a8b3ae93 100644
---
a/streampipes-extensions/streampipes-connect-adapters-iiot/src/test/java/org/apache/streampipes/connect/iiot/protocol/stream/FileReplayAdapterTest.java
+++
b/streampipes-extensions/streampipes-connect-adapters-iiot/src/test/java/org/apache/streampipes/connect/iiot/protocol/stream/FileReplayAdapterTest.java
@@ -19,7 +19,9 @@
package org.apache.streampipes.connect.iiot.protocol.stream;
import org.apache.streampipes.commons.exceptions.connect.AdapterException;
+import
org.apache.streampipes.connect.shared.preprocessing.transform.value.TimestampTranformationRuleMode;
import org.apache.streampipes.extensions.api.connect.IEventCollector;
+import
org.apache.streampipes.model.connect.rules.value.TimestampTranfsformationRuleDescription;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
@@ -102,4 +104,67 @@ class FileReplayAdapterTest {
}
+ @Test
+ void getTimestampFromEvent_returnsLongTimestamp() throws AdapterException {
+ event.put(TIMESTAMP, TIMESTAMP_VALUE);
+
+ long actualEventTimestamp = fileReplayAdapter.getTimestampFromEvent(event);
+
+ assertEquals(TIMESTAMP_VALUE, actualEventTimestamp);
+ }
+
+ @Test
+ void getTimestampFromEvent_forTimestampRuleInSecondsAsLong() throws
AdapterException {
+ setupEventAndRule(
+ TIMESTAMP,
+ 1622544682L,
+ TimestampTranformationRuleMode.TIME_UNIT,
+ 1000L
+ );
+ assertEventTimestamp(TIMESTAMP_VALUE);
+ }
+
+ @Test
+ void getTimestampFromEvent_forTimestampRuleInSecondsAsInteger() throws
AdapterException {
+ setupEventAndRule(
+ TIMESTAMP,
+ 1622544682,
+ TimestampTranformationRuleMode.TIME_UNIT,
+ 1000L
+ );
+ assertEventTimestamp(TIMESTAMP_VALUE);
+ }
+
+ @Test
+ void getTimestampFromEvent_forTimestampRuleAsString() throws
AdapterException {
+ setupEventAndRule(
+ TIMESTAMP,
+ "2024-07-01T12:00:00.000Z",
+ TimestampTranformationRuleMode.FORMAT_STRING,
+ "yyyy-MM-dd'T'HH:mm:ss.SSS'Z'"
+ );
+ assertEventTimestamp(1719828000000L);
+ }
+
+
+ private void setupEventAndRule(
+ String key,
+ Object value,
+ TimestampTranformationRuleMode mode,
+ Object additional
+ ) {
+ event.put(key, value);
+ var rule = new TimestampTranfsformationRuleDescription();
+ rule.setRuntimeKey(key);
+ rule.setMode(mode.internalName());
+ if (additional instanceof Long) {rule.setMultiplier((Long) additional);}
+ if (additional instanceof String) {rule.setFormatString((String)
additional);}
+ fileReplayAdapter.setTimestampTranfsformationRuleDescription(rule);
+ }
+
+ private void assertEventTimestamp(long expected) throws AdapterException {
+ assertEquals(expected, fileReplayAdapter.getTimestampFromEvent(event));
+ }
+
+
}
\ No newline at end of file
diff --git a/ui/cypress/tests/adapter/rules/valueRules.spec.ts
b/ui/cypress/tests/adapter/rules/valueRules.spec.ts
index d6d1ccdd1d..2163f454d2 100644
--- a/ui/cypress/tests/adapter/rules/valueRules.spec.ts
+++ b/ui/cypress/tests/adapter/rules/valueRules.spec.ts
@@ -37,14 +37,14 @@ describe('Connect value rule transformations', () => {
);
// Number transformation
- // ConnectEventSchemaUtils.numberTransformation('value', '10');
+ ConnectEventSchemaUtils.numberTransformation('value', '10');
// Unit transformation
- // ConnectEventSchemaUtils.unitTransformation(
- // 'temperature',
- // 'Degree Celsius',
- // 'Degree Fahrenheit',
- // );
+ ConnectEventSchemaUtils.unitTransformation(
+ 'temperature',
+ 'Degree Celsius',
+ 'Degree Fahrenheit',
+ );
ConnectEventSchemaUtils.finishEventSchemaConfiguration();