This is an automated email from the ASF dual-hosted git repository. dominikriemer pushed a commit to branch migrate-configuration-storage in repository https://gitbox.apache.org/repos/asf/streampipes.git
commit 1c0222127b05170d5efaec29a66f5ca804b58e82 Author: Dominik Riemer <[email protected]> AuthorDate: Tue Jun 23 23:24:17 2026 +0200 Finish migration of configuration storage --- .../api/IDataExplorerQueryManagement.java | 2 + .../export/ConfiguredOutputWriter.java | 10 ---- .../export/ConfiguredOutputWriterFactory.java | 54 ++++++++++++++++++++++ .../dataexplorer/export/OutputFormat.java | 21 ++------- .../influx/DataExplorerQueryManagementInflux.java | 4 +- .../iotdb/DataExplorerQueryManagementIotDb.java | 7 ++- .../dataexplorer/StreamedQueryResultProvider.java | 7 ++- .../streampipes/export/DataLakeExportManager.java | 8 +++- .../apache/streampipes/mail/utils/MailUtils.java | 5 -- .../resource/management/SpResourceManager.java | 10 +++- .../rest/core/base/impl/AbstractRestResource.java | 5 -- .../streampipes/rest/impl/Authentication.java | 9 ++-- .../streampipes/rest/impl/EmailResource.java | 12 ++++- .../apache/streampipes/rest/impl/UserResource.java | 10 +++- .../impl/admin/EmailConfigurationResource.java | 13 +++--- .../admin/ExportProviderConfigurationResource.java | 25 ++++++---- .../impl/admin/ExtensionsInstallationResource.java | 6 +-- .../impl/admin/GeneralConfigurationResource.java | 14 ++++-- .../impl/admin/LocationConfigurationResource.java | 14 ++++-- .../rest/impl/datalake/DataLakeResource.java | 9 +++- .../ExtensionServiceRequestConfiguration.java | 7 ++- .../core/oauth2/CustomOAuth2UserService.java | 10 ++-- .../service/core/oauth2/CustomOidcUserService.java | 10 ++-- .../service/core/oauth2/UserService.java | 10 ++-- .../service/core/scheduler/DataLakeScheduler.java | 3 +- .../core/storage/StorageApiConfiguration.java | 7 +++ .../storage/api/core/INoSqlStorage.java | 3 -- .../storage/couchdb/CouchDbStorageManager.java | 7 --- 28 files changed, 195 insertions(+), 107 deletions(-) diff --git a/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataExplorerQueryManagement.java b/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataExplorerQueryManagement.java index 2f5f3eb6c3..d114cb2886 100644 --- a/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataExplorerQueryManagement.java +++ b/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataExplorerQueryManagement.java @@ -18,6 +18,7 @@ package org.apache.streampipes.dataexplorer.api; +import org.apache.streampipes.dataexplorer.export.ConfiguredOutputWriterFactory; import org.apache.streampipes.dataexplorer.export.OutputFormat; import org.apache.streampipes.model.datalake.SpQueryResult; import org.apache.streampipes.model.datalake.param.ProvidedRestQueryParams; @@ -34,6 +35,7 @@ public interface IDataExplorerQueryManagement { void getDataAsStream(ProvidedRestQueryParams params, OutputFormat format, + ConfiguredOutputWriterFactory outputWriterFactory, boolean ignoreMissingValues, OutputStream outputStream) throws IOException; diff --git a/streampipes-data-explorer-export/src/main/java/org/apache/streampipes/dataexplorer/export/ConfiguredOutputWriter.java b/streampipes-data-explorer-export/src/main/java/org/apache/streampipes/dataexplorer/export/ConfiguredOutputWriter.java index 2b4344a477..217f3016f9 100644 --- a/streampipes-data-explorer-export/src/main/java/org/apache/streampipes/dataexplorer/export/ConfiguredOutputWriter.java +++ b/streampipes-data-explorer-export/src/main/java/org/apache/streampipes/dataexplorer/export/ConfiguredOutputWriter.java @@ -32,16 +32,6 @@ public abstract class ConfiguredOutputWriter { private final DecimalFormat df = new DecimalFormat("#"); - public static ConfiguredOutputWriter getConfiguredWriter(DataLakeMeasure schema, - OutputFormat format, - ProvidedRestQueryParams params, - boolean ignoreMissingValues) { - var writer = format.getWriter(); - writer.configure(schema, params, ignoreMissingValues); - - return writer; - } - protected String getHeaderName(DataLakeMeasure schema, String runtimeName, String headerColumnNameStrategy) { diff --git a/streampipes-data-explorer-export/src/main/java/org/apache/streampipes/dataexplorer/export/ConfiguredOutputWriterFactory.java b/streampipes-data-explorer-export/src/main/java/org/apache/streampipes/dataexplorer/export/ConfiguredOutputWriterFactory.java new file mode 100644 index 0000000000..98edbb310a --- /dev/null +++ b/streampipes-data-explorer-export/src/main/java/org/apache/streampipes/dataexplorer/export/ConfiguredOutputWriterFactory.java @@ -0,0 +1,54 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.streampipes.dataexplorer.export; + +import org.apache.streampipes.model.datalake.DataLakeMeasure; +import org.apache.streampipes.model.datalake.param.ProvidedRestQueryParams; +import org.apache.streampipes.storage.api.system.IFileMetadataStorage; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; + +public class ConfiguredOutputWriterFactory { + + private final IFileMetadataStorage fileMetadataStorage; + private final ISpCoreConfigurationStorage coreConfigurationStorage; + + public ConfiguredOutputWriterFactory(IFileMetadataStorage fileMetadataStorage, + ISpCoreConfigurationStorage coreConfigurationStorage) { + this.fileMetadataStorage = fileMetadataStorage; + this.coreConfigurationStorage = coreConfigurationStorage; + } + + public ConfiguredOutputWriter getConfiguredWriter(DataLakeMeasure schema, + OutputFormat format, + ProvidedRestQueryParams params, + boolean ignoreMissingValues) { + var writer = createWriter(format); + writer.configure(schema, params, ignoreMissingValues); + + return writer; + } + + private ConfiguredOutputWriter createWriter(OutputFormat format) { + return switch (format) { + case JSON -> new ConfiguredJsonOutputWriter(); + case CSV -> new ConfiguredCsvOutputWriter(); + case XLSX -> new ConfiguredExcelOutputWriter(fileMetadataStorage, coreConfigurationStorage); + }; + } +} diff --git a/streampipes-data-explorer-export/src/main/java/org/apache/streampipes/dataexplorer/export/OutputFormat.java b/streampipes-data-explorer-export/src/main/java/org/apache/streampipes/dataexplorer/export/OutputFormat.java index f36409b105..ad3dfefe2f 100644 --- a/streampipes-data-explorer-export/src/main/java/org/apache/streampipes/dataexplorer/export/OutputFormat.java +++ b/streampipes-data-explorer-export/src/main/java/org/apache/streampipes/dataexplorer/export/OutputFormat.java @@ -18,27 +18,12 @@ package org.apache.streampipes.dataexplorer.export; -import org.apache.streampipes.storage.management.StorageDispatcher; - import java.util.Arrays; -import java.util.function.Supplier; public enum OutputFormat { - JSON(ConfiguredJsonOutputWriter::new), - CSV(ConfiguredCsvOutputWriter::new), - XLSX(() -> new ConfiguredExcelOutputWriter( - StorageDispatcher.INSTANCE.getNoSqlStore().getFileMetadataStorage()) - ); - - private final Supplier<ConfiguredOutputWriter> writerSupplier; - - OutputFormat(Supplier<ConfiguredOutputWriter> writerSupplier) { - this.writerSupplier = writerSupplier; - } - - public ConfiguredOutputWriter getWriter() { - return writerSupplier.get(); - } + JSON, + CSV, + XLSX; public static OutputFormat fromString(String desiredFormat) { return Arrays.stream( diff --git a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerQueryManagementInflux.java b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerQueryManagementInflux.java index b6bf82e659..eee432d3e2 100644 --- a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerQueryManagementInflux.java +++ b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerQueryManagementInflux.java @@ -22,6 +22,7 @@ import org.apache.streampipes.dataexplorer.QueryResultProvider; import org.apache.streampipes.dataexplorer.StreamedQueryResultProvider; import org.apache.streampipes.dataexplorer.api.IDataExplorerQueryManagement; import org.apache.streampipes.dataexplorer.api.IDataExplorerSchemaManagement; +import org.apache.streampipes.dataexplorer.export.ConfiguredOutputWriterFactory; import org.apache.streampipes.dataexplorer.export.OutputFormat; import org.apache.streampipes.dataexplorer.param.DeleteQueryParams; import org.apache.streampipes.dataexplorer.param.ProvidedRestQueryParamConverter; @@ -58,10 +59,11 @@ public class DataExplorerQueryManagementInflux implements IDataExplorerQueryMana @Override public void getDataAsStream(ProvidedRestQueryParams params, OutputFormat format, + ConfiguredOutputWriterFactory outputWriterFactory, boolean ignoreMissingValues, OutputStream outputStream) throws IOException { - new StreamedQueryResultProvider(params, format, + new StreamedQueryResultProvider(params, format, outputWriterFactory, this, new DataExplorerInfluxQueryExecutor(), dataExplorerSchemaManagement, diff --git a/streampipes-data-explorer-iotdb/src/main/java/org/apache/streampipes/dataexplorer/iotdb/DataExplorerQueryManagementIotDb.java b/streampipes-data-explorer-iotdb/src/main/java/org/apache/streampipes/dataexplorer/iotdb/DataExplorerQueryManagementIotDb.java index 9ce1901fae..cc5d153704 100644 --- a/streampipes-data-explorer-iotdb/src/main/java/org/apache/streampipes/dataexplorer/iotdb/DataExplorerQueryManagementIotDb.java +++ b/streampipes-data-explorer-iotdb/src/main/java/org/apache/streampipes/dataexplorer/iotdb/DataExplorerQueryManagementIotDb.java @@ -20,6 +20,7 @@ package org.apache.streampipes.dataexplorer.iotdb; import org.apache.streampipes.dataexplorer.api.IDataExplorerQueryManagement; import org.apache.streampipes.dataexplorer.api.IDataExplorerSchemaManagement; +import org.apache.streampipes.dataexplorer.export.ConfiguredOutputWriterFactory; import org.apache.streampipes.dataexplorer.export.OutputFormat; import org.apache.streampipes.model.datalake.SpQueryResult; import org.apache.streampipes.model.datalake.param.ProvidedRestQueryParams; @@ -50,7 +51,11 @@ public class DataExplorerQueryManagementIotDb implements IDataExplorerQueryManag } @Override - public void getDataAsStream(ProvidedRestQueryParams params, OutputFormat format, boolean ignoreMissingValues, OutputStream outputStream) throws IOException { + public void getDataAsStream(ProvidedRestQueryParams params, + OutputFormat format, + ConfiguredOutputWriterFactory outputWriterFactory, + boolean ignoreMissingValues, + OutputStream outputStream) throws IOException { } diff --git a/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/StreamedQueryResultProvider.java b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/StreamedQueryResultProvider.java index 24cc28d6ea..42e27d65be 100644 --- a/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/StreamedQueryResultProvider.java +++ b/streampipes-data-explorer/src/main/java/org/apache/streampipes/dataexplorer/StreamedQueryResultProvider.java @@ -20,7 +20,7 @@ package org.apache.streampipes.dataexplorer; import org.apache.streampipes.dataexplorer.api.IDataExplorerQueryManagement; import org.apache.streampipes.dataexplorer.api.IDataExplorerSchemaManagement; -import org.apache.streampipes.dataexplorer.export.ConfiguredOutputWriter; +import org.apache.streampipes.dataexplorer.export.ConfiguredOutputWriterFactory; import org.apache.streampipes.dataexplorer.export.OutputFormat; import org.apache.streampipes.dataexplorer.query.DataExplorerQueryExecutor; import org.apache.streampipes.model.datalake.DataLakeMeasure; @@ -39,22 +39,25 @@ public class StreamedQueryResultProvider extends QueryResultProvider { private static final String TIME_FIELD = "time"; private final OutputFormat format; + private final ConfiguredOutputWriterFactory outputWriterFactory; public StreamedQueryResultProvider(ProvidedRestQueryParams params, OutputFormat format, + ConfiguredOutputWriterFactory outputWriterFactory, IDataExplorerQueryManagement dataExplorerQueryManagement, DataExplorerQueryExecutor<?, ?> queryExecutor, IDataExplorerSchemaManagement schemaManagement, boolean ignoreMissingValues) { super(params, dataExplorerQueryManagement, queryExecutor, schemaManagement, ignoreMissingValues); this.format = format; + this.outputWriterFactory = outputWriterFactory; } public void getDataAsStream(OutputStream outputStream) throws IOException { var usesLimit = queryParams.has(SupportedRestQueryParams.QP_LIMIT); var measurement = findByMeasurementName(queryParams.getMeasurementId()).get(); - var configuredWriter = ConfiguredOutputWriter + var configuredWriter = outputWriterFactory .getConfiguredWriter(measurement, format, queryParams, ignoreMissingData); if (!queryParams.has(SupportedRestQueryParams.QP_LIMIT)) { 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 0d129c9c5a..2d9de976cf 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 @@ -21,6 +21,7 @@ import org.apache.streampipes.commons.environment.Environment; import org.apache.streampipes.commons.environment.Environments; import org.apache.streampipes.dataexplorer.api.IDataExplorerQueryManagement; import org.apache.streampipes.dataexplorer.api.IDataExplorerSchemaManagement; +import org.apache.streampipes.dataexplorer.export.ConfiguredOutputWriterFactory; import org.apache.streampipes.dataexplorer.export.OutputFormat; import org.apache.streampipes.dataexplorer.export.objectstorage.ExportProviderFactory; import org.apache.streampipes.dataexplorer.export.objectstorage.IObjectStorage; @@ -30,6 +31,7 @@ import org.apache.streampipes.model.datalake.DataLakeMeasure; import org.apache.streampipes.model.datalake.RetentionAction; import org.apache.streampipes.model.datalake.RetentionLog; import org.apache.streampipes.model.datalake.param.ProvidedRestQueryParams; +import org.apache.streampipes.storage.api.system.IFileMetadataStorage; import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.slf4j.Logger; @@ -51,13 +53,16 @@ public class DataLakeExportManager { private final IDataExplorerSchemaManagement dataExplorerSchemaManagement; private final IDataExplorerQueryManagement dataExplorerQueryManagement; private final ISpCoreConfigurationStorage coreConfigurationStorage; + private final ConfiguredOutputWriterFactory outputWriterFactory; public DataLakeExportManager(IDataExplorerSchemaManagement dataLakeSchemaManagement, IDataExplorerQueryManagement dataLakeQueryManagement, - ISpCoreConfigurationStorage coreConfigurationStorage) { + ISpCoreConfigurationStorage coreConfigurationStorage, + IFileMetadataStorage fileMetadataStorage) { this.dataExplorerSchemaManagement = dataLakeSchemaManagement; this.dataExplorerQueryManagement = dataLakeQueryManagement; this.coreConfigurationStorage = coreConfigurationStorage; + this.outputWriterFactory = new ConfiguredOutputWriterFactory(fileMetadataStorage, coreConfigurationStorage); } private String savePath = ""; @@ -88,6 +93,7 @@ public class DataLakeExportManager { StreamingResponseBody streamingOutput = output -> dataExplorerQueryManagement.getDataAsStream( sanitizedParams, outputFormat, + outputWriterFactory, "ignore".equals( dataLakeMeasure.getRetentionTime().getRetentionExportConfig().getExportConfig() .missingValueBehaviour()), diff --git a/streampipes-mail/src/main/java/org/apache/streampipes/mail/utils/MailUtils.java b/streampipes-mail/src/main/java/org/apache/streampipes/mail/utils/MailUtils.java index 507765f9f5..06c75428f6 100644 --- a/streampipes-mail/src/main/java/org/apache/streampipes/mail/utils/MailUtils.java +++ b/streampipes-mail/src/main/java/org/apache/streampipes/mail/utils/MailUtils.java @@ -20,7 +20,6 @@ package org.apache.streampipes.mail.utils; import org.apache.streampipes.commons.resources.Resources; import org.apache.streampipes.model.configuration.GeneralConfig; import org.apache.streampipes.model.configuration.SpCoreConfiguration; -import org.apache.streampipes.storage.management.StorageDispatcher; import java.io.IOException; import java.nio.charset.StandardCharsets; @@ -40,8 +39,4 @@ public class MailUtils { public static String readResourceFileToString(String filename) throws IOException { return Resources.asString(filename, StandardCharsets.UTF_8); } - - public static SpCoreConfiguration getSpCoreConfiguration() { - return StorageDispatcher.INSTANCE.getNoSqlStore().getSpCoreConfigurationStorage().get(); - } } 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 03cb4d3ffb..45742f3a2b 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 @@ -23,6 +23,7 @@ import org.apache.streampipes.storage.api.explorer.IDashboardStorage; import org.apache.streampipes.storage.api.explorer.IDataLakeMeasureStorage; import org.apache.streampipes.storage.api.pipeline.IPipelineStorage; import org.apache.streampipes.storage.api.system.IAssetStorage; +import org.apache.streampipes.storage.api.system.IFileMetadataStorage; import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.streampipes.storage.api.user.IPermissionStorage; import org.apache.streampipes.storage.management.StorageDispatcher; @@ -37,6 +38,7 @@ public class SpResourceManager { private final IPipelineStorage pipelineStorage; private final IDataLakeMeasureStorage datasetStorage; private final ISpCoreConfigurationStorage coreConfigurationStorage; + private final IFileMetadataStorage fileMetadataStorage; public SpResourceManager(IPermissionStorage permissionStorage, IChartStorage chartStorage, @@ -45,7 +47,8 @@ public class SpResourceManager { IDashboardStorage dashboardStorage, IPipelineStorage pipelineStorage, IDataLakeMeasureStorage datasetStorage, - ISpCoreConfigurationStorage coreConfigurationStorage) { + ISpCoreConfigurationStorage coreConfigurationStorage, + IFileMetadataStorage fileMetadataStorage) { this.permissionStorage = permissionStorage; this.chartStorage = chartStorage; this.adapterStorage = adapterStorage; @@ -54,6 +57,7 @@ public class SpResourceManager { this.pipelineStorage = pipelineStorage; this.datasetStorage = datasetStorage; this.coreConfigurationStorage = coreConfigurationStorage; + this.fileMetadataStorage = fileMetadataStorage; } public AdapterDescriptionResourceManager manageAdapterDescriptions() { @@ -106,6 +110,10 @@ public class SpResourceManager { return coreConfigurationStorage; } + public IFileMetadataStorage getFileMetadataStorage() { + return fileMetadataStorage; + } + public UserResourceManager manageUsers() { return new UserResourceManager(coreConfigurationStorage); } diff --git a/streampipes-rest-core-base/src/main/java/org/apache/streampipes/rest/core/base/impl/AbstractRestResource.java b/streampipes-rest-core-base/src/main/java/org/apache/streampipes/rest/core/base/impl/AbstractRestResource.java index 4b7e97d941..b0c6cf8ffa 100644 --- a/streampipes-rest-core-base/src/main/java/org/apache/streampipes/rest/core/base/impl/AbstractRestResource.java +++ b/streampipes-rest-core-base/src/main/java/org/apache/streampipes/rest/core/base/impl/AbstractRestResource.java @@ -26,7 +26,6 @@ import org.apache.streampipes.rest.shared.impl.AbstractSharedRestInterface; import org.apache.streampipes.storage.api.core.INoSqlStorage; import org.apache.streampipes.storage.api.pipeline.IPipelineElementDescriptionStorage; import org.apache.streampipes.storage.api.pipeline.IPipelineElementTemplateStorage; -import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.streampipes.storage.api.user.IUserStorage; import org.apache.streampipes.storage.management.StorageDispatcher; @@ -34,10 +33,6 @@ import org.springframework.http.ResponseEntity; public class AbstractRestResource extends AbstractSharedRestInterface { - protected ISpCoreConfigurationStorage getSpCoreConfigurationStorage() { - return getNoSqlStorage().getSpCoreConfigurationStorage(); - } - protected IPipelineElementDescriptionStorage getPipelineElementStorage() { return getNoSqlStorage().getPipelineElementDescriptionStorage(); } diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/Authentication.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/Authentication.java index 43b35ded93..0ff0792fda 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/Authentication.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/Authentication.java @@ -78,11 +78,10 @@ public class Authentication extends AbstractRestResource { private final ISpCoreConfigurationStorage coreConfigurationStorage; public Authentication(AuthenticationManager authenticationManager, - SpResourceManager resourceManager, - ISpCoreConfigurationStorage coreConfigurationStorage) { + SpResourceManager resourceManager) { this.authenticationManager = authenticationManager; this.resourceManager = resourceManager; - this.coreConfigurationStorage = coreConfigurationStorage; + this.coreConfigurationStorage = resourceManager.getCoreConfigurationStorage(); } @PostMapping( @@ -167,7 +166,7 @@ public class Authentication extends AbstractRestResource { public synchronized ResponseEntity<SuccessMessage> doRegister( @RequestBody UserRegistrationData userRegistrationData ) { - GeneralConfig config = getSpCoreConfigurationStorage().get().getGeneralConfig(); + GeneralConfig config = coreConfigurationStorage.get().getGeneralConfig(); if (!config.isAllowSelfRegistration()) { return ResponseEntity.status(HttpStatus.FORBIDDEN).build(); } @@ -208,7 +207,7 @@ public class Authentication extends AbstractRestResource { path = "settings", produces = org.springframework.http.MediaType.APPLICATION_JSON_VALUE) public ResponseEntity<Map<String, Object>> getAuthSettings() { - GeneralConfig config = getSpCoreConfigurationStorage().get().getGeneralConfig(); + GeneralConfig config = coreConfigurationStorage.get().getGeneralConfig(); var termsAcknowledgmentRequired = config.getUserAcknowledgment() != null && config.getUserAcknowledgment().required(); Map<String, Object> response = new HashMap<>(); diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/EmailResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/EmailResource.java index 5d592c6182..78dec7ee41 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/EmailResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/EmailResource.java @@ -20,6 +20,7 @@ package org.apache.streampipes.rest.impl; import org.apache.streampipes.mail.MailSender; import org.apache.streampipes.model.mail.SpEmail; import org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.springframework.http.MediaType; import org.springframework.http.ResponseEntity; @@ -32,11 +33,18 @@ import org.springframework.web.bind.annotation.RestController; @RequestMapping("/api/v2/mail") public class EmailResource extends AbstractAuthGuardedRestResource { + private final ISpCoreConfigurationStorage configurationStorage; + + public EmailResource(ISpCoreConfigurationStorage configurationStorage) { + this.configurationStorage = configurationStorage; + } + @PostMapping(consumes = MediaType.APPLICATION_JSON_VALUE) public ResponseEntity<?> sendEmail(@RequestBody SpEmail email) { - if (getSpCoreConfigurationStorage().get().getEmailConfig().isEmailConfigured()) { + var configuration = configurationStorage.get(); + if (configuration.getEmailConfig().isEmailConfigured()) { try { - new MailSender().sendEmail(email); + new MailSender(configuration).sendEmail(email); return ok(); } catch (Exception e) { return badRequest(e); diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/UserResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/UserResource.java index 91ff8f4ed6..5920b34105 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/UserResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/UserResource.java @@ -30,6 +30,7 @@ import org.apache.streampipes.model.message.Notifications; import org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource; import org.apache.streampipes.rest.security.AuthConstants; import org.apache.streampipes.rest.utils.Utils; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.streampipes.user.management.encryption.SecretEncryptionManager; import org.apache.streampipes.user.management.service.TokenService; import org.apache.streampipes.user.management.util.PasswordUtil; @@ -60,6 +61,12 @@ import java.util.stream.Collectors; @RequestMapping("/api/v2/users") public class UserResource extends AbstractAuthGuardedRestResource { + private final ISpCoreConfigurationStorage configurationStorage; + + public UserResource(ISpCoreConfigurationStorage configurationStorage) { + this.configurationStorage = configurationStorage; + } + @GetMapping(produces = MediaType.APPLICATION_JSON_VALUE) public ResponseEntity<List<ShortUserInfo>> listUsers( @RequestParam("includeServiceAccounts") boolean includeServiceAccounts @@ -166,7 +173,8 @@ public class UserResource extends AbstractAuthGuardedRestResource { } else { String generatedProperty = PasswordUtil.generateRandomPassword(); encryptAndStore(userAccount, generatedProperty); - new MailSender().sendInitialPasswordMail(userAccount.getUsername(), generatedProperty); + new MailSender(configurationStorage.get()) + .sendInitialPasswordMail(userAccount.getUsername(), generatedProperty); } return ok(); } else { diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/EmailConfigurationResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/EmailConfigurationResource.java index 43f521c288..4bd6cf8a63 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/EmailConfigurationResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/EmailConfigurationResource.java @@ -54,21 +54,21 @@ public class EmailConfigurationResource extends AbstractAuthGuardedRestResource @GetMapping(produces = MediaType.APPLICATION_JSON_VALUE) @PreAuthorize(AuthConstants.IS_ADMIN_ROLE) public ResponseEntity<EmailConfig> getMailConfiguration() { - return ok(getSpCoreConfigurationStorage().get().getEmailConfig()); + return ok(configurationStorage.get().getEmailConfig()); } @GetMapping(path = "templates", produces = MediaType.APPLICATION_JSON_VALUE) @PreAuthorize(AuthConstants.IS_ADMIN_ROLE) public ResponseEntity<EmailTemplateConfig> getMailTemplates() { - return ok(getSpCoreConfigurationStorage().get().getEmailTemplateConfig()); + return ok(configurationStorage.get().getEmailTemplateConfig()); } @PutMapping(path = "templates", consumes = MediaType.APPLICATION_JSON_VALUE) @PreAuthorize(AuthConstants.IS_ADMIN_ROLE) public ResponseEntity<Void> updateMailTemplate(@RequestBody EmailTemplateConfig templateConfig) { - var config = getSpCoreConfigurationStorage().get(); + var config = configurationStorage.get(); config.setEmailTemplateConfig(templateConfig); - getSpCoreConfigurationStorage().updateElement(config); + configurationStorage.updateElement(config); return ok(); } @@ -85,10 +85,9 @@ public class EmailConfigurationResource extends AbstractAuthGuardedRestResource config.setSmtpPassword(SecretEncryptionManager.encrypt(config.getSmtpPassword())); config.setSmtpPassEncrypted(true); } - var storage = getSpCoreConfigurationStorage(); - var cfg = storage.get(); + var cfg = configurationStorage.get(); cfg.setEmailConfig(config); - storage.updateElement(cfg); + configurationStorage.updateElement(cfg); return ok(); } diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/ExportProviderConfigurationResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/ExportProviderConfigurationResource.java index 792ec3054b..84e5459fd9 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/ExportProviderConfigurationResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/ExportProviderConfigurationResource.java @@ -22,6 +22,7 @@ import org.apache.streampipes.model.configuration.ExportProviderSettings; import org.apache.streampipes.model.monitoring.SpLogMessage; import org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource; import org.apache.streampipes.rest.security.AuthConstants; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.streampipes.user.management.encryption.SecretEncryptionManager; import org.springframework.http.MediaType; @@ -45,16 +46,22 @@ import java.util.stream.Collectors; @RequestMapping("/api/v2/admin/exportprovider-config") public class ExportProviderConfigurationResource extends AbstractAuthGuardedRestResource { + private final ISpCoreConfigurationStorage configurationStorage; + + public ExportProviderConfigurationResource(ISpCoreConfigurationStorage configurationStorage) { + this.configurationStorage = configurationStorage; + } + @GetMapping(produces = MediaType.APPLICATION_JSON_VALUE) @PreAuthorize(AuthConstants.IS_ADMIN_ROLE) public ResponseEntity<List<ExportProviderSettings>> getExportProviderConfiguration() { - return ok(getSpCoreConfigurationStorage().get().getExportProviderSettings()); + return ok(configurationStorage.get().getExportProviderSettings()); } @GetMapping(value = "/{providerId}", produces = MediaType.APPLICATION_JSON_VALUE) @PreAuthorize(AuthConstants.IS_ADMIN_ROLE) public ResponseEntity<ExportProviderSettings> getExportProviderSettingById(@PathVariable String providerId) { - return getSpCoreConfigurationStorage().get().getExportProviderSettings().stream() + return configurationStorage.get().getExportProviderSettings().stream() .filter(setting -> setting.getProviderId().equalsIgnoreCase(providerId)) .findFirst() .map(ResponseEntity::ok) @@ -65,7 +72,7 @@ public class ExportProviderConfigurationResource extends AbstractAuthGuardedRest @PreAuthorize(AuthConstants.IS_ADMIN_ROLE) public ResponseEntity<?> testExportProviderSettingById(@PathVariable String providerId) { // Get Export Provider Settings - Optional<ExportProviderSettings> exportProviderSetting = getSpCoreConfigurationStorage().get() + Optional<ExportProviderSettings> exportProviderSetting = configurationStorage.get() .getExportProviderSettings().stream() .filter(setting -> setting.getProviderId().equalsIgnoreCase(providerId)) .findFirst(); @@ -95,8 +102,7 @@ public class ExportProviderConfigurationResource extends AbstractAuthGuardedRest config.setSecretKey(SecretEncryptionManager.encrypt(config.getSecretKey())); config.setSecretEncrypted(true); } - var storage = getSpCoreConfigurationStorage(); - var cfg = storage.get(); + var cfg = configurationStorage.get(); List<ExportProviderSettings> providerSettings = cfg.getExportProviderSettings(); if (providerSettings == null) { @@ -110,7 +116,7 @@ public class ExportProviderConfigurationResource extends AbstractAuthGuardedRest providerSettingsWithoutExisting.add(config); cfg.setExportProviderSettings(providerSettingsWithoutExisting); - storage.updateElement(cfg); + configurationStorage.updateElement(cfg); return ok(); } @@ -119,16 +125,15 @@ public class ExportProviderConfigurationResource extends AbstractAuthGuardedRest @PreAuthorize(AuthConstants.IS_ADMIN_ROLE) public ResponseEntity<Void> deleteExportProviderConfiguration(@PathVariable String providerId) { - List<ExportProviderSettings> allProviders = getSpCoreConfigurationStorage().get().getExportProviderSettings(); + List<ExportProviderSettings> allProviders = configurationStorage.get().getExportProviderSettings(); List<ExportProviderSettings> filteredProviders = allProviders.stream() .filter(provider -> !provider.getProviderId().equals(providerId)) .collect(Collectors.toList()); - var storage = getSpCoreConfigurationStorage(); - var cfg = storage.get(); + var cfg = configurationStorage.get(); cfg.setExportProviderSettings(filteredProviders); - storage.updateElement(cfg); + configurationStorage.updateElement(cfg); return ok(); } diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/ExtensionsInstallationResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/ExtensionsInstallationResource.java index b62931e4da..3f1a5e3d66 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/ExtensionsInstallationResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/ExtensionsInstallationResource.java @@ -33,7 +33,6 @@ import org.apache.streampipes.resource.management.AdapterDescriptionResourceMana import org.apache.streampipes.resource.management.DataProcessorResourceManager; import org.apache.streampipes.resource.management.DataSinkResourceManager; import org.apache.streampipes.resource.management.DataStreamResourceManager; -import org.apache.streampipes.resource.management.PermissionResourceManager; import org.apache.streampipes.resource.management.SpResourceManager; import org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource; import org.apache.streampipes.rest.security.AuthConstants; @@ -65,16 +64,17 @@ public class ExtensionsInstallationResource extends AbstractAuthGuardedRestResou private final AdapterDescriptionResourceManager adapterDescriptionResourceManager; private final DataStreamResourceManager dataStreamResourceManager; private final SpResourceManager resourceManager; + private final AssetManager assetManager; public ExtensionsInstallationResource(ExtensionServiceRequestManager extensionServiceRequestManager, SpResourceManager resourceManager) { this.resourceManager = resourceManager; this.extensionServiceRequestManager = extensionServiceRequestManager; - PermissionResourceManager permissionResourceManager = resourceManager.managePermissions(); this.dataSinkResourceManager = resourceManager.manageDataSinks(); this.dataProcessorResourceManager = resourceManager.manageDataProcessors(); this.adapterDescriptionResourceManager = resourceManager.manageAdapterDescriptions(); this.dataStreamResourceManager = resourceManager.manageDataStreams(); + this.assetManager = new AssetManager(resourceManager.getCoreConfigurationStorage()); } @PostMapping( @@ -124,7 +124,7 @@ public class ExtensionsInstallationResource extends AbstractAuthGuardedRestResou return constructErrorMessage(new Notification(NotificationType.STORAGE_ERROR.title(), NotificationType.STORAGE_ERROR.description())); } - AssetManager.deleteAsset(appId); + assetManager.deleteAsset(appId); } catch (IOException e) { return constructErrorMessage(new Notification(NotificationType.STORAGE_ERROR.title(), NotificationType.STORAGE_ERROR.description())); diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/GeneralConfigurationResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/GeneralConfigurationResource.java index 7231989716..759efc506a 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/GeneralConfigurationResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/GeneralConfigurationResource.java @@ -20,6 +20,7 @@ package org.apache.streampipes.rest.impl.admin; import org.apache.streampipes.model.configuration.GeneralConfig; import org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource; import org.apache.streampipes.rest.security.AuthConstants; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.springframework.http.HttpHeaders; import org.springframework.http.MediaType; @@ -41,20 +42,25 @@ import java.util.Base64; @RequestMapping("/api/v2/admin/general-config") public class GeneralConfigurationResource extends AbstractAuthGuardedRestResource { + private final ISpCoreConfigurationStorage configurationStorage; + + public GeneralConfigurationResource(ISpCoreConfigurationStorage configurationStorage) { + this.configurationStorage = configurationStorage; + } + @GetMapping(produces = MediaType.APPLICATION_JSON_VALUE) @PreAuthorize(AuthConstants.IS_ADMIN_ROLE) public GeneralConfig getGeneralConfiguration() { - return getSpCoreConfigurationStorage().get().getGeneralConfig(); + return configurationStorage.get().getGeneralConfig(); } @PutMapping(consumes = MediaType.APPLICATION_JSON_VALUE) @PreAuthorize(AuthConstants.IS_ADMIN_ROLE) public ResponseEntity<Void> updateGeneralConfiguration(@RequestBody GeneralConfig config) { config.setConfigured(true); - var storage = getSpCoreConfigurationStorage(); - var cfg = storage.get(); + var cfg = configurationStorage.get(); cfg.setGeneralConfig(config); - storage.updateElement(cfg); + configurationStorage.updateElement(cfg); return ok(); } diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/LocationConfigurationResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/LocationConfigurationResource.java index 7872efed0b..2980d6c277 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/LocationConfigurationResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/LocationConfigurationResource.java @@ -21,6 +21,7 @@ package org.apache.streampipes.rest.impl.admin; import org.apache.streampipes.model.configuration.LocationConfig; import org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource; import org.apache.streampipes.rest.security.AuthConstants; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.springframework.http.MediaType; import org.springframework.http.ResponseEntity; @@ -35,18 +36,23 @@ import org.springframework.web.bind.annotation.RestController; @RequestMapping("/api/v2/admin/location-config") public class LocationConfigurationResource extends AbstractAuthGuardedRestResource { + private final ISpCoreConfigurationStorage configurationStorage; + + public LocationConfigurationResource(ISpCoreConfigurationStorage configurationStorage) { + this.configurationStorage = configurationStorage; + } + @GetMapping(produces = MediaType.APPLICATION_JSON_VALUE) public LocationConfig getLocationConfig() { - return getSpCoreConfigurationStorage().get().getLocationConfig(); + return configurationStorage.get().getLocationConfig(); } @PutMapping(consumes = MediaType.APPLICATION_JSON_VALUE) @PreAuthorize(AuthConstants.IS_ADMIN_ROLE) public ResponseEntity<Void> updateGeneralConfiguration(@RequestBody LocationConfig config) { - var storage = getSpCoreConfigurationStorage(); - var cfg = storage.get(); + var cfg = configurationStorage.get(); cfg.setLocationConfig(config); - storage.updateElement(cfg); + configurationStorage.updateElement(cfg); return ok(); } 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 41bc5ee660..fa5d606826 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 @@ -20,6 +20,7 @@ package org.apache.streampipes.rest.impl.datalake; import org.apache.streampipes.commons.exceptions.SpRuntimeException; import org.apache.streampipes.dataexplorer.api.IDataExplorerQueryManagement; +import org.apache.streampipes.dataexplorer.export.ConfiguredOutputWriterFactory; import org.apache.streampipes.dataexplorer.export.OutputFormat; import org.apache.streampipes.dataexplorer.management.DataExplorerDispatcher; import org.apache.streampipes.export.DataLakeExportManager; @@ -99,6 +100,7 @@ public class DataLakeResource extends AbstractDataLakeResource { private final IDataExplorerQueryManagement dataExplorerQueryManagement; private final DataLakeExportManager dataLakeExportManager; private final IDataLakeMeasureStorage datasetStorage; + private final ConfiguredOutputWriterFactory outputWriterFactory; public DataLakeResource(IChartStorage chartStorage, SpResourceManager resourceManager) { @@ -107,10 +109,14 @@ public class DataLakeResource extends AbstractDataLakeResource { this.dataExplorerQueryManagement = new DataExplorerDispatcher() .getDataExplorerManager() .getQueryManagement(this.dataLakeMeasureManagement); + this.outputWriterFactory = new ConfiguredOutputWriterFactory( + resourceManager.getFileMetadataStorage(), + resourceManager.getCoreConfigurationStorage()); this.dataLakeExportManager = new DataLakeExportManager( this.dataLakeMeasureManagement, dataExplorerQueryManagement, - resourceManager.getCoreConfigurationStorage()); + resourceManager.getCoreConfigurationStorage(), + resourceManager.getFileMetadataStorage()); } @DeleteMapping(path = "/measurements/{measurementName}") @@ -275,6 +281,7 @@ public class DataLakeResource extends AbstractDataLakeResource { StreamingResponseBody streamingOutput = output -> dataExplorerQueryManagement.getDataAsStream( sanitizedParams, outputFormat, + outputWriterFactory, isIgnoreMissingValues(missingValueBehaviour), output); diff --git a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/extensions/ExtensionServiceRequestConfiguration.java b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/extensions/ExtensionServiceRequestConfiguration.java index 28f916b637..0808c1305c 100644 --- a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/extensions/ExtensionServiceRequestConfiguration.java +++ b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/extensions/ExtensionServiceRequestConfiguration.java @@ -29,6 +29,7 @@ import org.apache.streampipes.storage.api.explorer.IDashboardStorage; import org.apache.streampipes.storage.api.explorer.IDataLakeMeasureStorage; import org.apache.streampipes.storage.api.pipeline.IPipelineStorage; import org.apache.streampipes.storage.api.system.IAssetStorage; +import org.apache.streampipes.storage.api.system.IFileMetadataStorage; import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.streampipes.storage.api.user.IPermissionStorage; @@ -98,7 +99,8 @@ public class ExtensionServiceRequestConfiguration { IAssetStorage assetStorage, IPipelineStorage pipelineStorage, IDataLakeMeasureStorage datasetStorage, - ISpCoreConfigurationStorage coreConfigurationStorage) { + ISpCoreConfigurationStorage coreConfigurationStorage, + IFileMetadataStorage fileMetadataStorage) { return new SpResourceManager( permissionStorage, chartStorage, @@ -107,7 +109,8 @@ public class ExtensionServiceRequestConfiguration { dashboardStorage, pipelineStorage, datasetStorage, - coreConfigurationStorage + coreConfigurationStorage, + fileMetadataStorage ); } diff --git a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/CustomOAuth2UserService.java b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/CustomOAuth2UserService.java index 0177814cb9..d263258ec1 100755 --- a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/CustomOAuth2UserService.java +++ b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/CustomOAuth2UserService.java @@ -18,8 +18,8 @@ package org.apache.streampipes.service.core.oauth2; +import org.apache.streampipes.resource.management.SpResourceManager; import org.apache.streampipes.rest.security.OAuth2AuthenticationProcessingException; -import org.apache.streampipes.storage.api.user.IPermissionStorage; import org.springframework.security.core.AuthenticationException; import org.springframework.security.oauth2.client.userinfo.DefaultOAuth2UserService; @@ -33,10 +33,10 @@ import java.util.HashMap; @Service public class CustomOAuth2UserService extends DefaultOAuth2UserService { - private final IPermissionStorage permissionStorage; + private final SpResourceManager resourceManager; - public CustomOAuth2UserService(IPermissionStorage permissionStorage) { - this.permissionStorage = permissionStorage; + public CustomOAuth2UserService(SpResourceManager resourceManager) { + this.resourceManager = resourceManager; } @Override @@ -45,7 +45,7 @@ public class CustomOAuth2UserService extends DefaultOAuth2UserService { try { var attributes = new HashMap<>(oAuth2User.getAttributes()); var provider = oAuth2UserRequest.getClientRegistration().getRegistrationId(); - return new UserService(permissionStorage).processUserRegistration(provider, attributes); + return new UserService(resourceManager).processUserRegistration(provider, attributes); } catch (AuthenticationException e) { throw e; } catch (Exception e) { diff --git a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/CustomOidcUserService.java b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/CustomOidcUserService.java index eb9371402d..c8e413e6ed 100755 --- a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/CustomOidcUserService.java +++ b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/CustomOidcUserService.java @@ -20,8 +20,8 @@ package org.apache.streampipes.service.core.oauth2; import org.apache.streampipes.commons.environment.Environments; +import org.apache.streampipes.resource.management.SpResourceManager; import org.apache.streampipes.rest.security.OAuth2AuthenticationProcessingException; -import org.apache.streampipes.storage.api.user.IPermissionStorage; import org.springframework.security.core.AuthenticationException; import org.springframework.security.oauth2.client.oidc.userinfo.OidcUserRequest; @@ -35,10 +35,10 @@ import java.util.Objects; @Service public class CustomOidcUserService extends OidcUserService { - private final IPermissionStorage permissionStorage; + private final SpResourceManager resourceManager; - public CustomOidcUserService(IPermissionStorage permissionStorage) { - this.permissionStorage = permissionStorage; + public CustomOidcUserService(SpResourceManager resourceManager) { + this.resourceManager = resourceManager; var env = Environments.getEnvironment(); this.setRetrieveUserInfo(req -> { var config = env.getOAuthConfigurations() @@ -56,7 +56,7 @@ public class CustomOidcUserService extends OidcUserService { OidcUser oidcUser = super.loadUser(userRequest); try { var provider = userRequest.getClientRegistration().getRegistrationId(); - return new UserService(permissionStorage).processUserRegistration( + return new UserService(resourceManager).processUserRegistration( provider, oidcUser.getAttributes(), oidcUser.getIdToken(), diff --git a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/UserService.java b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/UserService.java index b2c2991134..c3754be18a 100755 --- a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/UserService.java +++ b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/UserService.java @@ -24,8 +24,10 @@ import org.apache.streampipes.commons.environment.model.OAuthConfiguration; import org.apache.streampipes.model.client.user.Group; import org.apache.streampipes.model.client.user.Role; import org.apache.streampipes.model.client.user.UserAccount; +import org.apache.streampipes.resource.management.SpResourceManager; import org.apache.streampipes.resource.management.UserResourceManager; import org.apache.streampipes.rest.security.OAuth2AuthenticationProcessingException; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.streampipes.storage.api.user.IPermissionStorage; import org.apache.streampipes.storage.api.user.IRoleStorage; import org.apache.streampipes.storage.api.user.IUserGroupStorage; @@ -55,15 +57,17 @@ public class UserService { private List<Role> allRoles; private List<Group> allGroups; private final IPermissionStorage permissionStorage; + private final ISpCoreConfigurationStorage configurationStorage; - public UserService(IPermissionStorage permissionStorage) { + public UserService(SpResourceManager resourceManager) { this.userStorage = StorageDispatcher.INSTANCE.getNoSqlStore().getUserStorageAPI(); this.roleStorage = StorageDispatcher.INSTANCE.getNoSqlStore().getRoleStorage(); this.groupStorage = StorageDispatcher.INSTANCE.getNoSqlStore().getUserGroupStorage(); this.allGroups = this.groupStorage.findAll(); this.allRoles = this.roleStorage.findAll(); this.env = Environments.getEnvironment(); - this.permissionStorage = permissionStorage; + this.permissionStorage = resourceManager.managePermissions().getDb(); + this.configurationStorage = resourceManager.getCoreConfigurationStorage(); } public OidcUserAccountDetails processUserRegistration(String registrationId, @@ -102,7 +106,7 @@ public class UserService { user = toUserAccount(registrationId, principalId, email, fullName); user.setLastLoginAtMillis(System.currentTimeMillis()); applyRoles(user, oAuthConfig, attributes, true); - new UserResourceManager().storeUser(user); + new UserResourceManager(configurationStorage).storeUser(user); } user = (UserAccount) userStorage.getUserById(principalId); 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 b7fc06ea82..b9d1d17552 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 @@ -56,7 +56,8 @@ public class DataLakeScheduler implements SchedulingConfigurer { dataExplorerSchemaManagement, new DataExplorerDispatcher().getDataExplorerManager() .getQueryManagement(dataExplorerSchemaManagement), - resourceManager.getCoreConfigurationStorage()); + resourceManager.getCoreConfigurationStorage(), + resourceManager.getFileMetadataStorage()); } public void cleanupMeasurements() { 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 fb1c15e30d..6aea14cbf7 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 @@ -26,6 +26,7 @@ import org.apache.streampipes.storage.api.explorer.IDataLakeMeasureStorage; import org.apache.streampipes.storage.api.function.IFunctionStateStorage; import org.apache.streampipes.storage.api.pipeline.IPipelineStorage; import org.apache.streampipes.storage.api.system.IAssetStorage; +import org.apache.streampipes.storage.api.system.IFileMetadataStorage; import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.streampipes.storage.api.user.IPermissionStorage; import org.apache.streampipes.storage.couchdb.impl.connect.AdapterInstanceStorageImpl; @@ -36,6 +37,7 @@ import org.apache.streampipes.storage.couchdb.impl.function.FunctionStateStorage import org.apache.streampipes.storage.couchdb.impl.pipeline.PipelineStorageImpl; import org.apache.streampipes.storage.couchdb.impl.system.AssetStorageImpl; import org.apache.streampipes.storage.couchdb.impl.system.CoreConfigurationStorageImpl; +import org.apache.streampipes.storage.couchdb.impl.system.FileMetadataStorageImpl; import org.apache.streampipes.storage.couchdb.impl.user.PermissionStorageImpl; import org.apache.streampipes.storage.couchdb.utils.Utils; @@ -110,6 +112,11 @@ public class StorageApiConfiguration { return new CoreConfigurationStorageImpl(); } + @Bean + public IFileMetadataStorage fileMetadataStorage() { + return new FileMetadataStorageImpl(); + } + @Bean public IPipelineStorage pipelineStorage(CacheManager cacheManager) { IPipelineStorage delegate = new PipelineStorageImpl(); 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 faaa5a5c1e..90fb043409 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 @@ -31,7 +31,6 @@ import org.apache.streampipes.storage.api.system.IExtensionsServiceStorage; import org.apache.streampipes.storage.api.system.IFileMetadataStorage; import org.apache.streampipes.storage.api.system.IGenericStorage; import org.apache.streampipes.storage.api.system.IImageStorage; -import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.streampipes.storage.api.system.ITransformationScriptTemplateStorage; import org.apache.streampipes.storage.api.user.IPasswordRecoveryTokenStorage; import org.apache.streampipes.storage.api.user.IPrivilegeStorage; @@ -77,8 +76,6 @@ public interface INoSqlStorage { IExtensionsServiceConfigurationStorage getExtensionsServiceConfigurationStorage(); - ISpCoreConfigurationStorage getSpCoreConfigurationStorage(); - IRoleStorage getRoleStorage(); IPrivilegeStorage getPrivilegeStorage(); 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 8bd1cc086c..561eed5cca 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 @@ -32,7 +32,6 @@ import org.apache.streampipes.storage.api.system.IExtensionsServiceStorage; import org.apache.streampipes.storage.api.system.IFileMetadataStorage; import org.apache.streampipes.storage.api.system.IGenericStorage; import org.apache.streampipes.storage.api.system.IImageStorage; -import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.streampipes.storage.api.system.ITransformationScriptTemplateStorage; import org.apache.streampipes.storage.api.user.IPasswordRecoveryTokenStorage; import org.apache.streampipes.storage.api.user.IPrivilegeStorage; @@ -50,7 +49,6 @@ import org.apache.streampipes.storage.couchdb.impl.pipeline.PipelineCanvasMetada import org.apache.streampipes.storage.couchdb.impl.pipeline.PipelineElementDescriptionStorageImpl; import org.apache.streampipes.storage.couchdb.impl.pipeline.PipelineElementTemplateStorageImpl; import org.apache.streampipes.storage.couchdb.impl.system.CertificateStorageImpl; -import org.apache.streampipes.storage.couchdb.impl.system.CoreConfigurationStorageImpl; import org.apache.streampipes.storage.couchdb.impl.system.ExtensionsServiceConfigurationStorageImpl; import org.apache.streampipes.storage.couchdb.impl.system.ExtensionsServiceStorageImpl; import org.apache.streampipes.storage.couchdb.impl.system.FileMetadataStorageImpl; @@ -152,11 +150,6 @@ public class CouchDbStorageManager implements INoSqlStorage { return new ExtensionsServiceConfigurationStorageImpl(); } - @Override - public ISpCoreConfigurationStorage getSpCoreConfigurationStorage() { - return new CoreConfigurationStorageImpl(); - } - @Override public IRoleStorage getRoleStorage() { return new RoleStorageImpl();
