This is an automated email from the ASF dual-hosted git repository.

SvenO3 pushed a commit to branch 
4766-unify-adapter-and-data-stream-lifecycle-handling
in repository https://gitbox.apache.org/repos/asf/streampipes.git


The following commit(s) were added to 
refs/heads/4766-unify-adapter-and-data-stream-lifecycle-handling by this push:
     new 028b82af38 Adjust tests
028b82af38 is described below

commit 028b82af38ce43d8b745a1858b9292a45e88c900
Author: Sven Oehler <[email protected]>
AuthorDate: Fri Jul 24 16:26:43 2026 +0200

    Adjust tests
---
 .../update/PipelineUpdateCoordinatorTest.java      | 139 +++++++--------------
 1 file changed, 47 insertions(+), 92 deletions(-)

diff --git 
a/streampipes-pipeline-management/src/test/java/org/apache/streampipes/manager/pipeline/update/PipelineUpdateCoordinatorTest.java
 
b/streampipes-pipeline-management/src/test/java/org/apache/streampipes/manager/pipeline/update/PipelineUpdateCoordinatorTest.java
index 5a0ff1b0f0..b557ba34c6 100644
--- 
a/streampipes-pipeline-management/src/test/java/org/apache/streampipes/manager/pipeline/update/PipelineUpdateCoordinatorTest.java
+++ 
b/streampipes-pipeline-management/src/test/java/org/apache/streampipes/manager/pipeline/update/PipelineUpdateCoordinatorTest.java
@@ -22,10 +22,7 @@ import 
org.apache.streampipes.manager.api.extensions.ExtensionServiceRequestMana
 import org.apache.streampipes.manager.execution.PipelineExecutor;
 import org.apache.streampipes.manager.matching.PipelineVerificationHandlerV2;
 import 
org.apache.streampipes.manager.matching.v2.pipeline.MeasurementChangeValidationStep;
-import org.apache.streampipes.manager.pipeline.PipelineManager;
 import org.apache.streampipes.model.SpDataStream;
-import org.apache.streampipes.model.connect.adapter.AdapterDescription;
-import org.apache.streampipes.model.connect.adapter.ChartSchemaUpdateInfo;
 import org.apache.streampipes.model.connect.adapter.PipelineUpdateInfo;
 import org.apache.streampipes.model.graph.DataProcessorInvocation;
 import org.apache.streampipes.model.graph.DataSinkInvocation;
@@ -38,8 +35,10 @@ import 
org.apache.streampipes.model.pipeline.PipelineModificationResult;
 import org.apache.streampipes.model.schema.EventPropertyPrimitive;
 import org.apache.streampipes.model.schema.EventSchema;
 import org.apache.streampipes.model.schema.PropertyScope;
+import org.apache.streampipes.resource.management.ChartResourceManager;
 import org.apache.streampipes.resource.management.PipelineResourceManager;
 import org.apache.streampipes.resource.management.SpResourceManager;
+import org.apache.streampipes.storage.api.explorer.IChartStorage;
 import org.apache.streampipes.storage.api.pipeline.IPipelineStorage;
 import org.apache.streampipes.vocabulary.XSD;
 
@@ -58,7 +57,6 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.mockConstruction;
 import static org.mockito.Mockito.verify;
