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());
}