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 742ccd31b62c3302aa495549cc96b37385174836 Author: Dominik Riemer <[email protected]> AuthorDate: Wed Jun 17 23:14:37 2026 +0200 refactor: Prepare chart storage for DI --- .../management/AdapterUpdateManagement.java | 6 ++- .../dataexplorer/api/IDataExplorerManager.java | 3 +- .../influx/DataExplorerManagerInflux.java | 9 +++- .../iotdb/DataExplorerManagerIotDb.java | 6 ++- .../management/DataExplorerDispatcher.java | 1 + .../dataexplorer/DataExplorerSchemaManagement.java | 5 -- .../streampipes/export/AssetLinkResolver.java | 9 +++- .../streampipes/export/DataLakeExportManager.java | 14 +++--- .../apache/streampipes/export/ExportManager.java | 14 ++++-- .../apache/streampipes/export/ImportManager.java | 12 +++-- .../export/dataimport/PerformImportGenerator.java | 11 +++-- .../export/dataimport/PreviewImportGenerator.java | 8 +++- .../export/generator/ExportPackageGenerator.java | 12 +++-- .../streampipes/export/resolver/ChartResolver.java | 13 ++++-- .../sinks/internal/jvm/datalake/DataLakeSink.java | 3 +- .../update/ChartSchemaUpdateCoordinator.java | 7 +-- .../update/DataStreamUpdateManagement.java | 5 +- .../update/MeasurementUpdateManagement.java | 6 +-- .../pipeline/update/PipelineUpdateCoordinator.java | 6 +-- .../management/DataExplorerResourceManager.java | 19 ++++++-- .../resource/management/SpResourceManager.java | 4 -- .../apache/streampipes/rest/ResetManagement.java | 53 ++++++++++++++-------- .../streampipes/rest/impl/PipelineResource.java | 11 ++++- .../streampipes/rest/impl/ResetResource.java | 22 ++++----- .../rest/impl/admin/DataExportResource.java | 11 +++-- .../rest/impl/admin/DataImportResource.java | 11 +++-- .../rest/impl/connect/AdapterResource.java | 16 +++++-- .../rest/impl/connect/CompactAdapterResource.java | 8 +++- .../impl/dashboard/DataLakeDashboardResource.java | 7 ++- .../impl/datalake/AbstractDataLakeResource.java | 9 ++-- .../impl/datalake/DataLakeMeasureResource.java | 6 ++- .../rest/impl/datalake/DataLakeResource.java | 14 +++--- .../rest/impl/datalake/DataLakeWidgetResource.java | 7 +-- .../datalake/KioskDashboardDataLakeResource.java | 14 +++--- .../datalake/importer/DataLakeImportResource.java | 6 ++- .../rest/impl/pe/DataStreamResource.java | 9 +++- .../service/core/StreamPipesCoreApplication.java | 6 ++- .../core/migrations/AvailableMigrations.java | 9 +++- .../v0980/FixImportedPermissionsMigration.java | 12 +++-- .../service/core/scheduler/DataLakeScheduler.java | 21 +++++++-- .../core/storage/StorageApiConfiguration.java | 7 +++ .../storage/api/core/INoSqlStorage.java | 3 -- .../storage/couchdb/CouchDbStorageManager.java | 5 -- 43 files changed, 275 insertions(+), 165 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 9a68ba38e3..094272aa00 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 @@ -20,6 +20,7 @@ package org.apache.streampipes.connect.management.management; import org.apache.streampipes.commons.exceptions.connect.AdapterException; import org.apache.streampipes.manager.api.extensions.ExtensionServiceRequestManager; +import org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator; import org.apache.streampipes.manager.pipeline.update.PipelineUpdateCoordinator; import org.apache.streampipes.model.SpDataStream; import org.apache.streampipes.model.connect.adapter.AdapterDescription; @@ -38,11 +39,12 @@ public class AdapterUpdateManagement { private final PipelineUpdateCoordinator pipelineUpdateCoordinator; public AdapterUpdateManagement(AdapterMasterManagement adapterMasterManagement, - ExtensionServiceRequestManager requestManager) { + ExtensionServiceRequestManager requestManager, + ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator) { this.adapterMasterManagement = adapterMasterManagement; this.adapterResourceManager = new SpResourceManager().manageAdapters(); this.dataStreamResourceManager = new SpResourceManager().manageDataStreams(); - this.pipelineUpdateCoordinator = new PipelineUpdateCoordinator(requestManager); + this.pipelineUpdateCoordinator = new PipelineUpdateCoordinator(requestManager, chartSchemaUpdateCoordinator); } public void updateAdapter(AdapterDescription ad) diff --git a/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataExplorerManager.java b/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataExplorerManager.java index 5711aade33..3de231ca5f 100644 --- a/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataExplorerManager.java +++ b/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataExplorerManager.java @@ -19,6 +19,7 @@ package org.apache.streampipes.dataexplorer.api; import org.apache.streampipes.client.api.IStreamPipesClient; +import org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator; import org.apache.streampipes.model.datalake.DataLakeMeasure; import java.util.List; @@ -41,7 +42,7 @@ public interface IDataExplorerManager { IDataExplorerQueryManagement getQueryManagement(IDataExplorerSchemaManagement dataExplorerSchemaManagement); - IDataExplorerSchemaManagement getSchemaManagement(); + IDataExplorerSchemaManagement getSchemaManagement(ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator); default ITimeSeriesStorage getTimeseriesStorage(DataLakeMeasure measure) { return getTimeseriesStorage(measure, false); diff --git a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerManagerInflux.java b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerManagerInflux.java index 5ec67f0db9..8d600847d5 100644 --- a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerManagerInflux.java +++ b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerManagerInflux.java @@ -30,6 +30,7 @@ import org.apache.streampipes.dataexplorer.api.ITimeSeriesStorage; import org.apache.streampipes.dataexplorer.influx.client.InfluxClientProvider; import org.apache.streampipes.dataexplorer.influx.sanitize.DataLakeMeasurementSanitizerInflux; import org.apache.streampipes.manager.permission.DataLakePermissionManager; +import org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator; import org.apache.streampipes.model.datalake.DataLakeMeasure; import org.apache.streampipes.storage.management.StorageDispatcher; @@ -37,6 +38,9 @@ import java.util.List; public class DataExplorerManagerInflux implements IDataExplorerManager { + public DataExplorerManagerInflux() { + } + @Override public IDataLakeMeasurementCounter getMeasurementCounter( List<DataLakeMeasure> allMeasurements, @@ -53,13 +57,14 @@ public class DataExplorerManagerInflux implements IDataExplorerManager { } @Override - public IDataExplorerSchemaManagement getSchemaManagement() { + public IDataExplorerSchemaManagement getSchemaManagement(ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator) { return new DataExplorerSchemaManagement(StorageDispatcher.INSTANCE .getNoSqlStore() .getDataLakeStorage(), new DataLakePermissionManager( StorageDispatcher.INSTANCE.getNoSqlStore().getPermissionStorage() - ) + ), + chartSchemaUpdateCoordinator ); } diff --git a/streampipes-data-explorer-iotdb/src/main/java/org/apache/streampipes/dataexplorer/iotdb/DataExplorerManagerIotDb.java b/streampipes-data-explorer-iotdb/src/main/java/org/apache/streampipes/dataexplorer/iotdb/DataExplorerManagerIotDb.java index ba0f060b94..4db64664d7 100644 --- a/streampipes-data-explorer-iotdb/src/main/java/org/apache/streampipes/dataexplorer/iotdb/DataExplorerManagerIotDb.java +++ b/streampipes-data-explorer-iotdb/src/main/java/org/apache/streampipes/dataexplorer/iotdb/DataExplorerManagerIotDb.java @@ -29,6 +29,7 @@ import org.apache.streampipes.dataexplorer.api.IDataLakeMeasurementSanitizer; import org.apache.streampipes.dataexplorer.api.ITimeSeriesStorage; import org.apache.streampipes.dataexplorer.iotdb.sanitize.DataLakeMeasurementSanitizerIotDb; import org.apache.streampipes.manager.permission.DataLakePermissionManager; +import org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator; import org.apache.streampipes.model.datalake.DataLakeMeasure; import org.apache.streampipes.storage.management.StorageDispatcher; @@ -52,13 +53,14 @@ public class DataExplorerManagerIotDb implements IDataExplorerManager { } @Override - public IDataExplorerSchemaManagement getSchemaManagement() { + public IDataExplorerSchemaManagement getSchemaManagement(ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator) { return new DataExplorerSchemaManagement(StorageDispatcher.INSTANCE .getNoSqlStore() .getDataLakeStorage(), new DataLakePermissionManager( StorageDispatcher.INSTANCE.getNoSqlStore().getPermissionStorage() - ) + ), + chartSchemaUpdateCoordinator ); } diff --git a/streampipes-data-explorer-management/src/main/java/org/apache/streampipes/dataexplorer/management/DataExplorerDispatcher.java b/streampipes-data-explorer-management/src/main/java/org/apache/streampipes/dataexplorer/management/DataExplorerDispatcher.java index 8614c15b57..10b810134f 100644 --- a/streampipes-data-explorer-management/src/main/java/org/apache/streampipes/dataexplorer/management/DataExplorerDispatcher.java +++ b/streampipes-data-explorer-management/src/main/java/org/apache/streampipes/dataexplorer/management/DataExplorerDispatcher.java @@ -22,6 +22,7 @@ import org.apache.streampipes.commons.environment.Environments; import org.apache.streampipes.dataexplorer.api.IDataExplorerManager; import org.apache.streampipes.dataexplorer.influx.DataExplorerManagerInflux; import org.apache.streampipes.dataexplorer.iotdb.DataExplorerManagerIotDb; +import org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator; public class DataExplorerDispatcher { diff --git a/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/DataExplorerSchemaManagement.java b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/DataExplorerSchemaManagement.java index a833e41860..fa44455a03 100644 --- a/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/DataExplorerSchemaManagement.java +++ b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/DataExplorerSchemaManagement.java @@ -44,11 +44,6 @@ public class DataExplorerSchemaManagement implements IDataExplorerSchemaManageme private final ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator; public DataExplorerSchemaManagement(CRUDStorage<DataLakeMeasure> dataLakeStorage, - DataLakePermissionManager permissionManager) { - this(dataLakeStorage, permissionManager, new ChartSchemaUpdateCoordinator()); - } - - DataExplorerSchemaManagement(CRUDStorage<DataLakeMeasure> dataLakeStorage, DataLakePermissionManager permissionManager, ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator) { this.dataLakeStorage = dataLakeStorage; diff --git a/streampipes-data-export/src/main/java/org/apache/streampipes/export/AssetLinkResolver.java b/streampipes-data-export/src/main/java/org/apache/streampipes/export/AssetLinkResolver.java index e767967ed9..708e680938 100644 --- a/streampipes-data-export/src/main/java/org/apache/streampipes/export/AssetLinkResolver.java +++ b/streampipes-data-export/src/main/java/org/apache/streampipes/export/AssetLinkResolver.java @@ -31,6 +31,7 @@ import org.apache.streampipes.model.assets.AssetLink; import org.apache.streampipes.model.assets.SpAssetModel; import org.apache.streampipes.model.export.AssetExportConfiguration; import org.apache.streampipes.serializers.json.JacksonSerializer; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import org.apache.streampipes.storage.management.StorageDispatcher; import com.fasterxml.jackson.databind.DeserializationFeature; @@ -51,11 +52,14 @@ public class AssetLinkResolver { private final String assetId; private final ObjectMapper mapper; private final ExtensionServiceRequestManager extensionServiceRequestManager; + private final IDataExplorerWidgetStorage chartStorage; public AssetLinkResolver(String assetId, - ExtensionServiceRequestManager extensionServiceRequestManager) { + ExtensionServiceRequestManager extensionServiceRequestManager, + IDataExplorerWidgetStorage chartStorage) { this.assetId = assetId; this.extensionServiceRequestManager = extensionServiceRequestManager; + this.chartStorage = chartStorage; this.mapper = JacksonSerializer.getObjectMapper(Map.of( DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, true, SerializationFeature.INDENT_OUTPUT, false @@ -72,7 +76,8 @@ public class AssetLinkResolver { exportConfig.setAssetName(asset.getAssetName()); exportConfig.setAdapters(new AdapterResolver(extensionServiceRequestManager) .resolve(getLinks(assetLinks, ResolvableAssetLinks.ADAPTER))); - exportConfig.setDataViews(new ChartResolver().resolve(getLinks(assetLinks, ResolvableAssetLinks.CHART))); + exportConfig.setDataViews(new ChartResolver(chartStorage) + .resolve(getLinks(assetLinks, ResolvableAssetLinks.CHART))); exportConfig.setDashboards(new DashboardResolver().resolve(getLinks(assetLinks, ResolvableAssetLinks.DASHBOARD))); exportConfig.setDataSources( new DataSourceResolver().resolve(getLinks(assetLinks, ResolvableAssetLinks.DATA_SOURCE))); diff --git a/streampipes-data-export/src/main/java/org/apache/streampipes/export/DataLakeExportManager.java b/streampipes-data-export/src/main/java/org/apache/streampipes/export/DataLakeExportManager.java index c865bad32b..a0274f5800 100644 --- a/streampipes-data-export/src/main/java/org/apache/streampipes/export/DataLakeExportManager.java +++ b/streampipes-data-export/src/main/java/org/apache/streampipes/export/DataLakeExportManager.java @@ -24,7 +24,6 @@ import org.apache.streampipes.dataexplorer.api.IDataExplorerSchemaManagement; import org.apache.streampipes.dataexplorer.export.OutputFormat; import org.apache.streampipes.dataexplorer.export.objectstorage.ExportProviderFactory; import org.apache.streampipes.dataexplorer.export.objectstorage.IObjectStorage; -import org.apache.streampipes.dataexplorer.management.DataExplorerDispatcher; import org.apache.streampipes.model.configuration.ExportProviderSettings; import org.apache.streampipes.model.configuration.ProviderType; import org.apache.streampipes.model.datalake.DataLakeMeasure; @@ -49,13 +48,14 @@ public class DataLakeExportManager { private static final Logger LOG = LoggerFactory.getLogger(DataLakeExportManager.class); private static final Environment env = Environments.getEnvironment(); - private final IDataExplorerSchemaManagement dataExplorerSchemaManagement = new DataExplorerDispatcher() - .getDataExplorerManager() - .getSchemaManagement(); + private final IDataExplorerSchemaManagement dataExplorerSchemaManagement; + private final IDataExplorerQueryManagement dataExplorerQueryManagement; - private final IDataExplorerQueryManagement dataExplorerQueryManagement = new DataExplorerDispatcher() - .getDataExplorerManager() - .getQueryManagement(this.dataExplorerSchemaManagement); + public DataLakeExportManager(IDataExplorerSchemaManagement dataLakeSchemaManagement, + IDataExplorerQueryManagement dataLakeQueryManagement) { + this.dataExplorerSchemaManagement = dataLakeSchemaManagement; + this.dataExplorerQueryManagement = dataLakeQueryManagement; + } private String savePath = ""; diff --git a/streampipes-data-export/src/main/java/org/apache/streampipes/export/ExportManager.java b/streampipes-data-export/src/main/java/org/apache/streampipes/export/ExportManager.java index 37ccdeaef8..4254ee693a 100644 --- a/streampipes-data-export/src/main/java/org/apache/streampipes/export/ExportManager.java +++ b/streampipes-data-export/src/main/java/org/apache/streampipes/export/ExportManager.java @@ -21,6 +21,7 @@ package org.apache.streampipes.export; import org.apache.streampipes.export.generator.ExportPackageGenerator; import org.apache.streampipes.manager.api.extensions.ExtensionServiceRequestManager; import org.apache.streampipes.model.export.ExportConfiguration; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import java.io.IOException; import java.util.List; @@ -29,11 +30,13 @@ import java.util.stream.Collectors; public class ExportManager { public static ExportConfiguration getExportPreview(List<String> selectedAssetIds, - ExtensionServiceRequestManager extensionServiceRequestManager) { + ExtensionServiceRequestManager extensionServiceRequestManager, + IDataExplorerWidgetStorage chartStorage) { var exportConfig = new ExportConfiguration(); var assetExportConfigurations = selectedAssetIds .stream() - .map(assetId -> new AssetLinkResolver(assetId, extensionServiceRequestManager).resolveResources()) + .map(assetId -> new AssetLinkResolver(assetId, extensionServiceRequestManager, chartStorage) + .resolveResources()) .collect(Collectors.toList()); exportConfig.setAssetExportConfiguration(assetExportConfigurations); @@ -42,9 +45,10 @@ public class ExportManager { } public static byte[] getExportPackage(ExportConfiguration exportConfiguration, - ExtensionServiceRequestManager extensionServiceRequestManager) - throws IOException { - return new ExportPackageGenerator(exportConfiguration, extensionServiceRequestManager).generateExportPackage(); + ExtensionServiceRequestManager extensionServiceRequestManager, + IDataExplorerWidgetStorage chartStorage) throws IOException { + return new ExportPackageGenerator(exportConfiguration, extensionServiceRequestManager, chartStorage) + .generateExportPackage(); } } diff --git a/streampipes-data-export/src/main/java/org/apache/streampipes/export/ImportManager.java b/streampipes-data-export/src/main/java/org/apache/streampipes/export/ImportManager.java index b98c4638f4..7802525cd4 100644 --- a/streampipes-data-export/src/main/java/org/apache/streampipes/export/ImportManager.java +++ b/streampipes-data-export/src/main/java/org/apache/streampipes/export/ImportManager.java @@ -22,6 +22,7 @@ import org.apache.streampipes.export.dataimport.PerformImportGenerator; import org.apache.streampipes.export.dataimport.PreviewImportGenerator; import org.apache.streampipes.manager.api.extensions.ExtensionServiceRequestManager; import org.apache.streampipes.model.export.AssetExportConfiguration; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import java.io.IOException; import java.io.InputStream; @@ -29,15 +30,18 @@ import java.io.InputStream; public class ImportManager { public static AssetExportConfiguration getImportPreview(InputStream packageZipStream, - ExtensionServiceRequestManager extensionServiceRequestManager) + ExtensionServiceRequestManager extensionServiceRequestManager, + IDataExplorerWidgetStorage chartStorage) throws IOException { - return new PreviewImportGenerator(extensionServiceRequestManager).generate(packageZipStream); + return new PreviewImportGenerator(extensionServiceRequestManager, chartStorage).generate(packageZipStream); } public static void performImport(InputStream packageZipStream, AssetExportConfiguration exportConfiguration, String ownerSid, - ExtensionServiceRequestManager extensionServiceRequestManager) throws IOException { - new PerformImportGenerator(exportConfiguration, ownerSid, extensionServiceRequestManager).generate(packageZipStream); + ExtensionServiceRequestManager extensionServiceRequestManager, + IDataExplorerWidgetStorage chartStorage) throws IOException { + new PerformImportGenerator(exportConfiguration, ownerSid, extensionServiceRequestManager, chartStorage) + .generate(packageZipStream); } } diff --git a/streampipes-data-export/src/main/java/org/apache/streampipes/export/dataimport/PerformImportGenerator.java b/streampipes-data-export/src/main/java/org/apache/streampipes/export/dataimport/PerformImportGenerator.java index 0e4532deec..b9b63c2e67 100644 --- a/streampipes-data-export/src/main/java/org/apache/streampipes/export/dataimport/PerformImportGenerator.java +++ b/streampipes-data-export/src/main/java/org/apache/streampipes/export/dataimport/PerformImportGenerator.java @@ -40,6 +40,7 @@ import org.apache.streampipes.model.export.ExportItem; import org.apache.streampipes.model.pipeline.Pipeline; import org.apache.streampipes.resource.management.PermissionResourceManager; import org.apache.streampipes.storage.api.core.INoSqlStorage; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import org.apache.streampipes.storage.management.StorageDispatcher; import com.fasterxml.jackson.core.JsonProcessingException; @@ -57,14 +58,17 @@ public class PerformImportGenerator extends ImportGenerator<Void> { private final Set<PermissionInfo> permissionsToStore = new HashSet<>(); private final String ownerSid; private final ExtensionServiceRequestManager extensionServiceRequestManager; + private final IDataExplorerWidgetStorage chartStorage; public PerformImportGenerator(AssetExportConfiguration config, String ownerSid, - ExtensionServiceRequestManager extensionServiceRequestManager) { + ExtensionServiceRequestManager extensionServiceRequestManager, + IDataExplorerWidgetStorage chartStorage) { this.config = config; this.storage = StorageDispatcher.INSTANCE.getNoSqlStore(); this.ownerSid = ownerSid; this.extensionServiceRequestManager = extensionServiceRequestManager; + this.chartStorage = chartStorage; } @Override @@ -93,8 +97,9 @@ public class PerformImportGenerator extends ImportGenerator<Void> { @Override protected void handleChart(String document, String chartId) throws JsonProcessingException { if (shouldStore(chartId, config.getDataViews())) { - writeDocument(document, new ChartResolver()); - var chart = new ChartResolver().deserializeDocument(document); + var chartResolver = new ChartResolver(chartStorage); + writeDocument(document, chartResolver); + var chart = chartResolver.deserializeDocument(document); permissionsToStore.add(new PermissionInfo(chart.getElementId(), DataExplorerWidgetModel.class)); } } diff --git a/streampipes-data-export/src/main/java/org/apache/streampipes/export/dataimport/PreviewImportGenerator.java b/streampipes-data-export/src/main/java/org/apache/streampipes/export/dataimport/PreviewImportGenerator.java index de8a0c976c..d440467d8a 100644 --- a/streampipes-data-export/src/main/java/org/apache/streampipes/export/dataimport/PreviewImportGenerator.java +++ b/streampipes-data-export/src/main/java/org/apache/streampipes/export/dataimport/PreviewImportGenerator.java @@ -29,6 +29,7 @@ import org.apache.streampipes.export.resolver.PipelineResolver; import org.apache.streampipes.manager.api.extensions.ExtensionServiceRequestManager; import org.apache.streampipes.model.export.AssetExportConfiguration; import org.apache.streampipes.model.export.ExportItem; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.core.type.TypeReference; @@ -43,11 +44,14 @@ public class PreviewImportGenerator extends ImportGenerator<AssetExportConfigura private static final Logger LOG = LoggerFactory.getLogger(PreviewImportGenerator.class); private final AssetExportConfiguration importConfig; private final ExtensionServiceRequestManager extensionServiceRequestManager; + private final IDataExplorerWidgetStorage chartStorage; - public PreviewImportGenerator(ExtensionServiceRequestManager extensionServiceRequestManager) { + public PreviewImportGenerator(ExtensionServiceRequestManager extensionServiceRequestManager, + IDataExplorerWidgetStorage chartStorage) { super(); this.importConfig = new AssetExportConfiguration(); this.extensionServiceRequestManager = extensionServiceRequestManager; + this.chartStorage = chartStorage; } @@ -85,7 +89,7 @@ public class PreviewImportGenerator extends ImportGenerator<AssetExportConfigura protected void handleChart(String document, String dataViewId) throws JsonProcessingException { addExportItem( dataViewId, - new ChartResolver().readDocument(document).getBaseAppearanceConfig().get("widgetTitle").toString(), + new ChartResolver(chartStorage).readDocument(document).getBaseAppearanceConfig().get("widgetTitle").toString(), importConfig::addDataView ); } diff --git a/streampipes-data-export/src/main/java/org/apache/streampipes/export/generator/ExportPackageGenerator.java b/streampipes-data-export/src/main/java/org/apache/streampipes/export/generator/ExportPackageGenerator.java index 1bc2829977..99944160b7 100644 --- a/streampipes-data-export/src/main/java/org/apache/streampipes/export/generator/ExportPackageGenerator.java +++ b/streampipes-data-export/src/main/java/org/apache/streampipes/export/generator/ExportPackageGenerator.java @@ -35,6 +35,7 @@ import org.apache.streampipes.model.export.ExportConfiguration; import org.apache.streampipes.model.export.ExportItem; import org.apache.streampipes.model.export.StreamPipesApplicationPackage; import org.apache.streampipes.serializers.json.JacksonSerializer; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import org.apache.streampipes.storage.management.StorageDispatcher; import com.fasterxml.jackson.core.JsonProcessingException; @@ -57,12 +58,15 @@ public class ExportPackageGenerator { private final ExportConfiguration exportConfiguration; private final ExtensionServiceRequestManager extensionServiceRequestManager; - private ObjectMapper defaultMapper; + private final ObjectMapper defaultMapper; + private final IDataExplorerWidgetStorage chartStorage; public ExportPackageGenerator(ExportConfiguration exportConfiguration, - ExtensionServiceRequestManager extensionServiceRequestManager) { + ExtensionServiceRequestManager extensionServiceRequestManager, + IDataExplorerWidgetStorage chartStorage) { this.exportConfiguration = exportConfiguration; this.extensionServiceRequestManager = extensionServiceRequestManager; + this.chartStorage = chartStorage; this.defaultMapper = JacksonSerializer.getObjectMapper(Map.of( DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, true, SerializationFeature.INDENT_OUTPUT, false @@ -107,14 +111,14 @@ public class ExportPackageGenerator { new DashboardResolver(), manifest::addDashboard); var charts = resolver.getCharts(item.getResourceId()); - var chartResolver = new ChartResolver(); + var chartResolver = new ChartResolver(chartStorage); charts.forEach(widgetId -> addDoc(builder, widgetId, chartResolver, manifest::addDataViewWidget)); }); config.getDataViews().forEach(item -> { addDoc(builder, item, - new ChartResolver(), + new ChartResolver(chartStorage), manifest::addDataView); }); diff --git a/streampipes-data-export/src/main/java/org/apache/streampipes/export/resolver/ChartResolver.java b/streampipes-data-export/src/main/java/org/apache/streampipes/export/resolver/ChartResolver.java index 74e417ac5d..d26996f856 100644 --- a/streampipes-data-export/src/main/java/org/apache/streampipes/export/resolver/ChartResolver.java +++ b/streampipes-data-export/src/main/java/org/apache/streampipes/export/resolver/ChartResolver.java @@ -22,6 +22,7 @@ import org.apache.streampipes.model.datalake.DataExplorerWidgetModel; import org.apache.streampipes.model.export.AssetExportConfiguration; import org.apache.streampipes.model.export.ExportItem; import org.apache.streampipes.serializers.json.JacksonSerializer; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.DeserializationFeature; @@ -30,9 +31,15 @@ import java.util.Map; public class ChartResolver extends AbstractResolver<DataExplorerWidgetModel> { + private final IDataExplorerWidgetStorage chartStorage; + + public ChartResolver(IDataExplorerWidgetStorage chartStorage) { + this.chartStorage = chartStorage; + } + @Override public DataExplorerWidgetModel findDocument(String resourceId) { - return getNoSqlStore().getDataExplorerWidgetStorage().getElementById(resourceId); + return chartStorage.getElementById(resourceId); } @Override @@ -58,7 +65,7 @@ public class ChartResolver extends AbstractResolver<DataExplorerWidgetModel> { @Override public void writeDocument(String document, AssetExportConfiguration config) throws JsonProcessingException { - getNoSqlStore().getDataExplorerWidgetStorage().persist(deserializeDocument(document)); + chartStorage.persist(deserializeDocument(document)); } @Override @@ -70,7 +77,7 @@ public class ChartResolver extends AbstractResolver<DataExplorerWidgetModel> { public void deleteDocument(String document) throws JsonProcessingException { var chart = readDocument(document); var resourceId = chart.getElementId(); - getNoSqlStore().getDataExplorerWidgetStorage().deleteElementById(resourceId); + chartStorage.deleteElementById(resourceId); } } 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 40230d9e5a..f0e1812556 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 @@ -166,8 +166,9 @@ public class DataLakeSink implements IStreamPipesDataSink, SupportsRuntimeConfig private RetentionTimeConfig getRetentionTime(String measureName, IStreamPipesClient client){ + // TODO IDataExplorerSchemaManagement dataExplorerSchemaManagement = new DataExplorerDispatcher().getDataExplorerManager() - .getSchemaManagement(); + .getSchemaManagement(null); var originalMeasure = dataExplorerSchemaManagement.getExistingMeasureByName(measureName); diff --git a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/ChartSchemaUpdateCoordinator.java b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/ChartSchemaUpdateCoordinator.java index b7d2272f50..1e26ea6e20 100644 --- a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/ChartSchemaUpdateCoordinator.java +++ b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/ChartSchemaUpdateCoordinator.java @@ -27,7 +27,6 @@ import org.apache.streampipes.model.schema.EventProperty; import org.apache.streampipes.model.schema.EventSchema; import org.apache.streampipes.serializers.json.JacksonSerializer; import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; -import org.apache.streampipes.storage.management.StorageDispatcher; import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.databind.ObjectMapper; @@ -56,11 +55,7 @@ public class ChartSchemaUpdateCoordinator { private final IDataExplorerWidgetStorage widgetStorage; private final ObjectMapper objectMapper; - public ChartSchemaUpdateCoordinator() { - this(StorageDispatcher.INSTANCE.getNoSqlStore().getDataExplorerWidgetStorage()); - } - - ChartSchemaUpdateCoordinator(IDataExplorerWidgetStorage widgetStorage) { + public ChartSchemaUpdateCoordinator(IDataExplorerWidgetStorage widgetStorage) { this.widgetStorage = widgetStorage; this.objectMapper = JacksonSerializer.getObjectMapper(); } 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 f4f51d2ee1..452b0e7678 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 @@ -31,9 +31,10 @@ public class DataStreamUpdateManagement { private final DataStreamResourceManager dataStreamResourceManager; private final PipelineUpdateCoordinator pipelineUpdateCoordinator; - public DataStreamUpdateManagement(ExtensionServiceRequestManager requestManager) { + public DataStreamUpdateManagement(ExtensionServiceRequestManager requestManager, + ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator) { this.dataStreamResourceManager = new SpResourceManager().manageDataStreams(); - this.pipelineUpdateCoordinator = new PipelineUpdateCoordinator(requestManager); + this.pipelineUpdateCoordinator = new PipelineUpdateCoordinator(requestManager, chartSchemaUpdateCoordinator); } public void updateDataStream(SpDataStream dataStream) { diff --git a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/MeasurementUpdateManagement.java b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/MeasurementUpdateManagement.java index 5d253a6738..e5cf7a078c 100644 --- a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/MeasurementUpdateManagement.java +++ b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/pipeline/update/MeasurementUpdateManagement.java @@ -36,11 +36,7 @@ public class MeasurementUpdateManagement { private final ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator; - public MeasurementUpdateManagement() { - this(new ChartSchemaUpdateCoordinator()); - } - - MeasurementUpdateManagement(ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator) { + public MeasurementUpdateManagement(ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator) { this.chartSchemaUpdateCoordinator = chartSchemaUpdateCoordinator; } 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 c71d241607..5a33381e50 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 @@ -52,11 +52,7 @@ public class PipelineUpdateCoordinator { private final ExtensionServiceRequestManager requestManager; private final ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator; - public PipelineUpdateCoordinator(ExtensionServiceRequestManager requestManager) { - this(requestManager, new ChartSchemaUpdateCoordinator()); - } - - PipelineUpdateCoordinator(ExtensionServiceRequestManager requestManager, + public PipelineUpdateCoordinator(ExtensionServiceRequestManager requestManager, ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator) { this.requestManager = requestManager; this.chartSchemaUpdateCoordinator = chartSchemaUpdateCoordinator; diff --git a/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/DataExplorerResourceManager.java b/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/DataExplorerResourceManager.java index 2a5d2fc65d..f2ca5eff70 100644 --- a/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/DataExplorerResourceManager.java +++ b/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/DataExplorerResourceManager.java @@ -22,6 +22,7 @@ import org.apache.streampipes.model.dashboard.DashboardModel; import org.apache.streampipes.model.dashboard.DashboardSummaryDto; import org.apache.streampipes.model.datalake.DataExplorerWidgetModel; import org.apache.streampipes.model.resource.ResourceSummaryDto; +import org.apache.streampipes.storage.api.explorer.IDataExplorerDashboardStorage; import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import org.apache.streampipes.storage.api.explorer.IDataLakeMeasureStorage; import org.apache.streampipes.storage.management.StorageDispatcher; @@ -36,10 +37,20 @@ public class DataExplorerResourceManager extends CrudResourceManager<DashboardMo private final IDataExplorerWidgetStorage widgetStorage; private final IDataLakeMeasureStorage dataLakeMeasureStorage; - public DataExplorerResourceManager() { - super(StorageDispatcher.INSTANCE.getNoSqlStore().getDataExplorerDashboardStorage(), DashboardModel.class); - this.widgetStorage = StorageDispatcher.INSTANCE.getNoSqlStore().getDataExplorerWidgetStorage(); - this.dataLakeMeasureStorage = StorageDispatcher.INSTANCE.getNoSqlStore().getDataLakeStorage(); + public DataExplorerResourceManager(IDataExplorerWidgetStorage widgetStorage) { + this( + StorageDispatcher.INSTANCE.getNoSqlStore().getDataExplorerDashboardStorage(), + widgetStorage, + StorageDispatcher.INSTANCE.getNoSqlStore().getDataLakeStorage() + ); + } + + private DataExplorerResourceManager(IDataExplorerDashboardStorage dashboardStorage, + IDataExplorerWidgetStorage widgetStorage, + IDataLakeMeasureStorage dataLakeMeasureStorage) { + super(dashboardStorage, DashboardModel.class); + this.widgetStorage = widgetStorage; + this.dataLakeMeasureStorage = dataLakeMeasureStorage; } public ResourceSummaryDto<DashboardSummaryDto> getSummary(Authentication auth) { diff --git a/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/SpResourceManager.java b/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/SpResourceManager.java index 294f9ddb1b..4c17a82856 100644 --- a/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/SpResourceManager.java +++ b/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/SpResourceManager.java @@ -33,10 +33,6 @@ public class SpResourceManager { return new AdapterDescriptionResourceManager(); } - public DataExplorerResourceManager manageDataExplorer() { - return new DataExplorerResourceManager(); - } - public DataExplorerWidgetResourceManager manageDataExplorerWidget(DataExplorerResourceManager dashboardManager, IDataExplorerWidgetStorage db) { return new DataExplorerWidgetResourceManager(dashboardManager, db); 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 61eaa3bbad..7d09e19c6a 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 @@ -29,12 +29,14 @@ import org.apache.streampipes.manager.file.FileManager; import org.apache.streampipes.manager.pipeline.PipelineCacheManager; import org.apache.streampipes.manager.pipeline.PipelineCanvasMetadataCacheManager; import org.apache.streampipes.manager.pipeline.PipelineManager; +import org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator; import org.apache.streampipes.model.connect.adapter.AdapterDescription; import org.apache.streampipes.model.datalake.DataLakeMeasure; import org.apache.streampipes.model.file.FileMetadata; import org.apache.streampipes.model.pipeline.Pipeline; import org.apache.streampipes.resource.management.SpResourceManager; import org.apache.streampipes.resource.management.UserResourceManager; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import org.apache.streampipes.storage.api.system.IExtensionsServiceStorage; import org.apache.streampipes.storage.api.system.IGenericStorage; import org.apache.streampipes.storage.management.StorageDispatcher; @@ -51,6 +53,23 @@ public class ResetManagement { // dependency between this package and streampipes-pipeline-management // See in issue [STREAMPIPES-405] + private final IDataExplorerWidgetStorage widgetStorage; + private final WorkerRestClient workerRestClient; + private final IExtensionsServiceStorage extensionsServiceStorage; + private final ExtensionServiceRequestManager requestManager; + private final ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator; + + public ResetManagement(IDataExplorerWidgetStorage widgetStorage, + WorkerRestClient workerRestClient, + IExtensionsServiceStorage extensionsServiceStorage, + ExtensionServiceRequestManager requestManager) { + this.widgetStorage = widgetStorage; + this.workerRestClient = workerRestClient; + this.extensionsServiceStorage = extensionsServiceStorage; + this.requestManager = requestManager; + this.chartSchemaUpdateCoordinator = new ChartSchemaUpdateCoordinator(widgetStorage); + } + private static final Logger logger = LoggerFactory.getLogger(ResetManagement.class); /** @@ -59,10 +78,7 @@ public class ResetManagement { * * @param username of the user to delte the resources */ - public static void reset(String username, - WorkerRestClient workerRestClient, - IExtensionsServiceStorage extensionsServiceStorage, - ExtensionServiceRequestManager requestManager) { + public void reset(String username) { logger.info("Start resetting the system"); setHideTutorialToFalse(username); @@ -77,7 +93,7 @@ public class ResetManagement { removeAllDataInDataLake(); - removeAllDataViewWidgets(); + removeAllDataViewWidgets(widgetStorage); removeAllDataViews(); @@ -90,16 +106,16 @@ public class ResetManagement { logger.info("Resetting the system was completed"); } - private static void setHideTutorialToFalse(String username) { + private void setHideTutorialToFalse(String username) { UserResourceManager.setHideTutorial(username, true); } - private static void clearPipelineAssemblyCache(String username) { + private void clearPipelineAssemblyCache(String username) { PipelineCacheManager.removeCachedPipeline(username); PipelineCanvasMetadataCacheManager.removeCanvasMetadataFromCache(username); } - private static void stopAndDeleteAllPipelines(ExtensionServiceRequestManager requestManager) { + private void stopAndDeleteAllPipelines(ExtensionServiceRequestManager requestManager) { List<Pipeline> allPipelines = PipelineManager.getAllPipelines(); allPipelines.forEach(pipeline -> { PipelineManager.stopPipeline(pipeline.getPipelineId(), true, requestManager); @@ -107,7 +123,7 @@ public class ResetManagement { }); } - private static void stopAndDeleteAllAdapters(WorkerRestClient workerRestClient, + private void stopAndDeleteAllAdapters(WorkerRestClient workerRestClient, IExtensionsServiceStorage extensionsServiceStorage, ExtensionServiceRequestManager requestManager) { AdapterMasterManagement adapterMasterManagement = new AdapterMasterManagement( @@ -131,16 +147,16 @@ public class ResetManagement { }); } - private static void deleteAllFiles() { + private void deleteAllFiles() { var fileManager = new FileManager(); List<FileMetadata> allFiles = fileManager.getAllFiles(); allFiles.forEach(fileMetadata -> fileManager.deleteFile(fileMetadata.getFileId())); } - private static void removeAllDataInDataLake() { + private void removeAllDataInDataLake() { var dataLakeMeasureManagement = new DataExplorerDispatcher() .getDataExplorerManager() - .getSchemaManagement(); + .getSchemaManagement(chartSchemaUpdateCoordinator); var dataExplorerQueryManagement = new DataExplorerDispatcher() .getDataExplorerManager() .getQueryManagement(dataLakeMeasureManagement); @@ -154,15 +170,12 @@ public class ResetManagement { }); } - private static void removeAllDataViewWidgets() { - var widgetStorage = - StorageDispatcher.INSTANCE.getNoSqlStore() - .getDataExplorerWidgetStorage(); + private void removeAllDataViewWidgets(IDataExplorerWidgetStorage widgetStorage) { widgetStorage.findAll() .forEach(widget -> widgetStorage.deleteElementById(widget.getElementId())); } - private static void removeAllDataViews() { + private void removeAllDataViews() { var dataLakeDashboardStorage = StorageDispatcher.INSTANCE.getNoSqlStore() .getDataExplorerDashboardStorage(); @@ -170,7 +183,7 @@ public class ResetManagement { .forEach(dashboard -> dataLakeDashboardStorage.deleteElementById(dashboard.getElementId())); } - private static void removeAllAssets(String username) { + private void removeAllAssets(String username) { IGenericStorage genericStorage = StorageDispatcher.INSTANCE.getNoSqlStore() .getGenericStorage(); try { @@ -182,7 +195,7 @@ public class ResetManagement { } } - private static void removeAllPipelineTemplates() { + private void removeAllPipelineTemplates() { var pipelineElementTemplateStorage = StorageDispatcher .INSTANCE .getNoSqlStore() @@ -194,7 +207,7 @@ public class ResetManagement { } - private static void clearGenericStorage() { + private void clearGenericStorage() { var appDocTypesToDelete = List.of( "asset-management", "asset-sites", diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/PipelineResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/PipelineResource.java index 1ccfa70fb3..f88a28c4ea 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/PipelineResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/PipelineResource.java @@ -24,6 +24,7 @@ import org.apache.streampipes.manager.execution.status.PipelineStatusManager; import org.apache.streampipes.manager.matching.PipelineVerificationHandlerV2; 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.MeasurementUpdateManagement; import org.apache.streampipes.manager.recommender.ElementRecommender; import org.apache.streampipes.manager.storage.PipelineStorageService; @@ -46,6 +47,7 @@ import org.apache.streampipes.resource.management.PipelineResourceManager; import org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource; import org.apache.streampipes.rest.shared.exception.SpMessageException; import org.apache.streampipes.rest.shared.exception.SpNotificationException; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import org.apache.streampipes.storage.management.StorageDispatcher; import com.google.gson.JsonSyntaxException; @@ -89,14 +91,19 @@ public class PipelineResource extends AbstractAuthGuardedRestResource { private final CompactPipelineManagement compactPipelineManagement; private final ExtensionServiceRequestManager requestManager; private final MeasurementUpdateManagement measurementUpdateManagement; + private final IDataExplorerWidgetStorage chartStorage; - public PipelineResource(ExtensionServiceRequestManager requestManager) { + public PipelineResource(ExtensionServiceRequestManager requestManager, + IDataExplorerWidgetStorage chartStorage) { this.compactPipelineManagement = new CompactPipelineManagement( StorageDispatcher.INSTANCE.getNoSqlStore().getPipelineElementDescriptionStorage(), requestManager ); this.requestManager = requestManager; - this.measurementUpdateManagement = new MeasurementUpdateManagement(); + this.chartStorage = chartStorage; + this.measurementUpdateManagement = new MeasurementUpdateManagement( + new ChartSchemaUpdateCoordinator(chartStorage) + ); } @GetMapping(produces = MediaType.APPLICATION_JSON_VALUE) diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/ResetResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/ResetResource.java index 26e8431326..210f82595e 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/ResetResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/ResetResource.java @@ -20,13 +20,13 @@ package org.apache.streampipes.rest.impl; import org.apache.streampipes.connect.management.management.WorkerRestClient; import org.apache.streampipes.manager.api.extensions.ExtensionServiceRequestManager; -import org.apache.streampipes.model.client.user.Principal; import org.apache.streampipes.model.client.user.PrincipalType; import org.apache.streampipes.model.message.Notifications; import org.apache.streampipes.model.message.SuccessMessage; import org.apache.streampipes.rest.ResetManagement; import org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource; import org.apache.streampipes.rest.security.AuthConstants; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import org.apache.streampipes.storage.api.system.IExtensionsServiceStorage; import org.apache.streampipes.storage.management.StorageDispatcher; @@ -47,30 +47,30 @@ import java.util.ArrayList; @PreAuthorize(AuthConstants.IS_ADMIN_ROLE) public class ResetResource extends AbstractAuthGuardedRestResource { - private final WorkerRestClient workerRestClient; - private final IExtensionsServiceStorage extensionsServiceStorage; - private final ExtensionServiceRequestManager requestManager; + private final ResetManagement resetManagement; public ResetResource(WorkerRestClient workerRestClient, - ExtensionServiceRequestManager requestManager) { - this.workerRestClient = workerRestClient; - this.extensionsServiceStorage = StorageDispatcher.INSTANCE.getNoSqlStore().getExtensionsServiceStorage(); - this.requestManager = requestManager; + ExtensionServiceRequestManager requestManager, + IDataExplorerWidgetStorage widgetStorage) { + IExtensionsServiceStorage extensionsServiceStorage = StorageDispatcher.INSTANCE.getNoSqlStore().getExtensionsServiceStorage(); + this.resetManagement = new ResetManagement( + widgetStorage, workerRestClient, extensionsServiceStorage, requestManager + ); } @PostMapping(produces = MediaType.APPLICATION_JSON_VALUE) @Operation(summary = "Resets StreamPipes instance") public ResponseEntity<SuccessMessage> reset() { - ResetManagement.reset(getAuthenticatedUsername(), workerRestClient, extensionsServiceStorage, requestManager); + resetManagement.reset(getAuthenticatedUsername()); var userStorage = getUserStorage(); // Delete all users other than current user (admin) and their resources - var allUsers = new ArrayList<Principal>(userStorage.getAllUsers()); + var allUsers = new ArrayList<>(userStorage.getAllUsers()); for (var user : allUsers) { if (user.getPrincipalType() == PrincipalType.USER_ACCOUNT && !user.getPrincipalId().equals(getAuthenticatedUserSid())) { - ResetManagement.reset(user.getUsername(), workerRestClient, extensionsServiceStorage, requestManager); + resetManagement.reset(user.getUsername()); userStorage.deleteUser(user.getPrincipalId()); } } diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/DataExportResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/DataExportResource.java index f9d1b44905..e6a3fe9ce0 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/DataExportResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/DataExportResource.java @@ -23,6 +23,7 @@ import org.apache.streampipes.manager.api.extensions.ExtensionServiceRequestMana import org.apache.streampipes.model.export.ExportConfiguration; import org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource; import org.apache.streampipes.rest.security.AuthConstants; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import org.springframework.http.MediaType; import org.springframework.http.ResponseEntity; @@ -41,9 +42,12 @@ import java.util.List; public class DataExportResource extends AbstractAuthGuardedRestResource { private final ExtensionServiceRequestManager extensionServiceRequestManager; + private final IDataExplorerWidgetStorage chartStorage; - public DataExportResource(ExtensionServiceRequestManager extensionServiceRequestManager) { + public DataExportResource(ExtensionServiceRequestManager extensionServiceRequestManager, + IDataExplorerWidgetStorage chartStorage) { this.extensionServiceRequestManager = extensionServiceRequestManager; + this.chartStorage = chartStorage; } @PostMapping( @@ -51,7 +55,7 @@ public class DataExportResource extends AbstractAuthGuardedRestResource { consumes = MediaType.APPLICATION_JSON_VALUE, produces = MediaType.APPLICATION_JSON_VALUE) public ResponseEntity<ExportConfiguration> getExportPreview(@RequestBody List<String> selectedAssetIds) { - var exportConfig = ExportManager.getExportPreview(selectedAssetIds, extensionServiceRequestManager); + var exportConfig = ExportManager.getExportPreview(selectedAssetIds, extensionServiceRequestManager, chartStorage); return ok(exportConfig); } @@ -60,7 +64,8 @@ public class DataExportResource extends AbstractAuthGuardedRestResource { consumes = MediaType.APPLICATION_JSON_VALUE, produces = MediaType.APPLICATION_OCTET_STREAM_VALUE) public ResponseEntity<byte[]> download(@RequestBody ExportConfiguration exportConfiguration) throws IOException { - var applicationPackage = ExportManager.getExportPackage(exportConfiguration, extensionServiceRequestManager); + var applicationPackage = ExportManager + .getExportPackage(exportConfiguration, extensionServiceRequestManager, chartStorage); return ok(applicationPackage); } diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/DataImportResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/DataImportResource.java index fff47a710c..baadde3df1 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/DataImportResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/DataImportResource.java @@ -23,6 +23,7 @@ import org.apache.streampipes.manager.api.extensions.ExtensionServiceRequestMana import org.apache.streampipes.model.export.AssetExportConfiguration; import org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource; import org.apache.streampipes.rest.security.AuthConstants; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -44,9 +45,12 @@ public class DataImportResource extends AbstractAuthGuardedRestResource { private static final Logger LOG = LoggerFactory.getLogger(DataImportResource.class); private final ExtensionServiceRequestManager extensionServiceRequestManager; + private final IDataExplorerWidgetStorage chartStorage; - public DataImportResource(ExtensionServiceRequestManager extensionServiceRequestManager) { + public DataImportResource(ExtensionServiceRequestManager extensionServiceRequestManager, + IDataExplorerWidgetStorage chartStorage) { this.extensionServiceRequestManager = extensionServiceRequestManager; + this.chartStorage = chartStorage; } @PostMapping( @@ -55,7 +59,7 @@ public class DataImportResource extends AbstractAuthGuardedRestResource { produces = MediaType.APPLICATION_JSON_VALUE) public ResponseEntity<AssetExportConfiguration> getImportPreview(@RequestPart("file_upload") MultipartFile fileDetail) throws IOException { - var importConfig = ImportManager.getImportPreview(fileDetail.getInputStream(), extensionServiceRequestManager); + var importConfig = ImportManager.getImportPreview(fileDetail.getInputStream(), extensionServiceRequestManager, chartStorage); return ok(importConfig); } @@ -69,7 +73,8 @@ public class DataImportResource extends AbstractAuthGuardedRestResource { fileDetail.getInputStream(), exportConfiguration, getAuthenticatedUserSid(), - extensionServiceRequestManager + extensionServiceRequestManager, + chartStorage ); return ok(); } catch (IOException e) { 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 689a739033..038e724b21 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,6 +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.model.client.user.DefaultRole; import org.apache.streampipes.model.client.user.Permission; import org.apache.streampipes.model.connect.adapter.AdapterDescription; @@ -45,6 +46,7 @@ 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.explorer.IDataExplorerWidgetStorage; import org.apache.streampipes.storage.api.pipeline.IPipelineStorage; import org.apache.streampipes.storage.management.StorageDispatcher; @@ -79,16 +81,19 @@ public class AdapterResource extends AbstractAdapterResource<AdapterMasterManage private static final Logger LOG = LoggerFactory.getLogger(AdapterResource.class); private final ExtensionServiceRequestManager requestManager; private final ApplicationEventPublisher eventPublisher; + private final ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator; public AdapterResource(WorkerRestClient workerRestClient, - ExtensionServiceRequestManager requestManager) { - this(workerRestClient, requestManager, null); + ExtensionServiceRequestManager requestManager, + IDataExplorerWidgetStorage chartStorage) { + this(workerRestClient, requestManager, null, chartStorage); } @Autowired public AdapterResource(WorkerRestClient workerRestClient, ExtensionServiceRequestManager requestManager, - ApplicationEventPublisher eventPublisher) { + ApplicationEventPublisher eventPublisher, + IDataExplorerWidgetStorage chartStorage) { super(() -> new AdapterMasterManagement( StorageDispatcher.INSTANCE.getNoSqlStore() .getAdapterInstanceStorage(), @@ -100,6 +105,7 @@ public class AdapterResource extends AbstractAdapterResource<AdapterMasterManage requestManager)); this.requestManager = requestManager; this.eventPublisher = eventPublisher; + this.chartSchemaUpdateCoordinator = new ChartSchemaUpdateCoordinator(chartStorage); } @PostMapping(consumes = MediaType.APPLICATION_JSON_VALUE) @@ -136,7 +142,7 @@ 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, requestManager); + var updateManager = new AdapterUpdateManagement(managementService, requestManager, chartSchemaUpdateCoordinator); try { updateManager.updateAdapter(adapterDescription); publishEvent(new AdapterUpdatedEvent(adapterDescription)); @@ -152,7 +158,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, requestManager); + var updateManager = new AdapterUpdateManagement(managementService, requestManager, chartSchemaUpdateCoordinator); var migrations = updateManager.checkPipelineMigrations(adapterDescription); return ok(migrations); 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 cb52424c9b..ddff2a5ca4 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 @@ -30,6 +30,7 @@ import org.apache.streampipes.connect.management.management.WorkerRestClient; import org.apache.streampipes.manager.api.extensions.ExtensionServiceRequestManager; import org.apache.streampipes.manager.execution.endpoint.ExtensionsServiceEndpointGenerator; import org.apache.streampipes.manager.pipeline.compact.CompactPipelineManagement; +import org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator; import org.apache.streampipes.model.connect.adapter.AdapterDescription; import org.apache.streampipes.model.connect.adapter.compact.CompactAdapter; import org.apache.streampipes.model.message.Notifications; @@ -37,6 +38,7 @@ import org.apache.streampipes.resource.management.SpResourceManager; import org.apache.streampipes.rest.shared.constants.SpMediaType; import org.apache.streampipes.rest.shared.exception.BadRequestException; import org.apache.streampipes.rest.shared.exception.SpMessageException; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import org.apache.streampipes.storage.management.StorageDispatcher; import org.slf4j.Logger; @@ -62,7 +64,8 @@ public class CompactAdapterResource extends AbstractAdapterResource<AdapterMaste private final ExtensionServiceRequestManager requestManager; public CompactAdapterResource(WorkerRestClient workerRestClient, - ExtensionServiceRequestManager requestManager) { + ExtensionServiceRequestManager requestManager, + IDataExplorerWidgetStorage chartStorage) { super(() -> new AdapterMasterManagement( StorageDispatcher.INSTANCE.getNoSqlStore() .getAdapterInstanceStorage(), @@ -81,7 +84,8 @@ public class CompactAdapterResource extends AbstractAdapterResource<AdapterMaste this.compactAdapterManagement = new CompactAdapterManagement( new AdapterGenerationSteps(guessManagement).getGenerators() ); - this.adapterUpdateManagement = new AdapterUpdateManagement(managementService, requestManager); + var chartCoordinator = new ChartSchemaUpdateCoordinator(chartStorage); + this.adapterUpdateManagement = new AdapterUpdateManagement(managementService, requestManager, chartCoordinator); } @PostMapping( diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/dashboard/DataLakeDashboardResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/dashboard/DataLakeDashboardResource.java index f940c21a0e..4e45709926 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/dashboard/DataLakeDashboardResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/dashboard/DataLakeDashboardResource.java @@ -25,6 +25,7 @@ import org.apache.streampipes.model.dashboard.DashboardSummaryDto; import org.apache.streampipes.model.resource.ResourceSummaryDto; import org.apache.streampipes.resource.management.DataExplorerResourceManager; import org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import org.apache.streampipes.storage.api.user.IPermissionStorage; import org.springframework.http.CacheControl; @@ -51,9 +52,11 @@ import java.util.Objects; public class DataLakeDashboardResource extends AbstractAuthGuardedRestResource { private final IPermissionStorage permissionStorage; + private final IDataExplorerWidgetStorage chartStorage; - public DataLakeDashboardResource() { + public DataLakeDashboardResource(IDataExplorerWidgetStorage chartStorage) { this.permissionStorage = getNoSqlStorage().getPermissionStorage(); + this.chartStorage = chartStorage; } @GetMapping(produces = MediaType.APPLICATION_JSON_VALUE) @@ -119,7 +122,7 @@ public class DataLakeDashboardResource extends AbstractAuthGuardedRestResource { private DataExplorerResourceManager getResourceManager() { - return getSpResourceManager().manageDataExplorer(); + return new DataExplorerResourceManager(chartStorage); } public boolean hasReadAuthority() { diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/AbstractDataLakeResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/AbstractDataLakeResource.java index 87e3a8aa5b..b24717c1bd 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/AbstractDataLakeResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/AbstractDataLakeResource.java @@ -19,6 +19,7 @@ package org.apache.streampipes.rest.impl.datalake; import org.apache.streampipes.dataexplorer.api.IDataExplorerSchemaManagement; import org.apache.streampipes.dataexplorer.management.DataExplorerDispatcher; +import org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator; import org.apache.streampipes.model.client.user.DefaultPrivilege; import org.apache.streampipes.resource.management.permission.SpPermissionEvaluator; import org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource; @@ -34,12 +35,12 @@ public class AbstractDataLakeResource extends AbstractAuthGuardedRestResource { final IDataExplorerSchemaManagement dataLakeMeasureManagement; private final IDataLakeMeasureStorage dataLakeMeasureStorage = StorageDispatcher.INSTANCE.getNoSqlStore().getDataLakeStorage(); + protected final ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator; - - - public AbstractDataLakeResource() { + public AbstractDataLakeResource(ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator) { + this.chartSchemaUpdateCoordinator = chartSchemaUpdateCoordinator; this.dataLakeMeasureManagement = new DataExplorerDispatcher().getDataExplorerManager() - .getSchemaManagement(); + .getSchemaManagement(chartSchemaUpdateCoordinator); } /** diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeMeasureResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeMeasureResource.java index 5f94ed50c3..41d9b19993 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeMeasureResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeMeasureResource.java @@ -19,11 +19,13 @@ package org.apache.streampipes.rest.impl.datalake; import org.apache.streampipes.dataexplorer.management.DataExplorerDispatcher; +import org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator; import org.apache.streampipes.model.datalake.DataLakeMeasure; import org.apache.streampipes.model.datalake.DatasetSummaryDto; import org.apache.streampipes.model.monitoring.SpLogMessage; import org.apache.streampipes.model.resource.ResourceSummaryDto; import org.apache.streampipes.resource.management.DataLakeMeasureResourceManager; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import io.swagger.v3.oas.annotations.Operation; import io.swagger.v3.oas.annotations.Parameter; @@ -49,8 +51,8 @@ import java.util.Objects; @RequestMapping("/api/v4/datalake/measure") public class DataLakeMeasureResource extends AbstractDataLakeResource { - public DataLakeMeasureResource() { - super(); + public DataLakeMeasureResource(IDataExplorerWidgetStorage chartStorage) { + super(new ChartSchemaUpdateCoordinator(chartStorage)); } @GetMapping(path = "/summary", produces = MediaType.APPLICATION_JSON_VALUE) diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeResource.java index df2ff38dee..1d6dd96d83 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeResource.java @@ -23,6 +23,7 @@ import org.apache.streampipes.dataexplorer.api.IDataExplorerQueryManagement; import org.apache.streampipes.dataexplorer.export.OutputFormat; import org.apache.streampipes.dataexplorer.management.DataExplorerDispatcher; import org.apache.streampipes.export.DataLakeExportManager; +import org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator; import org.apache.streampipes.model.datalake.DataLakeMeasure; import org.apache.streampipes.model.datalake.DataSeries; import org.apache.streampipes.model.datalake.RetentionTimeConfig; @@ -32,6 +33,7 @@ import org.apache.streampipes.model.message.Notifications; import org.apache.streampipes.model.monitoring.SpLogMessage; import org.apache.streampipes.rest.security.AuthConstants; import org.apache.streampipes.rest.shared.exception.SpMessageException; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import io.swagger.v3.oas.annotations.Operation; import io.swagger.v3.oas.annotations.Parameter; @@ -93,18 +95,14 @@ public class DataLakeResource extends AbstractDataLakeResource { private static final Logger LOG = LoggerFactory.getLogger(DataLakeResource.class); private final IDataExplorerQueryManagement dataExplorerQueryManagement; - private static DataLakeExportManager dataLakeExportManager = new DataLakeExportManager(); + private final DataLakeExportManager dataLakeExportManager; - public DataLakeResource() { - super(); + public DataLakeResource(IDataExplorerWidgetStorage chartStorage) { + super(new ChartSchemaUpdateCoordinator(chartStorage)); this.dataExplorerQueryManagement = new DataExplorerDispatcher() .getDataExplorerManager() .getQueryManagement(this.dataLakeMeasureManagement); - } - - public DataLakeResource(IDataExplorerQueryManagement dataExplorerQueryManagement) { - super(); - this.dataExplorerQueryManagement = dataExplorerQueryManagement; + this.dataLakeExportManager = new DataLakeExportManager(this.dataLakeMeasureManagement, dataExplorerQueryManagement); } @DeleteMapping(path = "/measurements/{measurementName}") diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeWidgetResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeWidgetResource.java index ca2ebe4222..0e8501d954 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeWidgetResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeWidgetResource.java @@ -28,6 +28,7 @@ import org.apache.streampipes.resource.management.SpResourceManager; import org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource; import org.apache.streampipes.rest.security.AuthConstants; import org.apache.streampipes.rest.shared.exception.BadRequestException; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import org.springframework.http.MediaType; import org.springframework.http.ResponseEntity; @@ -50,10 +51,10 @@ public class DataLakeWidgetResource extends AbstractAuthGuardedRestResource { private final DataExplorerWidgetResourceManager resourceManager; - public DataLakeWidgetResource() { + public DataLakeWidgetResource(IDataExplorerWidgetStorage dataExplorerWidgetStorage) { this.resourceManager = new SpResourceManager().manageDataExplorerWidget( - new DataExplorerResourceManager(), - getNoSqlStorage().getDataExplorerWidgetStorage() + new DataExplorerResourceManager(dataExplorerWidgetStorage), + dataExplorerWidgetStorage ); } diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/KioskDashboardDataLakeResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/KioskDashboardDataLakeResource.java index 35f82f08f4..ada7f5a0d8 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/KioskDashboardDataLakeResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/KioskDashboardDataLakeResource.java @@ -21,6 +21,7 @@ package org.apache.streampipes.rest.impl.datalake; import org.apache.streampipes.dataexplorer.api.IDataExplorerQueryManagement; import org.apache.streampipes.dataexplorer.api.IDataExplorerSchemaManagement; import org.apache.streampipes.dataexplorer.management.DataExplorerDispatcher; +import org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator; import org.apache.streampipes.model.client.user.DefaultPrivilege; import org.apache.streampipes.model.datalake.DataExplorerWidgetModel; import org.apache.streampipes.model.datalake.SpQueryResult; @@ -49,22 +50,19 @@ import java.util.Map; public class KioskDashboardDataLakeResource extends AbstractAuthGuardedRestResource { private final IDataExplorerQueryManagement dataExplorerQueryManagement; - private final IDataExplorerSchemaManagement dataExplorerSchemaManagement; private final IDataExplorerDashboardStorage dashboardStorage = StorageDispatcher.INSTANCE.getNoSqlStore().getDataExplorerDashboardStorage(); private final IDataExplorerWidgetStorage dataExplorerWidgetStorage; private final IPermissionStorage permissionStorage; - public KioskDashboardDataLakeResource() { - this.dataExplorerSchemaManagement = new DataExplorerDispatcher() + public KioskDashboardDataLakeResource(IDataExplorerWidgetStorage dataExplorerWidgetStorage) { + IDataExplorerSchemaManagement dataExplorerSchemaManagement = new DataExplorerDispatcher() .getDataExplorerManager() - .getSchemaManagement(); + .getSchemaManagement(new ChartSchemaUpdateCoordinator(dataExplorerWidgetStorage)); this.dataExplorerQueryManagement = new DataExplorerDispatcher() .getDataExplorerManager() - .getQueryManagement(this.dataExplorerSchemaManagement); - this.dataExplorerWidgetStorage = StorageDispatcher.INSTANCE - .getNoSqlStore() - .getDataExplorerWidgetStorage(); + .getQueryManagement(dataExplorerSchemaManagement); + this.dataExplorerWidgetStorage = dataExplorerWidgetStorage; this.permissionStorage = getNoSqlStorage().getPermissionStorage(); } diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/DataLakeImportResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/DataLakeImportResource.java index 7fde0ca48a..54a9840f3c 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/DataLakeImportResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/DataLakeImportResource.java @@ -18,6 +18,7 @@ package org.apache.streampipes.rest.impl.datalake.importer; +import org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator; import org.apache.streampipes.model.datalake.importer.CsvImportPreviewRequest; import org.apache.streampipes.model.datalake.importer.CsvImportPreviewResult; import org.apache.streampipes.model.datalake.importer.CsvImportRequest; @@ -26,6 +27,7 @@ import org.apache.streampipes.model.datalake.importer.CsvImportSchemaValidationR import org.apache.streampipes.model.datalake.importer.CsvImportSchemaValidationResult; import org.apache.streampipes.model.datalake.importer.CsvImportTargetMode; import org.apache.streampipes.rest.impl.datalake.AbstractDataLakeResource; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import org.springframework.http.HttpStatus; import org.springframework.http.MediaType; @@ -46,8 +48,8 @@ public class DataLakeImportResource extends AbstractDataLakeResource { private final CsvDataLakeImportService importService; - public DataLakeImportResource() { - super(); + public DataLakeImportResource(IDataExplorerWidgetStorage chartStorage) { + super(new ChartSchemaUpdateCoordinator(chartStorage)); this.importService = new CsvDataLakeImportService(getDataLakeMeasureManagement()); } 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 3ed8bd5bcd..0ac52a394a 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,6 +19,7 @@ package org.apache.streampipes.rest.impl.pe; import org.apache.streampipes.manager.api.extensions.ExtensionServiceRequestManager; +import org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator; import org.apache.streampipes.manager.pipeline.update.DataStreamUpdateManagement; import org.apache.streampipes.model.SpDataStream; import org.apache.streampipes.model.connect.adapter.PipelineUpdateInfo; @@ -30,6 +31,7 @@ import org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResourc 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.IDataExplorerWidgetStorage; import org.apache.http.client.HttpResponseException; import org.springframework.context.ApplicationEventPublisher; @@ -56,8 +58,11 @@ public class DataStreamResource extends AbstractAuthGuardedRestResource { private final ApplicationEventPublisher eventPublisher; public DataStreamResource(ExtensionServiceRequestManager requestManager, - ApplicationEventPublisher eventPublisher) { - this.dataStreamUpdateManagement = new DataStreamUpdateManagement(requestManager); + ApplicationEventPublisher eventPublisher, + IDataExplorerWidgetStorage chartStorage) { + this.dataStreamUpdateManagement = new DataStreamUpdateManagement( + requestManager, + new ChartSchemaUpdateCoordinator(chartStorage)); this.eventPublisher = eventPublisher; } diff --git a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/StreamPipesCoreApplication.java b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/StreamPipesCoreApplication.java index 33e8726054..abf883b470 100644 --- a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/StreamPipesCoreApplication.java +++ b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/StreamPipesCoreApplication.java @@ -59,6 +59,7 @@ import org.apache.streampipes.service.core.migrations.AvailableMigrations; import org.apache.streampipes.service.core.migrations.Migration; import org.apache.streampipes.service.core.migrations.MigrationsHandler; import org.apache.streampipes.service.core.storage.StorageApiConfiguration; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import org.apache.streampipes.storage.api.function.IFunctionStateStorage; import org.apache.streampipes.storage.api.pipeline.IPipelineStorage; import org.apache.streampipes.storage.api.system.IExtensionsServiceStorage; @@ -114,6 +115,9 @@ public class StreamPipesCoreApplication extends StreamPipesServiceBase { @Autowired private WorkerRestClient workerRestClient; + @Autowired + protected IDataExplorerWidgetStorage chartStorage; + private final IExtensionsServiceStorage extensionsServiceStorage = StorageDispatcher.INSTANCE.getNoSqlStore().getExtensionsServiceStorage(); @@ -237,7 +241,7 @@ public class StreamPipesCoreApplication extends StreamPipesServiceBase { } protected List<Migration> getMigrations() { - return new AvailableMigrations().getAvailableMigrations(); + return new AvailableMigrations(chartStorage).getAvailableMigrations(); } private boolean isConfigured() { diff --git a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/AvailableMigrations.java b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/AvailableMigrations.java index ad093e185e..4c89b7e0be 100644 --- a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/AvailableMigrations.java +++ b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/AvailableMigrations.java @@ -39,19 +39,26 @@ import org.apache.streampipes.service.core.migrations.v099.RemoveInternalNotific import org.apache.streampipes.service.core.migrations.v099.RemoveObsoletePrivilegesMigration; import org.apache.streampipes.service.core.migrations.v099.UniqueDashboardIdMigration; import org.apache.streampipes.service.core.migrations.v099.connect.MigrateAdaptersToUseScript; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import java.util.Arrays; import java.util.List; public class AvailableMigrations { + private final IDataExplorerWidgetStorage chartStorage; + + public AvailableMigrations(IDataExplorerWidgetStorage chartStorage) { + this.chartStorage = chartStorage; + } + public List<Migration> getAvailableMigrations() { return Arrays.asList( new ModifyAssetLinksMigration(), new ModifyAssetLinkTypesMigration(), new AddDataLakeMeasureViewMigration(), new AddDefaultExportProviderMigration(), - new FixImportedPermissionsMigration(), + new FixImportedPermissionsMigration(chartStorage), new AddAssetManagementViewMigration(), new MoveAssetContentMigration(), new CreateAssetPermissionMigration(), diff --git a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/v0980/FixImportedPermissionsMigration.java b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/v0980/FixImportedPermissionsMigration.java index 27e654f027..9fc98e05b0 100644 --- a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/v0980/FixImportedPermissionsMigration.java +++ b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/v0980/FixImportedPermissionsMigration.java @@ -20,6 +20,7 @@ package org.apache.streampipes.service.core.migrations.v0980; import org.apache.streampipes.model.shared.api.Storable; import org.apache.streampipes.service.core.migrations.Migration; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import org.apache.streampipes.storage.api.user.IPermissionStorage; import org.apache.streampipes.storage.management.StorageDispatcher; @@ -35,6 +36,12 @@ import java.util.List; */ public class FixImportedPermissionsMigration implements Migration { + private final IDataExplorerWidgetStorage chartStorage; + + public FixImportedPermissionsMigration(IDataExplorerWidgetStorage chartStorage) { + this.chartStorage = chartStorage; + } + private static final Logger LOG = LoggerFactory.getLogger(FixImportedPermissionsMigration.class); private final IPermissionStorage permissionStorage = @@ -65,10 +72,7 @@ public class FixImportedPermissionsMigration implements Migration { private void migrateChartsPermissions() { LOG.debug("Start migrate permissions for charts"); - var dataExplorerWidgetStorage = StorageDispatcher.INSTANCE - .getNoSqlStore() - .getDataExplorerWidgetStorage(); - var charts = dataExplorerWidgetStorage.findAll(); + var charts = chartStorage.findAll(); migrateResourcePermissions(charts); LOG.debug("Finished migrate permissions for charts"); } diff --git a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/scheduler/DataLakeScheduler.java b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/scheduler/DataLakeScheduler.java index 229ea35eaf..3c2522234d 100644 --- a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/scheduler/DataLakeScheduler.java +++ b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/scheduler/DataLakeScheduler.java @@ -21,7 +21,9 @@ import org.apache.streampipes.commons.environment.Environments; import org.apache.streampipes.dataexplorer.api.IDataExplorerSchemaManagement; import org.apache.streampipes.dataexplorer.management.DataExplorerDispatcher; import org.apache.streampipes.export.DataLakeExportManager; +import org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator; import org.apache.streampipes.model.datalake.DataLakeMeasure; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -35,12 +37,21 @@ import java.util.List; @Configuration public class DataLakeScheduler implements SchedulingConfigurer { - private static final DataLakeExportManager dataLakeExportManager = new DataLakeExportManager(); - private static final Logger LOG = LoggerFactory.getLogger(DataLakeExportManager.class); + private final DataLakeExportManager dataLakeExportManager; + private static final Logger LOG = LoggerFactory.getLogger(DataLakeScheduler.class); - private final IDataExplorerSchemaManagement dataExplorerSchemaManagement = new DataExplorerDispatcher() + private final IDataExplorerSchemaManagement dataExplorerSchemaManagement; + + public DataLakeScheduler(IDataExplorerWidgetStorage chartStorage) { + var chartSchemaUpdateCoordinator = new ChartSchemaUpdateCoordinator(chartStorage); + dataExplorerSchemaManagement = new DataExplorerDispatcher() .getDataExplorerManager() - .getSchemaManagement(); + .getSchemaManagement(chartSchemaUpdateCoordinator); + this.dataLakeExportManager = new DataLakeExportManager( + dataExplorerSchemaManagement, + new DataExplorerDispatcher().getDataExplorerManager() + .getQueryManagement(dataExplorerSchemaManagement)); + } public void cleanupMeasurements() { LOG.info("Retention CRON Job triggered."); @@ -72,4 +83,4 @@ public class DataLakeScheduler implements SchedulingConfigurer { } -} \ No newline at end of file +} diff --git a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/storage/StorageApiConfiguration.java b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/storage/StorageApiConfiguration.java index ae16f2e6ea..8945921de1 100644 --- a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/storage/StorageApiConfiguration.java +++ b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/storage/StorageApiConfiguration.java @@ -18,7 +18,9 @@ package org.apache.streampipes.service.core.storage; +import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import org.apache.streampipes.storage.api.function.IFunctionStateStorage; +import org.apache.streampipes.storage.couchdb.impl.explorer.DataExplorerWidgetStorageImpl; import org.apache.streampipes.storage.couchdb.impl.function.FunctionStateStorageImpl; import org.springframework.context.annotation.Bean; @@ -31,4 +33,9 @@ public class StorageApiConfiguration { public IFunctionStateStorage functionStateStorage() { return new FunctionStateStorageImpl(); } + + @Bean + public IDataExplorerWidgetStorage dataExplorerWidgetStorage() { + return new DataExplorerWidgetStorageImpl(); + } } diff --git a/streampipes-storage-api/src/main/java/org/apache/streampipes/storage/api/core/INoSqlStorage.java b/streampipes-storage-api/src/main/java/org/apache/streampipes/storage/api/core/INoSqlStorage.java index 3d616f1449..de4ef4b597 100644 --- a/streampipes-storage-api/src/main/java/org/apache/streampipes/storage/api/core/INoSqlStorage.java +++ b/streampipes-storage-api/src/main/java/org/apache/streampipes/storage/api/core/INoSqlStorage.java @@ -19,7 +19,6 @@ package org.apache.streampipes.storage.api.core; import org.apache.streampipes.storage.api.connect.IAdapterStorage; import org.apache.streampipes.storage.api.explorer.IDataExplorerDashboardStorage; -import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import org.apache.streampipes.storage.api.explorer.IDataLakeMeasureStorage; import org.apache.streampipes.storage.api.pipeline.ICompactPipelineTemplateStorage; import org.apache.streampipes.storage.api.pipeline.IDataProcessorStorage; @@ -69,8 +68,6 @@ public interface INoSqlStorage { IDataExplorerDashboardStorage getDataExplorerDashboardStorage(); - IDataExplorerWidgetStorage getDataExplorerWidgetStorage(); - IPipelineElementTemplateStorage getPipelineElementTemplateStorage(); IPipelineCanvasMetadataStorage getPipelineCanvasMetadataStorage(); diff --git a/streampipes-storage-couchdb/src/main/java/org/apache/streampipes/storage/couchdb/CouchDbStorageManager.java b/streampipes-storage-couchdb/src/main/java/org/apache/streampipes/storage/couchdb/CouchDbStorageManager.java index b75019c472..6b964775b4 100644 --- a/streampipes-storage-couchdb/src/main/java/org/apache/streampipes/storage/couchdb/CouchDbStorageManager.java +++ b/streampipes-storage-couchdb/src/main/java/org/apache/streampipes/storage/couchdb/CouchDbStorageManager.java @@ -135,11 +135,6 @@ public class CouchDbStorageManager implements INoSqlStorage { return new DataExplorerDashboardStorageImpl(); } - @Override - public IDataExplorerWidgetStorage getDataExplorerWidgetStorage() { - return new DataExplorerWidgetStorageImpl(); - } - @Override public IPipelineElementTemplateStorage getPipelineElementTemplateStorage() { return new PipelineElementTemplateStorageImpl();
