This is an automated email from the ASF dual-hosted git repository. SvenO3 pushed a commit to branch fix-adapter-script-migration-for-file-stream-adapter in repository https://gitbox.apache.org/repos/asf/streampipes.git
commit 0abf7f13693f1701875dd1c41a6e413dcda279ac Author: Sven Oehler <[email protected]> AuthorDate: Wed Jul 22 15:40:14 2026 +0200 Improve script migration --- .../compact/generator/AdapterSchemaGenerator.java | 31 +------- .../model/connect/TransformationConfig.java | 19 +++++ .../v099/connect/MigrateAdaptersToUseScript.java | 24 ++++-- .../v099/MigrateAdaptersToUseScriptTest.java | 85 ++++++++++++++++++++++ 4 files changed, 123 insertions(+), 36 deletions(-) diff --git a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/compact/generator/AdapterSchemaGenerator.java b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/compact/generator/AdapterSchemaGenerator.java index ff240bd30d..6b22c4abd9 100644 --- a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/compact/generator/AdapterSchemaGenerator.java +++ b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/compact/generator/AdapterSchemaGenerator.java @@ -62,8 +62,8 @@ public class AdapterSchemaGenerator implements AdapterModelGenerator { adapterDescription.getTransformationConfig() .setInputs(sampleData.getSamples()); - setDefaultScriptIfNotSet(adapterDescription); - setDefaultScriptLanguageIfNotSet(adapterDescription); + adapterDescription.getTransformationConfig() + .applyScriptDefaults(); guessManagement.transformSampleData(adapterDescription, userId); @@ -88,31 +88,4 @@ public class AdapterSchemaGenerator implements AdapterModelGenerator { } } - private void setDefaultScriptIfNotSet(AdapterDescription adapterDescription) { - if (adapterDescription.getTransformationConfig() - .getScript() == null - || adapterDescription.getTransformationConfig() - .getScript() - .isEmpty()) { - adapterDescription.getTransformationConfig().setScriptActive(true); - adapterDescription.getTransformationConfig() - .setScript(""" - function transform(event, out, ctx) { - out.collect(event); - } - """); - - } - } - - private void setDefaultScriptLanguageIfNotSet(AdapterDescription adapterDescription) { - if (adapterDescription.getTransformationConfig() - .getLanguage() == null - || adapterDescription.getTransformationConfig() - .getLanguage() - .isEmpty()) { - adapterDescription.getTransformationConfig() - .setLanguage("javascript"); - } - } } diff --git a/streampipes-model/src/main/java/org/apache/streampipes/model/connect/TransformationConfig.java b/streampipes-model/src/main/java/org/apache/streampipes/model/connect/TransformationConfig.java index fb8596c326..798d6a3d92 100644 --- a/streampipes-model/src/main/java/org/apache/streampipes/model/connect/TransformationConfig.java +++ b/streampipes-model/src/main/java/org/apache/streampipes/model/connect/TransformationConfig.java @@ -23,6 +23,12 @@ import java.util.List; import java.util.Map; public class TransformationConfig { + public static final String DEFAULT_LANGUAGE = "javascript"; + public static final String DEFAULT_SCRIPT = """ + function transform(event, out, ctx) { + out.collect(event); + }"""; + private boolean scriptActive; private String language; private String script; @@ -35,6 +41,19 @@ public class TransformationConfig { public TransformationConfig() { this.inputs = new ArrayList<>(); this.outputs = new ArrayList<>(); + this.scriptActive = false; + this.language = DEFAULT_LANGUAGE; + this.script = DEFAULT_SCRIPT; + } + + public void applyScriptDefaults() { + if (language == null || language.isBlank()) { + language = DEFAULT_LANGUAGE; + } + + if (script == null || script.isBlank()) { + script = DEFAULT_SCRIPT; + } } public String getLanguage() { diff --git a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/v099/connect/MigrateAdaptersToUseScript.java b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/v099/connect/MigrateAdaptersToUseScript.java index 6803d36856..ee1b7ff218 100644 --- a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/v099/connect/MigrateAdaptersToUseScript.java +++ b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/v099/connect/MigrateAdaptersToUseScript.java @@ -42,7 +42,6 @@ import java.util.Optional; public class MigrateAdaptersToUseScript implements Migration { private static final Logger LOG = LoggerFactory.getLogger(MigrateAdaptersToUseScript.class); - private static final String SCRIPT_LANGUAGE = "javascript"; private final IAdapterStorage adapterStorage; @@ -52,17 +51,17 @@ public class MigrateAdaptersToUseScript implements Migration { } @Override - // Execute if there is an adapter with no transformation config or script + // Execute if an adapter still has legacy transformation rules or misses its script config public boolean shouldExecute() { List<AdapterDescription> adapters = adapterStorage.findAll(); - return adapters != null && adapters.stream().anyMatch(this::hasEmptyTransformationConfig); + return adapters != null && adapters.stream().anyMatch(this::shouldMigrate); } @Override public void executeMigration() throws IOException { adapterStorage.findAll() .stream() - .filter(this::hasEmptyTransformationConfig) + .filter(this::shouldMigrate) .forEach(this::migrateAndUpdateAdapter); } @@ -71,8 +70,19 @@ public class MigrateAdaptersToUseScript implements Migration { return "Changes the rules based adapters to use script based transformations instead."; } - private boolean hasEmptyTransformationConfig(AdapterDescription adapter) { - return adapter.getTransformationConfig() == null || adapter.getTransformationConfig().getScript() == null; + private boolean shouldMigrate(AdapterDescription adapter) { + return adapter.getTransformationConfig() == null + || hasNoScript(adapter) + || hasLegacyRules(adapter); + } + + private boolean hasNoScript(AdapterDescription adapter) { + var script = adapter.getTransformationConfig().getScript(); + return script == null || script.isBlank(); + } + + private boolean hasLegacyRules(AdapterDescription adapter) { + return adapter.getRules() != null && !adapter.getRules().isEmpty(); } private void migrateAndUpdateAdapter(AdapterDescription adapterDescription) { @@ -125,7 +135,7 @@ public class MigrateAdaptersToUseScript implements Migration { private TransformationConfig initializeTransformationConfig(TransformationConfig oldConfig) { var config = new TransformationConfig(); - config.setLanguage(SCRIPT_LANGUAGE); + config.setLanguage(TransformationConfig.DEFAULT_LANGUAGE); if (oldConfig == null) { return config; diff --git a/streampipes-service-core/src/test/java/org/apache/streampipes/service/core/migrations/v099/MigrateAdaptersToUseScriptTest.java b/streampipes-service-core/src/test/java/org/apache/streampipes/service/core/migrations/v099/MigrateAdaptersToUseScriptTest.java index f2d901137e..ea72318a01 100644 --- a/streampipes-service-core/src/test/java/org/apache/streampipes/service/core/migrations/v099/MigrateAdaptersToUseScriptTest.java +++ b/streampipes-service-core/src/test/java/org/apache/streampipes/service/core/migrations/v099/MigrateAdaptersToUseScriptTest.java @@ -19,6 +19,9 @@ package org.apache.streampipes.service.core.migrations.v099; import org.apache.streampipes.model.SpDataStream; +import org.apache.streampipes.model.connect.ReduceEventRateRule; +import org.apache.streampipes.model.connect.RemoveDuplicateRule; +import org.apache.streampipes.model.connect.TransformationConfig; import org.apache.streampipes.model.connect.adapter.AdapterDescription; import org.apache.streampipes.model.connect.rules.TransformationRuleDescription; import org.apache.streampipes.model.connect.rules.schema.DeleteRuleDescription; @@ -47,10 +50,12 @@ import java.util.ArrayList; import java.util.Collections; import java.util.HashMap; import java.util.List; +import java.util.Map; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNotSame; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.verify; @@ -80,10 +85,22 @@ class MigrateAdaptersToUseScriptTest { assertTrue(result); } + @Test + void shouldExecute_ReturnsTrue_WhenAdapterHasLegacyRulesAndDefaultScript() { + var adapter = createBaseAdapter(new RenameRuleDescription("old", "new")); + + when(mockStorage.findAll()).thenReturn(List.of(adapter)); + + boolean result = migration.shouldExecute(); + + assertTrue(result); + } + @Test void executeMigration_RemoveAdditionalMetadata() throws IOException { // Arrange var adapter = new AdapterDescription(); + adapter.getTransformationConfig().setScript(null); var eventPropertyPrimitive = new EventPropertyPrimitive(); eventPropertyPrimitive.setAdditionalMetadata(Collections.singletonMap("key", "value")); var eventSchema = new EventSchema(); @@ -113,6 +130,66 @@ class MigrateAdaptersToUseScriptTest { verify(mockStorage).updateElement(adapter); } + @Test + void executeMigration_KeepsInputsAndOutputsNonNull_WhenExistingConfigHasNullLists() throws IOException { + var adapter = createBaseAdapterWithoutRules(); + var oldConfig = new TransformationConfig(); + oldConfig.setScript(null); + oldConfig.setInputs(null); + oldConfig.setOutputs(null); + adapter.setTransformationConfig(oldConfig); + + when(mockStorage.findAll()).thenReturn(List.of(adapter)); + + migration.executeMigration(); + + var resultConfig = adapter.getTransformationConfig(); + assertNotNull(resultConfig.getInputs()); + assertTrue(resultConfig.getInputs().isEmpty()); + assertNotNull(resultConfig.getOutputs()); + assertTrue(resultConfig.getOutputs().isEmpty()); + verify(mockStorage).updateElement(adapter); + } + + @Test + void executeMigration_CopiesExistingTransformationConfigValues() throws IOException { + var input = new HashMap<String, Object>(); + input.put("runtimeName", "temperature"); + var output = new HashMap<String, Object>(); + output.put("runtimeName", "temperature_celsius"); + + var inputs = new ArrayList<Map<String, Object>>(); + inputs.add(input); + var outputs = new ArrayList<Map<String, Object>>(); + outputs.add(output); + + var reduceEventRateRule = new ReduceEventRateRule(10, "mean"); + var removeDuplicateRule = new RemoveDuplicateRule("500"); + + var oldConfig = new TransformationConfig(); + oldConfig.setScript(null); + oldConfig.setInputs(inputs); + oldConfig.setOutputs(outputs); + oldConfig.setReduceEventRateRule(reduceEventRateRule); + oldConfig.setRemoveDuplicateRule(removeDuplicateRule); + + var adapter = createBaseAdapterWithoutRules(); + adapter.setTransformationConfig(oldConfig); + + when(mockStorage.findAll()).thenReturn(List.of(adapter)); + + migration.executeMigration(); + + var resultConfig = adapter.getTransformationConfig(); + assertEquals(TransformationConfig.DEFAULT_LANGUAGE, resultConfig.getLanguage()); + assertEquals(inputs, resultConfig.getInputs()); + assertEquals(outputs, resultConfig.getOutputs()); + assertNotSame(inputs, resultConfig.getInputs()); + assertNotSame(outputs, resultConfig.getOutputs()); + assertEquals(reduceEventRateRule, resultConfig.getReduceEventRateRule()); + assertEquals(removeDuplicateRule, resultConfig.getRemoveDuplicateRule()); + } + @Test void executeMigration_TransformsRenameRuleToScript() throws IOException { // Arrange @@ -424,4 +501,12 @@ class MigrateAdaptersToUseScriptTest { adapter.getDataStream().setEventSchema(new EventSchema()); return adapter; } + + private AdapterDescription createBaseAdapterWithoutRules() { + AdapterDescription adapter = new AdapterDescription(); + adapter.setRules(new ArrayList<>()); + adapter.setDataStream(new SpDataStream()); + adapter.getDataStream().setEventSchema(new EventSchema()); + return adapter; + } }
