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 e2a7b71d21 Unify adapter and data stream lifecycle handling
e2a7b71d21 is described below

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

    Unify adapter and data stream lifecycle handling
---
 .../management/AdapterUpdateManagement.java        | 52 +++++++++++++++-------
 .../pipeline/update}/DataStreamDeletedEvent.java   |  2 +-
 .../update/DataStreamUpdateManagement.java         | 10 +++--
 .../pipeline/update}/DataStreamUpdatedEvent.java   |  2 +-
 .../pipeline/update/PipelineUpdateCoordinator.java | 26 ++---------
 .../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 +++------
 10 files changed, 67 insertions(+), 135 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-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