-import static org.mockito.Mockito.verifyNoInteractions;
 import static org.mockito.Mockito.when;
 
 class PipelineUpdateCoordinatorTest {
@@ -68,12 +66,9 @@ class PipelineUpdateCoordinatorTest {
   @Test
   void updatePipelines_ShouldRestartRunningPipelinesForDataStreamUpdates() {
     var requestManager = mock(ExtensionServiceRequestManager.class);
-    var chartSchemaUpdateCoordinator = 
mock(ChartSchemaUpdateCoordinator.class);
     var resourceManager = mock(SpResourceManager.class);
-    var pipelineManager = mock(PipelineManager.class);
     var pipelineStorage = mock(IPipelineStorage.class);
-    var coordinator = makeCoordinator(
-        requestManager, resourceManager, chartSchemaUpdateCoordinator, 
pipelineManager, pipelineStorage);
+    var coordinator = makeCoordinator(requestManager, resourceManager, 
pipelineStorage);
     var dataStream = makeDataStream("stream-1", "Updated stream");
     var affectedPipeline = makePipeline("pipeline-1", "Pipeline", true, 
"stream-1", "Old stream");
     var storedPipeline = makePipeline("pipeline-1", "Pipeline", true, 
"stream-1", "Old stream");
@@ -81,8 +76,8 @@ class PipelineUpdateCoordinatorTest {
 
     var modificationMessage = new 
PipelineModificationMessage(List.of(validModification("sepa-1")));
     var verifiedPipelines = new ArrayList<Pipeline>();
-    
when(pipelineManager.getPipelinesContainingElements("stream-1")).thenReturn(List.of(affectedPipeline));
-    when(pipelineManager.getPipeline("pipeline-1")).thenReturn(storedPipeline, 
modifiedPipeline);
+    when(pipelineStorage.findAll()).thenReturn(List.of(affectedPipeline));
+    
when(pipelineStorage.getElementById("pipeline-1")).thenReturn(storedPipeline, 
modifiedPipeline);
 
     try (MockedConstruction<PipelineVerificationHandlerV2> 
verificationHandlerConstruction =
              mockConstruction(PipelineVerificationHandlerV2.class, (mock, 
context) -> {
@@ -105,29 +100,25 @@ class PipelineUpdateCoordinatorTest {
       verify(executorConstruction.constructed().get(0)).stopPipeline(true);
       verify(executorConstruction.constructed().get(1)).startPipeline();
       verify(pipelineStorage).updateElement(modifiedPipeline);
-      verifyNoInteractions(chartSchemaUpdateCoordinator);
     }
   }
 
   @Test
-  void 
updatePipelines_ShouldMarkPipelinesRequiringAttentionForAdapterUpdates() {
+  void 
updatePipelines_ShouldMarkPipelinesRequiringAttentionForDataStreamUpdates() {
     var requestManager = mock(ExtensionServiceRequestManager.class);
-    var chartSchemaUpdateCoordinator = 
mock(ChartSchemaUpdateCoordinator.class);
     var resourceManager = mock(SpResourceManager.class);
-    var pipelineManager = mock(PipelineManager.class);
     var pipelineStorage = mock(IPipelineStorage.class);
-    var coordinator = makeCoordinator(
-        requestManager, resourceManager, chartSchemaUpdateCoordinator, 
pipelineManager, pipelineStorage);
-    var adapterDescription = makeAdapter("stream-1", "Updated adapter");
+    var coordinator = makeCoordinator(requestManager, resourceManager, 
pipelineStorage);
+    var dataStream = makeDataStream("stream-1", "Updated stream");
     var storedPipeline = makePipeline("pipeline-1", "Pipeline", false, 
"stream-1", "Old stream");
-    var modifiedPipeline = makePipeline("pipeline-1", "Pipeline", false, 
"stream-1", "Updated adapter");
+    var modifiedPipeline = makePipeline("pipeline-1", "Pipeline", false, 
"stream-1", "Updated stream");
     modifiedPipeline.setSepas(List.of(makeSepa("sepa-1", "Processor")));
 
     var warning = PipelineElementValidationInfo.error("Schema mismatch");
     var modificationMessage = new 
PipelineModificationMessage(List.of(invalidModification("sepa-1", warning)));
     var verifiedPipelines = new ArrayList<Pipeline>();
-    
when(pipelineManager.getPipelinesContainingElements("stream-1")).thenReturn(List.of(storedPipeline));
-    when(pipelineManager.getPipeline("pipeline-1")).thenReturn(storedPipeline);
+    when(pipelineStorage.findAll()).thenReturn(List.of(storedPipeline));
+    
when(pipelineStorage.getElementById("pipeline-1")).thenReturn(storedPipeline);
 
     try (MockedConstruction<PipelineVerificationHandlerV2> 
verificationHandlerConstruction =
              mockConstruction(PipelineVerificationHandlerV2.class, (mock, 
context) -> {
@@ -139,19 +130,19 @@ class PipelineUpdateCoordinatorTest {
          MockedConstruction<PipelineExecutor> executorConstruction =
              mockConstruction(PipelineExecutor.class)) {
 
-      coordinator.updatePipelines(adapterDescription);
+      coordinator.updatePipelines(dataStream);
 
       assertEquals(1, verificationHandlerConstruction.constructed().size());
       var updatedPipeline = verifiedPipelines.get(0);
-      assertEquals("Updated adapter", 
updatedPipeline.getStreams().get(0).getName());
-      assertSame(adapterDescription.getEventSchema(), 
updatedPipeline.getStreams().get(0).getEventSchema());
+      assertEquals("Updated stream", 
updatedPipeline.getStreams().get(0).getName());
+      assertSame(dataStream.getEventSchema(), 
updatedPipeline.getStreams().get(0).getEventSchema());
 
       assertEquals(0, executorConstruction.constructed().size());
       var pipelineCaptor = ArgumentCaptor.forClass(Pipeline.class);
       verify(pipelineStorage).updateElement(pipelineCaptor.capture());
       assertEquals(PipelineHealthStatus.REQUIRES_ATTENTION, 
pipelineCaptor.getValue().getHealthStatus());
       assertFalse(pipelineCaptor.getValue().isValid());
-      assertEquals(List.of("Adapter modification: Processor: [Schema 
mismatch]"),
+      assertEquals(List.of("Data stream modification: Processor: [Schema 
mismatch]"),
           pipelineCaptor.getValue().getPipelineNotifications());
     }
   }
@@ -159,25 +150,22 @@ class PipelineUpdateCoordinatorTest {
   @Test
   void 
updatePipelines_ShouldMarkPipelineRequiringAttentionForCriticalMeasurementFieldChange()
 {
     var requestManager = mock(ExtensionServiceRequestManager.class);
-    var chartSchemaUpdateCoordinator = 
mock(ChartSchemaUpdateCoordinator.class);
     var resourceManager = mock(SpResourceManager.class);
-    var pipelineManager = mock(PipelineManager.class);
     var pipelineStorage = mock(IPipelineStorage.class);
-    var coordinator = makeCoordinator(
-        requestManager, resourceManager, chartSchemaUpdateCoordinator, 
pipelineManager, pipelineStorage);
-    var adapterDescription = makeAdapter("stream-1", "Updated adapter");
-    
adapterDescription.getDataStream().setEventSchema(makeSchema(makeMeasurementProperty("temperature",
 XSD.STRING)));
+    var coordinator = makeCoordinator(requestManager, resourceManager, 
pipelineStorage);
+    var dataStream = makeDataStream("stream-1", "Updated stream");
+    
dataStream.setEventSchema(makeSchema(makeMeasurementProperty("temperature", 
XSD.STRING)));
 
     var storedPipeline = makePipeline("pipeline-1", "Pipeline", true, 
"stream-1", "Old stream");
     
storedPipeline.getStreams().get(0).setEventSchema(makeSchema(makeMeasurementProperty("temperature",
 XSD.INTEGER)));
     storedPipeline.setActions(List.of(makeDataLakeSink()));
 
-    var modifiedPipeline = makePipeline("pipeline-1", "Pipeline", true, 
"stream-1", "Updated adapter");
+    var modifiedPipeline = makePipeline("pipeline-1", "Pipeline", true, 
"stream-1", "Updated stream");
     var measurementUpdateInfo = PipelineElementValidationInfo.info(
         measurementUpdateRequiredMessage());
     var modificationMessage = new 
PipelineModificationMessage(List.of(validModification("sepa-1", 
measurementUpdateInfo)));
-    
when(pipelineManager.getPipelinesContainingElements("stream-1")).thenReturn(List.of(storedPipeline));
-    when(pipelineManager.getPipeline("pipeline-1")).thenReturn(storedPipeline);
+    when(pipelineStorage.findAll()).thenReturn(List.of(storedPipeline));
+    
when(pipelineStorage.getElementById("pipeline-1")).thenReturn(storedPipeline);
 
     try (MockedConstruction<PipelineVerificationHandlerV2> 
verificationHandlerConstruction =
              mockConstruction(PipelineVerificationHandlerV2.class, (mock, 
context) -> {
@@ -188,7 +176,7 @@ class PipelineUpdateCoordinatorTest {
          MockedConstruction<PipelineExecutor> executorConstruction =
              mockConstruction(PipelineExecutor.class)) {
 
-      coordinator.updatePipelines(adapterDescription);
+      coordinator.updatePipelines(dataStream);
 
       assertEquals(1, verificationHandlerConstruction.constructed().size());
       var pipelineCaptor = ArgumentCaptor.forClass(Pipeline.class);
@@ -204,24 +192,14 @@ class PipelineUpdateCoordinatorTest {
   @Test
   void checkPipelineMigrations_ShouldUseUpdatedDataStreamValues() {
     var requestManager = mock(ExtensionServiceRequestManager.class);
-    var chartSchemaUpdateCoordinator = 
mock(ChartSchemaUpdateCoordinator.class);
     var resourceManager = mock(SpResourceManager.class);
-    var pipelineManager = mock(PipelineManager.class);
-    var coordinator = makeCoordinator(
-        requestManager,
-        resourceManager,
-        chartSchemaUpdateCoordinator,
-        pipelineManager,
-        mock(IPipelineStorage.class)
-    );
+    var pipelineStorage = mock(IPipelineStorage.class);
+    var coordinator = makeCoordinator(requestManager, resourceManager, 
pipelineStorage);
     var dataStream = makeDataStream("stream-1", "Updated stream");
     var pipeline = makePipeline("pipeline-1", "Pipeline", false, "stream-1", 
"Old stream");
     var modificationMessage = new 
PipelineModificationMessage(List.of(validModification("sepa-1")));
     var verifiedPipelines = new ArrayList<Pipeline>();
-    var chartUpdateInfo = new ChartSchemaUpdateInfo();
-    
when(pipelineManager.getPipelinesContainingElements("stream-1")).thenReturn(List.of(pipeline));
-    when(chartSchemaUpdateCoordinator.checkChartMigrations(pipeline, 
dataStream.getEventSchema()))
-        .thenReturn(List.of(chartUpdateInfo));
+    when(pipelineStorage.findAll()).thenReturn(List.of(pipeline));
 
     try (MockedConstruction<PipelineVerificationHandlerV2> 
verificationHandlerConstruction =
              mockConstruction(PipelineVerificationHandlerV2.class, (mock, 
context) -> {
@@ -237,7 +215,7 @@ class PipelineUpdateCoordinatorTest {
       assertEquals("pipeline-1", result.get(0).getPipelineId());
       assertEquals("Pipeline", result.get(0).getPipelineName());
       assertTrue(result.get(0).isCanAutoMigrate());
-      assertEquals(List.of(chartUpdateInfo), 
result.get(0).getChartSchemaUpdateInfos());
+      assertEquals(List.of(), result.get(0).getChartSchemaUpdateInfos());
 
       assertEquals(1, verificationHandlerConstruction.constructed().size());
       var updatedPipeline = verifiedPipelines.get(0);
@@ -247,25 +225,18 @@ class PipelineUpdateCoordinatorTest {
   }
 
   @Test
-  void checkPipelineMigrations_ShouldReportWarningsForAdapterUpdates() {
+  void checkPipelineMigrations_ShouldReportWarningsForDataStreamUpdates() {
     var requestManager = mock(ExtensionServiceRequestManager.class);
-    var chartSchemaUpdateCoordinator = 
mock(ChartSchemaUpdateCoordinator.class);
     var resourceManager = mock(SpResourceManager.class);
-    var pipelineManager = mock(PipelineManager.class);
-    var coordinator = makeCoordinator(
-        requestManager,
-        resourceManager,
-        chartSchemaUpdateCoordinator,
-        pipelineManager,
-        mock(IPipelineStorage.class)
-    );
-    var adapterDescription = makeAdapter("stream-1", "Updated adapter");
+    var pipelineStorage = mock(IPipelineStorage.class);
+    var coordinator = makeCoordinator(requestManager, resourceManager, 
pipelineStorage);
+    var dataStream = makeDataStream("stream-1", "Updated stream");
     var pipeline = makePipeline("pipeline-1", "Pipeline", false, "stream-1", 
"Old stream");
     pipeline.setSepas(List.of(makeSepa("sepa-1", "Processor")));
     var warning = PipelineElementValidationInfo.error("Schema mismatch");
     var modificationMessage = new 
PipelineModificationMessage(List.of(invalidModification("sepa-1", warning)));
     var verifiedPipelines = new ArrayList<Pipeline>();
-    
when(pipelineManager.getPipelinesContainingElements("stream-1")).thenReturn(List.of(pipeline));
+    when(pipelineStorage.findAll()).thenReturn(List.of(pipeline));
 
     try (MockedConstruction<PipelineVerificationHandlerV2> 
verificationHandlerConstruction =
              mockConstruction(PipelineVerificationHandlerV2.class, (mock, 
context) -> {
@@ -275,7 +246,7 @@ class PipelineUpdateCoordinatorTest {
                    .thenReturn(new PipelineModificationResult((Pipeline) 
context.arguments().get(0), List.of()));
              })) {
 
-      var result = coordinator.checkPipelineMigrations(adapterDescription);
+      var result = coordinator.checkPipelineMigrations(dataStream);
 
       assertEquals(1, result.size());
       PipelineUpdateInfo updateInfo = result.get(0);
@@ -285,26 +256,19 @@ class PipelineUpdateCoordinatorTest {
 
       assertEquals(1, verificationHandlerConstruction.constructed().size());
       var updatedPipeline = verifiedPipelines.get(0);
-      assertEquals("Updated adapter", 
updatedPipeline.getStreams().get(0).getName());
-      assertSame(adapterDescription.getEventSchema(), 
updatedPipeline.getStreams().get(0).getEventSchema());
+      assertEquals("Updated stream", 
updatedPipeline.getStreams().get(0).getName());
+      assertSame(dataStream.getEventSchema(), 
updatedPipeline.getStreams().get(0).getEventSchema());
     }
   }
 
   @Test
   void 
checkPipelineMigrations_ShouldDisableAutoMigrationForCriticalMeasurementFieldChange()
 {
     var requestManager = mock(ExtensionServiceRequestManager.class);
-    var chartSchemaUpdateCoordinator = 
mock(ChartSchemaUpdateCoordinator.class);
     var resourceManager = mock(SpResourceManager.class);
-    var pipelineManager = mock(PipelineManager.class);
-    var coordinator = makeCoordinator(
-        requestManager,
-        resourceManager,
-        chartSchemaUpdateCoordinator,
-        pipelineManager,
-        mock(IPipelineStorage.class)
-    );
-    var adapterDescription = makeAdapter("stream-1", "Updated adapter");
-    
adapterDescription.getDataStream().setEventSchema(makeSchema(makeMeasurementProperty("temperature",
 XSD.STRING)));
+    var pipelineStorage = mock(IPipelineStorage.class);
+    var coordinator = makeCoordinator(requestManager, resourceManager, 
pipelineStorage);
+    var dataStream = makeDataStream("stream-1", "Updated stream");
+    
dataStream.setEventSchema(makeSchema(makeMeasurementProperty("temperature", 
XSD.STRING)));
 
     var pipeline = makePipeline("pipeline-1", "Pipeline", false, "stream-1", 
"Old stream");
     
pipeline.getStreams().get(0).setEventSchema(makeSchema(makeMeasurementProperty("temperature",
 XSD.INTEGER)));
@@ -313,7 +277,7 @@ class PipelineUpdateCoordinatorTest {
     var measurementUpdateInfo = PipelineElementValidationInfo.info(
         measurementUpdateRequiredMessage());
     var modificationMessage = new 
PipelineModificationMessage(List.of(validModification("sepa-1", 
measurementUpdateInfo)));
-    
when(pipelineManager.getPipelinesContainingElements("stream-1")).thenReturn(List.of(pipeline));
+    when(pipelineStorage.findAll()).thenReturn(List.of(pipeline));
 
     try (MockedConstruction<PipelineVerificationHandlerV2> 
verificationHandlerConstruction =
              mockConstruction(PipelineVerificationHandlerV2.class, (mock, 
context) -> {
@@ -322,7 +286,7 @@ class PipelineUpdateCoordinatorTest {
                    .thenReturn(new PipelineModificationResult((Pipeline) 
context.arguments().get(0), List.of()));
              })) {
 
-      var result = coordinator.checkPipelineMigrations(adapterDescription);
+      var result = coordinator.checkPipelineMigrations(dataStream);
 
       assertEquals(1, verificationHandlerConstruction.constructed().size());
       assertEquals(1, result.size());
@@ -332,18 +296,16 @@ class PipelineUpdateCoordinatorTest {
 
   private PipelineUpdateCoordinator 
makeCoordinator(ExtensionServiceRequestManager requestManager,
                                                     SpResourceManager 
resourceManager,
-                                                    
ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator,
-                                                    PipelineManager 
pipelineManager,
                                                     IPipelineStorage 
pipelineStorage) {
     var pipelineResourceManager = mock(PipelineResourceManager.class);
+    var chartResourceManager = mock(ChartResourceManager.class);
+    var chartStorage = mock(IChartStorage.class);
     
when(resourceManager.managePipelines()).thenReturn(pipelineResourceManager);
     when(pipelineResourceManager.getDb()).thenReturn(pipelineStorage);
-    return new PipelineUpdateCoordinator(
-        requestManager,
-        resourceManager,
-        chartSchemaUpdateCoordinator,
-        pipelineManager
-    );
+    when(resourceManager.manageCharts()).thenReturn(chartResourceManager);
+    when(chartResourceManager.getDb()).thenReturn(chartStorage);
+    when(chartStorage.findAll()).thenReturn(List.of());
+    return new PipelineUpdateCoordinator(requestManager, resourceManager);
   }
 
   private SpDataStream makeDataStream(String elementId, String name) {
@@ -354,14 +316,6 @@ class PipelineUpdateCoordinatorTest {
     return dataStream;
   }
 
-  private AdapterDescription makeAdapter(String 
correspondingDataStreamElementId, String name) {
-    var adapterDescription = new AdapterDescription();
-    
adapterDescription.setCorrespondingDataStreamElementId(correspondingDataStreamElementId);
-    adapterDescription.setName(name);
-    adapterDescription.getDataStream().setEventSchema(new EventSchema());
-    return adapterDescription;
-  }
-
   private Pipeline makePipeline(String pipelineId,
                                 String name,
                                 boolean running,
@@ -373,6 +327,7 @@ class PipelineUpdateCoordinatorTest {
     pipeline.setRunning(running);
     pipeline.setStreams(List.of(makeDataStream(streamElementId, streamName)));
     pipeline.setSepas(List.of(makeSepa("sepa-1", "Processor")));
+    pipeline.setActions(List.of());
     return pipeline;
   }
 

Reply via email to