This is an automated email from the ASF dual-hosted git repository. dominikriemer pushed a commit to branch introduce-storage-cache in repository https://gitbox.apache.org/repos/asf/streampipes.git
commit b90aaf1888895d53a852c8820cc8c0f21bba0b10 Author: Dominik Riemer <[email protected]> AuthorDate: Thu Jun 18 18:46:30 2026 +0200 Fix tests --- .../sinks/internal/jvm/datalake/DataLakeSink.java | 20 +++-- .../update/PipelineUpdateCoordinatorTest.java | 87 +++++++++------------- .../management/AdapterResourceManagerTest.java | 3 +- .../apache/streampipes/rest/ResetManagement.java | 4 +- 4 files changed, 50 insertions(+), 64 deletions(-) diff --git a/streampipes-extensions/streampipes-sinks-internal-jvm/src/main/java/org/apache/streampipes/sinks/internal/jvm/datalake/DataLakeSink.java b/streampipes-extensions/streampipes-sinks-internal-jvm/src/main/java/org/apache/streampipes/sinks/internal/jvm/datalake/DataLakeSink.java index f0e1812556..1fc1fe3c23 100644 --- a/streampipes-extensions/streampipes-sinks-internal-jvm/src/main/java/org/apache/streampipes/sinks/internal/jvm/datalake/DataLakeSink.java +++ b/streampipes-extensions/streampipes-sinks-internal-jvm/src/main/java/org/apache/streampipes/sinks/internal/jvm/datalake/DataLakeSink.java @@ -22,7 +22,6 @@ import org.apache.streampipes.client.api.IStreamPipesClient; import org.apache.streampipes.commons.environment.Environments; import org.apache.streampipes.commons.exceptions.SpRuntimeException; import org.apache.streampipes.dataexplorer.TimeSeriesStore; -import org.apache.streampipes.dataexplorer.api.IDataExplorerSchemaManagement; import org.apache.streampipes.dataexplorer.management.DataExplorerDispatcher; import org.apache.streampipes.extensions.api.extractor.IStaticPropertyExtractor; import org.apache.streampipes.extensions.api.pe.IStreamPipesDataSink; @@ -166,18 +165,17 @@ public class DataLakeSink implements IStreamPipesDataSink, SupportsRuntimeConfig private RetentionTimeConfig getRetentionTime(String measureName, IStreamPipesClient client){ - // TODO - IDataExplorerSchemaManagement dataExplorerSchemaManagement = new DataExplorerDispatcher().getDataExplorerManager() - .getSchemaManagement(null); + try { + var originalMeasure = client.dataLakeMeasureApi().getByDatasetName(measureName); + RetentionTimeConfig retentionTime = null; - var originalMeasure = dataExplorerSchemaManagement.getExistingMeasureByName(measureName); - - RetentionTimeConfig retentionTime = null; - - if (originalMeasure.isPresent()){ - retentionTime = originalMeasure.get().getRetentionTime(); + if (originalMeasure.isPresent()) { + retentionTime = originalMeasure.get().getRetentionTime(); + } + return retentionTime; + } catch (Exception e) { + return null; } - return retentionTime; } /** 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 22ac557411..3b7440ac5b 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 @@ -39,7 +39,6 @@ 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.SpResourceManager; -import org.apache.streampipes.storage.api.core.INoSqlStorage; import org.apache.streampipes.storage.api.pipeline.IPipelineStorage; import org.apache.streampipes.storage.couchdb.CouchDbStorageManager; import org.apache.streampipes.vocabulary.XSD; @@ -47,7 +46,6 @@ import org.apache.streampipes.vocabulary.XSD; import org.junit.jupiter.api.Test; import org.mockito.ArgumentCaptor; import org.mockito.MockedConstruction; -import org.mockito.MockedStatic; import java.net.URI; import java.util.ArrayList; @@ -59,7 +57,6 @@ import static org.junit.jupiter.api.Assertions.assertSame; 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.mockStatic; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.verifyNoInteractions; import static org.mockito.Mockito.when; @@ -82,13 +79,12 @@ class PipelineUpdateCoordinatorTest { var modifiedPipeline = makePipeline("pipeline-1", "Pipeline", true, "stream-1", "Updated stream"); var modificationMessage = new PipelineModificationMessage(List.of(validModification("sepa-1"))); - var noSqlStorage = mock(INoSqlStorage.class); var pipelineStorage = mock(IPipelineStorage.class); var verifiedPipelines = new ArrayList<Pipeline>(); - when(noSqlStorage.getPipelineStorageAPI()).thenReturn(pipelineStorage); + when(pipelineManager.getPipelinesContainingElements("stream-1")).thenReturn(List.of(affectedPipeline)); + when(pipelineManager.getPipeline("pipeline-1")).thenReturn(storedPipeline, modifiedPipeline); - try (MockedStatic<PipelineManager> pipelineManager = mockStatic(PipelineManager.class); - MockedConstruction<PipelineVerificationHandlerV2> verificationHandlerConstruction = + try (MockedConstruction<PipelineVerificationHandlerV2> verificationHandlerConstruction = mockConstruction(PipelineVerificationHandlerV2.class, (mock, context) -> { verifiedPipelines.add((Pipeline) context.arguments().get(0)); when(mock.verifyPipeline()).thenReturn(modificationMessage); @@ -101,11 +97,6 @@ class PipelineUpdateCoordinatorTest { MockedConstruction<PipelineExecutor> executorConstruction = mockConstruction(PipelineExecutor.class)) { - pipelineManager.when(() -> PipelineManager.getPipelinesContainingElements("stream-1")) - .thenReturn(List.of(affectedPipeline)); - pipelineManager.when(() -> PipelineManager.getPipeline("pipeline-1")) - .thenReturn(storedPipeline, modifiedPipeline); - coordinator.updatePipelines(dataStream); assertEquals(1, verificationHandlerConstruction.constructed().size()); @@ -126,7 +117,10 @@ class PipelineUpdateCoordinatorTest { void updatePipelines_ShouldMarkPipelinesRequiringAttentionForAdapterUpdates() { var requestManager = mock(ExtensionServiceRequestManager.class); var chartSchemaUpdateCoordinator = mock(ChartSchemaUpdateCoordinator.class); - var coordinator = new PipelineUpdateCoordinator(requestManager, chartSchemaUpdateCoordinator); + var resourceManager = mock(SpResourceManager.class); + var pipelineManager = mock(PipelineManager.class); + var coordinator = new PipelineUpdateCoordinator( + requestManager, resourceManager, chartSchemaUpdateCoordinator, pipelineManager); var adapterDescription = makeAdapter("stream-1", "Updated adapter"); var storedPipeline = makePipeline("pipeline-1", "Pipeline", false, "stream-1", "Old stream"); var modifiedPipeline = makePipeline("pipeline-1", "Pipeline", false, "stream-1", "Updated adapter"); @@ -134,13 +128,12 @@ class PipelineUpdateCoordinatorTest { var warning = PipelineElementValidationInfo.error("Schema mismatch"); var modificationMessage = new PipelineModificationMessage(List.of(invalidModification("sepa-1", warning))); - var noSqlStorage = mock(INoSqlStorage.class); var pipelineStorage = mock(IPipelineStorage.class); var verifiedPipelines = new ArrayList<Pipeline>(); - when(noSqlStorage.getPipelineStorageAPI()).thenReturn(pipelineStorage); + when(pipelineManager.getPipelinesContainingElements("stream-1")).thenReturn(List.of(storedPipeline)); + when(pipelineManager.getPipeline("pipeline-1")).thenReturn(storedPipeline); - try (MockedStatic<PipelineManager> pipelineManager = mockStatic(PipelineManager.class); - MockedConstruction<PipelineVerificationHandlerV2> verificationHandlerConstruction = + try (MockedConstruction<PipelineVerificationHandlerV2> verificationHandlerConstruction = mockConstruction(PipelineVerificationHandlerV2.class, (mock, context) -> { verifiedPipelines.add((Pipeline) context.arguments().get(0)); when(mock.verifyPipeline()).thenReturn(modificationMessage); @@ -153,11 +146,6 @@ class PipelineUpdateCoordinatorTest { MockedConstruction<PipelineExecutor> executorConstruction = mockConstruction(PipelineExecutor.class)) { - pipelineManager.when(() -> PipelineManager.getPipelinesContainingElements("stream-1")) - .thenReturn(List.of(storedPipeline)); - pipelineManager.when(() -> PipelineManager.getPipeline("pipeline-1")) - .thenReturn(storedPipeline); - coordinator.updatePipelines(adapterDescription); assertEquals(1, verificationHandlerConstruction.constructed().size()); @@ -180,7 +168,10 @@ class PipelineUpdateCoordinatorTest { void updatePipelines_ShouldMarkPipelineRequiringAttentionForCriticalMeasurementFieldChange() { var requestManager = mock(ExtensionServiceRequestManager.class); var chartSchemaUpdateCoordinator = mock(ChartSchemaUpdateCoordinator.class); - var coordinator = new PipelineUpdateCoordinator(requestManager, chartSchemaUpdateCoordinator); + var resourceManager = mock(SpResourceManager.class); + var pipelineManager = mock(PipelineManager.class); + var coordinator = new PipelineUpdateCoordinator( + requestManager, resourceManager, chartSchemaUpdateCoordinator, pipelineManager); var adapterDescription = makeAdapter("stream-1", "Updated adapter"); adapterDescription.getDataStream().setEventSchema(makeSchema(makeMeasurementProperty("temperature", XSD.STRING))); @@ -193,9 +184,10 @@ class PipelineUpdateCoordinatorTest { measurementUpdateRequiredMessage()); var modificationMessage = new PipelineModificationMessage(List.of(validModification("sepa-1", measurementUpdateInfo))); var pipelineStorage = mock(IPipelineStorage.class); + when(pipelineManager.getPipelinesContainingElements("stream-1")).thenReturn(List.of(storedPipeline)); + when(pipelineManager.getPipeline("pipeline-1")).thenReturn(storedPipeline); - try (MockedStatic<PipelineManager> pipelineManager = mockStatic(PipelineManager.class); - MockedConstruction<PipelineVerificationHandlerV2> verificationHandlerConstruction = + try (MockedConstruction<PipelineVerificationHandlerV2> verificationHandlerConstruction = mockConstruction(PipelineVerificationHandlerV2.class, (mock, context) -> { when(mock.verifyPipeline()).thenReturn(modificationMessage); when(mock.makeModifiedPipeline(modificationMessage)) @@ -207,11 +199,6 @@ class PipelineUpdateCoordinatorTest { MockedConstruction<PipelineExecutor> executorConstruction = mockConstruction(PipelineExecutor.class)) { - pipelineManager.when(() -> PipelineManager.getPipelinesContainingElements("stream-1")) - .thenReturn(List.of(storedPipeline)); - pipelineManager.when(() -> PipelineManager.getPipeline("pipeline-1")) - .thenReturn(storedPipeline); - coordinator.updatePipelines(adapterDescription); assertEquals(1, verificationHandlerConstruction.constructed().size()); @@ -231,15 +218,20 @@ class PipelineUpdateCoordinatorTest { void checkPipelineMigrations_ShouldUseUpdatedDataStreamValues() { var requestManager = mock(ExtensionServiceRequestManager.class); var chartSchemaUpdateCoordinator = mock(ChartSchemaUpdateCoordinator.class); - var coordinator = new PipelineUpdateCoordinator(requestManager, chartSchemaUpdateCoordinator); + var resourceManager = mock(SpResourceManager.class); + var pipelineManager = mock(PipelineManager.class); + var coordinator = new PipelineUpdateCoordinator( + requestManager, resourceManager, chartSchemaUpdateCoordinator, pipelineManager); 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)); - try (MockedStatic<PipelineManager> pipelineManager = mockStatic(PipelineManager.class); - MockedConstruction<PipelineVerificationHandlerV2> verificationHandlerConstruction = + try (MockedConstruction<PipelineVerificationHandlerV2> verificationHandlerConstruction = mockConstruction(PipelineVerificationHandlerV2.class, (mock, context) -> { verifiedPipelines.add((Pipeline) context.arguments().get(0)); when(mock.verifyPipeline()).thenReturn(modificationMessage); @@ -247,11 +239,6 @@ class PipelineUpdateCoordinatorTest { .thenReturn(new PipelineModificationResult((Pipeline) context.arguments().get(0), List.of())); })) { - pipelineManager.when(() -> PipelineManager.getPipelinesContainingElements("stream-1")) - .thenReturn(List.of(pipeline)); - when(chartSchemaUpdateCoordinator.checkChartMigrations(pipeline, dataStream.getEventSchema())) - .thenReturn(List.of(chartUpdateInfo)); - var result = coordinator.checkPipelineMigrations(dataStream); assertEquals(1, result.size()); @@ -271,16 +258,19 @@ class PipelineUpdateCoordinatorTest { void checkPipelineMigrations_ShouldReportWarningsForAdapterUpdates() { var requestManager = mock(ExtensionServiceRequestManager.class); var chartSchemaUpdateCoordinator = mock(ChartSchemaUpdateCoordinator.class); - var coordinator = new PipelineUpdateCoordinator(requestManager, chartSchemaUpdateCoordinator); + var resourceManager = mock(SpResourceManager.class); + var pipelineManager = mock(PipelineManager.class); + var coordinator = new PipelineUpdateCoordinator( + requestManager, resourceManager, chartSchemaUpdateCoordinator, pipelineManager); var adapterDescription = makeAdapter("stream-1", "Updated adapter"); 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)); - try (MockedStatic<PipelineManager> pipelineManager = mockStatic(PipelineManager.class); - MockedConstruction<PipelineVerificationHandlerV2> verificationHandlerConstruction = + try (MockedConstruction<PipelineVerificationHandlerV2> verificationHandlerConstruction = mockConstruction(PipelineVerificationHandlerV2.class, (mock, context) -> { verifiedPipelines.add((Pipeline) context.arguments().get(0)); when(mock.verifyPipeline()).thenReturn(modificationMessage); @@ -288,9 +278,6 @@ class PipelineUpdateCoordinatorTest { .thenReturn(new PipelineModificationResult((Pipeline) context.arguments().get(0), List.of())); })) { - pipelineManager.when(() -> PipelineManager.getPipelinesContainingElements("stream-1")) - .thenReturn(List.of(pipeline)); - var result = coordinator.checkPipelineMigrations(adapterDescription); assertEquals(1, result.size()); @@ -310,7 +297,10 @@ class PipelineUpdateCoordinatorTest { void checkPipelineMigrations_ShouldDisableAutoMigrationForCriticalMeasurementFieldChange() { var requestManager = mock(ExtensionServiceRequestManager.class); var chartSchemaUpdateCoordinator = mock(ChartSchemaUpdateCoordinator.class); - var coordinator = new PipelineUpdateCoordinator(requestManager, chartSchemaUpdateCoordinator); + var resourceManager = mock(SpResourceManager.class); + var pipelineManager = mock(PipelineManager.class); + var coordinator = new PipelineUpdateCoordinator( + requestManager, resourceManager, chartSchemaUpdateCoordinator, pipelineManager); var adapterDescription = makeAdapter("stream-1", "Updated adapter"); adapterDescription.getDataStream().setEventSchema(makeSchema(makeMeasurementProperty("temperature", XSD.STRING))); @@ -321,18 +311,15 @@ 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)); - try (MockedStatic<PipelineManager> pipelineManager = mockStatic(PipelineManager.class); - MockedConstruction<PipelineVerificationHandlerV2> verificationHandlerConstruction = + try (MockedConstruction<PipelineVerificationHandlerV2> verificationHandlerConstruction = mockConstruction(PipelineVerificationHandlerV2.class, (mock, context) -> { when(mock.verifyPipeline()).thenReturn(modificationMessage); when(mock.makeModifiedPipeline(modificationMessage)) .thenReturn(new PipelineModificationResult((Pipeline) context.arguments().get(0), List.of())); })) { - pipelineManager.when(() -> PipelineManager.getPipelinesContainingElements("stream-1")) - .thenReturn(List.of(pipeline)); - var result = coordinator.checkPipelineMigrations(adapterDescription); assertEquals(1, verificationHandlerConstruction.constructed().size()); diff --git a/streampipes-resource-management/src/test/java/org/apache/streampipes/resource/management/AdapterResourceManagerTest.java b/streampipes-resource-management/src/test/java/org/apache/streampipes/resource/management/AdapterResourceManagerTest.java index ec63f239b9..8e6fbf767f 100644 --- a/streampipes-resource-management/src/test/java/org/apache/streampipes/resource/management/AdapterResourceManagerTest.java +++ b/streampipes-resource-management/src/test/java/org/apache/streampipes/resource/management/AdapterResourceManagerTest.java @@ -43,7 +43,8 @@ public class AdapterResourceManagerTest { @BeforeEach void setUp() { storage = mock(IAdapterStorage.class); - adapterResourceManager = new AdapterResourceManager(storage, null, null); + PermissionResourceManager permissionResourceManager = mock(PermissionResourceManager.class); + adapterResourceManager = new AdapterResourceManager(storage, null, permissionResourceManager); } @Test diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/ResetManagement.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/ResetManagement.java index 50065f2b22..22c98de554 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/ResetManagement.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/ResetManagement.java @@ -127,8 +127,8 @@ public class ResetManagement { } private void stopAndDeleteAllAdapters(WorkerRestClient workerRestClient, - IExtensionsServiceStorage extensionsServiceStorage, - ExtensionServiceRequestManager requestManager) { + IExtensionsServiceStorage extensionsServiceStorage, + ExtensionServiceRequestManager requestManager) { AdapterMasterManagement adapterMasterManagement = new AdapterMasterManagement( StorageDispatcher.INSTANCE.getNoSqlStore() .getAdapterInstanceStorage(),
