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

SvenO3 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 34efc0105b refactor(#4766): Unify adapter and data stream lifecycle 
handling (#4767)
34efc0105b is described below

commit 34efc0105be228bc36499078393247049552206b
Author: Sven Oehler <[email protected]>
AuthorDate: Mon Jul 27 15:04:30 2026 +0200

    refactor(#4766): Unify adapter and data stream lifecycle handling (#4767)
---
 .../management/AdapterUpdateManagement.java        |  52 +++++---
 .../pipeline/update}/DataStreamDeletedEvent.java   |   2 +-
 .../update/DataStreamUpdateManagement.java         |  10 +-
 .../pipeline/update}/DataStreamUpdatedEvent.java   |   2 +-
 .../pipeline/update/PipelineUpdateCoordinator.java |  26 +---
 .../update/PipelineUpdateCoordinatorTest.java      | 139 +++++++--------------
 .../rest/event/AdapterDeletedEvent.java            |  24 ----
 .../rest/event/AdapterUpdatedEvent.java            |  24 ----
 .../rest/impl/connect/AdapterResource.java         |  28 ++---
 .../rest/impl/connect/CompactAdapterResource.java  |  14 +--
 .../rest/impl/pe/DataStreamResource.java           |  20 +--
 11 files changed, 114 insertions(+), 227 deletions(-)

diff --git 
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterUpdateManagement.java
 
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterUpdateManagement.java
index 66512001d3..ed529161d9 100644
--- 
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterUpdateManagement.java
+++ 
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterUpdateManagement.java
@@ -19,7 +19,9 @@
 package org.apache.streampipes.connect.management.management;
 
 import org.apache.streampipes.commons.exceptions.connect.AdapterException;
-import 
org.apache.streampipes.manager.pipeline.update.PipelineUpdateCoordinator;
+import 
org.apache.streampipes.manager.api.extensions.ExtensionServiceRequestManager;
+import 
org.apache.streampipes.manager.pipeline.update.DataStreamUpdateManagement;
+import org.apache.streampipes.manager.pipeline.update.DataStreamUpdatedEvent;
 import org.apache.streampipes.model.SpDataStream;
 import org.apache.streampipes.model.connect.adapter.AdapterDescription;
 import org.apache.streampipes.model.connect.adapter.PipelineUpdateInfo;
@@ -27,6 +29,8 @@ import 
org.apache.streampipes.resource.management.AdapterResourceManager;
 import org.apache.streampipes.resource.management.DataStreamResourceManager;
 import org.apache.streampipes.resource.management.SpResourceManager;
 
+import org.springframework.context.ApplicationEventPublisher;
+
 import java.util.List;
 
 public class AdapterUpdateManagement {
@@ -34,20 +38,26 @@ public class AdapterUpdateManagement {
   private final AdapterMasterManagement adapterMasterManagement;
   private final AdapterResourceManager adapterResourceManager;
   private final DataStreamResourceManager dataStreamResourceManager;
-  private final PipelineUpdateCoordinator pipelineUpdateCoordinator;
+  private final DataStreamUpdateManagement dataStreamUpdateManagement;
+  private final ApplicationEventPublisher eventPublisher;
 
   public AdapterUpdateManagement(AdapterMasterManagement 
adapterMasterManagement,
-                                 PipelineUpdateCoordinator 
pipelineUpdateCoordinator,
-                                 SpResourceManager resourceManager) {
+                                 ExtensionServiceRequestManager requestManager,
+                                 SpResourceManager resourceManager,
+                                 ApplicationEventPublisher eventPublisher) {
     this.adapterMasterManagement = adapterMasterManagement;
     this.adapterResourceManager = resourceManager.manageAdapters();
     this.dataStreamResourceManager = resourceManager.manageDataStreams();
-    this.pipelineUpdateCoordinator = pipelineUpdateCoordinator;
+    this.dataStreamUpdateManagement = new DataStreamUpdateManagement(
+        requestManager,
+        resourceManager
+    );
+    this.eventPublisher = eventPublisher;
   }
 
   public void updateAdapter(AdapterDescription ad)
       throws AdapterException {
-    // update adapter in database 
+    // update adapter in database
     AdapterTransformationConfigDefaults.applyTo(ad);
     this.adapterResourceManager.encryptAndUpdate(ad);
     boolean shouldRestart = ad.isRunning();
@@ -56,10 +66,10 @@ public class AdapterUpdateManagement {
       this.adapterMasterManagement.stopAdapter(ad.getElementId(), true);
     }
 
-    // update data source in database
-    this.updateDataSource(ad);
-
-    pipelineUpdateCoordinator.updatePipelines(ad);
+    // update data source
+    var updatedDataStream = this.updateDataSource(ad);
+    dataStreamUpdateManagement.updateDataStream(updatedDataStream);
+    publishEvent(new DataStreamUpdatedEvent(updatedDataStream));
 
     if (shouldRestart) {
       this.adapterMasterManagement.startAdapter(ad.getElementId());
@@ -67,16 +77,28 @@ public class AdapterUpdateManagement {
   }
 
   public List<PipelineUpdateInfo> checkPipelineMigrations(AdapterDescription 
adapterDescription) {
-    return 
pipelineUpdateCoordinator.checkPipelineMigrations(adapterDescription);
+    return 
dataStreamUpdateManagement.checkPipelineMigrations(toDataStreamUpdate(adapterDescription));
   }
 
-  private void updateDataSource(AdapterDescription ad) {
+  private SpDataStream updateDataSource(AdapterDescription ad) {
     // get data source
     SpDataStream dataStream = 
this.dataStreamResourceManager.find(ad.getCorrespondingDataStreamElementId());
 
-    SourcesManagement.updateDataStream(ad, dataStream);
+    return SourcesManagement.updateDataStream(ad, dataStream);
+  }
+
+  private SpDataStream toDataStreamUpdate(AdapterDescription 
adapterDescription) {
+    // create a new data stream because adapterDescription.getStream() misses 
the real id and name
+    var dataStream = new SpDataStream();
+    
dataStream.setElementId(adapterDescription.getCorrespondingDataStreamElementId());
+    dataStream.setName(adapterDescription.getName());
+    dataStream.setEventSchema(adapterDescription.getEventSchema());
+    return dataStream;
+  }
 
-    // Update data source in database
-    this.dataStreamResourceManager.update(dataStream);
+  private void publishEvent(Object event) {
+    if (eventPublisher != null) {
+      eventPublisher.publishEvent(event);
+    }
   }
 }
diff --git 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/event/DataStreamDeletedEvent.java
 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/DataStreamDeletedEvent.java
similarity index 93%
rename from 
streampipes-rest/src/main/java/org/apache/streampipes/rest/event/DataStreamDeletedEvent.java
rename to 
streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/DataStreamDeletedEvent.java
index bcbbd36c8e..47fbf9f0e4 100644
--- 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/event/DataStreamDeletedEvent.java
+++ 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/DataStreamDeletedEvent.java
@@ -16,7 +16,7 @@
  *
  */
 
-package org.apache.streampipes.rest.event;
+package org.apache.streampipes.manager.pipeline.update;
 
 public record DataStreamDeletedEvent(String elementId) {
 }
diff --git 
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/DataStreamUpdateManagement.java
 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/DataStreamUpdateManagement.java
index a6454d78b4..9fe004577c 100644
--- 
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/DataStreamUpdateManagement.java
+++ 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/DataStreamUpdateManagement.java
@@ -18,9 +18,11 @@
 
 package org.apache.streampipes.manager.pipeline.update;
 
+import 
org.apache.streampipes.manager.api.extensions.ExtensionServiceRequestManager;
 import org.apache.streampipes.model.SpDataStream;
 import org.apache.streampipes.model.connect.adapter.PipelineUpdateInfo;
 import org.apache.streampipes.resource.management.DataStreamResourceManager;
+import org.apache.streampipes.resource.management.SpResourceManager;
 
 import java.util.List;
 
@@ -29,10 +31,10 @@ public class DataStreamUpdateManagement {
   private final DataStreamResourceManager dataStreamResourceManager;
   private final PipelineUpdateCoordinator pipelineUpdateCoordinator;
 
-  public DataStreamUpdateManagement(PipelineUpdateCoordinator 
pipelineUpdateCoordinator,
-                                    DataStreamResourceManager 
dataStreamResourceManager) {
-    this.dataStreamResourceManager = dataStreamResourceManager;
-    this.pipelineUpdateCoordinator = pipelineUpdateCoordinator;
+  public DataStreamUpdateManagement(ExtensionServiceRequestManager 
requestManager,
+                                    SpResourceManager resourceManager) {
+    this.pipelineUpdateCoordinator = new 
PipelineUpdateCoordinator(requestManager, resourceManager);
+    this.dataStreamResourceManager = resourceManager.manageDataStreams();
   }
 
   public void updateDataStream(SpDataStream dataStream) {
diff --git 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/event/DataStreamUpdatedEvent.java
 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/DataStreamUpdatedEvent.java
similarity index 94%
rename from 
streampipes-rest/src/main/java/org/apache/streampipes/rest/event/DataStreamUpdatedEvent.java
rename to 
streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/DataStreamUpdatedEvent.java
index ac20a0ed11..76bb4622c6 100644
--- 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/event/DataStreamUpdatedEvent.java
+++ 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/DataStreamUpdatedEvent.java
@@ -16,7 +16,7 @@
  *
  */
 
-package org.apache.streampipes.rest.event;
+package org.apache.streampipes.manager.pipeline.update;
 
 import org.apache.streampipes.model.SpDataStream;
 
diff --git 
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/PipelineUpdateCoordinator.java
 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/PipelineUpdateCoordinator.java
index e89d9476eb..6b95a49770 100644
--- 
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/PipelineUpdateCoordinator.java
+++ 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/PipelineUpdateCoordinator.java
@@ -27,7 +27,6 @@ import 
org.apache.streampipes.manager.pipeline.PipelineElementUserCleaner;
 import org.apache.streampipes.manager.pipeline.PipelineManager;
 import org.apache.streampipes.model.SpDataStream;
 import org.apache.streampipes.model.base.NamedStreamPipesEntity;
-import org.apache.streampipes.model.connect.adapter.AdapterDescription;
 import org.apache.streampipes.model.connect.adapter.PipelineUpdateInfo;
 import org.apache.streampipes.model.message.PipelineModificationMessage;
 import org.apache.streampipes.model.pipeline.Pipeline;
@@ -58,13 +57,11 @@ public class PipelineUpdateCoordinator {
   private final IPipelineStorage pipelineStorage;
 
   public PipelineUpdateCoordinator(ExtensionServiceRequestManager 
requestManager,
-                                   SpResourceManager resourceManager,
-                                   ChartSchemaUpdateCoordinator 
chartSchemaUpdateCoordinator,
-                                   PipelineManager pipelineManager) {
+                                   SpResourceManager resourceManager) {
     this.requestManager = requestManager;
     this.resourceManager = resourceManager;
-    this.chartSchemaUpdateCoordinator = chartSchemaUpdateCoordinator;
-    this.pipelineManager = pipelineManager;
+    this.chartSchemaUpdateCoordinator = new 
ChartSchemaUpdateCoordinator(resourceManager.manageCharts().getDb());
+    this.pipelineManager = new PipelineManager(resourceManager);
     this.pipelineStorage = resourceManager.managePipelines().getDb();
   }
 
@@ -77,15 +74,6 @@ public class PipelineUpdateCoordinator {
     );
   }
 
-  public void updatePipelines(AdapterDescription adapterDescription) {
-    updatePipelines(
-        adapterDescription.getCorrespondingDataStreamElementId(),
-        adapterDescription.getName(),
-        adapterDescription.getEventSchema(),
-        "Adapter"
-    );
-  }
-
   public List<PipelineUpdateInfo> checkPipelineMigrations(SpDataStream 
dataStream) {
     return checkPipelineMigrations(
         dataStream.getElementId(),
@@ -94,14 +82,6 @@ public class PipelineUpdateCoordinator {
     );
   }
 
-  public List<PipelineUpdateInfo> checkPipelineMigrations(AdapterDescription 
adapterDescription) {
-    return checkPipelineMigrations(
-        adapterDescription.getCorrespondingDataStreamElementId(),
-        adapterDescription.getName(),
-        adapterDescription.getEventSchema()
-    );
-  }
-
   private void updatePipelines(String affectedElementId,
                                String updatedStreamName,
                                EventSchema updatedEventSchema,
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;
   }
 
diff --git 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/event/AdapterDeletedEvent.java
 
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/event/AdapterDeletedEvent.java
deleted file mode 100644
index d6bb23574e..0000000000
--- 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/event/AdapterDeletedEvent.java
+++ /dev/null
@@ -1,24 +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.rest.event;
-
-import org.apache.streampipes.model.connect.adapter.AdapterDescription;
-
-public record AdapterDeletedEvent(AdapterDescription adapterDescription) {
-}
diff --git 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/event/AdapterUpdatedEvent.java
 
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/event/AdapterUpdatedEvent.java
deleted file mode 100644
index 0e67fdbe67..0000000000
--- 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/event/AdapterUpdatedEvent.java
+++ /dev/null
@@ -1,24 +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.rest.event;
-
-import org.apache.streampipes.model.connect.adapter.AdapterDescription;
-
-public record AdapterUpdatedEvent(AdapterDescription adapterDescription) {
-}
diff --git 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/AdapterResource.java
 
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/AdapterResource.java
index 84d86461d0..b44a5ab574 100644
--- 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/AdapterResource.java
+++ 
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/AdapterResource.java
@@ -26,8 +26,7 @@ import 
org.apache.streampipes.connect.management.management.CompactAdapterManage
 import org.apache.streampipes.connect.management.management.WorkerRestClient;
 import 
org.apache.streampipes.manager.api.extensions.ExtensionServiceRequestManager;
 import org.apache.streampipes.manager.pipeline.PipelineManager;
-import 
org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator;
-import 
org.apache.streampipes.manager.pipeline.update.PipelineUpdateCoordinator;
+import org.apache.streampipes.manager.pipeline.update.DataStreamDeletedEvent;
 import org.apache.streampipes.model.client.user.DefaultRole;
 import org.apache.streampipes.model.client.user.Permission;
 import org.apache.streampipes.model.connect.adapter.AdapterDescription;
@@ -42,8 +41,6 @@ import org.apache.streampipes.model.util.ElementIdGenerator;
 import org.apache.streampipes.resource.management.PermissionResourceManager;
 import org.apache.streampipes.resource.management.SpResourceManager;
 import 
org.apache.streampipes.resource.management.permission.SpPermissionEvaluator;
-import org.apache.streampipes.rest.event.AdapterDeletedEvent;
-import org.apache.streampipes.rest.event.AdapterUpdatedEvent;
 import org.apache.streampipes.rest.security.AuthConstants;
 import org.apache.streampipes.rest.shared.constants.SpMediaType;
 import org.apache.streampipes.storage.api.pipeline.IPipelineStorage;
@@ -82,7 +79,7 @@ public class AdapterResource extends 
AbstractAdapterResource<AdapterMasterManage
   private final ApplicationEventPublisher eventPublisher;
   private final PermissionResourceManager permissionResourceManager;
   private final PipelineManager pipelineManager;
-  private final PipelineUpdateCoordinator pipelineUpdateCoordinator;
+  private final AdapterUpdateManagement adapterUpdateManagement;
   private final SpResourceManager resourceManager;
 
   public AdapterResource(WorkerRestClient workerRestClient,
@@ -109,11 +106,11 @@ public class AdapterResource extends 
AbstractAdapterResource<AdapterMasterManage
     this.pipelineManager = new PipelineManager(
         resourceManager
     );
-    this.pipelineUpdateCoordinator = new PipelineUpdateCoordinator(
+    this.adapterUpdateManagement = new AdapterUpdateManagement(
+        managementService,
         requestManager,
         resourceManager,
-        new 
ChartSchemaUpdateCoordinator(resourceManager.manageCharts().getDb()),
-        pipelineManager
+        eventPublisher
     );
   }
 
@@ -151,12 +148,8 @@ public class AdapterResource extends 
AbstractAdapterResource<AdapterMasterManage
   @PutMapping(produces = MediaType.APPLICATION_JSON_VALUE, consumes = 
MediaType.APPLICATION_JSON_VALUE)
   @PreAuthorize("this.hasWriteAuthority() and 
hasPermission(#adapterDescription.correspondingDataStreamElementId, 'WRITE')")
   public ResponseEntity<? extends Message> updateAdapter(@RequestBody 
AdapterDescription adapterDescription) {
-    var updateManager = new AdapterUpdateManagement(
-        managementService, pipelineUpdateCoordinator, resourceManager
-    );
     try {
-      updateManager.updateAdapter(adapterDescription);
-      publishEvent(new AdapterUpdatedEvent(adapterDescription));
+      adapterUpdateManagement.updateAdapter(adapterDescription);
     } catch (AdapterException e) {
       LOG.error("Error while updating adapter with id {}", 
adapterDescription.getElementId(), e);
       return ok(Notifications.error(e.getMessage(), 
ExceptionUtils.getStackTrace(e)));
@@ -169,10 +162,7 @@ public class AdapterResource extends 
AbstractAdapterResource<AdapterMasterManage
   @PreAuthorize(AuthConstants.HAS_WRITE_ADAPTER_PRIVILEGE)
   public ResponseEntity<List<PipelineUpdateInfo>> 
performPipelineMigrationPreflight(
       @RequestBody AdapterDescription adapterDescription) {
-    var updateManager = new AdapterUpdateManagement(
-        managementService, pipelineUpdateCoordinator, resourceManager
-    );
-    var migrations = updateManager.checkPipelineMigrations(adapterDescription);
+    var migrations = 
adapterUpdateManagement.checkPipelineMigrations(adapterDescription);
 
     return ok(migrations);
   }
@@ -272,7 +262,7 @@ public class AdapterResource extends 
AbstractAdapterResource<AdapterMasterManage
         if (pipelinesUsingAdapter.isEmpty()) {
           try {
             managementService.deleteAdapter(elementId);
-            publishEvent(new AdapterDeletedEvent(adapter));
+            publishEvent(new 
DataStreamDeletedEvent(adapter.getCorrespondingDataStreamElementId()));
 
             return ok(Notifications.success("Adapter with id: " + elementId + 
" is deleted."));
           } catch (AdapterException e) {
@@ -325,7 +315,7 @@ public class AdapterResource extends 
AbstractAdapterResource<AdapterMasterManage
                 pipelineManager.deletePipeline(pipelineId);
               }
               managementService.deleteAdapter(elementId);
-              publishEvent(new AdapterDeletedEvent(adapter));
+              publishEvent(new 
DataStreamDeletedEvent(adapter.getCorrespondingDataStreamElementId()));
 
               return ok(Notifications.success("Adapter with id: " + elementId
                   + " and all pipelines using the adapter are deleted."));
diff --git 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/CompactAdapterResource.java
 
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/CompactAdapterResource.java
index e8f88ae5e9..7479446b4c 100644
--- 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/CompactAdapterResource.java
+++ 
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/CompactAdapterResource.java
@@ -31,8 +31,6 @@ import 
org.apache.streampipes.manager.api.extensions.ExtensionServiceRequestMana
 import 
org.apache.streampipes.manager.execution.endpoint.ExtensionsServiceEndpointGenerator;
 import org.apache.streampipes.manager.pipeline.PipelineManager;
 import 
org.apache.streampipes.manager.pipeline.compact.CompactPipelineManagement;
-import 
org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator;
-import 
org.apache.streampipes.manager.pipeline.update.PipelineUpdateCoordinator;
 import org.apache.streampipes.model.connect.adapter.AdapterDescription;
 import org.apache.streampipes.model.connect.adapter.compact.CompactAdapter;
 import org.apache.streampipes.model.message.Notifications;
@@ -44,6 +42,7 @@ import 
org.apache.streampipes.storage.management.StorageDispatcher;
 
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
+import org.springframework.context.ApplicationEventPublisher;
 import org.springframework.http.HttpStatus;
 import org.springframework.http.MediaType;
 import org.springframework.http.ResponseEntity;
@@ -67,6 +66,7 @@ public class CompactAdapterResource extends 
AbstractAdapterResource<AdapterMaste
 
   public CompactAdapterResource(WorkerRestClient workerRestClient,
                                 ExtensionServiceRequestManager requestManager,
+                                ApplicationEventPublisher eventPublisher,
                                 SpResourceManager resourceManager) {
     super(() -> new AdapterMasterManagement(
         resourceManager,
@@ -87,16 +87,12 @@ public class CompactAdapterResource extends 
AbstractAdapterResource<AdapterMaste
     this.pipelineManager = new PipelineManager(
         resourceManager
     );
-    var pipelineUpdateCoordinator = new PipelineUpdateCoordinator(
+    this.adapterUpdateManagement = new AdapterUpdateManagement(
+        managementService,
         requestManager,
         resourceManager,
-        new 
ChartSchemaUpdateCoordinator(resourceManager.manageCharts().getDb()),
-        pipelineManager
+        eventPublisher
     );
-    this.adapterUpdateManagement = new AdapterUpdateManagement(
-        managementService,
-        pipelineUpdateCoordinator,
-        resourceManager);
   }
 
   @PostMapping(
diff --git 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/pe/DataStreamResource.java
 
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/pe/DataStreamResource.java
index b106a39cc0..1080571485 100644
--- 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/pe/DataStreamResource.java
+++ 
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/pe/DataStreamResource.java
@@ -19,10 +19,9 @@
 package org.apache.streampipes.rest.impl.pe;
 
 import 
org.apache.streampipes.manager.api.extensions.ExtensionServiceRequestManager;
-import org.apache.streampipes.manager.pipeline.PipelineManager;
-import 
org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator;
+import org.apache.streampipes.manager.pipeline.update.DataStreamDeletedEvent;
 import 
org.apache.streampipes.manager.pipeline.update.DataStreamUpdateManagement;
-import 
org.apache.streampipes.manager.pipeline.update.PipelineUpdateCoordinator;
+import org.apache.streampipes.manager.pipeline.update.DataStreamUpdatedEvent;
 import org.apache.streampipes.model.SpDataStream;
 import org.apache.streampipes.model.connect.adapter.PipelineUpdateInfo;
 import org.apache.streampipes.model.message.Message;
@@ -31,10 +30,7 @@ import org.apache.streampipes.model.monitoring.SpLogMessage;
 import org.apache.streampipes.resource.management.DataStreamResourceManager;
 import org.apache.streampipes.resource.management.SpResourceManager;
 import 
org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource;
-import org.apache.streampipes.rest.event.DataStreamDeletedEvent;
-import org.apache.streampipes.rest.event.DataStreamUpdatedEvent;
 import org.apache.streampipes.rest.security.AuthConstants;
-import org.apache.streampipes.storage.api.explorer.IChartStorage;
 
 import org.apache.http.client.HttpResponseException;
 import org.springframework.context.ApplicationEventPublisher;
@@ -63,17 +59,11 @@ public class DataStreamResource extends 
AbstractAuthGuardedRestResource {
 
   public DataStreamResource(ExtensionServiceRequestManager requestManager,
                             ApplicationEventPublisher eventPublisher,
-                            IChartStorage chartStorage,
                             SpResourceManager resourceManager) {
-    var pipelineUpdateCoordinator = new PipelineUpdateCoordinator(
-        requestManager,
-        resourceManager,
-        new ChartSchemaUpdateCoordinator(chartStorage),
-        new PipelineManager(resourceManager)
-    );
     this.dataStreamResourceManager = resourceManager.manageDataStreams();
     this.dataStreamUpdateManagement = new DataStreamUpdateManagement(
-        pipelineUpdateCoordinator, dataStreamResourceManager
+        requestManager,
+        resourceManager
     );
     this.eventPublisher = eventPublisher;
   }
@@ -95,8 +85,8 @@ public class DataStreamResource extends 
AbstractAuthGuardedRestResource {
   @DeleteMapping(path = "/{elementId}", produces = 
MediaType.APPLICATION_JSON_VALUE)
   @PreAuthorize(AuthConstants.HAS_WRITE_PIPELINE_ELEMENT_PRIVILEGE)
   public ResponseEntity<Message> delete(@PathVariable("elementId") String 
elementId) {
-    publishEvent(new DataStreamDeletedEvent(elementId));
     dataStreamResourceManager.delete(elementId);
+    publishEvent(new DataStreamDeletedEvent(elementId));
     return 
constructSuccessMessage(NotificationType.STORAGE_SUCCESS.uiNotification());
   }
 


Reply via email to