This is an automated email from the ASF dual-hosted git repository.
SvenO3 pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/streampipes.git
The following commit(s) were added to refs/heads/dev by this push:
new 90c926c53f refactor: Improve rest endpoints for different resources
(#4679)
90c926c53f is described below
commit 90c926c53f6a0fca6f5edc30b0a4ed12a7a58abe
Author: Sven Oehler <[email protected]>
AuthorDate: Wed Jul 8 10:54:13 2026 +0200
refactor: Improve rest endpoints for different resources (#4679)
---
.../apache/streampipes/commons/constants/Envs.java | 1 +
.../commons/environment/DefaultEnvironment.java | 5 ++
.../commons/environment/Environment.java | 2 +
.../streampipes/rest/impl/EmailResource.java | 7 +++
.../FunctionStateEndpointsEnabledCondition.java | 34 +++++++++++
...onsResource.java => FunctionStateResource.java} | 49 ++--------------
.../streampipes/rest/impl/FunctionsResource.java | 61 +++----------------
.../rest/impl/PipelineCanvasMetadataResource.java | 68 +++++++++-------------
.../streampipes/rest/impl/PipelineTemplate.java | 13 +++++
.../impl/pe/PipelineElementTemplateResource.java | 27 ++++++++-
.../lib/apis/pipeline-canvas-metadata.service.ts | 12 +---
.../save-pipeline/save-pipeline.component.ts | 17 ++----
.../functions-logs/functions-logs.component.html | 4 +-
.../functions-logs/functions-logs.component.ts | 2 +
.../functions-metrics.component.html | 22 ++++---
.../functions-metrics.component.ts | 2 +
ui/src/app/pipelines/pipelines.component.html | 28 +++++----
ui/src/app/pipelines/pipelines.component.ts | 10 ++--
18 files changed, 171 insertions(+), 193 deletions(-)
diff --git
a/streampipes-commons/src/main/java/org/apache/streampipes/commons/constants/Envs.java
b/streampipes-commons/src/main/java/org/apache/streampipes/commons/constants/Envs.java
index d73cdfbda0..ec5ca22205 100644
---
a/streampipes-commons/src/main/java/org/apache/streampipes/commons/constants/Envs.java
+++
b/streampipes-commons/src/main/java/org/apache/streampipes/commons/constants/Envs.java
@@ -48,6 +48,7 @@ public enum Envs {
SP_OAUTH_ENABLED("SP_OAUTH_ENABLED", "false"),
SP_OAUTH_REDIRECT_URI("SP_OAUTH_REDIRECT_URI"),
SP_RESET_ENDPOINT_ENABLED("SP_RESET_ENDPOINT_ENABLED", "false"),
+ SP_FUNCTION_STATE_ENDPOINTS_ENABLED("SP_FUNCTION_STATE_ENDPOINTS_ENABLED",
"false"),
SP_DEBUG("SP_DEBUG", "false"),
SP_MAX_WAIT_TIME_AT_SHUTDOWN("SP_MAX_WAIT_TIME_AT_SHUTDOWN"),
diff --git
a/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/DefaultEnvironment.java
b/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/DefaultEnvironment.java
index 09f3877eb7..9fbdf39100 100644
---
a/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/DefaultEnvironment.java
+++
b/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/DefaultEnvironment.java
@@ -197,6 +197,11 @@ public class DefaultEnvironment implements Environment {
return new BooleanEnvironmentVariable(Envs.SP_RESET_ENDPOINT_ENABLED);
}
+ @Override
+ public BooleanEnvironmentVariable getFunctionStateEndpointsEnabled() {
+ return new
BooleanEnvironmentVariable(Envs.SP_FUNCTION_STATE_ENDPOINTS_ENABLED);
+ }
+
@Override
public List<OAuthConfiguration> getOAuthConfigurations() {
return new OAuthConfigurationParser().parse(System.getenv());
diff --git
a/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/Environment.java
b/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/Environment.java
index 2460a2d10b..94b9757739 100644
---
a/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/Environment.java
+++
b/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/Environment.java
@@ -104,6 +104,8 @@ public interface Environment {
BooleanEnvironmentVariable getResetEndpointEnabled();
+ BooleanEnvironmentVariable getFunctionStateEndpointsEnabled();
+
List<OAuthConfiguration> getOAuthConfigurations();
// Messaging
diff --git
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/EmailResource.java
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/EmailResource.java
index 78dec7ee41..c443613ce0 100644
---
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/EmailResource.java
+++
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/EmailResource.java
@@ -18,12 +18,14 @@
package org.apache.streampipes.rest.impl;
import org.apache.streampipes.mail.MailSender;
+import org.apache.streampipes.model.client.user.DefaultPrivilege;
import org.apache.streampipes.model.mail.SpEmail;
import
org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource;
import org.apache.streampipes.storage.api.system.ISpCoreConfigurationStorage;
import org.springframework.http.MediaType;
import org.springframework.http.ResponseEntity;
+import org.springframework.security.access.prepost.PreAuthorize;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
@@ -40,6 +42,7 @@ public class EmailResource extends
AbstractAuthGuardedRestResource {
}
@PostMapping(consumes = MediaType.APPLICATION_JSON_VALUE)
+ @PreAuthorize("this.hasWriteAuthority()")
public ResponseEntity<?> sendEmail(@RequestBody SpEmail email) {
var configuration = configurationStorage.get();
if (configuration.getEmailConfig().isEmailConfigured()) {
@@ -54,4 +57,8 @@ public class EmailResource extends
AbstractAuthGuardedRestResource {
"Could not send email - no valid mail configuration provided in
StreamPipes (go to settings -> mail)");
}
}
+
+ public boolean hasWriteAuthority() {
+ return
isAdminOrHasAnyAuthority(DefaultPrivilege.Constants.PRIVILEGE_WRITE_PIPELINE_VALUE);
+ }
}
diff --git
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/FunctionStateEndpointsEnabledCondition.java
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/FunctionStateEndpointsEnabledCondition.java
new file mode 100644
index 0000000000..f56fef23bd
--- /dev/null
+++
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/FunctionStateEndpointsEnabledCondition.java
@@ -0,0 +1,34 @@
+/*
+ * 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.rest.impl;
+
+import org.apache.streampipes.commons.environment.Environments;
+
+import org.springframework.context.annotation.Condition;
+import org.springframework.context.annotation.ConditionContext;
+import org.springframework.core.type.AnnotatedTypeMetadata;
+
+public class FunctionStateEndpointsEnabledCondition implements Condition {
+
+ @Override
+ public boolean matches(ConditionContext context,
+ AnnotatedTypeMetadata metadata) {
+ return
Environments.getEnvironment().getFunctionStateEndpointsEnabled().getValueOrDefault();
+ }
+}
diff --git
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/FunctionsResource.java
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/FunctionStateResource.java
similarity index 62%
copy from
streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/FunctionsResource.java
copy to
streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/FunctionStateResource.java
index a2dcc7d8aa..933b591731 100644
---
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/FunctionsResource.java
+++
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/FunctionStateResource.java
@@ -18,79 +18,38 @@
package org.apache.streampipes.rest.impl;
-import org.apache.streampipes.loadbalance.pipeline.ExtensionsLogProvider;
-import org.apache.streampipes.manager.function.FunctionRegistrationService;
-import org.apache.streampipes.model.function.FunctionDefinition;
import org.apache.streampipes.model.function.FunctionState;
import org.apache.streampipes.model.message.Notifications;
import org.apache.streampipes.model.message.SuccessMessage;
-import org.apache.streampipes.model.monitoring.SpLogEntry;
-import org.apache.streampipes.model.monitoring.SpMetricsEntry;
import
org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource;
import org.apache.streampipes.rest.shared.exception.SpMessageException;
import org.apache.streampipes.storage.api.function.IFunctionStateStorage;
+import org.springframework.context.annotation.Conditional;
import org.springframework.http.HttpStatus;
import org.springframework.http.MediaType;
import org.springframework.http.ResponseEntity;
import org.springframework.web.bind.annotation.DeleteMapping;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
-import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.PutMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
-import java.util.Collection;
-import java.util.List;
import java.util.Map;
@RestController
@RequestMapping("/api/v2/functions")
-public class FunctionsResource extends AbstractAuthGuardedRestResource {
+@Conditional(FunctionStateEndpointsEnabledCondition.class)
+public class FunctionStateResource extends AbstractAuthGuardedRestResource {
private final IFunctionStateStorage functionStateStorage;
- public FunctionsResource(IFunctionStateStorage functionStateStorage) {
+ public FunctionStateResource(IFunctionStateStorage functionStateStorage) {
this.functionStateStorage = functionStateStorage;
}
- @GetMapping(produces = MediaType.APPLICATION_JSON_VALUE)
- public ResponseEntity<Collection<FunctionDefinition>> getActiveFunctions() {
- return ok(FunctionRegistrationService.INSTANCE.getAllFunctions());
- }
-
- @GetMapping(path = "{functionId}", produces =
MediaType.APPLICATION_JSON_VALUE)
- public ResponseEntity<FunctionDefinition>
getFunction(@PathVariable("functionId") String functionId) {
- return ok(FunctionRegistrationService.INSTANCE.getFunction(functionId));
- }
-
- @PostMapping(
- produces = MediaType.APPLICATION_JSON_VALUE,
- consumes = MediaType.APPLICATION_JSON_VALUE
- )
- public ResponseEntity<SuccessMessage> registerFunctions(@RequestBody
List<FunctionDefinition> functions) {
- functions.forEach(FunctionRegistrationService.INSTANCE::registerFunction);
- return ok(Notifications.success("Function successfully registered"));
- }
-
- @DeleteMapping(path = "{functionId}", produces =
MediaType.APPLICATION_JSON_VALUE)
- public ResponseEntity<SuccessMessage>
deregisterFunction(@PathVariable("functionId") String functionId) {
- FunctionRegistrationService.INSTANCE.deregisterFunction(functionId);
- return ok(Notifications.success("Function successfully deregistered"));
- }
-
- @GetMapping(path = "{functionId}/metrics", produces =
MediaType.APPLICATION_JSON_VALUE)
- public ResponseEntity<SpMetricsEntry>
getFunctionMetrics(@PathVariable("functionId") String functionId) {
- return
ok(ExtensionsLogProvider.INSTANCE.getMetricInfosForResource(functionId));
- }
-
- @GetMapping(path = "{functionId}/logs")
- public ResponseEntity<List<SpLogEntry>>
getFunctionLogs(@PathVariable("functionId") String functionId) {
- return
ok(ExtensionsLogProvider.INSTANCE.getLogInfosForResource(functionId));
- }
-
@GetMapping(path = "{functionId}/state", produces =
MediaType.APPLICATION_JSON_VALUE)
public ResponseEntity<Map<String, Object>>
getFunctionState(@PathVariable("functionId") String functionId) {
var functionState = functionStateStorage.getElementById(functionId);
diff --git
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/FunctionsResource.java
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/FunctionsResource.java
index a2dcc7d8aa..8f4dd48cca 100644
---
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/FunctionsResource.java
+++
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/FunctionsResource.java
@@ -21,47 +21,39 @@ package org.apache.streampipes.rest.impl;
import org.apache.streampipes.loadbalance.pipeline.ExtensionsLogProvider;
import org.apache.streampipes.manager.function.FunctionRegistrationService;
import org.apache.streampipes.model.function.FunctionDefinition;
-import org.apache.streampipes.model.function.FunctionState;
import org.apache.streampipes.model.message.Notifications;
import org.apache.streampipes.model.message.SuccessMessage;
import org.apache.streampipes.model.monitoring.SpLogEntry;
import org.apache.streampipes.model.monitoring.SpMetricsEntry;
import
org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource;
-import org.apache.streampipes.rest.shared.exception.SpMessageException;
-import org.apache.streampipes.storage.api.function.IFunctionStateStorage;
+import org.apache.streampipes.rest.security.AuthConstants;
-import org.springframework.http.HttpStatus;
import org.springframework.http.MediaType;
import org.springframework.http.ResponseEntity;
+import org.springframework.security.access.prepost.PreAuthorize;
import org.springframework.web.bind.annotation.DeleteMapping;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.PostMapping;
-import org.springframework.web.bind.annotation.PutMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import java.util.Collection;
import java.util.List;
-import java.util.Map;
@RestController
@RequestMapping("/api/v2/functions")
public class FunctionsResource extends AbstractAuthGuardedRestResource {
- private final IFunctionStateStorage functionStateStorage;
-
- public FunctionsResource(IFunctionStateStorage functionStateStorage) {
- this.functionStateStorage = functionStateStorage;
- }
-
@GetMapping(produces = MediaType.APPLICATION_JSON_VALUE)
+ @PreAuthorize(AuthConstants.IS_ADMIN_ROLE)
public ResponseEntity<Collection<FunctionDefinition>> getActiveFunctions() {
return ok(FunctionRegistrationService.INSTANCE.getAllFunctions());
}
@GetMapping(path = "{functionId}", produces =
MediaType.APPLICATION_JSON_VALUE)
+ @PreAuthorize(AuthConstants.IS_ADMIN_ROLE)
public ResponseEntity<FunctionDefinition>
getFunction(@PathVariable("functionId") String functionId) {
return ok(FunctionRegistrationService.INSTANCE.getFunction(functionId));
}
@@ -70,65 +62,28 @@ public class FunctionsResource extends
AbstractAuthGuardedRestResource {
produces = MediaType.APPLICATION_JSON_VALUE,
consumes = MediaType.APPLICATION_JSON_VALUE
)
+ @PreAuthorize(AuthConstants.IS_ADMIN_ROLE)
public ResponseEntity<SuccessMessage> registerFunctions(@RequestBody
List<FunctionDefinition> functions) {
functions.forEach(FunctionRegistrationService.INSTANCE::registerFunction);
return ok(Notifications.success("Function successfully registered"));
}
@DeleteMapping(path = "{functionId}", produces =
MediaType.APPLICATION_JSON_VALUE)
+ @PreAuthorize(AuthConstants.IS_ADMIN_ROLE)
public ResponseEntity<SuccessMessage>
deregisterFunction(@PathVariable("functionId") String functionId) {
FunctionRegistrationService.INSTANCE.deregisterFunction(functionId);
return ok(Notifications.success("Function successfully deregistered"));
}
@GetMapping(path = "{functionId}/metrics", produces =
MediaType.APPLICATION_JSON_VALUE)
+ @PreAuthorize(AuthConstants.IS_ADMIN_ROLE)
public ResponseEntity<SpMetricsEntry>
getFunctionMetrics(@PathVariable("functionId") String functionId) {
return
ok(ExtensionsLogProvider.INSTANCE.getMetricInfosForResource(functionId));
}
@GetMapping(path = "{functionId}/logs")
+ @PreAuthorize(AuthConstants.IS_ADMIN_ROLE)
public ResponseEntity<List<SpLogEntry>>
getFunctionLogs(@PathVariable("functionId") String functionId) {
return
ok(ExtensionsLogProvider.INSTANCE.getLogInfosForResource(functionId));
}
-
- @GetMapping(path = "{functionId}/state", produces =
MediaType.APPLICATION_JSON_VALUE)
- public ResponseEntity<Map<String, Object>>
getFunctionState(@PathVariable("functionId") String functionId) {
- var functionState = functionStateStorage.getElementById(functionId);
- if (functionState != null) {
- return ok(functionState.getState());
- } else {
- throw new SpMessageException(HttpStatus.NOT_FOUND,
Notifications.error("Function state not found"));
- }
- }
-
- @PutMapping(
- path = "{functionId}/state",
- produces = MediaType.APPLICATION_JSON_VALUE,
- consumes = MediaType.APPLICATION_JSON_VALUE
- )
- public ResponseEntity<SuccessMessage>
persistFunctionState(@PathVariable("functionId") String functionId,
- @RequestBody
Map<String, Object> state) {
- var existingFunctionState =
functionStateStorage.getElementById(functionId);
-
- if (existingFunctionState != null) {
- existingFunctionState.setState(state);
- functionStateStorage.updateElement(existingFunctionState);
- } else {
- functionStateStorage.persist(new FunctionState(functionId, state));
- }
-
- return ok(Notifications.success("Function state successfully persisted"));
- }
-
- @DeleteMapping(path = "{functionId}/state", produces =
MediaType.APPLICATION_JSON_VALUE)
- public ResponseEntity<SuccessMessage>
deleteFunctionState(@PathVariable("functionId") String functionId) {
- var existingFunctionState =
functionStateStorage.getElementById(functionId);
-
- if (existingFunctionState == null) {
- throw new SpMessageException(HttpStatus.NOT_FOUND,
Notifications.error("Function state not found"));
- }
-
- functionStateStorage.deleteElementById(functionId);
- return ok(Notifications.success("Function state successfully deleted"));
- }
}
diff --git
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/PipelineCanvasMetadataResource.java
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/PipelineCanvasMetadataResource.java
index f081e4b579..ee86ad9d14 100644
---
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/PipelineCanvasMetadataResource.java
+++
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/PipelineCanvasMetadataResource.java
@@ -18,18 +18,19 @@
package org.apache.streampipes.rest.impl;
import org.apache.streampipes.model.canvas.PipelineCanvasMetadata;
+import org.apache.streampipes.model.client.user.DefaultPrivilege;
import org.apache.streampipes.model.message.Notifications;
-import org.apache.streampipes.rest.core.base.impl.AbstractRestResource;
+import
org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource;
import org.apache.streampipes.rest.shared.exception.SpMessageException;
import
org.apache.streampipes.storage.api.pipeline.IPipelineCanvasMetadataStorage;
import org.springframework.http.HttpStatus;
import org.springframework.http.MediaType;
import org.springframework.http.ResponseEntity;
+import org.springframework.security.access.prepost.PreAuthorize;
import org.springframework.web.bind.annotation.DeleteMapping;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
-import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.PutMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
@@ -37,9 +38,10 @@ import
org.springframework.web.bind.annotation.RestController;
@RestController
@RequestMapping("/api/v2/pipeline-canvas-metadata")
-public class PipelineCanvasMetadataResource extends AbstractRestResource {
+public class PipelineCanvasMetadataResource extends
AbstractAuthGuardedRestResource {
@GetMapping(path = "/pipeline/{pipelineId}", produces =
MediaType.APPLICATION_JSON_VALUE)
+ @PreAuthorize("this.hasReadAuthority() and hasPermission(#pipelineId,
'READ')")
public ResponseEntity<PipelineCanvasMetadata>
getPipelineCanvasMetadataForPipeline(
@PathVariable("pipelineId") String pipelineId) {
try {
@@ -50,65 +52,49 @@ public class PipelineCanvasMetadataResource extends
AbstractRestResource {
}
}
- @GetMapping(path = "{canvasId}", produces = MediaType.APPLICATION_JSON_VALUE)
- public ResponseEntity<PipelineCanvasMetadata> getPipelineCanvasMetadata(
- @PathVariable("canvasId") String pipelineCanvasId) {
- try {
- return ok(getPipelineCanvasMetadataStorage()
- .getElementById(pipelineCanvasId));
- } catch (IllegalArgumentException e) {
- throw new SpMessageException(HttpStatus.BAD_REQUEST,
Notifications.error(e.getMessage()));
- }
- }
-
- @PostMapping(
- consumes = MediaType.APPLICATION_JSON_VALUE,
- produces = MediaType.APPLICATION_JSON_VALUE)
- public ResponseEntity<Void> storePipelineCanvasMetadata(@RequestBody
PipelineCanvasMetadata pipelineCanvasMetadata) {
- getPipelineCanvasMetadataStorage().persist(pipelineCanvasMetadata);
- return ok();
- }
-
- @DeleteMapping(
- path = "{canvasId}",
- produces = MediaType.APPLICATION_JSON_VALUE)
- public ResponseEntity<Void>
deletePipelineCanvasMetadata(@PathVariable("canvasId") String pipelineCanvasId)
{
- PipelineCanvasMetadata metadata = find(pipelineCanvasId);
- getPipelineCanvasMetadataStorage().deleteElement(metadata);
- return ok();
- }
-
@DeleteMapping(
path = "/pipeline/{pipelineId}",
produces = MediaType.APPLICATION_JSON_VALUE)
+ @PreAuthorize("this.hasWriteAuthority() and hasPermission(#pipelineId,
'WRITE')")
public ResponseEntity<Void>
deletePipelineCanvasMetadataForPipeline(@PathVariable("pipelineId") String
pipelineId) {
PipelineCanvasMetadata metadata =
getPipelineCanvasMetadataStorage().getPipelineCanvasMetadataForPipeline(pipelineId);
- getPipelineCanvasMetadataStorage().deleteElement(metadata);
+ if (metadata != null) {
+ getPipelineCanvasMetadataStorage().deleteElement(metadata);
+ }
return ok();
}
@PutMapping(
- path = "{canvasId}",
+ path = "/pipeline/{pipelineId}",
consumes = MediaType.APPLICATION_JSON_VALUE,
produces = MediaType.APPLICATION_JSON_VALUE)
- public ResponseEntity<Void>
updatePipelineCanvasMetadata(@PathVariable("canvasId") String pipelineCanvasId,
+ @PreAuthorize("this.hasWriteAuthority() and hasPermission(#pipelineId,
'WRITE')")
+ public ResponseEntity<Void>
updatePipelineCanvasMetadata(@PathVariable("pipelineId") String pipelineId,
@RequestBody
PipelineCanvasMetadata pipelineCanvasMetadata) {
- try {
- var existing =
getPipelineCanvasMetadataStorage().getElementById(pipelineCanvasMetadata.getId());
+ var existing =
getPipelineCanvasMetadataStorage().getPipelineCanvasMetadataForPipeline(pipelineId);
+ pipelineCanvasMetadata.setPipelineId(pipelineId);
+ if (existing != null) {
+ pipelineCanvasMetadata.setId(existing.getId());
pipelineCanvasMetadata.setRev(existing.getRev());
getPipelineCanvasMetadataStorage().updateElement(pipelineCanvasMetadata);
- } catch (IllegalArgumentException e) {
+ } else {
+ pipelineCanvasMetadata.setId(null);
+ pipelineCanvasMetadata.setRev(null);
getPipelineCanvasMetadataStorage().persist(pipelineCanvasMetadata);
}
return ok();
}
- private PipelineCanvasMetadata find(String canvasId) {
- return getPipelineCanvasMetadataStorage().getElementById(canvasId);
- }
-
private IPipelineCanvasMetadataStorage getPipelineCanvasMetadataStorage() {
return getNoSqlStorage().getPipelineCanvasMetadataStorage();
}
+
+ public boolean hasWriteAuthority() {
+ return
isAdminOrHasAnyAuthority(DefaultPrivilege.Constants.PRIVILEGE_WRITE_PIPELINE_VALUE);
+ }
+
+ public boolean hasReadAuthority() {
+ return
isAdminOrHasAnyAuthority(DefaultPrivilege.Constants.PRIVILEGE_READ_PIPELINE_VALUE);
+ }
}
diff --git
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/PipelineTemplate.java
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/PipelineTemplate.java
index d35ec77377..28408cf9a8 100644
---
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/PipelineTemplate.java
+++
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/PipelineTemplate.java
@@ -20,6 +20,7 @@ package org.apache.streampipes.rest.impl;
import
org.apache.streampipes.manager.api.extensions.ExtensionServiceRequestManager;
import
org.apache.streampipes.manager.template.compact.CompactPipelineTemplateManagement;
+import org.apache.streampipes.model.client.user.DefaultPrivilege;
import org.apache.streampipes.model.template.CompactPipelineTemplate;
import org.apache.streampipes.model.template.PipelineTemplateGenerationRequest;
import
org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource;
@@ -30,6 +31,7 @@ import
org.apache.streampipes.storage.management.StorageDispatcher;
import org.springframework.http.MediaType;
import org.springframework.http.ResponseEntity;
+import org.springframework.security.access.prepost.PreAuthorize;
import org.springframework.web.bind.annotation.DeleteMapping;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
@@ -61,6 +63,7 @@ public class PipelineTemplate extends
AbstractAuthGuardedRestResource {
@GetMapping(
produces = {MediaType.APPLICATION_JSON_VALUE, SpMediaType.YAML,
SpMediaType.YML})
+ @PreAuthorize("this.hasWriteAuthority()")
public List<CompactPipelineTemplate> findAll() {
return storage.findAll()
.stream()
@@ -71,6 +74,7 @@ public class PipelineTemplate extends
AbstractAuthGuardedRestResource {
@GetMapping(
path = "/{id}",
produces = {MediaType.APPLICATION_JSON_VALUE, SpMediaType.YAML,
SpMediaType.YML})
+ @PreAuthorize("this.hasWriteAuthority()")
public ResponseEntity<?> findById(@PathVariable("id") String id) {
return ok(storage.getElementById(id));
}
@@ -79,6 +83,7 @@ public class PipelineTemplate extends
AbstractAuthGuardedRestResource {
@PostMapping(
produces = {MediaType.APPLICATION_JSON_VALUE, SpMediaType.YAML,
SpMediaType.YML},
consumes = {MediaType.APPLICATION_JSON_VALUE, SpMediaType.YAML,
SpMediaType.YML})
+ @PreAuthorize("this.hasWriteAuthority()")
public void create(@RequestBody CompactPipelineTemplate entity) {
storage.persist(entity);
}
@@ -86,11 +91,13 @@ public class PipelineTemplate extends
AbstractAuthGuardedRestResource {
@PutMapping(path = "/{id}",
produces = {MediaType.APPLICATION_JSON_VALUE, SpMediaType.YAML,
SpMediaType.YML},
consumes = {MediaType.APPLICATION_JSON_VALUE, SpMediaType.YAML,
SpMediaType.YML})
+ @PreAuthorize("this.hasWriteAuthority()")
public void update(@PathVariable("id") String id, @RequestBody
CompactPipelineTemplate entity) {
storage.updateElement(entity);
}
@DeleteMapping(path = "/{id}")
+ @PreAuthorize("this.hasWriteAuthority()")
public void delete(@PathVariable("id") String id) {
storage.deleteElementById(id);
}
@@ -99,6 +106,7 @@ public class PipelineTemplate extends
AbstractAuthGuardedRestResource {
@PostMapping(path = "/{id}/pipeline",
produces = {MediaType.APPLICATION_JSON_VALUE, SpMediaType.YAML,
SpMediaType.YML},
consumes = {MediaType.APPLICATION_JSON_VALUE, SpMediaType.YAML,
SpMediaType.YML})
+ @PreAuthorize("this.hasWriteAuthority()")
public ResponseEntity<?> makePipelineFromTemplate(@RequestBody
PipelineTemplateGenerationRequest request) {
try {
return ok(templateManagement.makePipeline(request).pipeline());
@@ -112,6 +120,7 @@ public class PipelineTemplate extends
AbstractAuthGuardedRestResource {
@GetMapping(
path = "/{id}/streams",
produces = {MediaType.APPLICATION_JSON_VALUE, SpMediaType.YAML,
SpMediaType.YML})
+ @PreAuthorize("this.hasWriteAuthority()")
public ResponseEntity<Map<String, List<List<String>>>>
getAvailableStreamsForTemplate(
@PathVariable("id") String pipelineTemplateId) {
try {
@@ -122,4 +131,8 @@ public class PipelineTemplate extends
AbstractAuthGuardedRestResource {
throw new RuntimeException(e);
}
}
+
+ public boolean hasWriteAuthority() {
+ return
isAdminOrHasAnyAuthority(DefaultPrivilege.Constants.PRIVILEGE_WRITE_PIPELINE_VALUE);
+ }
}
diff --git
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/pe/PipelineElementTemplateResource.java
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/pe/PipelineElementTemplateResource.java
index de59ef28a5..33eaa6e154 100644
---
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/pe/PipelineElementTemplateResource.java
+++
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/pe/PipelineElementTemplateResource.java
@@ -21,11 +21,12 @@ package org.apache.streampipes.rest.impl.pe;
import org.apache.streampipes.manager.template.AdapterTemplateHandler;
import org.apache.streampipes.manager.template.DataProcessorTemplateHandler;
import org.apache.streampipes.manager.template.DataSinkTemplateHandler;
+import org.apache.streampipes.model.client.user.DefaultPrivilege;
import org.apache.streampipes.model.connect.adapter.AdapterDescription;
import org.apache.streampipes.model.graph.DataProcessorInvocation;
import org.apache.streampipes.model.graph.DataSinkInvocation;
import org.apache.streampipes.model.template.PipelineElementTemplate;
-import org.apache.streampipes.rest.core.base.impl.AbstractRestResource;
+import
org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource;
import io.swagger.v3.oas.annotations.Operation;
import io.swagger.v3.oas.annotations.Parameter;
@@ -36,6 +37,7 @@ import io.swagger.v3.oas.annotations.parameters.RequestBody;
import io.swagger.v3.oas.annotations.responses.ApiResponse;
import org.springframework.http.MediaType;
import org.springframework.http.ResponseEntity;
+import org.springframework.security.access.prepost.PreAuthorize;
import org.springframework.web.bind.annotation.DeleteMapping;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
@@ -49,7 +51,7 @@ import java.util.List;
@RestController
@RequestMapping("/api/v2/pipeline-element-templates")
-public class PipelineElementTemplateResource extends AbstractRestResource {
+public class PipelineElementTemplateResource extends
AbstractAuthGuardedRestResource {
@GetMapping(produces = MediaType.APPLICATION_JSON_VALUE)
@Operation(summary = "Get a list of all pipeline element templates",
@@ -61,6 +63,7 @@ public class PipelineElementTemplateResource extends
AbstractRestResource {
array = @ArraySchema(schema = @Schema(implementation
= PipelineElementTemplate.class)))
})
})
+ @PreAuthorize("this.hasWriteAuthority()")
public ResponseEntity<List<PipelineElementTemplate>> getAll(
@Parameter(description = "Filter all templates by this appId")
@RequestParam("appId") String appId
@@ -83,6 +86,7 @@ public class PipelineElementTemplateResource extends
AbstractRestResource {
}),
@ApiResponse(responseCode = "400", description = "Template
with given id not found")
})
+ @PreAuthorize("this.hasWriteAuthority()")
public ResponseEntity<?> getById(
@Parameter(description = "The id of the pipeline element template",
required = true)
@PathVariable("id") String s
@@ -100,6 +104,7 @@ public class PipelineElementTemplateResource extends
AbstractRestResource {
responses = {
@ApiResponse(responseCode = "200", description = "Template
successfully stored")
})
+ @PreAuthorize("this.hasWriteAuthority()")
public ResponseEntity<Void> create(
@RequestBody(description = "The pipeline element template to be stored",
content = @Content(schema = @Schema(implementation =
PipelineElementTemplate.class)))
@@ -120,6 +125,7 @@ public class PipelineElementTemplateResource extends
AbstractRestResource {
}, responseCode = "200", description = "Template successfully
updated"),
@ApiResponse(responseCode = "400", description = "Template
with given id not found")
})
+ @PreAuthorize("this.hasWriteAuthority()")
public ResponseEntity<?> update(
@Parameter(description = "The id of the pipeline element template",
required = true)
@PathVariable("id") String id,
@@ -165,6 +171,7 @@ public class PipelineElementTemplateResource extends
AbstractRestResource {
schema = @Schema(implementation =
DataSinkInvocation.class))
}, responseCode = "200", description = "The configured data
sink invocation model"),
})
+ @PreAuthorize("this.hasWriteAuthorityPipeline()")
public ResponseEntity<DataSinkInvocation> getPipelineElementForTemplate(
@Parameter(description = "The id of the pipeline element template",
required = true)
@PathVariable("id") String id,
@@ -196,6 +203,7 @@ public class PipelineElementTemplateResource extends
AbstractRestResource {
schema = @Schema(implementation =
DataProcessorInvocation.class))
}, responseCode = "200", description = "The configured data
processor invocation model"),
})
+ @PreAuthorize("this.hasWriteAuthorityPipeline()")
public ResponseEntity<DataProcessorInvocation> getPipelineElementForTemplate(
@Parameter(description = "The id of the pipeline element template",
required = true)
@PathVariable("id") String id,
@@ -227,6 +235,7 @@ public class PipelineElementTemplateResource extends
AbstractRestResource {
schema = @Schema(implementation =
AdapterDescription.class))
}, responseCode = "200", description = "The configured
adapter model"),
})
+ @PreAuthorize("this.hasWriteAuthorityAdapter()")
public ResponseEntity<AdapterDescription> getPipelineElementForTemplate(
@Parameter(description = "The id of the pipeline element template",
required = true)
@PathVariable("id") String id,
@@ -246,4 +255,18 @@ public class PipelineElementTemplateResource extends
AbstractRestResource {
.applyTemplateOnPipelineElement();
return ok(desc);
}
+
+ public boolean hasWriteAuthority() {
+ return isAdminOrHasAnyAuthority(
+ DefaultPrivilege.Constants.PRIVILEGE_WRITE_PIPELINE_VALUE,
+ DefaultPrivilege.Constants.PRIVILEGE_WRITE_ADAPTER_VALUE);
+ }
+
+ public boolean hasWriteAuthorityPipeline() {
+ return
isAdminOrHasAnyAuthority(DefaultPrivilege.Constants.PRIVILEGE_WRITE_PIPELINE_VALUE);
+ }
+
+ public boolean hasWriteAuthorityAdapter() {
+ return
isAdminOrHasAnyAuthority(DefaultPrivilege.Constants.PRIVILEGE_WRITE_ADAPTER_VALUE);
+ }
}
diff --git
a/ui/projects/streampipes/platform-services/src/lib/apis/pipeline-canvas-metadata.service.ts
b/ui/projects/streampipes/platform-services/src/lib/apis/pipeline-canvas-metadata.service.ts
index 4babb0052b..141b8e604f 100644
---
a/ui/projects/streampipes/platform-services/src/lib/apis/pipeline-canvas-metadata.service.ts
+++
b/ui/projects/streampipes/platform-services/src/lib/apis/pipeline-canvas-metadata.service.ts
@@ -30,13 +30,6 @@ export class PipelineCanvasMetadataService {
private http = inject(HttpClient);
private platformServicesCommons = inject(PlatformServicesCommons);
- addPipelineCanvasMetadata(pipelineCanvasMetadata: PipelineCanvasMetadata) {
- return this.http.post(
- this.pipelineCanvasMetadataBasePath,
- pipelineCanvasMetadata,
- );
- }
-
getPipelineCanvasMetadata(
pipelineId: string,
): Observable<PipelineCanvasMetadata> {
@@ -50,12 +43,11 @@ export class PipelineCanvasMetadataService {
}
updatePipelineCanvasMetadata(
+ pipelineId: string,
pipelineCanvasMetadata: PipelineCanvasMetadata,
) {
return this.http.put(
- this.pipelineCanvasMetadataBasePath +
- '/' +
- pipelineCanvasMetadata.pipelineId,
+ this.pipelineCanvasMetadataPipelinePath + pipelineId,
pipelineCanvasMetadata,
);
}
diff --git a/ui/src/app/editor/dialog/save-pipeline/save-pipeline.component.ts
b/ui/src/app/editor/dialog/save-pipeline/save-pipeline.component.ts
index 0853a11381..1cb1fa125c 100644
--- a/ui/src/app/editor/dialog/save-pipeline/save-pipeline.component.ts
+++ b/ui/src/app/editor/dialog/save-pipeline/save-pipeline.component.ts
@@ -333,20 +333,11 @@ export class SavePipelineComponent implements OnInit {
}
getPipelineCanvasMetadata$(pipelineId: string): Observable<object> {
- let request;
this.pipelineCanvasMetadata.pipelineId = pipelineId;
- if (this.storageOptions.updateModeActive) {
- request = this.pipelineCanvasService.updatePipelineCanvasMetadata(
- this.pipelineCanvasMetadata,
- );
- } else {
- this.pipelineCanvasMetadata._id = undefined;
- this.pipelineCanvasMetadata._rev = undefined;
- request = this.pipelineCanvasService.addPipelineCanvasMetadata(
- this.pipelineCanvasMetadata,
- );
- }
- return request;
+ return this.pipelineCanvasService.updatePipelineCanvasMetadata(
+ pipelineId,
+ this.pipelineCanvasMetadata,
+ );
}
addStatusIndicator(message: string, status: Status) {
diff --git
a/ui/src/app/pipelines/components/functions-overview/functions-logs/functions-logs.component.html
b/ui/src/app/pipelines/components/functions-overview/functions-logs/functions-logs.component.html
index 582cc3548e..ef06f6a17e 100644
---
a/ui/src/app/pipelines/components/functions-overview/functions-logs/functions-logs.component.html
+++
b/ui/src/app/pipelines/components/functions-overview/functions-logs/functions-logs.component.html
@@ -22,12 +22,12 @@
[showBackLink]="true"
[backLinkTarget]="['pipelines']"
>
- <div nav fxFlex="100" fxLayout="row" fxLayoutAlign="end center">
+ <div nav fxLayout="row" fxLayoutAlign="end center">
<button
mat-icon-button
color="accent"
class="mr-10"
- matTooltip="Refresh"
+ [matTooltip]="'Refresh' | translate"
(click)="triggerUpdate()"
>
<i class="material-icons">refresh</i>
diff --git
a/ui/src/app/pipelines/components/functions-overview/functions-logs/functions-logs.component.ts
b/ui/src/app/pipelines/components/functions-overview/functions-logs/functions-logs.component.ts
index b0ebd1ab0f..5bb37d61eb 100644
---
a/ui/src/app/pipelines/components/functions-overview/functions-logs/functions-logs.component.ts
+++
b/ui/src/app/pipelines/components/functions-overview/functions-logs/functions-logs.component.ts
@@ -28,6 +28,7 @@ import {
import { MatIconButton } from '@angular/material/button';
import { MatTooltip } from '@angular/material/tooltip';
import { SpSimpleLogsComponent } from
'../../../../core-ui/monitoring/simple-logs/simple-logs.component';
+import { TranslatePipe } from '@ngx-translate/core';
@Component({
selector: 'sp-functions-logs',
@@ -41,6 +42,7 @@ import { SpSimpleLogsComponent } from
'../../../../core-ui/monitoring/simple-log
MatIconButton,
MatTooltip,
SpSimpleLogsComponent,
+ TranslatePipe,
],
})
export class SpFunctionsLogsComponent
diff --git
a/ui/src/app/pipelines/components/functions-overview/functions-metrics/functions-metrics.component.html
b/ui/src/app/pipelines/components/functions-overview/functions-metrics/functions-metrics.component.html
index fceb4f3f14..f1019fc661 100644
---
a/ui/src/app/pipelines/components/functions-overview/functions-metrics/functions-metrics.component.html
+++
b/ui/src/app/pipelines/components/functions-overview/functions-metrics/functions-metrics.component.html
@@ -22,12 +22,12 @@
[showBackLink]="true"
[backLinkTarget]="['pipelines']"
>
- <div nav fxFlex="100" fxLayout="row" fxLayoutAlign="end center">
+ <div nav fxLayout="row" fxLayoutAlign="end center">
<button
mat-icon-button
color="accent"
class="mr-10"
- matTooltip="Refresh"
+ [matTooltip]="'Refresh' | translate"
(click)="triggerUpdate()"
>
<i class="material-icons">refresh</i>
@@ -39,16 +39,14 @@
messagesIn of metrics.messagesIn | keyvalue;
track messagesIn
) {
- <div>
- <sp-simple-metrics
- [elementName]="streamNames[messagesIn.key]"
- lastPublishedLabel="Last consumed message"
- statusValueLabel="Consumed messages"
- [lastTimestamp]="messagesIn.value.lastTimestamp"
- [statusValue]="messagesIn.value.counter"
- >
- </sp-simple-metrics>
- </div>
+ <sp-simple-metrics
+ [elementName]="streamNames[messagesIn.key]"
+ lastPublishedLabel="Last consumed message"
+ statusValueLabel="Consumed messages"
+ [lastTimestamp]="messagesIn.value.lastTimestamp"
+ [statusValue]="messagesIn.value.counter"
+ >
+ </sp-simple-metrics>
}
</div>
}
diff --git
a/ui/src/app/pipelines/components/functions-overview/functions-metrics/functions-metrics.component.ts
b/ui/src/app/pipelines/components/functions-overview/functions-metrics/functions-metrics.component.ts
index ca691fd644..3e6d19fa4f 100644
---
a/ui/src/app/pipelines/components/functions-overview/functions-metrics/functions-metrics.component.ts
+++
b/ui/src/app/pipelines/components/functions-overview/functions-metrics/functions-metrics.component.ts
@@ -29,6 +29,7 @@ import { MatIconButton } from '@angular/material/button';
import { MatTooltip } from '@angular/material/tooltip';
import { KeyValuePipe } from '@angular/common';
import { SpSimpleMetricsComponent } from
'../../../../core-ui/monitoring/simple-metrics/simple-metrics.component';
+import { TranslatePipe } from '@ngx-translate/core';
@Component({
selector: 'sp-functions-metrics',
@@ -43,6 +44,7 @@ import { SpSimpleMetricsComponent } from
'../../../../core-ui/monitoring/simple-
MatTooltip,
KeyValuePipe,
SpSimpleMetricsComponent,
+ TranslatePipe,
],
})
export class SpFunctionsMetricsComponent
diff --git a/ui/src/app/pipelines/pipelines.component.html
b/ui/src/app/pipelines/pipelines.component.html
index 4ff79306f7..34c0db5950 100644
--- a/ui/src/app/pipelines/pipelines.component.html
+++ b/ui/src/app/pipelines/pipelines.component.html
@@ -73,19 +73,25 @@
}
</div>
</div>
- <div fxFlex="100" fxLayout="column" style="margin-top: 20px">
- <sp-basic-header-title-component
- title="Functions"
- ></sp-basic-header-title-component>
- <div fxFlex="100" fxLayout="row" fxLayoutAlign="center start">
- <div fxFlex="100">
- @if (functionsReady) {
- <sp-functions-overview [functions]="functions">
- </sp-functions-overview>
- }
+ @if (isAdminRole) {
+ <div fxFlex="100" fxLayout="column" style="margin-top: 20px">
+ <sp-basic-header-title-component
+ title="Functions"
+ ></sp-basic-header-title-component>
+ <div
+ fxFlex="100"
+ fxLayout="row"
+ fxLayoutAlign="center start"
+ >
+ <div fxFlex="100">
+ @if (functionsReady) {
+ <sp-functions-overview [functions]="functions">
+ </sp-functions-overview>
+ }
+ </div>
</div>
</div>
- </div>
+ }
</div>
</div>
</sp-basic-view>
diff --git a/ui/src/app/pipelines/pipelines.component.ts
b/ui/src/app/pipelines/pipelines.component.ts
index 5f94bcc0e6..d755a95f39 100644
--- a/ui/src/app/pipelines/pipelines.component.ts
+++ b/ui/src/app/pipelines/pipelines.component.ts
@@ -132,10 +132,12 @@ export class PipelinesComponent implements OnInit,
OnDestroy {
}
getFunctions() {
- this.functionsService.getActiveFunctions().subscribe(functions => {
- this.functions = functions.map(f => f.functionId).sort();
- this.functionsReady = true;
- });
+ if (this.isAdminRole) {
+ this.functionsService.getActiveFunctions().subscribe(functions => {
+ this.functions = functions.map(f => f.functionId).sort();
+ this.functionsReady = true;
+ });
+ }
}
getPipelines() {