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 c39c564cd030448dc6ea54fe0e1595b91a90263a Author: Dominik Riemer <[email protected]> AuthorDate: Tue Jun 23 20:54:50 2026 +0200 Migrate configuration storage --- .../management/AdapterMasterManagement.java | 2 +- .../management/WorkerAdministrationManagement.java | 6 +- .../streampipes/connect/management/util/Utils.java | 32 ----- .../export/ConfiguredExcelOutputWriter.java | 9 +- .../streampipes/export/DataLakeExportManager.java | 11 +- .../export/dataimport/PerformImportGenerator.java | 3 +- .../export/generator/ExportPackageGenerator.java | 4 +- .../apache/streampipes/mail/AbstractMailer.java | 24 ++-- .../org/apache/streampipes/mail/MailSender.java | 22 ++-- .../mail/template/AbstractMailTemplate.java | 10 +- .../apache/streampipes/mail/utils/MailUtils.java | 8 +- .../api/extensions/ExtensionServiceRequests.java | 10 +- .../streampipes/manager/assets/AssetConstants.java | 31 ----- .../streampipes/manager/assets/AssetExtractor.java | 13 +- .../streampipes/manager/assets/AssetManager.java | 31 +++-- .../streampipes/manager/file/FileConstants.java | 30 ----- .../streampipes/manager/file/FileHandler.java | 10 +- .../streampipes/manager/file/FileManager.java | 5 +- .../migration/AbstractMigrationManager.java | 2 +- .../manager/setup/AutoInstallation.java | 2 +- .../manager/setup/InstallationConfiguration.java | 5 +- .../manager/setup/SpCoreConfigurationStep.java | 12 +- .../manager/setup/StreamPipesEnvChecker.java | 17 ++- .../streampipes/manager/util/AuthTokenUtils.java | 7 +- .../manager/verification/ElementVerifier.java | 15 ++- .../manager/verification/TypedElementVerifier.java | 13 +- .../resource/management/SpResourceManager.java | 10 +- .../apache/streampipes/rest/ResetManagement.java | 2 +- .../streampipes/rest/impl/Authentication.java | 10 +- .../apache/streampipes/rest/impl/FileResource.java | 5 +- .../rest/impl/PipelineElementAsset.java | 8 +- .../org/apache/streampipes/rest/impl/Setup.java | 8 +- .../admin/ExtensionsServiceEndpointResource.java | 8 +- .../rest/impl/admin/MigrationResource.java | 7 +- .../rest/impl/datalake/DataLakeResource.java | 5 +- .../service/core/StreamPipesCoreApplication.java | 12 +- .../service/core/WebSecurityConfig.java | 10 +- .../ExtensionServiceRequestConfiguration.java | 7 +- .../core/filter/TokenAuthenticationFilter.java | 6 +- .../core/migrations/AvailableMigrations.java | 5 +- .../v0980/AddDefaultExportProviderMigration.java | 8 +- .../oauth2/OAuth2AuthenticationSuccessHandler.java | 140 +++++++++++---------- .../service/core/scheduler/DataLakeScheduler.java | 3 +- .../CertificateExpiryEmailScheduler.java | 7 ++ .../core/storage/StorageApiConfiguration.java | 7 ++ .../user/management/jwt/JwtTokenProvider.java | 19 ++- .../user/management/jwt/SpKeyResolver.java | 11 +- 47 files changed, 323 insertions(+), 309 deletions(-) diff --git a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterMasterManagement.java b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterMasterManagement.java index f4780b0d48..53b7c66fad 100644 --- a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterMasterManagement.java +++ b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterMasterManagement.java @@ -225,7 +225,7 @@ public class AdapterMasterManagement { storageApi::update, SpServiceUrlProvider.DATA_STREAM, requestManager, - resourceManager.managePermissions() + resourceManager ); verifier.verifyAndAdd(principalSid, false); } diff --git a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/WorkerAdministrationManagement.java b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/WorkerAdministrationManagement.java index 63e23fd51c..8f71808fb5 100644 --- a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/WorkerAdministrationManagement.java +++ b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/WorkerAdministrationManagement.java @@ -45,6 +45,7 @@ public class WorkerAdministrationManagement { private final PermissionResourceManager permissionResourceManager; private final UserResourceManager userResourceManager; private final ExtensionServiceRequestManager requestManager; + private final AssetManager assetManager; public WorkerAdministrationManagement( IAdapterStorage adapterDescriptionStorage, @@ -55,6 +56,7 @@ public class WorkerAdministrationManagement { this.permissionStorage = resourceManager.managePermissions().getDb(); this.permissionResourceManager = resourceManager.managePermissions(); this.requestManager = requestManager; + this.assetManager = new AssetManager(resourceManager.getCoreConfigurationStorage()); } public void performAdapterMigrations(List<SpServiceTag> tags) { @@ -63,10 +65,10 @@ public class WorkerAdministrationManagement { installedAdapters.stream() .filter(adapter -> tags.stream().anyMatch(tag -> tag.getValue().equals(adapter.getAppId()))) .forEach(adapter -> { - if (!AssetManager.existsAssetDir(adapter.getAppId())) { + if (!assetManager.existsAssetDir(adapter.getAppId())) { try { LOG.info("Updating assets for adapter {}", adapter.getAppId()); - AssetManager.storeAsset(SpServiceUrlProvider.ADAPTER, adapter.getAppId(), requestManager); + assetManager.storeAsset(SpServiceUrlProvider.ADAPTER, adapter.getAppId(), requestManager); } catch (IOException | NoServiceEndpointsAvailableException e) { LOG.error( "Could not fetch asset for adapter {}, please try to manually update this adapter.", diff --git a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/util/Utils.java b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/util/Utils.java deleted file mode 100644 index ca6d05337d..0000000000 --- a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/util/Utils.java +++ /dev/null @@ -1,32 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - * - */ - -package org.apache.streampipes.connect.management.util; - -import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; -import org.apache.streampipes.storage.management.StorageDispatcher; - -public class Utils { - - public static ISpCoreConfigurationStorage getCoreConfigStorage() { - return StorageDispatcher - .INSTANCE - .getNoSqlStore() - .getSpCoreConfigurationStorage(); - } -} diff --git a/streampipes-data-explorer-export/src/main/java/org/apache/streampipes/dataexplorer/export/ConfiguredExcelOutputWriter.java b/streampipes-data-explorer-export/src/main/java/org/apache/streampipes/dataexplorer/export/ConfiguredExcelOutputWriter.java index 100ebf0b46..a20e2b064f 100644 --- a/streampipes-data-explorer-export/src/main/java/org/apache/streampipes/dataexplorer/export/ConfiguredExcelOutputWriter.java +++ b/streampipes-data-explorer-export/src/main/java/org/apache/streampipes/dataexplorer/export/ConfiguredExcelOutputWriter.java @@ -23,6 +23,7 @@ import org.apache.streampipes.model.datalake.DataLakeMeasure; import org.apache.streampipes.model.datalake.param.ProvidedRestQueryParams; import org.apache.streampipes.model.datalake.param.SupportedRestQueryParams; import org.apache.streampipes.storage.api.system.IFileMetadataStorage; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.poi.ss.usermodel.Sheet; import org.apache.poi.xssf.streaming.SXSSFWorkbook; @@ -38,6 +39,7 @@ import java.util.Objects; public class ConfiguredExcelOutputWriter extends ConfiguredOutputWriter { private final IFileMetadataStorage storage; + private final ISpCoreConfigurationStorage coreConfigurationStorage; private SXSSFWorkbook wb; private Sheet ws; @@ -47,8 +49,10 @@ public class ConfiguredExcelOutputWriter extends ConfiguredOutputWriter { private DataLakeMeasure schema; private String headerColumnNameStrategy; - public ConfiguredExcelOutputWriter(IFileMetadataStorage fileMetadataStorage) { + public ConfiguredExcelOutputWriter(IFileMetadataStorage fileMetadataStorage, + ISpCoreConfigurationStorage coreConfigurationStorage) { this.storage = fileMetadataStorage; + this.coreConfigurationStorage = coreConfigurationStorage; } @Override @@ -70,7 +74,8 @@ public class ConfiguredExcelOutputWriter extends ConfiguredOutputWriter { if (useTemplate && Objects.nonNull(templateId)) { var fileMetadata = storage.getElementById(templateId); if (fileMetadata != null) { - var path = new FileManager().getFile(fileMetadata.getFilename()).getAbsoluteFile().toPath(); + var path = new FileManager(coreConfigurationStorage) + .getFile(fileMetadata.getFilename()).getAbsoluteFile().toPath(); try (InputStream is = Files.newInputStream(path)) { XSSFWorkbook templateWorkbook = new XSSFWorkbook(is); wb = new SXSSFWorkbook(templateWorkbook); 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 a0274f5800..0d129c9c5a 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 @@ -30,7 +30,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.management.StorageDispatcher; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -50,11 +50,14 @@ public class DataLakeExportManager { private final IDataExplorerSchemaManagement dataExplorerSchemaManagement; private final IDataExplorerQueryManagement dataExplorerQueryManagement; + private final ISpCoreConfigurationStorage coreConfigurationStorage; public DataLakeExportManager(IDataExplorerSchemaManagement dataLakeSchemaManagement, - IDataExplorerQueryManagement dataLakeQueryManagement) { + IDataExplorerQueryManagement dataLakeQueryManagement, + ISpCoreConfigurationStorage coreConfigurationStorage) { this.dataExplorerSchemaManagement = dataLakeSchemaManagement; this.dataExplorerQueryManagement = dataLakeQueryManagement; + this.coreConfigurationStorage = coreConfigurationStorage; } private String savePath = ""; @@ -94,9 +97,7 @@ public class DataLakeExportManager { .getExportProviderId(); // FInd Item in Document - List<ExportProviderSettings> exportProviders = StorageDispatcher.INSTANCE - .getNoSqlStore() - .getSpCoreConfigurationStorage() + List<ExportProviderSettings> exportProviders = coreConfigurationStorage .get() .getExportProviderSettings(); 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 64c6ed766b..bd98673457 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 @@ -155,7 +155,8 @@ public class PerformImportGenerator extends ImportGenerator<Void> { writeDocument(document, resolver); byte[] file = zipContent.get( fileMetadata.getFilename().substring(0, fileMetadata.getFilename().lastIndexOf("."))); - new FileHandler().storeFile(fileMetadata.getFilename(), new ByteArrayInputStream(file)); + new FileHandler(resourceManager.getCoreConfigurationStorage()) + .storeFile(fileMetadata.getFilename(), new ByteArrayInputStream(file)); } @Override 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 86ad0027ed..1b0b6ee249 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 @@ -155,7 +155,9 @@ public class ExportPackageGenerator { String filename = fileResolver.findDocument(item.getResourceId()).getFilename(); addDoc(builder, item, new FileResolver(), manifest::addFile); try { - builder.addBinary(filename, Files.readAllBytes(new FileManager().getFile(filename).toPath())); + builder.addBinary(filename, Files.readAllBytes( + new FileManager(resourceManager.getCoreConfigurationStorage()).getFile(filename).toPath()) + ); } catch (IOException e) { LOG.warn("Could not add binary file to export package: {}", e.getMessage()); } diff --git a/streampipes-mail/src/main/java/org/apache/streampipes/mail/AbstractMailer.java b/streampipes-mail/src/main/java/org/apache/streampipes/mail/AbstractMailer.java index b974929c2b..34af63c277 100644 --- a/streampipes-mail/src/main/java/org/apache/streampipes/mail/AbstractMailer.java +++ b/streampipes-mail/src/main/java/org/apache/streampipes/mail/AbstractMailer.java @@ -18,8 +18,8 @@ package org.apache.streampipes.mail; import org.apache.streampipes.mail.config.MailConfigurationBuilder; -import org.apache.streampipes.mail.utils.MailUtils; import org.apache.streampipes.model.configuration.EmailConfig; +import org.apache.streampipes.model.configuration.SpCoreConfiguration; import org.apache.streampipes.user.management.encryption.SecretEncryptionManager; import org.simplejavamail.api.email.Email; @@ -35,17 +35,22 @@ import java.util.stream.Collectors; public class AbstractMailer { + protected final SpCoreConfiguration spCoreConfiguration; + + public AbstractMailer(SpCoreConfiguration configuration) { + this.spCoreConfiguration = configuration; + } + protected Mailer getMailer() { - return getMailer(getEmailConfig()); + return getMailer(getEmailConfig(spCoreConfiguration.getEmailConfig())); } protected Mailer getMailer(EmailConfig config) { return new MailConfigurationBuilder().buildMailerFromConfig(config); } - protected EmailConfig getEmailConfig() { - EmailConfig config = MailUtils.getSpCoreConfiguration().getEmailConfig(); - return getDecryptedEmailConfig(config); + protected EmailConfig getEmailConfig(EmailConfig emailConfig) { + return getDecryptedEmailConfig(emailConfig); } protected EmailConfig getDecryptedEmailConfig(EmailConfig config) { @@ -63,7 +68,10 @@ public class AbstractMailer { } protected void deliverMail(Email email) { - deliverMail(getEmailConfig(), email); + deliverMail( + getEmailConfig(spCoreConfiguration.getEmailConfig()), + email + ); } protected void deliverMail(EmailConfig config, Email email) { @@ -73,7 +81,9 @@ public class AbstractMailer { } protected EmailPopulatingBuilder baseEmail() { - return baseEmail(getEmailConfig()); + return baseEmail( + getEmailConfig(spCoreConfiguration.getEmailConfig()) + ); } protected EmailPopulatingBuilder baseEmail(EmailConfig config) { diff --git a/streampipes-mail/src/main/java/org/apache/streampipes/mail/MailSender.java b/streampipes-mail/src/main/java/org/apache/streampipes/mail/MailSender.java index 43f9d57803..145ef282d3 100644 --- a/streampipes-mail/src/main/java/org/apache/streampipes/mail/MailSender.java +++ b/streampipes-mail/src/main/java/org/apache/streampipes/mail/MailSender.java @@ -22,6 +22,7 @@ import org.apache.streampipes.mail.template.CustomMailTemplate; import org.apache.streampipes.mail.template.InitialPasswordMailTemplate; import org.apache.streampipes.mail.template.PasswordRecoveryMailTemplate; import org.apache.streampipes.mail.utils.MailUtils; +import org.apache.streampipes.model.configuration.SpCoreConfiguration; import org.apache.streampipes.model.mail.SpEmail; import org.simplejavamail.api.email.Email; @@ -30,6 +31,10 @@ import java.io.IOException; public class MailSender extends AbstractMailer { + public MailSender(SpCoreConfiguration configuration) { + super(configuration); + } + public void sendEmail(SpEmail mail) throws IOException { Email email = baseEmail() .withRecipients(toSimpleRecipientList(mail.getRecipients())) @@ -37,7 +42,7 @@ public class MailSender extends AbstractMailer { .appendTextHTML(new CustomMailTemplate( mail.getSubject(), mail.getPreheader(), - mail.getMessage()).generateTemplate()) + mail.getMessage()).generateTemplate(spCoreConfiguration.getEmailTemplateConfig())) .buildEmail(); deliverMail(email); @@ -46,8 +51,9 @@ public class MailSender extends AbstractMailer { public void sendAccountActivationMail(String recipientAddress, String activationCode) throws IOException { Email email = baseEmail() - .withSubject(MailUtils.extractAppName() + " - Account Activation") - .appendTextHTML(new AccountActiviationMailTemplate(activationCode).generateTemplate()) + .withSubject(MailUtils.extractAppName(spCoreConfiguration) + " - Account Activation") + .appendTextHTML(new AccountActiviationMailTemplate(activationCode) + .generateTemplate(spCoreConfiguration.getEmailTemplateConfig())) .to(recipientAddress) .buildEmail(); @@ -57,8 +63,9 @@ public class MailSender extends AbstractMailer { public void sendPasswordRecoveryMail(String recipientAddress, String recoveryCode) throws IOException { Email email = baseEmail() - .withSubject(MailUtils.extractAppName() + " - Password Recovery") - .appendTextHTML(new PasswordRecoveryMailTemplate(recoveryCode).generateTemplate()) + .withSubject(MailUtils.extractAppName(spCoreConfiguration) + " - Password Recovery") + .appendTextHTML(new PasswordRecoveryMailTemplate(recoveryCode) + .generateTemplate(spCoreConfiguration.getEmailTemplateConfig())) .to(recipientAddress) .buildEmail(); @@ -68,8 +75,9 @@ public class MailSender extends AbstractMailer { public void sendInitialPasswordMail(String recipientAddress, String generatedProperty) throws IOException { Email email = baseEmail() - .withSubject(MailUtils.extractAppName() + " - New Account") - .appendTextHTML(new InitialPasswordMailTemplate(generatedProperty).generateTemplate()) + .withSubject(MailUtils.extractAppName(spCoreConfiguration) + " - New Account") + .appendTextHTML(new InitialPasswordMailTemplate(generatedProperty) + .generateTemplate(spCoreConfiguration.getEmailTemplateConfig())) .to(recipientAddress) .buildEmail(); diff --git a/streampipes-mail/src/main/java/org/apache/streampipes/mail/template/AbstractMailTemplate.java b/streampipes-mail/src/main/java/org/apache/streampipes/mail/template/AbstractMailTemplate.java index 1b7d140952..51a6553ffe 100644 --- a/streampipes-mail/src/main/java/org/apache/streampipes/mail/template/AbstractMailTemplate.java +++ b/streampipes-mail/src/main/java/org/apache/streampipes/mail/template/AbstractMailTemplate.java @@ -21,7 +21,7 @@ import org.apache.streampipes.mail.template.generation.DefaultPlaceholders; import org.apache.streampipes.mail.template.generation.MailTemplateBuilder; import org.apache.streampipes.mail.template.part.BaseUrlPart; import org.apache.streampipes.mail.template.part.LogoPart; -import org.apache.streampipes.storage.management.StorageDispatcher; +import org.apache.streampipes.model.configuration.EmailTemplateConfig; import com.google.common.base.Charsets; @@ -40,15 +40,11 @@ public abstract class AbstractMailTemplate { protected abstract void configureTemplate(MailTemplateBuilder builder); - public String generateTemplate() throws IOException { + public String generateTemplate(EmailTemplateConfig emailTemplateConfig) throws IOException { Map<String, String> placeholders = new HashMap<>(); addPlaceholders(placeholders); - var template = StorageDispatcher.INSTANCE - .getNoSqlStore() - .getSpCoreConfigurationStorage() - .get() - .getEmailTemplateConfig() + var template = emailTemplateConfig .getTemplate(); var builder = MailTemplateBuilder.create(template) 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 4a66194365..507765f9f5 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 @@ -27,14 +27,14 @@ import java.nio.charset.StandardCharsets; public class MailUtils { - public static String extractBaseUrl() { - GeneralConfig config = getSpCoreConfiguration().getGeneralConfig(); + public static String extractBaseUrl(SpCoreConfiguration spCoreConfiguration) { + GeneralConfig config = spCoreConfiguration.getGeneralConfig(); return config.getProtocol() + "://" + config.getHostname() + ":" + config.getPort(); } - public static String extractAppName() { - return getSpCoreConfiguration().getGeneralConfig().getAppName(); + public static String extractAppName(SpCoreConfiguration spCoreConfiguration) { + return spCoreConfiguration.getGeneralConfig().getAppName(); } public static String readResourceFileToString(String filename) throws IOException { diff --git a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/api/extensions/ExtensionServiceRequests.java b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/api/extensions/ExtensionServiceRequests.java index 65aafe6dfe..aa411d839d 100644 --- a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/api/extensions/ExtensionServiceRequests.java +++ b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/api/extensions/ExtensionServiceRequests.java @@ -52,8 +52,14 @@ public final class ExtensionServiceRequests { return post(target, payload, authToken); } - public static ExtensionServiceRequest migration(ExtensionServiceRequestTarget target, String payload) { - return post(target, payload, AuthTokenUtils.getAuthTokenForCurrentUser()); + public static ExtensionServiceRequest migration(ExtensionServiceRequestTarget target, + String payload, + SpResourceManager resourceManager) { + return post( + target, + payload, + AuthTokenUtils.getAuthTokenForCurrentUser(resourceManager.getCoreConfigurationStorage()) + ); } public static ExtensionServiceRequest descriptionUpdate(ExtensionServiceRequestTarget target, diff --git a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/assets/AssetConstants.java b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/assets/AssetConstants.java deleted file mode 100644 index 67ad1e1986..0000000000 --- a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/assets/AssetConstants.java +++ /dev/null @@ -1,31 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - * - */ -package org.apache.streampipes.manager.assets; - -import org.apache.streampipes.storage.management.StorageDispatcher; - -public class AssetConstants { - - public static final String ASSET_BASE_DIR = StorageDispatcher - .INSTANCE - .getNoSqlStore() - .getSpCoreConfigurationStorage() - .get() - .getAssetDir(); - -} diff --git a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/assets/AssetExtractor.java b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/assets/AssetExtractor.java index 9bbdc74e0e..740c76b46f 100644 --- a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/assets/AssetExtractor.java +++ b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/assets/AssetExtractor.java @@ -19,6 +19,7 @@ package org.apache.streampipes.manager.assets; import org.apache.streampipes.commons.constants.GlobalStreamPipesConstants; import org.apache.streampipes.commons.zip.ZipFileExtractor; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import java.io.File; import java.io.IOException; @@ -26,12 +27,16 @@ import java.io.InputStream; public class AssetExtractor { - private InputStream zipInputStream; - private String appId; + private final InputStream zipInputStream; + private final String appId; + private final ISpCoreConfigurationStorage coreConfigurationStorage; - public AssetExtractor(InputStream zipInputStream, String appId) { + public AssetExtractor(InputStream zipInputStream, + String appId, + ISpCoreConfigurationStorage coreConfigurationStorage) { this.zipInputStream = zipInputStream; this.appId = appId; + this.coreConfigurationStorage = coreConfigurationStorage; } public void extractAssetContents() throws IOException { @@ -44,7 +49,7 @@ public class AssetExtractor { } private String makeAssetLocation(String appId) { - return AssetConstants.ASSET_BASE_DIR + return coreConfigurationStorage.get().getAssetDir() + File.separator + appId; } diff --git a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/assets/AssetManager.java b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/assets/AssetManager.java index 2733844bec..61d7ea5ca8 100644 --- a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/assets/AssetManager.java +++ b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/assets/AssetManager.java @@ -20,6 +20,7 @@ package org.apache.streampipes.manager.assets; import org.apache.streampipes.commons.constants.GlobalStreamPipesConstants; import org.apache.streampipes.commons.exceptions.NoServiceEndpointsAvailableException; import org.apache.streampipes.manager.api.extensions.ExtensionServiceRequestManager; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.streampipes.svcdiscovery.api.model.SpServiceUrlProvider; import org.apache.commons.io.FileUtils; @@ -33,52 +34,58 @@ import java.nio.file.Paths; public class AssetManager { - public static byte[] getAssetIcon(String appId) throws IOException { + private final ISpCoreConfigurationStorage coreConfigurationStorage; + + public AssetManager(ISpCoreConfigurationStorage coreConfigurationStorage) { + this.coreConfigurationStorage = coreConfigurationStorage; + } + + public byte[] getAssetIcon(String appId) throws IOException { return Files.readAllBytes(Paths.get(getAssetIconPath(appId))); } - public static String getAssetDocumentation(String appId) throws IOException { + public String getAssetDocumentation(String appId) throws IOException { return new String(Files.readAllBytes(Paths.get(getAssetDocumentationPath(appId)))); } - public static byte[] getAsset(String appId, String assetName) throws IOException { + public byte[] getAsset(String appId, String assetName) throws IOException { return Files.readAllBytes(Paths.get(getAssetPath(appId, assetName))); } - public static boolean existsAssetDir(String appId) { + public boolean existsAssetDir(String appId) { var directory = new File(getAssetDir(appId)); return directory.exists() && directory.isDirectory(); } - public static void storeAsset(SpServiceUrlProvider spServiceUrlProvider, + public void storeAsset(SpServiceUrlProvider spServiceUrlProvider, String appId, ExtensionServiceRequestManager requestManager) throws IOException, NoServiceEndpointsAvailableException { InputStream assetStream = new AssetFetcher(spServiceUrlProvider, appId, requestManager) .fetchPipelineElementAssets(); - new AssetExtractor(assetStream, appId).extractAssetContents(); + new AssetExtractor(assetStream, appId, coreConfigurationStorage).extractAssetContents(); } - public static void deleteAsset(String appId) throws IOException { + public void deleteAsset(String appId) throws IOException { Path path = Paths.get(getAssetDir(appId)); if (Files.exists(path)) { FileUtils.deleteDirectory(path.toFile()); } } - private static String getAssetPath(String appId, String assetName) { + private String getAssetPath(String appId, String assetName) { return getAssetDir(appId) + File.separator + assetName; } - private static String getAssetIconPath(String appId) { + private String getAssetIconPath(String appId) { return getAssetDir(appId) + File.separator + GlobalStreamPipesConstants.STD_ICON_NAME; } - private static String getAssetDocumentationPath(String appId) { + private String getAssetDocumentationPath(String appId) { return getAssetDir(appId) + File.separator + GlobalStreamPipesConstants.STD_DOCUMENTATION_NAME; } - private static String getAssetDir(String appId) { - return AssetConstants.ASSET_BASE_DIR + File.separator + appId; + private String getAssetDir(String appId) { + return this.coreConfigurationStorage.get().getAssetDir() + File.separator + appId; } } diff --git a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/file/FileConstants.java b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/file/FileConstants.java deleted file mode 100644 index 372ddda1d1..0000000000 --- a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/file/FileConstants.java +++ /dev/null @@ -1,30 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - * - */ -package org.apache.streampipes.manager.file; - -import org.apache.streampipes.storage.management.StorageDispatcher; - -public class FileConstants { - public static final String FILES_BASE_DIR = StorageDispatcher - .INSTANCE - .getNoSqlStore() - .getSpCoreConfigurationStorage() - .get() - .getFilesDir(); - -} diff --git a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/file/FileHandler.java b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/file/FileHandler.java index 943da3d95f..34fc98351f 100644 --- a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/file/FileHandler.java +++ b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/file/FileHandler.java @@ -17,6 +17,8 @@ */ package org.apache.streampipes.manager.file; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; + import org.apache.commons.io.FileUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -31,6 +33,12 @@ public class FileHandler { Logger logger = LoggerFactory.getLogger(FileHandler.class); + private final ISpCoreConfigurationStorage coreConfigurationStorage; + + public FileHandler(ISpCoreConfigurationStorage coreConfigurationStorage) { + this.coreConfigurationStorage = coreConfigurationStorage; + } + public void storeFile(String filename, InputStream fileInputStream) throws IOException { File targetFile = makeFile(filename); FileUtils.copyInputStreamToFile(fileInputStream, targetFile); @@ -68,7 +76,7 @@ public class FileHandler { } private String makeFileLocation() { - return FileConstants.FILES_BASE_DIR + return coreConfigurationStorage.get().getFilesDir() + File.separator; } } diff --git a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/file/FileManager.java b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/file/FileManager.java index f8d22a78b4..42f41ebd04 100644 --- a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/file/FileManager.java +++ b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/file/FileManager.java @@ -21,6 +21,7 @@ import org.apache.streampipes.commons.file.FileHasher; import org.apache.streampipes.model.file.FileMetadata; import org.apache.streampipes.sdk.helpers.Filetypes; import org.apache.streampipes.storage.api.system.IFileMetadataStorage; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.streampipes.storage.management.StorageDispatcher; import org.apache.commons.io.input.BOMInputStream; @@ -46,12 +47,12 @@ public class FileManager { this.fileHasher = fileHasher; } - public FileManager() { + public FileManager(ISpCoreConfigurationStorage coreConfigurationStorage) { this.fileMetadataStorage = StorageDispatcher .INSTANCE .getNoSqlStore() .getFileMetadataStorage(); - this.fileHandler = new FileHandler(); + this.fileHandler = new FileHandler(coreConfigurationStorage); this.fileHasher = new FileHasher(); } diff --git a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/migration/AbstractMigrationManager.java b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/migration/AbstractMigrationManager.java index 6edbf49142..7af2a5598b 100644 --- a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/migration/AbstractMigrationManager.java +++ b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/migration/AbstractMigrationManager.java @@ -94,7 +94,7 @@ public abstract class AbstractMigrationManager { String serializedRequest = JacksonSerializer.getObjectMapper().writeValueAsString(migrationRequest); var migrationResponse = requestManager.request( - ExtensionServiceRequests.migration(requestTarget, serializedRequest) + ExtensionServiceRequests.migration(requestTarget, serializedRequest, resourceManager) ); TypeReference<MigrationResult<T>> typeReference = new TypeReference<>() { diff --git a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/AutoInstallation.java b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/AutoInstallation.java index df7de44532..5a9096f263 100644 --- a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/AutoInstallation.java +++ b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/AutoInstallation.java @@ -53,7 +53,7 @@ public class AutoInstallation implements BackgroundTaskNotifier { public void startAutoInstallation() { InitialSettings settings = collectInitialSettings(); - List<InstallationStep> steps = InstallationConfiguration.getInstallationSteps(settings); + List<InstallationStep> steps = InstallationConfiguration.getInstallationSteps(settings, resourceManager); List<Runnable> backgroundSteps = InstallationConfiguration.getBackgroundInstallationSteps( settings, this, diff --git a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/InstallationConfiguration.java b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/InstallationConfiguration.java index f6ce2eb6c1..70e5fba099 100644 --- a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/InstallationConfiguration.java +++ b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/InstallationConfiguration.java @@ -28,10 +28,11 @@ import java.util.List; public class InstallationConfiguration { - public static List<InstallationStep> getInstallationSteps(InitialSettings settings) { + public static List<InstallationStep> getInstallationSteps(InitialSettings settings, + SpResourceManager resourceManager) { List<InstallationStep> steps = new ArrayList<>(); - steps.add(new SpCoreConfigurationStep()); + steps.add(new SpCoreConfigurationStep(resourceManager.getCoreConfigurationStorage())); steps.add(new CouchDbInstallationStep()); steps.add(new UserRegistrationInstallationStep( settings.getAdminEmail(), diff --git a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/SpCoreConfigurationStep.java b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/SpCoreConfigurationStep.java index b5e0a8bd3b..5605d29578 100644 --- a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/SpCoreConfigurationStep.java +++ b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/SpCoreConfigurationStep.java @@ -20,7 +20,7 @@ package org.apache.streampipes.manager.setup; import org.apache.streampipes.model.configuration.DefaultSpCoreConfiguration; import org.apache.streampipes.model.configuration.SpCoreConfigurationStatus; -import org.apache.streampipes.storage.management.StorageDispatcher; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -29,13 +29,19 @@ public class SpCoreConfigurationStep extends InstallationStep { private static final Logger LOG = LoggerFactory.getLogger(SpCoreConfigurationStep.class); + private final ISpCoreConfigurationStorage coreConfigurationStorage; + + public SpCoreConfigurationStep(ISpCoreConfigurationStorage coreConfigurationStorage) { + this.coreConfigurationStorage = coreConfigurationStorage; + } + @Override public void install() { var coreCfg = new DefaultSpCoreConfiguration().make(); coreCfg.setServiceStatus(SpCoreConfigurationStatus.INSTALLING); - StorageDispatcher.INSTANCE.getNoSqlStore().getSpCoreConfigurationStorage().createElement(coreCfg); + coreConfigurationStorage.createElement(coreCfg); LOG.info("Core is now in {} state", coreCfg.getServiceStatus()); - new StreamPipesEnvChecker().updateEnvironmentVariables(); + new StreamPipesEnvChecker(coreConfigurationStorage).updateEnvironmentVariables(); } @Override diff --git a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/StreamPipesEnvChecker.java b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/StreamPipesEnvChecker.java index 0062419eb4..4133d9e0d0 100644 --- a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/StreamPipesEnvChecker.java +++ b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/StreamPipesEnvChecker.java @@ -24,7 +24,7 @@ import org.apache.streampipes.model.configuration.DefaultSpCoreConfiguration; import org.apache.streampipes.model.configuration.JwtSigningMode; import org.apache.streampipes.model.configuration.LocalAuthConfig; import org.apache.streampipes.model.configuration.SpCoreConfiguration; -import org.apache.streampipes.storage.management.StorageDispatcher; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -40,26 +40,23 @@ public class StreamPipesEnvChecker { private SpCoreConfiguration coreConfig; private final Environment env; + private final ISpCoreConfigurationStorage coreConfigStorage; - public StreamPipesEnvChecker() { + public StreamPipesEnvChecker(ISpCoreConfigurationStorage coreConfigurationStorage) { this.env = Environments.getEnvironment(); + this.coreConfigStorage = coreConfigurationStorage; } public void updateEnvironmentVariables() { - var configStorage = StorageDispatcher - .INSTANCE - .getNoSqlStore() - .getSpCoreConfigurationStorage(); - - if (configStorage.exists()) { - this.coreConfig = configStorage.get(); + if (coreConfigStorage.exists()) { + this.coreConfig = coreConfigStorage.get(); LOG.info("Checking and updating environment variables..."); var shouldUpdateJwtConfig = updateJwtSettings(); var shouldUpdateDirectoryConfig = updateDirectorySettings(); if (shouldUpdateJwtConfig || shouldUpdateDirectoryConfig) { - configStorage.updateElement(coreConfig); + coreConfigStorage.updateElement(coreConfig); } } } diff --git a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/util/AuthTokenUtils.java b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/util/AuthTokenUtils.java index 8ba2b4387a..e34bd26f4e 100644 --- a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/util/AuthTokenUtils.java +++ b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/util/AuthTokenUtils.java @@ -23,6 +23,7 @@ import org.apache.streampipes.model.client.user.Principal; import org.apache.streampipes.resource.management.PermissionResourceManager; import org.apache.streampipes.resource.management.SpResourceManager; import org.apache.streampipes.resource.management.UserResourceManager; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.streampipes.user.management.jwt.JwtTokenProvider; import org.springframework.security.core.Authentication; @@ -30,15 +31,15 @@ import org.springframework.security.core.context.SecurityContextHolder; public class AuthTokenUtils { - public static String getAuthTokenForCurrentUser() { + public static String getAuthTokenForCurrentUser(ISpCoreConfigurationStorage coreConfigurationStorage) { Authentication auth = SecurityContextHolder.getContext().getAuthentication(); - return makeBearerToken(new JwtTokenProvider().createToken(auth)); + return makeBearerToken(new JwtTokenProvider(coreConfigurationStorage).createToken(auth)); } public static String getAuthToken(String resourceId, SpResourceManager resourceManager) { if (SecurityContextHolder.getContext().getAuthentication() != null) { - return getAuthTokenForCurrentUser(); + return getAuthTokenForCurrentUser(resourceManager.getCoreConfigurationStorage()); } else { if (resourceId != null) { String ownerSid = getOwnerSid(resourceId, resourceManager.managePermissions()); diff --git a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/verification/ElementVerifier.java b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/verification/ElementVerifier.java index d23540ddf5..b18fdb3c84 100644 --- a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/verification/ElementVerifier.java +++ b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/verification/ElementVerifier.java @@ -29,6 +29,7 @@ import org.apache.streampipes.model.message.Notification; import org.apache.streampipes.model.message.NotificationType; import org.apache.streampipes.model.message.SuccessMessage; import org.apache.streampipes.resource.management.PermissionResourceManager; +import org.apache.streampipes.resource.management.SpResourceManager; import org.apache.streampipes.serializers.json.JacksonSerializer; import org.apache.streampipes.storage.api.pipeline.IPipelineElementDescriptionStorage; @@ -51,30 +52,34 @@ public abstract class ElementVerifier<T extends NamedStreamPipesEntity> { protected T elementDescription; protected final IPipelineElementDescriptionStorage storageApi; + protected final AssetManager assetManager; private final PermissionResourceManager permissionResourceManager; public ElementVerifier( String graphData, Class<T> elementClass, IPipelineElementDescriptionStorage storageApi, - PermissionResourceManager permissionResourceManager + SpResourceManager resourceManager ) { this.elementClass = elementClass; this.graphData = graphData; this.storageApi = storageApi; this.shouldTransform = true; - this.permissionResourceManager = permissionResourceManager; + this.permissionResourceManager = resourceManager.managePermissions(); + this.assetManager = new AssetManager(resourceManager.getCoreConfigurationStorage()); } public ElementVerifier(T elementDescription, IPipelineElementDescriptionStorage storageApi, - PermissionResourceManager permissionResourceManager) { + SpResourceManager resourceManager) { this.elementDescription = elementDescription; this.storageApi = storageApi; this.graphData = null; this.elementClass = null; this.shouldTransform = false; - this.permissionResourceManager = permissionResourceManager; + this.permissionResourceManager = resourceManager.managePermissions(); + this.assetManager = new AssetManager(resourceManager.getCoreConfigurationStorage()); + } protected abstract StorageState store(); @@ -122,7 +127,7 @@ public abstract class ElementVerifier<T extends NamedStreamPipesEntity> { protected void updateAssets() throws IOException, NoServiceEndpointsAvailableException { if (elementDescription.isIncludesAssets()) { - AssetManager.deleteAsset(elementDescription.getAppId()); + assetManager.deleteAsset(elementDescription.getAppId()); storeAssets(); } } diff --git a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/verification/TypedElementVerifier.java b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/verification/TypedElementVerifier.java index 0a553bba42..f00bf3364c 100644 --- a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/verification/TypedElementVerifier.java +++ b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/verification/TypedElementVerifier.java @@ -20,9 +20,8 @@ package org.apache.streampipes.manager.verification; import org.apache.streampipes.commons.exceptions.NoServiceEndpointsAvailableException; import org.apache.streampipes.manager.api.extensions.ExtensionServiceRequestManager; -import org.apache.streampipes.manager.assets.AssetManager; import org.apache.streampipes.model.base.NamedStreamPipesEntity; -import org.apache.streampipes.resource.management.PermissionResourceManager; +import org.apache.streampipes.resource.management.SpResourceManager; import org.apache.streampipes.storage.api.pipeline.IPipelineElementDescriptionStorage; import org.apache.streampipes.svcdiscovery.api.model.SpServiceUrlProvider; @@ -47,9 +46,9 @@ public class TypedElementVerifier<T extends NamedStreamPipesEntity> extends Elem Consumer<T> updateOperation, SpServiceUrlProvider serviceUrlProvider, ExtensionServiceRequestManager requestManager, - PermissionResourceManager permissionResourceManager + SpResourceManager resourceManager ) { - super(graphData, elementClass, storageApi, permissionResourceManager); + super(graphData, elementClass, storageApi, resourceManager); this.existsChecker = existsChecker; this.storeOperation = storeOperation; this.updateOperation = updateOperation; @@ -65,9 +64,9 @@ public class TypedElementVerifier<T extends NamedStreamPipesEntity> extends Elem Consumer<T> updateOperation, SpServiceUrlProvider serviceUrlProvider, ExtensionServiceRequestManager requestManager, - PermissionResourceManager permissionResourceManager + SpResourceManager resourceManager ) { - super(elementDescription, storageApi, permissionResourceManager); + super(elementDescription, storageApi, resourceManager); this.existsChecker = existsChecker; this.storeOperation = storeOperation; this.updateOperation = updateOperation; @@ -92,7 +91,7 @@ public class TypedElementVerifier<T extends NamedStreamPipesEntity> extends Elem @Override protected void storeAssets() throws IOException, NoServiceEndpointsAvailableException { if (elementDescription.isIncludesAssets()) { - AssetManager.storeAsset(serviceUrlProvider, elementDescription.getAppId(), requestManager); + assetManager.storeAsset(serviceUrlProvider, elementDescription.getAppId(), requestManager); } } } 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 9811d5b1cc..8d2653b1cb 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.ISpCoreConfigurationStorage; import org.apache.streampipes.storage.api.user.IPermissionStorage; import org.apache.streampipes.storage.management.StorageDispatcher; @@ -35,6 +36,7 @@ public class SpResourceManager { private final IDashboardStorage dashboardStorage; private final IPipelineStorage pipelineStorage; private final IDataLakeMeasureStorage datasetStorage; + private final ISpCoreConfigurationStorage coreConfigurationStorage; public SpResourceManager(IPermissionStorage permissionStorage, IChartStorage chartStorage, @@ -42,7 +44,8 @@ public class SpResourceManager { IAssetStorage assetStorage, IDashboardStorage dashboardStorage, IPipelineStorage pipelineStorage, - IDataLakeMeasureStorage datasetStorage) { + IDataLakeMeasureStorage datasetStorage, + ISpCoreConfigurationStorage coreConfigurationStorage) { this.permissionStorage = permissionStorage; this.chartStorage = chartStorage; this.adapterStorage = adapterStorage; @@ -50,6 +53,7 @@ public class SpResourceManager { this.dashboardStorage = dashboardStorage; this.pipelineStorage = pipelineStorage; this.datasetStorage = datasetStorage; + this.coreConfigurationStorage = coreConfigurationStorage; } public AdapterDescriptionResourceManager manageAdapterDescriptions() { @@ -98,6 +102,10 @@ public class SpResourceManager { ); } + public ISpCoreConfigurationStorage getCoreConfigurationStorage() { + return coreConfigurationStorage; + } + public UserResourceManager manageUsers() { return new UserResourceManager(); } 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 f2ce2cabf8..45a6c38e50 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 @@ -148,7 +148,7 @@ public class ResetManagement { } private void deleteAllFiles() { - var fileManager = new FileManager(); + var fileManager = new FileManager(resourceManager.getCoreConfigurationStorage()); List<FileMetadata> allFiles = fileManager.getAllFiles(); allFiles.forEach(fileMetadata -> fileManager.deleteFile(fileMetadata.getFileId())); } 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 97826293b7..43b35ded93 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 @@ -33,6 +33,7 @@ import org.apache.streampipes.model.message.SuccessMessage; import org.apache.streampipes.resource.management.SpResourceManager; import org.apache.streampipes.rest.core.base.impl.AbstractRestResource; import org.apache.streampipes.rest.shared.exception.SpMessageException; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.streampipes.storage.management.StorageDispatcher; import org.apache.streampipes.user.management.jwt.JwtTokenProvider; import org.apache.streampipes.user.management.model.PrincipalUserDetails; @@ -74,11 +75,14 @@ public class Authentication extends AbstractRestResource { AuthenticationManager authenticationManager; private final SpResourceManager resourceManager; + private final ISpCoreConfigurationStorage coreConfigurationStorage; public Authentication(AuthenticationManager authenticationManager, - SpResourceManager resourceManager) { + SpResourceManager resourceManager, + ISpCoreConfigurationStorage coreConfigurationStorage) { this.authenticationManager = authenticationManager; this.resourceManager = resourceManager; + this.coreConfigurationStorage = coreConfigurationStorage; } @PostMapping( @@ -129,7 +133,7 @@ public class Authentication extends AbstractRestResource { setRefreshCookie(request, response, issuedRefreshToken); - String jwt = new JwtTokenProvider().createToken(userAccount); + String jwt = new JwtTokenProvider(coreConfigurationStorage).createToken(userAccount); return ok(new JwtAuthenticationResponse(jwt)); } @@ -241,7 +245,7 @@ public class Authentication extends AbstractRestResource { } private JwtAuthenticationResponse makeJwtResponse(org.springframework.security.core.Authentication auth) { - String jwt = new JwtTokenProvider().createToken(auth); + String jwt = new JwtTokenProvider(coreConfigurationStorage).createToken(auth); return new JwtAuthenticationResponse(jwt); } diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/FileResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/FileResource.java index 209d30f061..4bdd7baa9b 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/FileResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/FileResource.java @@ -24,6 +24,7 @@ import org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResourc import org.apache.streampipes.rest.security.AuthConstants; import org.apache.streampipes.rest.shared.exception.SpMessageException; import org.apache.streampipes.sdk.helpers.Filetypes; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import io.swagger.v3.oas.annotations.Operation; import io.swagger.v3.oas.annotations.Parameter; @@ -56,8 +57,8 @@ public class FileResource extends AbstractAuthGuardedRestResource { private final FileManager fileManager; - public FileResource() { - this.fileManager = new FileManager(); + public FileResource(ISpCoreConfigurationStorage coreConfigurationStorage) { + this.fileManager = new FileManager(coreConfigurationStorage); } @PostMapping( diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/PipelineElementAsset.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/PipelineElementAsset.java index f80ccf344b..e51ccced37 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/PipelineElementAsset.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/PipelineElementAsset.java @@ -41,15 +41,17 @@ public class PipelineElementAsset extends AbstractRestResource { private static final Logger LOG = LoggerFactory.getLogger(PipelineElementAsset.class); private final SpResourceManager resourceManager; + private final AssetManager assetManager; public PipelineElementAsset(SpResourceManager resourceManager) { this.resourceManager = resourceManager; + this.assetManager = new AssetManager(resourceManager.getCoreConfigurationStorage()); } @GetMapping(path = "/{appId}/assets/icon") public ResponseEntity<?> getIconAsset(@PathVariable("appId") String appId) { try { - byte[] icon = AssetManager.getAssetIcon(appId); + byte[] icon = assetManager.getAssetIcon(appId); return ResponseEntity.ok() .contentType(MediaType.parseMediaType(ImageMimeTypeDetector.detect(icon))) .body(icon); @@ -75,7 +77,7 @@ public class PipelineElementAsset extends AbstractRestResource { .getElementById(dataStream.getCorrespondingAdapterId()); appId = adapterDescription.getAppId(); } - return ok(AssetManager.getAssetDocumentation(appId)); + return ok(assetManager.getAssetDocumentation(appId)); } catch (IOException e) { return fail(); } @@ -85,7 +87,7 @@ public class PipelineElementAsset extends AbstractRestResource { public ResponseEntity<?> getAsset(@PathVariable("appId") String appId, @PathVariable("assetName") String assetName) { try { - byte[] asset = AssetManager.getAsset(appId, assetName); + byte[] asset = assetManager.getAsset(appId, assetName); return ok(asset); } catch (IOException e) { LOG.error("Could not find asset {}", assetName); diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/Setup.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/Setup.java index 778c09d897..a0236cf4ad 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/Setup.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/Setup.java @@ -22,7 +22,6 @@ package org.apache.streampipes.rest.impl; import org.apache.streampipes.manager.health.CoreServiceStatusManager; import org.apache.streampipes.rest.core.base.impl.AbstractRestResource; import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; -import org.apache.streampipes.storage.management.StorageDispatcher; import com.google.gson.JsonObject; import io.swagger.v3.oas.annotations.Operation; @@ -36,8 +35,11 @@ import org.springframework.web.bind.annotation.RestController; @RequestMapping("/api/v2/setup") public class Setup extends AbstractRestResource { - private final ISpCoreConfigurationStorage storage = StorageDispatcher - .INSTANCE.getNoSqlStore().getSpCoreConfigurationStorage(); + private final ISpCoreConfigurationStorage storage; + + public Setup(ISpCoreConfigurationStorage storage) { + this.storage = storage; + } @GetMapping(path = "/configured", produces = MediaType.APPLICATION_JSON_VALUE) @Operation(summary = "Endpoint is used by UI to determine whether the core is running upon startup", diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/ExtensionsServiceEndpointResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/ExtensionsServiceEndpointResource.java index f342f4b363..d6d629d4aa 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/ExtensionsServiceEndpointResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/ExtensionsServiceEndpointResource.java @@ -30,6 +30,7 @@ import org.apache.streampipes.model.extensions.ExtensionItemDescription; import org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource; import org.apache.streampipes.rest.security.AuthConstants; import org.apache.streampipes.rest.shared.exception.SpMessageException; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.streampipes.svcdiscovery.api.model.SpServiceUrlProvider; import org.springframework.http.HttpStatus; @@ -52,9 +53,12 @@ import java.util.Set; public class ExtensionsServiceEndpointResource extends AbstractAuthGuardedRestResource { private final ExtensionServiceRequestManager extensionServiceRequestManager; + private final AssetManager assetManager; - public ExtensionsServiceEndpointResource(ExtensionServiceRequestManager extensionServiceRequestManager) { + public ExtensionsServiceEndpointResource(ExtensionServiceRequestManager extensionServiceRequestManager, + ISpCoreConfigurationStorage coreConfigurationStorage) { this.extensionServiceRequestManager = extensionServiceRequestManager; + this.assetManager = new AssetManager(coreConfigurationStorage); } @GetMapping(produces = MediaType.APPLICATION_JSON_VALUE) @@ -77,7 +81,7 @@ public class ExtensionsServiceEndpointResource extends AbstractAuthGuardedRestRe private byte[] getIconImage(ExtensionItemDescription extensionItemDescription) throws IOException { if (extensionItemDescription.isInstalled()) { - return AssetManager.getAssetIcon(extensionItemDescription.getAppId()); + return assetManager.getAssetIcon(extensionItemDescription.getAppId()); } if (extensionItemDescription.getServiceTagPrefix() == null) { diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/MigrationResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/MigrationResource.java index 3fbdfe3d14..f08b9e6d19 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/MigrationResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/MigrationResource.java @@ -73,9 +73,7 @@ public class MigrationResource extends AbstractAuthGuardedRestResource { private final IDataSinkStorage dataSinkStorage = getNoSqlStorage().getDataSinkStorage(); private final IPipelineStorage pipelineStorage; - private final CoreServiceStatusManager coreServiceStatusManager = new CoreServiceStatusManager( - getNoSqlStorage().getSpCoreConfigurationStorage() - ); + private final CoreServiceStatusManager coreServiceStatusManager; private final ExtensionServiceRequestManager extensionServiceRequestManager; private final WorkerRestClient workerRestClient; private final SpResourceManager resourceManager; @@ -88,6 +86,9 @@ public class MigrationResource extends AbstractAuthGuardedRestResource { this.resourceManager = resourceManager; this.adapterStorage = resourceManager.manageAdapters().getDb(); this.pipelineStorage = resourceManager.managePipelines().getDb(); + this.coreServiceStatusManager = new CoreServiceStatusManager( + resourceManager.getCoreConfigurationStorage() + ); } @PostMapping(path = "{serviceId}", consumes = 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 13f439cc42..41bc5ee660 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 @@ -107,7 +107,10 @@ public class DataLakeResource extends AbstractDataLakeResource { this.dataExplorerQueryManagement = new DataExplorerDispatcher() .getDataExplorerManager() .getQueryManagement(this.dataLakeMeasureManagement); - this.dataLakeExportManager = new DataLakeExportManager(this.dataLakeMeasureManagement, dataExplorerQueryManagement); + this.dataLakeExportManager = new DataLakeExportManager( + this.dataLakeMeasureManagement, + dataExplorerQueryManagement, + resourceManager.getCoreConfigurationStorage()); } @DeleteMapping(path = "/measurements/{measurementName}") 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 2e463403bf..09d2109ecb 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 @@ -62,7 +62,6 @@ import org.apache.streampipes.service.core.storage.StorageApiConfiguration; import org.apache.streampipes.storage.api.function.IFunctionStateStorage; import org.apache.streampipes.storage.api.pipeline.IPipelineStorage; import org.apache.streampipes.storage.api.system.IExtensionsServiceStorage; -import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.streampipes.storage.couchdb.impl.user.UserStorage; import org.apache.streampipes.storage.couchdb.utils.CouchDbViewGenerator; import org.apache.streampipes.storage.management.StorageDispatcher; @@ -99,12 +98,6 @@ public class StreamPipesCoreApplication extends StreamPipesServiceBase { private static final Logger LOG = LoggerFactory.getLogger(StreamPipesCoreApplication.class.getCanonicalName()); - private final ISpCoreConfigurationStorage coreConfigStorage = - StorageDispatcher.INSTANCE.getNoSqlStore().getSpCoreConfigurationStorage(); - - private final CoreServiceStatusManager coreStatusManager = - new CoreServiceStatusManager(coreConfigStorage); - @Autowired private IFunctionStateStorage functionStateStorage; @@ -160,7 +153,7 @@ public class StreamPipesCoreApplication extends StreamPipesServiceBase { var executorService = Executors.newSingleThreadScheduledExecutor(); var logCheckExecutorService = Executors.newSingleThreadScheduledExecutor(); - new StreamPipesEnvChecker().updateEnvironmentVariables(); + new StreamPipesEnvChecker(resourceManager.getCoreConfigurationStorage()).updateEnvironmentVariables(); new CouchDbViewGenerator().createGenericDatabaseIfNotExists(); var env = Environments.getEnvironment(); @@ -179,6 +172,9 @@ public class StreamPipesCoreApplication extends StreamPipesServiceBase { if (env.getLoadManagerEnable().getValueOrDefault()) { LoadManager.initialize(resourceManager); } + + var coreConfigStorage = resourceManager.getCoreConfigurationStorage(); + var coreStatusManager = new CoreServiceStatusManager(coreConfigStorage); if (!isConfigured()) { CoreInitialInstallationProgress.INSTANCE.triggerInitiallyInstallingMode(); doInitialSetup(env.getInitialWaitTimeBeforeInstallationInMillis().getValueOrDefault()); diff --git a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/WebSecurityConfig.java b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/WebSecurityConfig.java index eb25ed6f3c..288cea95e5 100644 --- a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/WebSecurityConfig.java +++ b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/WebSecurityConfig.java @@ -30,6 +30,7 @@ import org.apache.streampipes.service.core.oauth2.OAuth2AccessTokenResponseConve import org.apache.streampipes.service.core.oauth2.OAuth2AuthenticationFailureHandler; import org.apache.streampipes.service.core.oauth2.OAuth2AuthenticationSuccessHandler; import org.apache.streampipes.service.core.oauth2.OAuthEnabledCondition; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.streampipes.storage.api.user.IPermissionStorage; import org.apache.streampipes.user.management.service.SpUserDetailsService; @@ -121,7 +122,8 @@ public class WebSecurityConfig { } @Bean - public SecurityFilterChain filterChain(HttpSecurity http) { + public SecurityFilterChain filterChain(HttpSecurity http, + ISpCoreConfigurationStorage coreConfigurationStorage) { http .cors(Customizer.withDefaults()) .sessionManagement(sm -> sm.sessionCreationPolicy(SessionCreationPolicy.STATELESS)) @@ -152,13 +154,13 @@ public class WebSecurityConfig { ); } - http.addFilterBefore(tokenAuthenticationFilter(), UsernamePasswordAuthenticationFilter.class); + http.addFilterBefore(tokenAuthenticationFilter(coreConfigurationStorage), UsernamePasswordAuthenticationFilter.class); return http.build(); } - public TokenAuthenticationFilter tokenAuthenticationFilter() { - return new TokenAuthenticationFilter(permissionStorage); + public TokenAuthenticationFilter tokenAuthenticationFilter(ISpCoreConfigurationStorage coreConfigurationStorage) { + return new TokenAuthenticationFilter(permissionStorage, coreConfigurationStorage); } @Bean(BeanIds.USER_DETAILS_SERVICE) 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 b1cc374707..28f916b637 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.ISpCoreConfigurationStorage; import org.apache.streampipes.storage.api.user.IPermissionStorage; import org.slf4j.Logger; @@ -96,7 +97,8 @@ public class ExtensionServiceRequestConfiguration { IDashboardStorage dashboardStorage, IAssetStorage assetStorage, IPipelineStorage pipelineStorage, - IDataLakeMeasureStorage datasetStorage) { + IDataLakeMeasureStorage datasetStorage, + ISpCoreConfigurationStorage coreConfigurationStorage) { return new SpResourceManager( permissionStorage, chartStorage, @@ -104,7 +106,8 @@ public class ExtensionServiceRequestConfiguration { assetStorage, dashboardStorage, pipelineStorage, - datasetStorage + datasetStorage, + coreConfigurationStorage ); } diff --git a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/filter/TokenAuthenticationFilter.java b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/filter/TokenAuthenticationFilter.java index e91c87e940..a79f6824ef 100644 --- a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/filter/TokenAuthenticationFilter.java +++ b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/filter/TokenAuthenticationFilter.java @@ -23,6 +23,7 @@ import org.apache.streampipes.model.client.user.DefaultRole; import org.apache.streampipes.model.client.user.Principal; import org.apache.streampipes.model.client.user.ServiceAccount; import org.apache.streampipes.model.client.user.UserAccount; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.streampipes.storage.api.user.IPermissionStorage; import org.apache.streampipes.storage.api.user.IUserStorage; import org.apache.streampipes.storage.management.StorageDispatcher; @@ -72,8 +73,9 @@ public class TokenAuthenticationFilter extends OncePerRequestFilter { private static final Logger logger = LoggerFactory.getLogger(TokenAuthenticationFilter.class); - public TokenAuthenticationFilter(IPermissionStorage permissionStorage) { - this.tokenProvider = new JwtTokenProvider(); + public TokenAuthenticationFilter(IPermissionStorage permissionStorage, + ISpCoreConfigurationStorage coreConfigurationStorage) { + this.tokenProvider = new JwtTokenProvider(coreConfigurationStorage); this.userStorage = StorageDispatcher.INSTANCE.getNoSqlStore().getUserStorageAPI(); this.permissionStorage = permissionStorage; } 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 9d2180a841..ea58e06ac3 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 @@ -46,6 +46,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.ISpCoreConfigurationStorage; import org.apache.streampipes.storage.api.user.IPermissionStorage; import java.util.Arrays; @@ -60,6 +61,7 @@ public class AvailableMigrations { private final IAssetStorage assetStorage; private final IPipelineStorage pipelineStorage; private final IDataLakeMeasureStorage datasetStorage; + private final ISpCoreConfigurationStorage coreConfigStorage; public AvailableMigrations(SpResourceManager resourceManager) { this.chartStorage = resourceManager.manageCharts().getDb(); @@ -69,6 +71,7 @@ public class AvailableMigrations { this.assetStorage = resourceManager.manageAssets().getDb(); this.pipelineStorage = resourceManager.managePipelines().getDb(); this.datasetStorage = resourceManager.manageDataLakeMeasures().getDb(); + this.coreConfigStorage = resourceManager.getCoreConfigurationStorage(); } public List<Migration> getAvailableMigrations() { @@ -76,7 +79,7 @@ public class AvailableMigrations { new ModifyAssetLinksMigration(), new ModifyAssetLinkTypesMigration(), new AddDataLakeMeasureViewMigration(), - new AddDefaultExportProviderMigration(), + new AddDefaultExportProviderMigration(coreConfigStorage), new FixImportedPermissionsMigration(chartStorage, dashboardStorage, permissionStorage), new AddAssetManagementViewMigration(), new MoveAssetContentMigration(), diff --git a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/v0980/AddDefaultExportProviderMigration.java b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/v0980/AddDefaultExportProviderMigration.java index 3fd7e57615..1037c9417e 100644 --- a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/v0980/AddDefaultExportProviderMigration.java +++ b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/v0980/AddDefaultExportProviderMigration.java @@ -21,14 +21,16 @@ package org.apache.streampipes.service.core.migrations.v0980; import org.apache.streampipes.model.configuration.DefaultExportProviderConfig; import org.apache.streampipes.service.core.migrations.Migration; import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; -import org.apache.streampipes.storage.management.StorageDispatcher; import java.io.IOException; public class AddDefaultExportProviderMigration implements Migration { - private final ISpCoreConfigurationStorage storage = StorageDispatcher.INSTANCE.getNoSqlStore() - .getSpCoreConfigurationStorage(); + private final ISpCoreConfigurationStorage storage; + + public AddDefaultExportProviderMigration(ISpCoreConfigurationStorage storage) { + this.storage = storage; + } @Override public boolean shouldExecute() { diff --git a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/OAuth2AuthenticationSuccessHandler.java b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/OAuth2AuthenticationSuccessHandler.java index 36cfd187ab..61e8d7b5cb 100755 --- a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/OAuth2AuthenticationSuccessHandler.java +++ b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/OAuth2AuthenticationSuccessHandler.java @@ -1,28 +1,29 @@ -/* - * 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.service.core.oauth2; - +/* + * 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.service.core.oauth2; + import org.apache.streampipes.commons.environment.Environment; import org.apache.streampipes.commons.environment.Environments; import org.apache.streampipes.model.client.user.Principal; import org.apache.streampipes.rest.shared.exception.BadRequestException; import org.apache.streampipes.service.core.oauth2.util.CookieUtils; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.streampipes.user.management.jwt.JwtTokenProvider; import org.apache.streampipes.user.management.model.PrincipalUserDetails; import org.apache.streampipes.user.management.service.RefreshTokenService; @@ -33,8 +34,8 @@ import org.springframework.http.ResponseCookie; import org.springframework.security.core.Authentication; import org.springframework.security.web.authentication.SimpleUrlAuthenticationSuccessHandler; import org.springframework.stereotype.Component; - -import jakarta.servlet.http.Cookie; + +import jakarta.servlet.http.Cookie; import jakarta.servlet.http.HttpServletRequest; import jakarta.servlet.http.HttpServletResponse; @@ -53,43 +54,44 @@ public class OAuth2AuthenticationSuccessHandler extends SimpleUrlAuthenticationS private static final long MIN_REFRESH_COOKIE_SECONDS = 1; private final JwtTokenProvider tokenProvider; - private final HttpCookieOAuth2AuthorizationRequestRepository httpCookieOAuth2AuthorizationRequestRepository; - private final Environment env; - - @Autowired - OAuth2AuthenticationSuccessHandler(HttpCookieOAuth2AuthorizationRequestRepository - httpCookieOAuth2AuthorizationRequestRepository) { - this.tokenProvider = new JwtTokenProvider(); - this.httpCookieOAuth2AuthorizationRequestRepository = httpCookieOAuth2AuthorizationRequestRepository; - this.env = Environments.getEnvironment(); - } - - @Override - public void onAuthenticationSuccess(HttpServletRequest request, - HttpServletResponse response, - Authentication authentication) throws IOException { - String targetUrl = determineTargetUrl(request, response, authentication); - - if (response.isCommitted()) { - return; - } - - clearAuthenticationAttributes(request, response); - getRedirectStrategy().sendRedirect(request, response, targetUrl); - } - - @Override + private final HttpCookieOAuth2AuthorizationRequestRepository httpCookieOAuth2AuthorizationRequestRepository; + private final Environment env; + + @Autowired + OAuth2AuthenticationSuccessHandler(HttpCookieOAuth2AuthorizationRequestRepository + httpCookieOAuth2AuthorizationRequestRepository, + ISpCoreConfigurationStorage coreConfigurationStorage) { + this.tokenProvider = new JwtTokenProvider(coreConfigurationStorage); + this.httpCookieOAuth2AuthorizationRequestRepository = httpCookieOAuth2AuthorizationRequestRepository; + this.env = Environments.getEnvironment(); + } + + @Override + public void onAuthenticationSuccess(HttpServletRequest request, + HttpServletResponse response, + Authentication authentication) throws IOException { + String targetUrl = determineTargetUrl(request, response, authentication); + + if (response.isCommitted()) { + return; + } + + clearAuthenticationAttributes(request, response); + getRedirectStrategy().sendRedirect(request, response, targetUrl); + } + + @Override protected String determineTargetUrl(HttpServletRequest request, HttpServletResponse response, Authentication authentication) { Optional<String> redirectUri = CookieUtils .getCookie(request, HttpCookieOAuth2AuthorizationRequestRepository.REDIRECT_URI_PARAM_COOKIE_NAME) .map(Cookie::getValue); - - if (redirectUri.isPresent() && !isAuthorizedRedirectUri(redirectUri.get())) { - throw new BadRequestException( - "Unauthorized redirect uri found - check the redirect uri in your OAuth config" - ); + + if (redirectUri.isPresent() && !isAuthorizedRedirectUri(redirectUri.get())) { + throw new BadRequestException( + "Unauthorized redirect uri found - check the redirect uri in your OAuth config" + ); } String targetUrl = redirectUri.orElse(getDefaultTargetUrl()); @@ -150,19 +152,19 @@ public class OAuth2AuthenticationSuccessHandler extends SimpleUrlAuthenticationS protected void clearAuthenticationAttributes(HttpServletRequest request, HttpServletResponse response) { - super.clearAuthenticationAttributes(request); - httpCookieOAuth2AuthorizationRequestRepository.removeAuthorizationRequestCookies(request, response); - } - - private boolean isAuthorizedRedirectUri(String uri) { - URI clientRedirectUri = URI.create(uri); - var authorizedRedirectUri = env.getOAuthRedirectUri(); - if (authorizedRedirectUri.exists()) { - URI authorizedURI = URI.create(authorizedRedirectUri.getValue()); - return authorizedURI.getHost().equalsIgnoreCase(clientRedirectUri.getHost()) - && authorizedURI.getPort() == clientRedirectUri.getPort(); - } else { - return false; - } - } -} + super.clearAuthenticationAttributes(request); + httpCookieOAuth2AuthorizationRequestRepository.removeAuthorizationRequestCookies(request, response); + } + + private boolean isAuthorizedRedirectUri(String uri) { + URI clientRedirectUri = URI.create(uri); + var authorizedRedirectUri = env.getOAuthRedirectUri(); + if (authorizedRedirectUri.exists()) { + URI authorizedURI = URI.create(authorizedRedirectUri.getValue()); + return authorizedURI.getHost().equalsIgnoreCase(clientRedirectUri.getHost()) + && authorizedURI.getPort() == clientRedirectUri.getPort(); + } else { + return false; + } + } +} 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 453a68105a..b7fc06ea82 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 @@ -55,7 +55,8 @@ public class DataLakeScheduler implements SchedulingConfigurer { this.dataLakeExportManager = new DataLakeExportManager( dataExplorerSchemaManagement, new DataExplorerDispatcher().getDataExplorerManager() - .getQueryManagement(dataExplorerSchemaManagement)); + .getQueryManagement(dataExplorerSchemaManagement), + resourceManager.getCoreConfigurationStorage()); } public void cleanupMeasurements() { diff --git a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/scheduler/certificates/CertificateExpiryEmailScheduler.java b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/scheduler/certificates/CertificateExpiryEmailScheduler.java index 8a234022cc..df8d376d17 100644 --- a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/scheduler/certificates/CertificateExpiryEmailScheduler.java +++ b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/scheduler/certificates/CertificateExpiryEmailScheduler.java @@ -24,6 +24,7 @@ import org.apache.streampipes.model.client.user.DefaultRole; import org.apache.streampipes.model.client.user.Principal; import org.apache.streampipes.model.mail.SpEmail; import org.apache.streampipes.model.opcua.Certificate; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.streampipes.storage.api.user.IUserStorage; import org.apache.streampipes.storage.management.StorageDispatcher; @@ -48,6 +49,12 @@ public class CertificateExpiryEmailScheduler implements SchedulingConfigurer { private static final String SUBJECT = "Upcoming certificate expirations — action required"; + private final ISpCoreConfigurationStorage coreConfigurationStorage; + + public CertificateExpiryEmailScheduler(ISpCoreConfigurationStorage coreConfigurationStorage) { + this.coreConfigurationStorage = coreConfigurationStorage; + } + public void checkForExpiringCertificates() { var certificateExpiryEmailDays = Environments.getEnvironment().getCertificateExpiryEmailDays().getValueOrDefault(); 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 1fdcadc25d..fb1c15e30d 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.ISpCoreConfigurationStorage; import org.apache.streampipes.storage.api.user.IPermissionStorage; import org.apache.streampipes.storage.couchdb.impl.connect.AdapterInstanceStorageImpl; import org.apache.streampipes.storage.couchdb.impl.explorer.ChartStorageImpl; @@ -34,6 +35,7 @@ import org.apache.streampipes.storage.couchdb.impl.explorer.DataLakeMeasureStora import org.apache.streampipes.storage.couchdb.impl.function.FunctionStateStorageImpl; 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.user.PermissionStorageImpl; import org.apache.streampipes.storage.couchdb.utils.Utils; @@ -103,6 +105,11 @@ public class StorageApiConfiguration { return new AssetStorageImpl(); } + @Bean + public ISpCoreConfigurationStorage coreConfigurationStorage() { + return new CoreConfigurationStorageImpl(); + } + @Bean public IPipelineStorage pipelineStorage(CacheManager cacheManager) { IPipelineStorage delegate = new PipelineStorageImpl(); diff --git a/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/jwt/JwtTokenProvider.java b/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/jwt/JwtTokenProvider.java index c4af4619eb..beb3bc76d3 100644 --- a/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/jwt/JwtTokenProvider.java +++ b/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/jwt/JwtTokenProvider.java @@ -23,11 +23,10 @@ import org.apache.streampipes.commons.environment.Environments; import org.apache.streampipes.model.client.user.Principal; import org.apache.streampipes.model.configuration.JwtSigningMode; import org.apache.streampipes.model.configuration.LocalAuthConfig; -import org.apache.streampipes.model.configuration.SpCoreConfiguration; import org.apache.streampipes.security.jwt.JwtTokenGenerator; import org.apache.streampipes.security.jwt.JwtTokenUtils; import org.apache.streampipes.security.jwt.JwtTokenValidator; -import org.apache.streampipes.storage.management.StorageDispatcher; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.streampipes.user.management.model.PrincipalUserDetails; import org.apache.streampipes.user.management.util.GrantedAuthoritiesBuilder; import org.apache.streampipes.user.management.util.UserInfoUtil; @@ -52,15 +51,11 @@ public class JwtTokenProvider { public static final String CLAIM_USER = "user"; private static final Logger LOG = LoggerFactory.getLogger(JwtTokenProvider.class); - private SpCoreConfiguration config; private Environment env; + private final ISpCoreConfigurationStorage coreConfigurationStorage; - public JwtTokenProvider() { - this.config = StorageDispatcher - .INSTANCE - .getNoSqlStore() - .getSpCoreConfigurationStorage() - .get(); + public JwtTokenProvider(ISpCoreConfigurationStorage coreConfigurationStorage) { + this.coreConfigurationStorage = coreConfigurationStorage; this.env = Environments.getEnvironment(); } @@ -109,11 +104,11 @@ public class JwtTokenProvider { } public String getUserIdFromToken(String token) { - return JwtTokenUtils.getUserIdFromToken(token, new SpKeyResolver(tokenSecret())); + return JwtTokenUtils.getUserIdFromToken(token, new SpKeyResolver(tokenSecret(), coreConfigurationStorage)); } public boolean validateJwtToken(String jwtToken) { - return JwtTokenValidator.validateJwtToken(jwtToken, new SpKeyResolver(tokenSecret())); + return JwtTokenValidator.validateJwtToken(jwtToken, new SpKeyResolver(tokenSecret(), coreConfigurationStorage)); } public boolean validateJwtToken(String tokenSecret, @@ -130,7 +125,7 @@ public class JwtTokenProvider { } private LocalAuthConfig authConfig() { - return this.config.getLocalAuthConfig(); + return this.coreConfigurationStorage.get().getLocalAuthConfig(); } private Date makeExpirationDate() { diff --git a/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/jwt/SpKeyResolver.java b/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/jwt/SpKeyResolver.java index d0d9a48858..9577c43d4a 100644 --- a/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/jwt/SpKeyResolver.java +++ b/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/jwt/SpKeyResolver.java @@ -21,6 +21,7 @@ import org.apache.streampipes.model.client.user.Principal; import org.apache.streampipes.model.client.user.ServiceAccount; import org.apache.streampipes.model.client.user.UserAccount; import org.apache.streampipes.security.jwt.KeyGenerator; +import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.streampipes.storage.api.user.IUserStorage; import org.apache.streampipes.storage.management.StorageDispatcher; import org.apache.streampipes.user.management.encryption.SecretEncryptionManager; @@ -35,9 +36,12 @@ public class SpKeyResolver implements SigningKeyResolver { private final String tokenSecret; private final IUserStorage userStorage; + private final ISpCoreConfigurationStorage coreConfigStorage; - public SpKeyResolver(String tokenSecret) { + public SpKeyResolver(String tokenSecret, + ISpCoreConfigurationStorage coreConfigurationStorage) { this.tokenSecret = tokenSecret; + this.coreConfigStorage = coreConfigurationStorage; this.userStorage = StorageDispatcher.INSTANCE.getNoSqlStore().getUserStorageAPI(); } @@ -69,10 +73,7 @@ public class SpKeyResolver implements SigningKeyResolver { } public String getPublicKeyFromConfig() { - return StorageDispatcher - .INSTANCE - .getNoSqlStore() - .getSpCoreConfigurationStorage() + return coreConfigStorage .get() .getLocalAuthConfig() .getPublicKey();
