This is an automated email from the ASF dual-hosted git repository. dominikriemer pushed a commit to branch introduce-storage-cache in repository https://gitbox.apache.org/repos/asf/streampipes.git
commit 8a0610629b742d1e06f5e84c02b6b0e3d24c9f7d Author: Dominik Riemer <[email protected]> AuthorDate: Thu Jun 18 17:30:26 2026 +0200 Add permission cache --- .../dataexplorer/api/IDataExplorerManager.java | 4 +- .../influx/DataExplorerManagerInflux.java | 6 +- .../iotdb/DataExplorerManagerIotDb.java | 8 +- .../loadbalance/unit/InvokeHttpRequest.java | 3 +- .../AbstractPipelineElementResourceManager.java | 3 +- .../management/AdapterResourceManager.java | 9 +- .../resource/management/AssetResourceManager.java | 4 +- .../resource/management/CrudResourceManager.java | 2 +- .../management/DataLakeMeasureResourceManager.java | 4 +- .../management/PermissionResourceManager.java | 5 - .../resource/management/SpResourceManager.java | 6 +- .../permission/SpPermissionEvaluator.java | 5 +- .../management/AdapterResourceManagerTest.java | 2 +- .../apache/streampipes/rest/ResetManagement.java | 2 +- .../rest/impl/AdapterMonitoringResource.java | 2 +- .../rest/impl/connect/AdapterResource.java | 2 +- .../impl/datalake/AbstractDataLakeResource.java | 10 +- .../impl/datalake/DataLakeMeasureResource.java | 2 +- .../rest/impl/datalake/DataLakeResource.java | 6 +- .../datalake/KioskDashboardDataLakeResource.java | 10 +- .../datalake/importer/DataLakeImportResource.java | 6 +- .../service/core/WebSecurityConfig.java | 11 +- .../core/filter/TokenAuthenticationFilter.java | 9 +- .../core/oauth2/CustomOAuth2UserService.java | 9 +- .../service/core/oauth2/CustomOidcUserService.java | 8 +- .../core/oauth2/OidcUserAccountDetails.java | 11 +- .../service/core/oauth2/UserService.java | 33 +-- .../service/core/scheduler/DataLakeScheduler.java | 6 +- .../core/storage/CachedPermissionStorage.java | 236 +++++++++++++++++++++ .../core/storage/StorageApiConfiguration.java | 7 +- .../src/main/resources/application.properties | 2 +- .../core/storage/CachedPermissionStorageTest.java | 122 +++++++++++ .../storage/api/core/INoSqlStorage.java | 3 - .../storage/couchdb/CouchDbStorageManager.java | 7 - .../management/model/PrincipalUserDetails.java | 6 +- .../management/model/ServiceAccountDetails.java | 6 +- .../user/management/model/UserAccountDetails.java | 6 +- .../management/service/SpUserDetailsService.java | 11 +- .../management/util/GrantedPermissionsBuilder.java | 12 +- 39 files changed, 504 insertions(+), 102 deletions(-) diff --git a/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataExplorerManager.java b/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataExplorerManager.java index 3de231ca5f..3f87b4b5f5 100644 --- a/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataExplorerManager.java +++ b/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataExplorerManager.java @@ -21,6 +21,7 @@ package org.apache.streampipes.dataexplorer.api; import org.apache.streampipes.client.api.IStreamPipesClient; import org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator; import org.apache.streampipes.model.datalake.DataLakeMeasure; +import org.apache.streampipes.storage.api.user.IPermissionStorage; import java.util.List; @@ -42,7 +43,8 @@ public interface IDataExplorerManager { IDataExplorerQueryManagement getQueryManagement(IDataExplorerSchemaManagement dataExplorerSchemaManagement); - IDataExplorerSchemaManagement getSchemaManagement(ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator); + IDataExplorerSchemaManagement getSchemaManagement(ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator, + IPermissionStorage permissionStorage); default ITimeSeriesStorage getTimeseriesStorage(DataLakeMeasure measure) { return getTimeseriesStorage(measure, false); diff --git a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerManagerInflux.java b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerManagerInflux.java index 8d600847d5..181ecb5be5 100644 --- a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerManagerInflux.java +++ b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerManagerInflux.java @@ -32,6 +32,7 @@ import org.apache.streampipes.dataexplorer.influx.sanitize.DataLakeMeasurementSa import org.apache.streampipes.manager.permission.DataLakePermissionManager; import org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator; import org.apache.streampipes.model.datalake.DataLakeMeasure; +import org.apache.streampipes.storage.api.user.IPermissionStorage; import org.apache.streampipes.storage.management.StorageDispatcher; import java.util.List; @@ -57,12 +58,13 @@ public class DataExplorerManagerInflux implements IDataExplorerManager { } @Override - public IDataExplorerSchemaManagement getSchemaManagement(ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator) { + public IDataExplorerSchemaManagement getSchemaManagement(ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator, + IPermissionStorage permissionStorage) { return new DataExplorerSchemaManagement(StorageDispatcher.INSTANCE .getNoSqlStore() .getDataLakeStorage(), new DataLakePermissionManager( - StorageDispatcher.INSTANCE.getNoSqlStore().getPermissionStorage() + permissionStorage ), chartSchemaUpdateCoordinator ); diff --git a/streampipes-data-explorer-iotdb/src/main/java/org/apache/streampipes/dataexplorer/iotdb/DataExplorerManagerIotDb.java b/streampipes-data-explorer-iotdb/src/main/java/org/apache/streampipes/dataexplorer/iotdb/DataExplorerManagerIotDb.java index 4db64664d7..a8c31313d9 100644 --- a/streampipes-data-explorer-iotdb/src/main/java/org/apache/streampipes/dataexplorer/iotdb/DataExplorerManagerIotDb.java +++ b/streampipes-data-explorer-iotdb/src/main/java/org/apache/streampipes/dataexplorer/iotdb/DataExplorerManagerIotDb.java @@ -31,6 +31,7 @@ import org.apache.streampipes.dataexplorer.iotdb.sanitize.DataLakeMeasurementSan import org.apache.streampipes.manager.permission.DataLakePermissionManager; import org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator; import org.apache.streampipes.model.datalake.DataLakeMeasure; +import org.apache.streampipes.storage.api.user.IPermissionStorage; import org.apache.streampipes.storage.management.StorageDispatcher; import java.util.List; @@ -53,13 +54,12 @@ public class DataExplorerManagerIotDb implements IDataExplorerManager { } @Override - public IDataExplorerSchemaManagement getSchemaManagement(ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator) { + public IDataExplorerSchemaManagement getSchemaManagement(ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator, + IPermissionStorage permissionStorage) { return new DataExplorerSchemaManagement(StorageDispatcher.INSTANCE .getNoSqlStore() .getDataLakeStorage(), - new DataLakePermissionManager( - StorageDispatcher.INSTANCE.getNoSqlStore().getPermissionStorage() - ), + new DataLakePermissionManager(permissionStorage), chartSchemaUpdateCoordinator ); } diff --git a/streampipes-load-balancer/src/main/java/org/apache/streampipes/loadbalance/unit/InvokeHttpRequest.java b/streampipes-load-balancer/src/main/java/org/apache/streampipes/loadbalance/unit/InvokeHttpRequest.java index ba7d6026c9..0af6c20bcd 100644 --- a/streampipes-load-balancer/src/main/java/org/apache/streampipes/loadbalance/unit/InvokeHttpRequest.java +++ b/streampipes-load-balancer/src/main/java/org/apache/streampipes/loadbalance/unit/InvokeHttpRequest.java @@ -23,6 +23,7 @@ import org.apache.streampipes.model.client.user.Permission; import org.apache.streampipes.model.client.user.Principal; import org.apache.streampipes.model.pipeline.PipelineElementStatus; import org.apache.streampipes.serializers.json.JacksonSerializer; +import org.apache.streampipes.storage.couchdb.impl.user.PermissionStorageImpl; import org.apache.streampipes.storage.management.StorageDispatcher; import org.apache.streampipes.user.management.jwt.JwtTokenProvider; @@ -121,7 +122,7 @@ public class InvokeHttpRequest{ } private static String getOwnerSid(String resourceId) { - return StorageDispatcher.INSTANCE.getNoSqlStore().getPermissionStorage().getUserPermissionsForObject(resourceId) + return new PermissionStorageImpl("users/permissions").getUserPermissionsForObject(resourceId) .stream() .findFirst() .map(Permission::getOwnerSid) diff --git a/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/AbstractPipelineElementResourceManager.java b/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/AbstractPipelineElementResourceManager.java index 2e4ee34f18..4ca033cf56 100644 --- a/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/AbstractPipelineElementResourceManager.java +++ b/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/AbstractPipelineElementResourceManager.java @@ -78,8 +78,7 @@ public abstract class AbstractPipelineElementResourceManager<T extends CRUDStora W existing = find(pipelineElement.getElementId()); if (existing == null) { this.db.persist(pipelineElement); - new PermissionResourceManager() - .createDefault(pipelineElement.getElementId(), SpDataStream.class, principalSid, false); + permissionResourceManager.createDefault(pipelineElement.getElementId(), SpDataStream.class, principalSid, false); } else { throw new IllegalArgumentException("This pipeline element already exists"); } diff --git a/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/AdapterResourceManager.java b/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/AdapterResourceManager.java index f3ba00fc39..b4ff3cbed2 100644 --- a/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/AdapterResourceManager.java +++ b/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/AdapterResourceManager.java @@ -36,17 +36,18 @@ public class AdapterResourceManager extends AbstractResourceManager<IAdapterStor private final SpPermissionEvaluator permissionEvaluator; public AdapterResourceManager(IAdapterStorage adapterStorage, - ICertificateStorage certificateStorage) { + ICertificateStorage certificateStorage, + PermissionResourceManager permissionResourceManager) { super(adapterStorage); this.certificateStorage = certificateStorage; - this.permissionEvaluator = new SpPermissionEvaluator(); + this.permissionEvaluator = new SpPermissionEvaluator(permissionResourceManager.getDb()); } - public AdapterResourceManager() { + public AdapterResourceManager(PermissionResourceManager permissionResourceManager) { super(StorageDispatcher.INSTANCE.getNoSqlStore() .getAdapterInstanceStorage()); this.certificateStorage = StorageDispatcher.INSTANCE.getNoSqlStore().getCertificateStorage(); - this.permissionEvaluator = new SpPermissionEvaluator(); + this.permissionEvaluator = new SpPermissionEvaluator(permissionResourceManager.getDb()); } public ResourceSummaryDto<AdapterSummaryDto> getSummary(Authentication auth) { diff --git a/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/AssetResourceManager.java b/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/AssetResourceManager.java index b8aa3a0ba4..ecd02ee85f 100644 --- a/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/AssetResourceManager.java +++ b/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/AssetResourceManager.java @@ -31,9 +31,9 @@ public class AssetResourceManager extends AbstractResourceManager<IAssetStorage> private final SpPermissionEvaluator permissionEvaluator; - public AssetResourceManager() { + public AssetResourceManager(PermissionResourceManager permissionResourceManager) { super(StorageDispatcher.INSTANCE.getNoSqlStore().getAssetStorage()); - this.permissionEvaluator = new SpPermissionEvaluator(); + this.permissionEvaluator = new SpPermissionEvaluator(permissionResourceManager.getDb()); } public ResourceSummaryDto<AssetSummaryDto> getSummary(Authentication auth) { diff --git a/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/CrudResourceManager.java b/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/CrudResourceManager.java index 8951bbb6ad..edb0586267 100644 --- a/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/CrudResourceManager.java +++ b/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/CrudResourceManager.java @@ -37,7 +37,7 @@ public class CrudResourceManager<T extends Storable> PermissionResourceManager permissionResourceManager) { super(db); this.elementClass = elementClass; - this.permissionEvaluator = new SpPermissionEvaluator(); + this.permissionEvaluator = new SpPermissionEvaluator(permissionResourceManager.getDb()); this.permissionResourceManager = permissionResourceManager; } diff --git a/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/DataLakeMeasureResourceManager.java b/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/DataLakeMeasureResourceManager.java index a80c7280e1..ee8892818f 100644 --- a/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/DataLakeMeasureResourceManager.java +++ b/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/DataLakeMeasureResourceManager.java @@ -46,10 +46,10 @@ public class DataLakeMeasureResourceManager extends AbstractResourceManager<IDat private final IPipelineStorage pipelineStorage; private final SpPermissionEvaluator permissionEvaluator; - public DataLakeMeasureResourceManager() { + public DataLakeMeasureResourceManager(PermissionResourceManager permissionResourceManager) { super(StorageDispatcher.INSTANCE.getNoSqlStore().getDataLakeStorage()); this.pipelineStorage = StorageDispatcher.INSTANCE.getNoSqlStore().getPipelineStorageAPI(); - this.permissionEvaluator = new SpPermissionEvaluator(); + this.permissionEvaluator = new SpPermissionEvaluator(permissionResourceManager.getDb()); } public ResourceSummaryDto<DatasetSummaryDto> getSummary(Authentication auth) { diff --git a/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/PermissionResourceManager.java b/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/PermissionResourceManager.java index b9b5751968..ea72f5b75c 100644 --- a/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/PermissionResourceManager.java +++ b/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/PermissionResourceManager.java @@ -21,7 +21,6 @@ import org.apache.streampipes.model.client.user.Permission; import org.apache.streampipes.model.client.user.PermissionBuilder; import org.apache.streampipes.model.dashboard.DashboardModel; import org.apache.streampipes.storage.api.user.IPermissionStorage; -import org.apache.streampipes.storage.management.StorageDispatcher; import java.util.List; @@ -31,10 +30,6 @@ public class PermissionResourceManager extends AbstractResourceManager<IPermissi DashboardModel.class.getCanonicalName() ); - public PermissionResourceManager() { - super(StorageDispatcher.INSTANCE.getNoSqlStore().getPermissionStorage()); - } - public PermissionResourceManager(IPermissionStorage permissionStorage) { super(permissionStorage); } 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 83107b42e1..11a9e08d96 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 @@ -49,15 +49,15 @@ public class SpResourceManager { } public AssetResourceManager manageAssets() { - return new AssetResourceManager(); + return new AssetResourceManager(managePermissions()); } public AdapterResourceManager manageAdapters() { - return new AdapterResourceManager(); + return new AdapterResourceManager(managePermissions()); } public DataLakeMeasureResourceManager manageDataLakeMeasures() { - return new DataLakeMeasureResourceManager(); + return new DataLakeMeasureResourceManager(managePermissions()); } public PermissionResourceManager managePermissions() { diff --git a/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/permission/SpPermissionEvaluator.java b/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/permission/SpPermissionEvaluator.java index 0c208847c4..ce6f451c7d 100644 --- a/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/permission/SpPermissionEvaluator.java +++ b/streampipes-resource-management/src/main/java/org/apache/streampipes/resource/management/permission/SpPermissionEvaluator.java @@ -22,7 +22,6 @@ import org.apache.streampipes.model.client.user.Permission; import org.apache.streampipes.model.pipeline.PipelineElementRecommendation; import org.apache.streampipes.model.pipeline.PipelineElementRecommendationMessage; import org.apache.streampipes.storage.api.user.IPermissionStorage; -import org.apache.streampipes.storage.management.StorageDispatcher; import org.apache.streampipes.user.management.model.PrincipalUserDetails; import org.springframework.security.access.PermissionEvaluator; @@ -40,8 +39,8 @@ public class SpPermissionEvaluator implements PermissionEvaluator { private final IPermissionStorage permissionStorage; - public SpPermissionEvaluator() { - this.permissionStorage = StorageDispatcher.INSTANCE.getNoSqlStore().getPermissionStorage(); + public SpPermissionEvaluator(IPermissionStorage permissionStorage) { + this.permissionStorage = permissionStorage; } /** diff --git a/streampipes-resource-management/src/test/java/org/apache/streampipes/resource/management/AdapterResourceManagerTest.java b/streampipes-resource-management/src/test/java/org/apache/streampipes/resource/management/AdapterResourceManagerTest.java index 3e0175c738..ec63f239b9 100644 --- a/streampipes-resource-management/src/test/java/org/apache/streampipes/resource/management/AdapterResourceManagerTest.java +++ b/streampipes-resource-management/src/test/java/org/apache/streampipes/resource/management/AdapterResourceManagerTest.java @@ -43,7 +43,7 @@ public class AdapterResourceManagerTest { @BeforeEach void setUp() { storage = mock(IAdapterStorage.class); - adapterResourceManager = new AdapterResourceManager(storage, null); + adapterResourceManager = new AdapterResourceManager(storage, null, null); } @Test diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/ResetManagement.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/ResetManagement.java index 9d3e764166..50065f2b22 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 @@ -158,7 +158,7 @@ public class ResetManagement { private void removeAllDataInDataLake() { var dataLakeMeasureManagement = new DataExplorerDispatcher() .getDataExplorerManager() - .getSchemaManagement(chartSchemaUpdateCoordinator); + .getSchemaManagement(chartSchemaUpdateCoordinator, resourceManager.managePermissions().getDb()); var dataExplorerQueryManagement = new DataExplorerDispatcher() .getDataExplorerManager() .getQueryManagement(dataLakeMeasureManagement); diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/AdapterMonitoringResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/AdapterMonitoringResource.java index cb13ff6e73..4fa53dbe19 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/AdapterMonitoringResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/AdapterMonitoringResource.java @@ -130,7 +130,7 @@ public class AdapterMonitoringResource extends AbstractMonitoringResource { */ private boolean checkAdapterPermission(AdapterDescription adapterDescription, String permission) { - var spPermissionEvaluator = new SpPermissionEvaluator(); + var spPermissionEvaluator = new SpPermissionEvaluator(resourceManager.managePermissions().getDb()); var authentication = SecurityContextHolder.getContext() .getAuthentication(); return spPermissionEvaluator.hasPermission( diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/AdapterResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/AdapterResource.java index da43b782af..86eae06dca 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/AdapterResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/AdapterResource.java @@ -218,7 +218,7 @@ public class AdapterResource extends AbstractAdapterResource<AdapterMasterManage */ private boolean checkAdapterPermission(AdapterDescription adapterDescription, String permission) { - var spPermissionEvaluator = new SpPermissionEvaluator(); + var spPermissionEvaluator = new SpPermissionEvaluator(resourceManager.managePermissions().getDb()); var authentication = SecurityContextHolder.getContext() .getAuthentication(); return spPermissionEvaluator.hasPermission( diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/AbstractDataLakeResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/AbstractDataLakeResource.java index b24717c1bd..ef9dab09b2 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/AbstractDataLakeResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/AbstractDataLakeResource.java @@ -21,6 +21,7 @@ import org.apache.streampipes.dataexplorer.api.IDataExplorerSchemaManagement; import org.apache.streampipes.dataexplorer.management.DataExplorerDispatcher; import org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator; import org.apache.streampipes.model.client.user.DefaultPrivilege; +import org.apache.streampipes.resource.management.SpResourceManager; import org.apache.streampipes.resource.management.permission.SpPermissionEvaluator; import org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource; import org.apache.streampipes.storage.api.explorer.IDataLakeMeasureStorage; @@ -36,11 +37,14 @@ public class AbstractDataLakeResource extends AbstractAuthGuardedRestResource { private final IDataLakeMeasureStorage dataLakeMeasureStorage = StorageDispatcher.INSTANCE.getNoSqlStore().getDataLakeStorage(); protected final ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator; + private final SpResourceManager resourceManager; - public AbstractDataLakeResource(ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator) { + public AbstractDataLakeResource(ChartSchemaUpdateCoordinator chartSchemaUpdateCoordinator, + SpResourceManager resourceManager) { this.chartSchemaUpdateCoordinator = chartSchemaUpdateCoordinator; + this.resourceManager = resourceManager; this.dataLakeMeasureManagement = new DataExplorerDispatcher().getDataExplorerManager() - .getSchemaManagement(chartSchemaUpdateCoordinator); + .getSchemaManagement(chartSchemaUpdateCoordinator, resourceManager.managePermissions().getDb()); } /** @@ -72,7 +76,7 @@ public class AbstractDataLakeResource extends AbstractAuthGuardedRestResource { var measure = dataLakeMeasureStorage.getByMeasureName(measurementName); if (Objects.nonNull(measure)) { - var spPermissionEvaluator = new SpPermissionEvaluator(); + var spPermissionEvaluator = new SpPermissionEvaluator(resourceManager.managePermissions().getDb()); var authentication = SecurityContextHolder.getContext() .getAuthentication(); return spPermissionEvaluator.hasPermission( diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeMeasureResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeMeasureResource.java index bb6b2176c0..b490c591ae 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeMeasureResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeMeasureResource.java @@ -55,7 +55,7 @@ public class DataLakeMeasureResource extends AbstractDataLakeResource { public DataLakeMeasureResource(IDataExplorerWidgetStorage chartStorage, SpResourceManager resourceManager) { - super(new ChartSchemaUpdateCoordinator(chartStorage)); + super(new ChartSchemaUpdateCoordinator(chartStorage), resourceManager); this.resourceManager = resourceManager; } 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 1d6dd96d83..d2db088eab 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 @@ -31,6 +31,7 @@ import org.apache.streampipes.model.datalake.SpQueryResult; import org.apache.streampipes.model.datalake.param.ProvidedRestQueryParams; import org.apache.streampipes.model.message.Notifications; import org.apache.streampipes.model.monitoring.SpLogMessage; +import org.apache.streampipes.resource.management.SpResourceManager; import org.apache.streampipes.rest.security.AuthConstants; import org.apache.streampipes.rest.shared.exception.SpMessageException; import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; @@ -97,8 +98,9 @@ public class DataLakeResource extends AbstractDataLakeResource { private final IDataExplorerQueryManagement dataExplorerQueryManagement; private final DataLakeExportManager dataLakeExportManager; - public DataLakeResource(IDataExplorerWidgetStorage chartStorage) { - super(new ChartSchemaUpdateCoordinator(chartStorage)); + public DataLakeResource(IDataExplorerWidgetStorage chartStorage, + SpResourceManager resourceManager) { + super(new ChartSchemaUpdateCoordinator(chartStorage), resourceManager); this.dataExplorerQueryManagement = new DataExplorerDispatcher() .getDataExplorerManager() .getQueryManagement(this.dataLakeMeasureManagement); diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/KioskDashboardDataLakeResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/KioskDashboardDataLakeResource.java index 84305ad078..1c5b44daad 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/KioskDashboardDataLakeResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/KioskDashboardDataLakeResource.java @@ -27,6 +27,7 @@ import org.apache.streampipes.model.datalake.DataExplorerWidgetModel; import org.apache.streampipes.model.datalake.SpQueryResult; import org.apache.streampipes.model.datalake.param.ProvidedRestQueryParams; import org.apache.streampipes.model.monitoring.SpLogMessage; +import org.apache.streampipes.resource.management.SpResourceManager; import org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource; import org.apache.streampipes.storage.api.explorer.IDataExplorerDashboardStorage; import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; @@ -56,15 +57,18 @@ public class KioskDashboardDataLakeResource extends AbstractAuthGuardedRestResou private final IPermissionStorage permissionStorage; public KioskDashboardDataLakeResource(IDataExplorerWidgetStorage dataExplorerWidgetStorage, - IPermissionStorage permissionStorage) { + SpResourceManager resourceManager) { IDataExplorerSchemaManagement dataExplorerSchemaManagement = new DataExplorerDispatcher() .getDataExplorerManager() - .getSchemaManagement(new ChartSchemaUpdateCoordinator(dataExplorerWidgetStorage)); + .getSchemaManagement( + new ChartSchemaUpdateCoordinator(dataExplorerWidgetStorage), + resourceManager.managePermissions().getDb() + ); this.dataExplorerQueryManagement = new DataExplorerDispatcher() .getDataExplorerManager() .getQueryManagement(dataExplorerSchemaManagement); this.dataExplorerWidgetStorage = dataExplorerWidgetStorage; - this.permissionStorage = permissionStorage; + this.permissionStorage = resourceManager.managePermissions().getDb(); } @PostMapping(path = "/{dashboardId}/{widgetId}/data", diff --git a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/DataLakeImportResource.java b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/DataLakeImportResource.java index 54a9840f3c..338e8e5762 100644 --- a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/DataLakeImportResource.java +++ b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/importer/DataLakeImportResource.java @@ -26,6 +26,7 @@ import org.apache.streampipes.model.datalake.importer.CsvImportResult; import org.apache.streampipes.model.datalake.importer.CsvImportSchemaValidationRequest; import org.apache.streampipes.model.datalake.importer.CsvImportSchemaValidationResult; import org.apache.streampipes.model.datalake.importer.CsvImportTargetMode; +import org.apache.streampipes.resource.management.SpResourceManager; import org.apache.streampipes.rest.impl.datalake.AbstractDataLakeResource; import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; @@ -48,8 +49,9 @@ public class DataLakeImportResource extends AbstractDataLakeResource { private final CsvDataLakeImportService importService; - public DataLakeImportResource(IDataExplorerWidgetStorage chartStorage) { - super(new ChartSchemaUpdateCoordinator(chartStorage)); + public DataLakeImportResource(IDataExplorerWidgetStorage chartStorage, + SpResourceManager resourceManager) { + super(new ChartSchemaUpdateCoordinator(chartStorage), resourceManager); this.importService = new CsvDataLakeImportService(getDataLakeMeasureManagement()); } 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 229d03e167..eb25ed6f3c 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.user.IPermissionStorage; import org.apache.streampipes.user.management.service.SpUserDetailsService; import org.slf4j.Logger; @@ -95,10 +96,14 @@ public class WebSecurityConfig { @Autowired private OAuth2AuthenticationFailureHandler oAuth2AuthenticationFailureHandler; - public WebSecurityConfig(StreamPipesPasswordEncoder passwordEncoder) { + private final IPermissionStorage permissionStorage; + + public WebSecurityConfig(StreamPipesPasswordEncoder passwordEncoder, + IPermissionStorage permissionStorage) { this.passwordEncoder = passwordEncoder; - this.userDetailsService = new SpUserDetailsService(); + this.userDetailsService = new SpUserDetailsService(permissionStorage); this.env = Environments.getEnvironment(); + this.permissionStorage = permissionStorage; } @Autowired @@ -153,7 +158,7 @@ public class WebSecurityConfig { } public TokenAuthenticationFilter tokenAuthenticationFilter() { - return new TokenAuthenticationFilter(); + return new TokenAuthenticationFilter(permissionStorage); } @Bean(BeanIds.USER_DETAILS_SERVICE) 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 6a14bd320f..e91c87e940 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.user.IPermissionStorage; import org.apache.streampipes.storage.api.user.IUserStorage; import org.apache.streampipes.storage.management.StorageDispatcher; import org.apache.streampipes.user.management.encryption.SecretEncryptionManager; @@ -61,6 +62,7 @@ public class TokenAuthenticationFilter extends OncePerRequestFilter { private final JwtTokenProvider tokenProvider; private final IUserStorage userStorage; + private final IPermissionStorage permissionStorage; private final List<String> supportedBasicAuthPaths = List.of( "/actuator/prometheus" @@ -70,9 +72,10 @@ public class TokenAuthenticationFilter extends OncePerRequestFilter { private static final Logger logger = LoggerFactory.getLogger(TokenAuthenticationFilter.class); - public TokenAuthenticationFilter() { + public TokenAuthenticationFilter(IPermissionStorage permissionStorage) { this.tokenProvider = new JwtTokenProvider(); this.userStorage = StorageDispatcher.INSTANCE.getNoSqlStore().getUserStorageAPI(); + this.permissionStorage = permissionStorage; } @Override @@ -155,8 +158,8 @@ public class TokenAuthenticationFilter extends OncePerRequestFilter { } private PrincipalUserDetails<?> makeDetails(Principal user) { - return user instanceof UserAccount ? new UserAccountDetails((UserAccount) user) : - new ServiceAccountDetails((ServiceAccount) user); + return user instanceof UserAccount ? new UserAccountDetails((UserAccount) user, permissionStorage) : + new ServiceAccountDetails((ServiceAccount) user, permissionStorage); } private boolean isAdminUser(PrincipalUserDetails<?> userDetails) { diff --git a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/CustomOAuth2UserService.java b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/CustomOAuth2UserService.java index 7715d1c668..0177814cb9 100755 --- a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/CustomOAuth2UserService.java +++ b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/CustomOAuth2UserService.java @@ -19,6 +19,7 @@ package org.apache.streampipes.service.core.oauth2; import org.apache.streampipes.rest.security.OAuth2AuthenticationProcessingException; +import org.apache.streampipes.storage.api.user.IPermissionStorage; import org.springframework.security.core.AuthenticationException; import org.springframework.security.oauth2.client.userinfo.DefaultOAuth2UserService; @@ -32,13 +33,19 @@ import java.util.HashMap; @Service public class CustomOAuth2UserService extends DefaultOAuth2UserService { + private final IPermissionStorage permissionStorage; + + public CustomOAuth2UserService(IPermissionStorage permissionStorage) { + this.permissionStorage = permissionStorage; + } + @Override public OAuth2User loadUser(OAuth2UserRequest oAuth2UserRequest) throws OAuth2AuthenticationException { OAuth2User oAuth2User = super.loadUser(oAuth2UserRequest); try { var attributes = new HashMap<>(oAuth2User.getAttributes()); var provider = oAuth2UserRequest.getClientRegistration().getRegistrationId(); - return new UserService().processUserRegistration(provider, attributes); + return new UserService(permissionStorage).processUserRegistration(provider, attributes); } catch (AuthenticationException e) { throw e; } catch (Exception e) { diff --git a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/CustomOidcUserService.java b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/CustomOidcUserService.java index 558fb870c3..eb9371402d 100755 --- a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/CustomOidcUserService.java +++ b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/CustomOidcUserService.java @@ -21,6 +21,7 @@ package org.apache.streampipes.service.core.oauth2; import org.apache.streampipes.commons.environment.Environments; import org.apache.streampipes.rest.security.OAuth2AuthenticationProcessingException; +import org.apache.streampipes.storage.api.user.IPermissionStorage; import org.springframework.security.core.AuthenticationException; import org.springframework.security.oauth2.client.oidc.userinfo.OidcUserRequest; @@ -34,7 +35,10 @@ import java.util.Objects; @Service public class CustomOidcUserService extends OidcUserService { - public CustomOidcUserService() { + private final IPermissionStorage permissionStorage; + + public CustomOidcUserService(IPermissionStorage permissionStorage) { + this.permissionStorage = permissionStorage; var env = Environments.getEnvironment(); this.setRetrieveUserInfo(req -> { var config = env.getOAuthConfigurations() @@ -52,7 +56,7 @@ public class CustomOidcUserService extends OidcUserService { OidcUser oidcUser = super.loadUser(userRequest); try { var provider = userRequest.getClientRegistration().getRegistrationId(); - return new UserService().processUserRegistration( + return new UserService(permissionStorage).processUserRegistration( provider, oidcUser.getAttributes(), oidcUser.getIdToken(), diff --git a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/OidcUserAccountDetails.java b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/OidcUserAccountDetails.java index 48b3fd5c4f..671adbad20 100755 --- a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/OidcUserAccountDetails.java +++ b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/OidcUserAccountDetails.java @@ -19,6 +19,7 @@ package org.apache.streampipes.service.core.oauth2; import org.apache.streampipes.model.client.user.UserAccount; +import org.apache.streampipes.storage.api.user.IPermissionStorage; import org.apache.streampipes.user.management.model.UserAccountDetails; import org.springframework.security.oauth2.core.oidc.OidcIdToken; @@ -36,8 +37,9 @@ public class OidcUserAccountDetails extends UserAccountDetails implements OAuth2 public OidcUserAccountDetails(UserAccount user, OidcIdToken idToken, - OidcUserInfo userInfo) { - super(user); + OidcUserInfo userInfo, + IPermissionStorage permissionStorage) { + super(user, permissionStorage); this.idToken = idToken; this.userInfo = userInfo; } @@ -45,8 +47,9 @@ public class OidcUserAccountDetails extends UserAccountDetails implements OAuth2 public static OidcUserAccountDetails create(UserAccount user, Map<String, Object> attributes, OidcIdToken idToken, - OidcUserInfo userInfo) { - OidcUserAccountDetails localUser = new OidcUserAccountDetails(user, idToken, userInfo); + OidcUserInfo userInfo, + IPermissionStorage permissionStorage) { + OidcUserAccountDetails localUser = new OidcUserAccountDetails(user, idToken, userInfo, permissionStorage); localUser.setAttributes(attributes); return localUser; } diff --git a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/UserService.java b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/UserService.java index bac55bb048..b2c2991134 100755 --- a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/UserService.java +++ b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/oauth2/UserService.java @@ -21,14 +21,15 @@ 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.commons.environment.model.OAuthConfiguration; -import org.apache.streampipes.model.client.user.Group; -import org.apache.streampipes.model.client.user.Role; -import org.apache.streampipes.model.client.user.UserAccount; -import org.apache.streampipes.resource.management.UserResourceManager; -import org.apache.streampipes.rest.security.OAuth2AuthenticationProcessingException; -import org.apache.streampipes.storage.api.user.IRoleStorage; -import org.apache.streampipes.storage.api.user.IUserGroupStorage; -import org.apache.streampipes.storage.api.user.IUserStorage; +import org.apache.streampipes.model.client.user.Group; +import org.apache.streampipes.model.client.user.Role; +import org.apache.streampipes.model.client.user.UserAccount; +import org.apache.streampipes.resource.management.UserResourceManager; +import org.apache.streampipes.rest.security.OAuth2AuthenticationProcessingException; +import org.apache.streampipes.storage.api.user.IPermissionStorage; +import org.apache.streampipes.storage.api.user.IRoleStorage; +import org.apache.streampipes.storage.api.user.IUserGroupStorage; +import org.apache.streampipes.storage.api.user.IUserStorage; import org.apache.streampipes.storage.management.StorageDispatcher; import org.slf4j.Logger; @@ -45,22 +46,24 @@ import java.util.stream.Collectors; public class UserService { - private static final Logger LOG = LoggerFactory.getLogger(UserService.class); - - private final IUserStorage userStorage; - private final IRoleStorage roleStorage; - private final IUserGroupStorage groupStorage; + private static final Logger LOG = LoggerFactory.getLogger(UserService.class); + + private final IUserStorage userStorage; + private final IRoleStorage roleStorage; + private final IUserGroupStorage groupStorage; private final Environment env; private List<Role> allRoles; private List<Group> allGroups; + private final IPermissionStorage permissionStorage; - public UserService() { + public UserService(IPermissionStorage permissionStorage) { this.userStorage = StorageDispatcher.INSTANCE.getNoSqlStore().getUserStorageAPI(); this.roleStorage = StorageDispatcher.INSTANCE.getNoSqlStore().getRoleStorage(); this.groupStorage = StorageDispatcher.INSTANCE.getNoSqlStore().getUserGroupStorage(); this.allGroups = this.groupStorage.findAll(); this.allRoles = this.roleStorage.findAll(); this.env = Environments.getEnvironment(); + this.permissionStorage = permissionStorage; } public OidcUserAccountDetails processUserRegistration(String registrationId, @@ -103,7 +106,7 @@ public class UserService { } user = (UserAccount) userStorage.getUserById(principalId); - return OidcUserAccountDetails.create(user, attributes, idToken, userInfo); + return OidcUserAccountDetails.create(user, attributes, idToken, userInfo, permissionStorage); } else { throw new OAuth2AuthenticationProcessingException( String.format("No config found for provider %s", registrationId) 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 3c2522234d..04c45cd1fa 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 @@ -23,6 +23,7 @@ import org.apache.streampipes.dataexplorer.management.DataExplorerDispatcher; import org.apache.streampipes.export.DataLakeExportManager; import org.apache.streampipes.manager.pipeline.update.ChartSchemaUpdateCoordinator; import org.apache.streampipes.model.datalake.DataLakeMeasure; +import org.apache.streampipes.resource.management.SpResourceManager; import org.apache.streampipes.storage.api.explorer.IDataExplorerWidgetStorage; import org.slf4j.Logger; @@ -42,11 +43,12 @@ public class DataLakeScheduler implements SchedulingConfigurer { private final IDataExplorerSchemaManagement dataExplorerSchemaManagement; - public DataLakeScheduler(IDataExplorerWidgetStorage chartStorage) { + public DataLakeScheduler(IDataExplorerWidgetStorage chartStorage, + SpResourceManager resourceManager) { var chartSchemaUpdateCoordinator = new ChartSchemaUpdateCoordinator(chartStorage); dataExplorerSchemaManagement = new DataExplorerDispatcher() .getDataExplorerManager() - .getSchemaManagement(chartSchemaUpdateCoordinator); + .getSchemaManagement(chartSchemaUpdateCoordinator, resourceManager.managePermissions().getDb()); this.dataLakeExportManager = new DataLakeExportManager( dataExplorerSchemaManagement, new DataExplorerDispatcher().getDataExplorerManager() diff --git a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/storage/CachedPermissionStorage.java b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/storage/CachedPermissionStorage.java new file mode 100644 index 0000000000..3076842bf2 --- /dev/null +++ b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/storage/CachedPermissionStorage.java @@ -0,0 +1,236 @@ +/* + * 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.storage; + +import org.apache.streampipes.model.Tuple2; +import org.apache.streampipes.model.client.user.Permission; +import org.apache.streampipes.serializers.json.JacksonSerializer; +import org.apache.streampipes.storage.api.user.IPermissionStorage; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.core.type.TypeReference; +import com.fasterxml.jackson.databind.ObjectMapper; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.cache.Cache; +import org.springframework.cache.CacheManager; + +import java.util.ArrayList; +import java.util.List; +import java.util.Objects; +import java.util.Set; + +public class CachedPermissionStorage implements IPermissionStorage { + + static final String CACHE_NAME = "permissions"; + static final String FIND_ALL_CACHE_NAME = "permissionsAll"; + static final String BY_OBJECT_CACHE_NAME = "permissionsByObject"; + static final String BY_PRINCIPALS_CACHE_NAME = "objectPermissionsByPrincipals"; + + private static final Logger LOG = LoggerFactory.getLogger(CachedPermissionStorage.class); + private static final String FIND_ALL_CACHE_KEY = "all"; + private static final TypeReference<List<Permission>> PERMISSION_LIST_TYPE = new TypeReference<>() { + }; + private static final TypeReference<Set<String>> OBJECT_PERMISSION_SET_TYPE = new TypeReference<>() { + }; + + private final IPermissionStorage delegate; + private final Cache cache; + private final Cache findAllCache; + private final Cache byObjectCache; + private final Cache byPrincipalsCache; + private final ObjectMapper objectMapper; + + public CachedPermissionStorage(IPermissionStorage delegate, + CacheManager cacheManager) { + this(delegate, cacheManager, JacksonSerializer.getObjectMapper()); + } + + CachedPermissionStorage(IPermissionStorage delegate, + CacheManager cacheManager, + ObjectMapper objectMapper) { + this.delegate = delegate; + this.cache = getCache(cacheManager, CACHE_NAME); + this.findAllCache = getCache(cacheManager, FIND_ALL_CACHE_NAME); + this.byObjectCache = getCache(cacheManager, BY_OBJECT_CACHE_NAME); + this.byPrincipalsCache = getCache(cacheManager, BY_PRINCIPALS_CACHE_NAME); + this.objectMapper = objectMapper; + } + + @Override + public List<Permission> findAll() { + var cachedPermissions = getPermissionList(findAllCache, FIND_ALL_CACHE_KEY); + if (cachedPermissions != null) { + return cachedPermissions; + } + + var permissions = delegate.findAll(); + put(findAllCache, FIND_ALL_CACHE_KEY, permissions); + return permissions; + } + + @Override + public Tuple2<Boolean, String> persist(Permission element) { + var result = delegate.persist(element); + if (Boolean.TRUE.equals(result.k)) { + clearCaches(); + } + return result; + } + + @Override + public Permission getElementById(String id) { + var cachedPermission = getPermission(id); + if (cachedPermission != null) { + return cachedPermission; + } + + var permission = delegate.getElementById(id); + if (permission != null) { + put(cache, id, permission); + } + return permission; + } + + @Override + public Permission updateElement(Permission element) { + var updatedElement = delegate.updateElement(element); + clearCaches(); + return updatedElement; + } + + @Override + public void deleteElement(Permission element) { + delegate.deleteElement(element); + clearCaches(); + } + + @Override + public void deleteElementById(String id) { + delegate.deleteElementById(id); + clearCaches(); + } + + @Override + public Set<String> getObjectPermissions(List<String> sids) { + var cacheKey = makePrincipalCacheKey(sids); + var cachedPermissions = getObjectPermissionSet(cacheKey); + if (cachedPermissions != null) { + return cachedPermissions; + } + + var permissions = delegate.getObjectPermissions(sids); + put(byPrincipalsCache, cacheKey, permissions); + return permissions; + } + + @Override + public List<Permission> getUserPermissionsForObject(String objectInstanceId) { + var cachedPermissions = getPermissionList(byObjectCache, objectInstanceId); + if (cachedPermissions != null) { + return cachedPermissions; + } + + var permissions = delegate.getUserPermissionsForObject(objectInstanceId); + put(byObjectCache, objectInstanceId, permissions); + return permissions; + } + + private Permission getPermission(String id) { + try { + var cachedValue = cache.get(id, String.class); + return cachedValue == null ? null : objectMapper.readValue(cachedValue, Permission.class); + } catch (JsonProcessingException | RuntimeException e) { + LOG.warn("Could not read permission {} from cache", id, e); + evict(cache, id); + return null; + } + } + + private List<Permission> getPermissionList(Cache targetCache, + String key) { + try { + var cachedValue = targetCache.get(key, String.class); + return cachedValue == null ? null : objectMapper.readValue(cachedValue, PERMISSION_LIST_TYPE); + } catch (JsonProcessingException | RuntimeException e) { + LOG.warn("Could not read permissions from cache {}", targetCache.getName(), e); + evict(targetCache, key); + return null; + } + } + + private Set<String> getObjectPermissionSet(String key) { + try { + var cachedValue = byPrincipalsCache.get(key, String.class); + return cachedValue == null ? null : objectMapper.readValue(cachedValue, OBJECT_PERMISSION_SET_TYPE); + } catch (JsonProcessingException | RuntimeException e) { + LOG.warn("Could not read object permissions from cache", e); + evict(byPrincipalsCache, key); + return null; + } + } + + private String makePrincipalCacheKey(List<String> sids) { + var sortedSids = new ArrayList<>(sids); + sortedSids.sort(String::compareTo); + try { + return objectMapper.writeValueAsString(sortedSids); + } catch (JsonProcessingException e) { + throw new IllegalStateException("Could not create permission cache key", e); + } + } + + private void put(Cache targetCache, + String key, + Object value) { + try { + targetCache.put(key, objectMapper.writeValueAsString(value)); + } catch (JsonProcessingException | RuntimeException e) { + LOG.warn("Could not write permissions to cache {}", targetCache.getName(), e); + } + } + + private void clearCaches() { + clear(cache); + clear(findAllCache); + clear(byObjectCache); + clear(byPrincipalsCache); + } + + private void clear(Cache targetCache) { + try { + targetCache.clear(); + } catch (RuntimeException e) { + LOG.warn("Could not clear permission cache {}", targetCache.getName(), e); + } + } + + private void evict(Cache targetCache, + String key) { + try { + targetCache.evict(key); + } catch (RuntimeException e) { + LOG.warn("Could not evict permission cache entry from {}", targetCache.getName(), e); + } + } + + private static Cache getCache(CacheManager cacheManager, + String cacheName) { + return Objects.requireNonNull(cacheManager.getCache(cacheName)); + } +} 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 cdfcffeada..002d6fcfae 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 @@ -48,7 +48,10 @@ public class StorageApiConfiguration { } @Bean - public IPermissionStorage permissionStorage() { - return new PermissionStorageImpl("users/permissions"); + public IPermissionStorage permissionStorage(CacheManager cacheManager) { + return new CachedPermissionStorage( + new PermissionStorageImpl("users/permissions"), + cacheManager + ); } } diff --git a/streampipes-service-core/src/main/resources/application.properties b/streampipes-service-core/src/main/resources/application.properties index ece15145a6..f48263acff 100644 --- a/streampipes-service-core/src/main/resources/application.properties +++ b/streampipes-service-core/src/main/resources/application.properties @@ -22,7 +22,7 @@ spring.servlet.multipart.max-request-size=5120MB server.tomcat.additional-tld-skip-patterns=*.jar spring.mvc.async.request-timeout=900000 spring.cache.type=caffeine -spring.cache.cache-names=dataExplorerWidgets,dataExplorerWidgetsAll +spring.cache.cache-names=dataExplorerWidgets,dataExplorerWidgetsAll,permissions,permissionsAll,permissionsByObject,objectPermissionsByPrincipals spring.cache.caffeine.spec=maximumSize=10000,expireAfterWrite=10m logging.config=classpath:logback.xml spring.output.ansi.enabled=always diff --git a/streampipes-service-core/src/test/java/org/apache/streampipes/service/core/storage/CachedPermissionStorageTest.java b/streampipes-service-core/src/test/java/org/apache/streampipes/service/core/storage/CachedPermissionStorageTest.java new file mode 100644 index 0000000000..b276a9149e --- /dev/null +++ b/streampipes-service-core/src/test/java/org/apache/streampipes/service/core/storage/CachedPermissionStorageTest.java @@ -0,0 +1,122 @@ +/* + * 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.storage; + +import org.apache.streampipes.model.client.user.Permission; +import org.apache.streampipes.storage.api.user.IPermissionStorage; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.springframework.cache.concurrent.ConcurrentMapCacheManager; + +import java.util.ArrayList; +import java.util.List; +import java.util.Set; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotSame; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +class CachedPermissionStorageTest { + + private static final String PERMISSION_ID = "permission-id"; + private static final String OBJECT_ID = "object-id"; + + private IPermissionStorage delegate; + private CachedPermissionStorage storage; + + @BeforeEach + void setUp() { + delegate = mock(IPermissionStorage.class); + var cacheManager = new ConcurrentMapCacheManager( + CachedPermissionStorage.CACHE_NAME, + CachedPermissionStorage.FIND_ALL_CACHE_NAME, + CachedPermissionStorage.BY_OBJECT_CACHE_NAME, + CachedPermissionStorage.BY_PRINCIPALS_CACHE_NAME + ); + storage = new CachedPermissionStorage(delegate, cacheManager); + } + + @Test + void getUserPermissionsForObjectCachesSerializedCopies() { + var permission = makePermission("owner"); + when(delegate.getUserPermissionsForObject(OBJECT_ID)).thenReturn(List.of(permission)); + + var firstResult = storage.getUserPermissionsForObject(OBJECT_ID); + firstResult.get(0).setOwnerSid("changed-owner"); + var secondResult = storage.getUserPermissionsForObject(OBJECT_ID); + + assertNotSame(firstResult, secondResult); + assertNotSame(firstResult.get(0), secondResult.get(0)); + assertEquals("owner", secondResult.get(0).getOwnerSid()); + verify(delegate, times(1)).getUserPermissionsForObject(OBJECT_ID); + } + + @Test + void getObjectPermissionsUsesOrderIndependentPrincipalKey() { + var firstOrder = new ArrayList<>(List.of("user", "group")); + var secondOrder = new ArrayList<>(List.of("group", "user")); + when(delegate.getObjectPermissions(firstOrder)).thenReturn(Set.of(OBJECT_ID)); + + var firstResult = storage.getObjectPermissions(firstOrder); + var secondResult = storage.getObjectPermissions(secondOrder); + + assertEquals(Set.of(OBJECT_ID), firstResult); + assertEquals(Set.of(OBJECT_ID), secondResult); + verify(delegate, times(1)).getObjectPermissions(firstOrder); + } + + @Test + void updateElementClearsAllQueryCaches() { + var permission = makePermission("owner"); + var updatedPermission = makePermission("updated-owner"); + var sids = List.of("user"); + when(delegate.getElementById(PERMISSION_ID)).thenReturn(permission, updatedPermission); + when(delegate.findAll()).thenReturn(List.of(permission), List.of(updatedPermission)); + when(delegate.getUserPermissionsForObject(OBJECT_ID)) + .thenReturn(List.of(permission), List.of(updatedPermission)); + when(delegate.getObjectPermissions(sids)).thenReturn(Set.of("old-object"), Set.of("new-object")); + when(delegate.updateElement(updatedPermission)).thenReturn(updatedPermission); + + storage.getElementById(PERMISSION_ID); + storage.findAll(); + storage.getUserPermissionsForObject(OBJECT_ID); + storage.getObjectPermissions(sids); + storage.updateElement(updatedPermission); + + assertEquals("updated-owner", storage.getElementById(PERMISSION_ID).getOwnerSid()); + assertEquals("updated-owner", storage.findAll().get(0).getOwnerSid()); + assertEquals("updated-owner", storage.getUserPermissionsForObject(OBJECT_ID).get(0).getOwnerSid()); + assertEquals(Set.of("new-object"), storage.getObjectPermissions(sids)); + verify(delegate, times(2)).getElementById(PERMISSION_ID); + verify(delegate, times(2)).findAll(); + verify(delegate, times(2)).getUserPermissionsForObject(OBJECT_ID); + verify(delegate, times(2)).getObjectPermissions(sids); + } + + private Permission makePermission(String ownerSid) { + var permission = new Permission(); + permission.setPermissionId(PERMISSION_ID); + permission.setObjectInstanceId(OBJECT_ID); + permission.setOwnerSid(ownerSid); + return permission; + } +} diff --git a/streampipes-storage-api/src/main/java/org/apache/streampipes/storage/api/core/INoSqlStorage.java b/streampipes-storage-api/src/main/java/org/apache/streampipes/storage/api/core/INoSqlStorage.java index de4ef4b597..c5891d7d85 100644 --- a/streampipes-storage-api/src/main/java/org/apache/streampipes/storage/api/core/INoSqlStorage.java +++ b/streampipes-storage-api/src/main/java/org/apache/streampipes/storage/api/core/INoSqlStorage.java @@ -38,7 +38,6 @@ import org.apache.streampipes.storage.api.system.IImageStorage; import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.streampipes.storage.api.system.ITransformationScriptTemplateStorage; import org.apache.streampipes.storage.api.user.IPasswordRecoveryTokenStorage; -import org.apache.streampipes.storage.api.user.IPermissionStorage; import org.apache.streampipes.storage.api.user.IPrivilegeStorage; import org.apache.streampipes.storage.api.user.IRefreshTokenStorage; import org.apache.streampipes.storage.api.user.IRoleStorage; @@ -74,8 +73,6 @@ public interface INoSqlStorage { IPipelineElementDescriptionStorage getPipelineElementDescriptionStorage(); - IPermissionStorage getPermissionStorage(); - IDataProcessorStorage getDataProcessorStorage(); IDataSinkStorage getDataSinkStorage(); diff --git a/streampipes-storage-couchdb/src/main/java/org/apache/streampipes/storage/couchdb/CouchDbStorageManager.java b/streampipes-storage-couchdb/src/main/java/org/apache/streampipes/storage/couchdb/CouchDbStorageManager.java index 577b82b51f..e4fb7b29b4 100644 --- a/streampipes-storage-couchdb/src/main/java/org/apache/streampipes/storage/couchdb/CouchDbStorageManager.java +++ b/streampipes-storage-couchdb/src/main/java/org/apache/streampipes/storage/couchdb/CouchDbStorageManager.java @@ -40,7 +40,6 @@ import org.apache.streampipes.storage.api.system.IImageStorage; import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage; import org.apache.streampipes.storage.api.system.ITransformationScriptTemplateStorage; import org.apache.streampipes.storage.api.user.IPasswordRecoveryTokenStorage; -import org.apache.streampipes.storage.api.user.IPermissionStorage; import org.apache.streampipes.storage.api.user.IPrivilegeStorage; import org.apache.streampipes.storage.api.user.IRefreshTokenStorage; import org.apache.streampipes.storage.api.user.IRoleStorage; @@ -69,7 +68,6 @@ import org.apache.streampipes.storage.couchdb.impl.system.GenericStorageImpl; import org.apache.streampipes.storage.couchdb.impl.system.ImageStorageImpl; import org.apache.streampipes.storage.couchdb.impl.system.TransformationScriptTemplateStorageImpl; import org.apache.streampipes.storage.couchdb.impl.user.PasswordRecoveryTokenStorageImpl; -import org.apache.streampipes.storage.couchdb.impl.user.PermissionStorageImpl; import org.apache.streampipes.storage.couchdb.impl.user.PrivilegeStorageImpl; import org.apache.streampipes.storage.couchdb.impl.user.RefreshTokenStorageImpl; import org.apache.streampipes.storage.couchdb.impl.user.RoleStorageImpl; @@ -148,11 +146,6 @@ public class CouchDbStorageManager implements INoSqlStorage { return new PipelineElementDescriptionStorageImpl(); } - @Override - public IPermissionStorage getPermissionStorage() { - return new PermissionStorageImpl("users/permissions"); - } - @Override public IDataProcessorStorage getDataProcessorStorage() { return new DataProcessorStorageImpl(); diff --git a/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/model/PrincipalUserDetails.java b/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/model/PrincipalUserDetails.java index a9c943822d..7a029b2161 100644 --- a/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/model/PrincipalUserDetails.java +++ b/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/model/PrincipalUserDetails.java @@ -18,6 +18,7 @@ package org.apache.streampipes.user.management.model; import org.apache.streampipes.model.client.user.Principal; +import org.apache.streampipes.storage.api.user.IPermissionStorage; import org.apache.streampipes.user.management.util.GrantedAuthoritiesBuilder; import org.apache.streampipes.user.management.util.GrantedPermissionsBuilder; @@ -35,10 +36,11 @@ public abstract class PrincipalUserDetails<T extends Principal> implements UserD private Set<String> allAuthorities; private Set<String> allObjectPermissions; - public PrincipalUserDetails(T details) { + public PrincipalUserDetails(T details, + IPermissionStorage permissionStorage) { this.details = details; this.allAuthorities = new GrantedAuthoritiesBuilder(details).buildAllAuthorities(); - this.allObjectPermissions = new GrantedPermissionsBuilder(details).buildAllPermissions(); + this.allObjectPermissions = new GrantedPermissionsBuilder(details, permissionStorage).buildAllPermissions(); } public T getDetails() { diff --git a/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/model/ServiceAccountDetails.java b/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/model/ServiceAccountDetails.java index fee4703089..a5ae3758a2 100644 --- a/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/model/ServiceAccountDetails.java +++ b/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/model/ServiceAccountDetails.java @@ -18,12 +18,14 @@ package org.apache.streampipes.user.management.model; import org.apache.streampipes.model.client.user.ServiceAccount; +import org.apache.streampipes.storage.api.user.IPermissionStorage; public class ServiceAccountDetails extends PrincipalUserDetails<ServiceAccount> { - public ServiceAccountDetails(ServiceAccount details) { - super(details); + public ServiceAccountDetails(ServiceAccount details, + IPermissionStorage permissionStorage) { + super(details, permissionStorage); } @Override diff --git a/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/model/UserAccountDetails.java b/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/model/UserAccountDetails.java index da2405066e..08f4452098 100644 --- a/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/model/UserAccountDetails.java +++ b/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/model/UserAccountDetails.java @@ -18,11 +18,13 @@ package org.apache.streampipes.user.management.model; import org.apache.streampipes.model.client.user.UserAccount; +import org.apache.streampipes.storage.api.user.IPermissionStorage; public class UserAccountDetails extends PrincipalUserDetails<UserAccount> { - public UserAccountDetails(UserAccount details) { - super(details); + public UserAccountDetails(UserAccount details, + IPermissionStorage permissionStorage) { + super(details, permissionStorage); } @Override diff --git a/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/service/SpUserDetailsService.java b/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/service/SpUserDetailsService.java index a113185810..5fd1bb74fb 100644 --- a/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/service/SpUserDetailsService.java +++ b/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/service/SpUserDetailsService.java @@ -20,6 +20,7 @@ package org.apache.streampipes.user.management.service; 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.user.IPermissionStorage; import org.apache.streampipes.storage.management.StorageDispatcher; import org.apache.streampipes.user.management.model.ServiceAccountDetails; import org.apache.streampipes.user.management.model.UserAccountDetails; @@ -30,10 +31,16 @@ import org.springframework.security.core.userdetails.UsernameNotFoundException; public class SpUserDetailsService implements UserDetailsService { + private final IPermissionStorage permissionStorage; + + public SpUserDetailsService(IPermissionStorage permissionStorage) { + this.permissionStorage = permissionStorage; + } + @Override public UserDetails loadUserByUsername(String s) throws UsernameNotFoundException { Principal user = StorageDispatcher.INSTANCE.getNoSqlStore().getUserStorageAPI().getUser(s); - return user instanceof UserAccount ? new UserAccountDetails((UserAccount) user) : - new ServiceAccountDetails((ServiceAccount) user); + return user instanceof UserAccount ? new UserAccountDetails((UserAccount) user, permissionStorage) : + new ServiceAccountDetails((ServiceAccount) user, permissionStorage); } } diff --git a/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/util/GrantedPermissionsBuilder.java b/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/util/GrantedPermissionsBuilder.java index 1674849767..90db32e7dd 100644 --- a/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/util/GrantedPermissionsBuilder.java +++ b/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/util/GrantedPermissionsBuilder.java @@ -18,7 +18,7 @@ package org.apache.streampipes.user.management.util; import org.apache.streampipes.model.client.user.Principal; -import org.apache.streampipes.storage.management.StorageDispatcher; +import org.apache.streampipes.storage.api.user.IPermissionStorage; import java.util.ArrayList; import java.util.Set; @@ -26,18 +26,18 @@ import java.util.Set; public class GrantedPermissionsBuilder { private final Principal principal; + private final IPermissionStorage permissionStorage; - public GrantedPermissionsBuilder(Principal principal) { + public GrantedPermissionsBuilder(Principal principal, + IPermissionStorage permissionStorage) { this.principal = principal; + this.permissionStorage = permissionStorage; } public Set<String> buildAllPermissions() { Set<String> sids = extractSids(); - return StorageDispatcher - .INSTANCE - .getNoSqlStore() - .getPermissionStorage() + return permissionStorage .getObjectPermissions(new ArrayList<>(sids)); }
