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