This is an automated email from the ASF dual-hosted git repository.
riemer 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 3ecdebb35 Improve and centralize HTTP requests to extensions (#1918)
3ecdebb35 is described below
commit 3ecdebb35500dd6ce497ed9f53949f33220ddb74
Author: Dominik Riemer <[email protected]>
AuthorDate: Sun Sep 17 10:38:17 2023 +0200
Improve and centralize HTTP requests to extensions (#1918)
* Improve and centralize HTTP requests to extensions
* Use current user token for guess schema request
* Add SP prefix to env variable
* Refactor fetching of runtime static property requests
---
.../apache/streampipes/commons/constants/Envs.java | 2 +-
.../management/health/AdapterHealthCheck.java | 7 ++-
.../management/management/GuessManagement.java | 9 +--
.../management/management/SourcesManagement.java | 17 ------
.../management/management/WorkerRestClient.java | 25 +++-----
.../manager/endpoint/EndpointFetcher.java | 6 +-
.../manager/endpoint/EndpointItemFetcher.java | 17 +++---
.../execution/ExtensionServiceExecutions.java | 67 ++++++++++++++++++++++
.../health/PipelineElementEndpointHealthCheck.java | 5 +-
.../manager/health/ServiceHealthCheck.java | 5 +-
.../pipeline/ExtensionsServiceLogExecutor.java | 5 +-
.../remote/ContainerProvidedOptionsHandler.java | 13 ++---
.../manager/setup/ExtensionsInstallationTask.java | 8 +--
.../setup/PipelineElementInstallationStep.java | 41 ++++++++++---
.../manager/setup/StreamPipesEnvChecker.java | 5 ++
.../streampipes/manager/util/AuthTokenUtils.java | 5 ++
.../svcdiscovery/SpServiceDiscoveryCore.java | 6 +-
.../security/UnauthenticatedInterfaces.java | 10 ++--
18 files changed, 166 insertions(+), 87 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 a56b1257c..6b424bfa7 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
@@ -40,7 +40,7 @@ public enum Envs {
SP_CONSUL_PORT("SP_CONSUL_PORT", DefaultEnvValues.CONSUL_PORT_DEFAULT),
SP_KAFKA_RETENTION_MS("SP_KAFKA_RETENTION_MS",
DefaultEnvValues.SP_KAFKA_RETENTION_MS_DEFAULT),
SP_PRIORITIZED_PROTOCOL("SP_PRIORITIZED_PROTOCOL", "kafka"),
- SP_JWT_SECRET("JWT_SECRET"),
+ SP_JWT_SECRET("SP_JWT_SECRET"),
SP_JWT_SIGNING_MODE("SP_JWT_SIGNING_MODE"),
SP_JWT_PRIVATE_KEY_LOC("SP_JWT_PRIVATE_KEY_LOC"),
SP_JWT_PUBLIC_KEY_LOC("SP_JWT_PUBLIC_KEY_LOC"),
diff --git
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/health/AdapterHealthCheck.java
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/health/AdapterHealthCheck.java
index c575e9a43..07d70cf7d 100644
---
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/health/AdapterHealthCheck.java
+++
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/health/AdapterHealthCheck.java
@@ -26,6 +26,9 @@ import
org.apache.streampipes.model.connect.adapter.AdapterDescription;
import org.apache.streampipes.storage.api.IAdapterStorage;
import org.apache.streampipes.storage.couchdb.CouchDbStorageManager;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
@@ -33,6 +36,8 @@ import java.util.Map;
public class AdapterHealthCheck {
+ private static final Logger LOG =
LoggerFactory.getLogger(AdapterHealthCheck.class);
+
private final IAdapterStorage adapterStorage;
private final AdapterMasterManagement adapterMasterManagement;
@@ -125,7 +130,7 @@ public class AdapterHealthCheck {
this.adapterMasterManagement.startStreamAdapter(adapterDescription.getElementId());
}
} catch (AdapterException e) {
- e.printStackTrace();
+ LOG.warn("Could not start adapter {}", adapterDescription.getName(),
e);
}
}
diff --git
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/GuessManagement.java
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/GuessManagement.java
index 383e889cb..a8cae5b72 100644
---
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/GuessManagement.java
+++
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/GuessManagement.java
@@ -24,6 +24,7 @@ import
org.apache.streampipes.commons.exceptions.connect.ParseException;
import org.apache.streampipes.connect.management.AdapterEventPreviewPipeline;
import org.apache.streampipes.connect.management.util.WorkerPaths;
import
org.apache.streampipes.extensions.api.connect.exception.WorkerAdapterException;
+import org.apache.streampipes.manager.execution.ExtensionServiceExecutions;
import org.apache.streampipes.model.connect.adapter.AdapterDescription;
import org.apache.streampipes.model.connect.guess.AdapterEventPreview;
import org.apache.streampipes.model.connect.guess.GuessSchema;
@@ -31,9 +32,7 @@ import
org.apache.streampipes.serializers.json.JacksonSerializer;
import com.fasterxml.jackson.core.JsonProcessingException;
import org.apache.http.HttpStatus;
-import org.apache.http.client.fluent.Request;
import org.apache.http.client.fluent.Response;
-import org.apache.http.entity.ContentType;
import org.apache.http.util.EntityUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -57,10 +56,8 @@ public class GuessManagement {
var objectMapper = JacksonSerializer.getObjectMapper();
var description = objectMapper.writeValueAsString(adapterDescription);
logger.info("Guess schema at: " + workerUrl);
- Response requestResponse = Request.Post(workerUrl)
- .bodyString(description, ContentType.APPLICATION_JSON)
- .connectTimeout(1000)
- .socketTimeout(100000)
+ Response requestResponse = ExtensionServiceExecutions
+ .extServicePostRequest(workerUrl, description)
.execute();
var httpResponse = requestResponse.returnResponse();
diff --git
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/SourcesManagement.java
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/SourcesManagement.java
index da2c96a01..74eefd734 100644
---
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/SourcesManagement.java
+++
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/SourcesManagement.java
@@ -21,26 +21,9 @@ package org.apache.streampipes.connect.management.management;
import org.apache.streampipes.model.SpDataStream;
import org.apache.streampipes.model.connect.adapter.AdapterDescription;
import org.apache.streampipes.model.grounding.EventGrounding;
-import org.apache.streampipes.storage.couchdb.impl.AdapterInstanceStorageImpl;
-
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
public class SourcesManagement {
- private final Logger logger =
LoggerFactory.getLogger(SourcesManagement.class);
-
- private final AdapterInstanceStorageImpl adapterInstanceStorage;
-
-
- public SourcesManagement(AdapterInstanceStorageImpl adapterStorage) {
- this.adapterInstanceStorage = adapterStorage;
- }
-
- public SourcesManagement() {
- this(new AdapterInstanceStorageImpl());
- }
-
public static SpDataStream updateDataStream(AdapterDescription
adapterDescription,
SpDataStream oldDataStream) {
diff --git
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/WorkerRestClient.java
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/WorkerRestClient.java
index adefb1f6e..58f6a49bc 100644
---
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/WorkerRestClient.java
+++
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/WorkerRestClient.java
@@ -21,6 +21,7 @@ package org.apache.streampipes.connect.management.management;
import org.apache.streampipes.commons.exceptions.SpConfigurationException;
import org.apache.streampipes.commons.exceptions.connect.AdapterException;
import org.apache.streampipes.connect.management.util.WorkerPaths;
+import org.apache.streampipes.manager.execution.ExtensionServiceExecutions;
import org.apache.streampipes.model.connect.adapter.AdapterDescription;
import org.apache.streampipes.model.runtime.RuntimeOptionsRequest;
import org.apache.streampipes.model.runtime.RuntimeOptionsResponse;
@@ -35,7 +36,6 @@ import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.commons.io.IOUtils;
import org.apache.http.HttpResponse;
import org.apache.http.client.fluent.Request;
-import org.apache.http.entity.ContentType;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -73,10 +73,8 @@ public class WorkerRestClient {
public static List<AdapterDescription>
getAllRunningAdapterInstanceDescriptions(String url) throws AdapterException {
try {
logger.info("Requesting all running adapter description instances: " +
url);
-
- var responseString = Request.Get(url)
- .connectTimeout(1000)
- .socketTimeout(100000)
+ var responseString = ExtensionServiceExecutions
+ .extServiceGetRequest(url)
.execute().returnContent().asString();
List<AdapterDescription> result =
JacksonSerializer.getObjectMapper().readValue(responseString, List.class);
@@ -108,7 +106,7 @@ public class WorkerRestClient {
try {
String adapterDescription =
JacksonSerializer.getObjectMapper().writeValueAsString(ad);
- var response = triggerPost(url, adapterDescription);
+ var response = triggerPost(url, ad.getElementId(), adapterDescription);
var responseString = getResponseBody(response);
if (response.getStatusLine().getStatusCode() != 200) {
@@ -116,7 +114,7 @@ public class WorkerRestClient {
throw new AdapterException(exception.getMessage(),
exception.getCause());
}
- logger.info("Adapter {} on endpoint: " + url + " with Response: " +
responseString);
+ logger.info("Adapter {} on endpoint: " + url + " with Response: ",
ad.getName() + responseString);
} catch (IOException e) {
logger.error("Adapter was not {} successfully", action, e);
@@ -129,12 +127,10 @@ public class WorkerRestClient {
}
private static HttpResponse triggerPost(String url,
+ String elementId,
String payload) throws IOException {
- return Request.Post(url)
- .bodyString(payload, ContentType.APPLICATION_JSON)
- .connectTimeout(1000)
- .socketTimeout(100000)
- .execute().returnResponse();
+ var request = ExtensionServiceExecutions.extServicePostRequest(url,
elementId, payload);
+ return request.execute().returnResponse();
}
public static RuntimeOptionsResponse getConfiguration(String workerEndpoint,
@@ -145,10 +141,7 @@ public class WorkerRestClient {
try {
String payload =
JacksonSerializer.getObjectMapper().writeValueAsString(runtimeOptionsRequest);
- var response = Request.Post(url)
- .bodyString(payload, ContentType.APPLICATION_JSON)
- .connectTimeout(1000)
- .socketTimeout(100000)
+ var response = ExtensionServiceExecutions.extServicePostRequest(url,
payload)
.execute()
.returnResponse();
diff --git
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/endpoint/EndpointFetcher.java
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/endpoint/EndpointFetcher.java
index 1aa54c870..1896f3319 100644
---
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/endpoint/EndpointFetcher.java
+++
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/endpoint/EndpointFetcher.java
@@ -31,19 +31,19 @@ public class EndpointFetcher {
public List<ExtensionsServiceEndpoint> getEndpoints() {
List<String> endpoints =
SpServiceDiscovery.getServiceDiscovery().getActivePipelineElementEndpoints();
- List<ExtensionsServiceEndpoint> servicerdExtensionsServiceEndpoints = new
LinkedList<>();
+ List<ExtensionsServiceEndpoint> serviceExtensionsServiceEndpoints = new
LinkedList<>();
for (String endpoint : endpoints) {
ExtensionsServiceEndpoint extensionsServiceEndpoint =
new ExtensionsServiceEndpoint(endpoint);
- servicerdExtensionsServiceEndpoints.add(extensionsServiceEndpoint);
+ serviceExtensionsServiceEndpoints.add(extensionsServiceEndpoint);
}
List<ExtensionsServiceEndpoint> databasedExtensionsServiceEndpoints =
StorageDispatcher.INSTANCE.getNoSqlStore()
.getRdfEndpointStorage()
.getExtensionsServiceEndpoints();
List<ExtensionsServiceEndpoint> concatList =
- Stream.of(databasedExtensionsServiceEndpoints,
servicerdExtensionsServiceEndpoints)
+ Stream.of(databasedExtensionsServiceEndpoints,
serviceExtensionsServiceEndpoints)
.flatMap(Collection::stream)
.collect(Collectors.toList());
diff --git
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/endpoint/EndpointItemFetcher.java
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/endpoint/EndpointItemFetcher.java
index 0cfda6e2b..2fd3a4ef9 100644
---
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/endpoint/EndpointItemFetcher.java
+++
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/endpoint/EndpointItemFetcher.java
@@ -18,24 +18,24 @@
package org.apache.streampipes.manager.endpoint;
+import org.apache.streampipes.manager.execution.ExtensionServiceExecutions;
import org.apache.streampipes.model.client.endpoint.ExtensionsServiceEndpoint;
import
org.apache.streampipes.model.client.endpoint.ExtensionsServiceEndpointItem;
import org.apache.streampipes.serializers.json.JacksonSerializer;
import com.fasterxml.jackson.core.type.TypeReference;
-import org.apache.http.client.fluent.Request;
-import org.apache.http.message.BasicHeader;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.util.ArrayList;
+import java.util.Collections;
import java.util.List;
public class EndpointItemFetcher {
Logger logger = LoggerFactory.getLogger(EndpointItemFetcher.class);
- private List<ExtensionsServiceEndpoint> extensionsServiceEndpoints;
+ private final List<ExtensionsServiceEndpoint> extensionsServiceEndpoints;
public EndpointItemFetcher(List<ExtensionsServiceEndpoint>
extensionsServiceEndpoints) {
this.extensionsServiceEndpoints = extensionsServiceEndpoints;
@@ -49,19 +49,18 @@ public class EndpointItemFetcher {
private List<ExtensionsServiceEndpointItem>
getEndpointItems(ExtensionsServiceEndpoint e) {
try {
- String result = Request.Get(e.getEndpointUrl())
- .addHeader(new BasicHeader("Accept", "application/json"))
- .connectTimeout(1000)
+ String result = ExtensionServiceExecutions
+ .extServiceGetRequest(e.getEndpointUrl())
.execute()
.returnContent()
.asString();
return JacksonSerializer.getObjectMapper()
- .readValue(result, new
TypeReference<List<ExtensionsServiceEndpointItem>>() {
+ .readValue(result, new TypeReference<>() {
});
} catch (IOException e1) {
- logger.warn("Processing Element Descriptions could not be fetched from
RDF endpoint: " + e.getEndpointUrl());
- return new ArrayList<>();
+ logger.warn("Processing Element Descriptions could not be fetched from
endpoint: " + e.getEndpointUrl(), e1);
+ return Collections.emptyList();
}
}
}
diff --git
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/execution/ExtensionServiceExecutions.java
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/execution/ExtensionServiceExecutions.java
new file mode 100644
index 000000000..70bd2e8ce
--- /dev/null
+++
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/execution/ExtensionServiceExecutions.java
@@ -0,0 +1,67 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ *
+ */
+
+package org.apache.streampipes.manager.execution;
+
+import org.apache.streampipes.manager.util.AuthTokenUtils;
+import org.apache.streampipes.resource.management.SpResourceManager;
+
+import org.apache.http.client.fluent.Request;
+import org.apache.http.entity.ContentType;
+
+public class ExtensionServiceExecutions {
+
+ public static Request extServiceGetRequest(String url) {
+ return Request
+ .Get(url)
+ .addHeader("Authorization",
AuthTokenUtils.getAuthTokenForUser(getServiceAdminSid()))
+ .addHeader("Accept", "application/json")
+ .connectTimeout(10000)
+ .socketTimeout(10000);
+ }
+
+
+ private static String getServiceAdminSid() {
+ return new
SpResourceManager().manageUsers().getServiceAdmin().getPrincipalId();
+ }
+
+ public static Request extServicePostRequest(String url,
+ String payload) {
+ return authenticatedPostRequest(url,
AuthTokenUtils.getAuthTokenForCurrentUser(), payload);
+ }
+
+ public static Request extServicePostRequest(String url,
+ String elementId,
+ String payload) {
+ return authenticatedPostRequest(
+ url,
+ AuthTokenUtils.getAuthToken(elementId),
+ payload
+ );
+ }
+
+ private static Request authenticatedPostRequest(String url,
+ String token,
+ String payload) {
+ return Request.Post(url)
+ .addHeader("Authorization", token)
+ .bodyString(payload, ContentType.APPLICATION_JSON)
+ .connectTimeout(1000)
+ .socketTimeout(100000);
+ }
+}
diff --git
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/health/PipelineElementEndpointHealthCheck.java
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/health/PipelineElementEndpointHealthCheck.java
index 8d9dc937e..9bec880a8 100644
---
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/health/PipelineElementEndpointHealthCheck.java
+++
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/health/PipelineElementEndpointHealthCheck.java
@@ -17,10 +17,10 @@
*/
package org.apache.streampipes.manager.health;
+import org.apache.streampipes.manager.execution.ExtensionServiceExecutions;
import org.apache.streampipes.serializers.json.JacksonSerializer;
import com.fasterxml.jackson.core.JsonProcessingException;
-import org.apache.http.client.fluent.Request;
import java.io.IOException;
import java.util.Arrays;
@@ -37,7 +37,8 @@ public class PipelineElementEndpointHealthCheck {
}
public List<String> checkRunningInstances() throws IOException {
- return
asList(Request.Get(makeRequestUrl()).execute().returnContent().toString());
+ var request =
ExtensionServiceExecutions.extServiceGetRequest(makeRequestUrl());
+ return asList(request.execute().returnContent().toString());
}
private List<String> asList(String json) throws JsonProcessingException {
diff --git
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/health/ServiceHealthCheck.java
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/health/ServiceHealthCheck.java
index 768430d5a..dadfe9c64 100644
---
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/health/ServiceHealthCheck.java
+++
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/health/ServiceHealthCheck.java
@@ -18,11 +18,11 @@
package org.apache.streampipes.manager.health;
+import org.apache.streampipes.manager.execution.ExtensionServiceExecutions;
import
org.apache.streampipes.model.extensions.svcdiscovery.SpServiceRegistration;
import org.apache.streampipes.storage.api.CRUDStorage;
import org.apache.streampipes.storage.management.StorageDispatcher;
-import org.apache.http.client.fluent.Request;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -51,7 +51,8 @@ public class ServiceHealthCheck implements Runnable {
String healthCheckUrl = makeHealthCheckUrl(service);
try {
- var response = Request.Get(healthCheckUrl).execute();
+ var request =
ExtensionServiceExecutions.extServiceGetRequest(healthCheckUrl);
+ var response = request.execute();
if (response.returnResponse().getStatusLine().getStatusCode() != 200) {
processUnhealthyService(service);
} else {
diff --git
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/monitoring/pipeline/ExtensionsServiceLogExecutor.java
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/monitoring/pipeline/ExtensionsServiceLogExecutor.java
index ee8e19885..124c216c2 100644
---
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/monitoring/pipeline/ExtensionsServiceLogExecutor.java
+++
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/monitoring/pipeline/ExtensionsServiceLogExecutor.java
@@ -19,7 +19,7 @@
package org.apache.streampipes.manager.monitoring.pipeline;
-import org.apache.streampipes.manager.util.AuthTokenUtils;
+import org.apache.streampipes.manager.execution.ExtensionServiceExecutions;
import org.apache.streampipes.model.client.user.Principal;
import org.apache.streampipes.model.monitoring.SpEndpointMonitoringInfo;
import org.apache.streampipes.resource.management.SpResourceManager;
@@ -61,8 +61,7 @@ public class ExtensionsServiceLogExecutor implements Runnable
{
}
private Request makeRequest(String serviceEndpointUrl) {
- return Request.Get(makeLogUrl(serviceEndpointUrl))
- .addHeader("Authorization",
AuthTokenUtils.getAuthTokenForUser(getServiceAdmin()));
+ return
ExtensionServiceExecutions.extServiceGetRequest(makeLogUrl(serviceEndpointUrl));
}
private Principal getServiceAdmin() {
diff --git
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/remote/ContainerProvidedOptionsHandler.java
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/remote/ContainerProvidedOptionsHandler.java
index 0fa9a66fb..fd1957f5c 100644
---
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/remote/ContainerProvidedOptionsHandler.java
+++
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/remote/ContainerProvidedOptionsHandler.java
@@ -18,6 +18,7 @@
package org.apache.streampipes.manager.remote;
import
org.apache.streampipes.commons.exceptions.NoServiceEndpointsAvailableException;
+import org.apache.streampipes.manager.execution.ExtensionServiceExecutions;
import
org.apache.streampipes.manager.execution.endpoint.ExtensionsServiceEndpointGenerator;
import
org.apache.streampipes.manager.execution.endpoint.ExtensionsServiceEndpointUtils;
import org.apache.streampipes.model.runtime.RuntimeOptionsRequest;
@@ -26,9 +27,7 @@ import
org.apache.streampipes.serializers.json.JacksonSerializer;
import org.apache.streampipes.svcdiscovery.api.model.SpServiceUrlProvider;
import com.google.gson.JsonSyntaxException;
-import org.apache.http.client.fluent.Request;
import org.apache.http.client.fluent.Response;
-import org.apache.http.entity.ContentType;
import java.io.IOException;
@@ -38,11 +37,11 @@ public class ContainerProvidedOptionsHandler {
public RuntimeOptionsResponse fetchRemoteOptions(RuntimeOptionsRequest
request) {
try {
- String httpRequestBody =
JacksonSerializer.getObjectMapper().writeValueAsString(request);
- Response httpResp =
-
Request.Post(getEndpointUrl(request.getAppId())).bodyString(httpRequestBody,
ContentType.APPLICATION_JSON)
- .execute();
- return handleResponse(httpResp);
+ var payload =
JacksonSerializer.getObjectMapper().writeValueAsString(request);
+ var url = getEndpointUrl(request.getAppId());
+ var resp = ExtensionServiceExecutions.extServicePostRequest(url,
payload).execute();
+
+ return handleResponse(resp);
} catch (Exception e) {
e.printStackTrace();
return new RuntimeOptionsResponse();
diff --git
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/ExtensionsInstallationTask.java
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/ExtensionsInstallationTask.java
index 51864794e..026b6eef1 100644
---
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/ExtensionsInstallationTask.java
+++
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/ExtensionsInstallationTask.java
@@ -34,8 +34,8 @@ public class ExtensionsInstallationTask implements Runnable {
private static final Logger LOG =
LoggerFactory.getLogger(ExtensionsInstallationTask.class);
- private static final int MAX_RETRIES = 4;
- private static final int SLEEP_TIME_SECONDS = 3;
+ private static final int MAX_RETRIES = 6;
+ private static final int SLEEP_TIME_SECONDS = 2;
private final InitialSettings settings;
private final BackgroundTaskNotifier callback;
@@ -55,7 +55,7 @@ public class ExtensionsInstallationTask implements Runnable {
do {
endpoints = new EndpointFetcher().getEndpoints();
numberOfAttempts++;
- if (endpoints.size() == 0) {
+ if (endpoints.isEmpty()) {
LOG.info("Found 0 endpoints - waiting {} seconds to make sure all
endpoints have properly started",
SLEEP_TIME_SECONDS);
try {
@@ -64,7 +64,7 @@ public class ExtensionsInstallationTask implements Runnable {
e.printStackTrace();
}
}
- } while (endpoints.size() == 0 && numberOfAttempts < MAX_RETRIES);
+ } while (endpoints.isEmpty() && numberOfAttempts < MAX_RETRIES);
LOG.info("Found {} endpoints from which we will install extensions.",
endpoints.size());
LOG.info(
"Further available extensions can always be installed by navigating
to the 'Install pipeline elements' view");
diff --git
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/PipelineElementInstallationStep.java
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/PipelineElementInstallationStep.java
index 319318eeb..568f2902f 100644
---
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/PipelineElementInstallationStep.java
+++
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/PipelineElementInstallationStep.java
@@ -23,14 +23,23 @@ import
org.apache.streampipes.model.client.endpoint.ExtensionsServiceEndpoint;
import
org.apache.streampipes.model.client.endpoint.ExtensionsServiceEndpointItem;
import org.apache.streampipes.model.message.Message;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
+import java.util.concurrent.TimeUnit;
public class PipelineElementInstallationStep extends InstallationStep {
+ private static final Logger LOG =
LoggerFactory.getLogger(PipelineElementInstallationStep.class);
+ private static final int MAX_RETRIES = 8;
+
private final ExtensionsServiceEndpoint endpoint;
private final String principalSid;
+ private int retries = 0;
+
public PipelineElementInstallationStep(ExtensionsServiceEndpoint endpoint,
String principalSid) {
@@ -42,15 +51,31 @@ public class PipelineElementInstallationStep extends
InstallationStep {
public void install() {
List<Message> statusMessages = new ArrayList<>();
List<ExtensionsServiceEndpointItem> items =
Operations.getEndpointUriContents(Collections.singletonList(endpoint));
- for (ExtensionsServiceEndpointItem item : items) {
- statusMessages.add(new
EndpointItemParser().parseAndAddEndpointItem(item.getUri(),
- principalSid, true));
- }
-
- if (statusMessages.stream().allMatch(Message::isSuccess)) {
- logSuccess(getTitle());
+ if (items.isEmpty() && retries <= MAX_RETRIES) {
+ retries++;
+ LOG.info(
+ "Endpoint available but no extensions yet found, so we will retry to
fetch pipeline elements ({}/{})",
+ retries,
+ MAX_RETRIES
+ );
+ try {
+ TimeUnit.SECONDS.sleep(2);
+ install();
+ } catch (InterruptedException e) {
+ throw new RuntimeException(e);
+ }
} else {
- logFailure(getTitle());
+ LOG.info("Found {} endpoint items for endpoint {}", items.size(),
endpoint.getEndpointUrl());
+ for (ExtensionsServiceEndpointItem item : items) {
+ statusMessages.add(new
EndpointItemParser().parseAndAddEndpointItem(item.getUri(),
+ principalSid, true));
+ }
+
+ if (statusMessages.stream().allMatch(Message::isSuccess)) {
+ logSuccess(getTitle());
+ } else {
+ logFailure(getTitle());
+ }
}
}
diff --git
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/StreamPipesEnvChecker.java
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/StreamPipesEnvChecker.java
index 17dc5072e..c4fc1bba0 100644
---
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/StreamPipesEnvChecker.java
+++
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/StreamPipesEnvChecker.java
@@ -70,7 +70,12 @@ public class StreamPipesEnvChecker {
if (signingMode.exists()) {
localAuthConfig.setJwtSigningMode(JwtSigningMode.valueOf(signingMode.getValue()));
+ } else {
+ if (localAuthConfig.getJwtSigningMode() != JwtSigningMode.HMAC) {
+ localAuthConfig.setJwtSigningMode(JwtSigningMode.HMAC);
+ }
}
+
if (jwtSecret.exists()) {
localAuthConfig.setTokenSecret(jwtSecret.getValue());
}
diff --git
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/util/AuthTokenUtils.java
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/util/AuthTokenUtils.java
index 43fb7089e..12c88359b 100644
---
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/util/AuthTokenUtils.java
+++
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/util/AuthTokenUtils.java
@@ -28,6 +28,11 @@ import
org.springframework.security.core.context.SecurityContextHolder;
public class AuthTokenUtils {
+ public static String getAuthTokenForCurrentUser() {
+ Authentication auth =
SecurityContextHolder.getContext().getAuthentication();
+ return makeBearerToken(new JwtTokenProvider().createToken(auth));
+ }
+
public static String getAuthToken(String resourceId) {
if (SecurityContextHolder.getContext().getAuthentication() != null) {
Authentication auth =
SecurityContextHolder.getContext().getAuthentication();
diff --git
a/streampipes-service-discovery/src/main/java/org/apache/streampipes/svcdiscovery/SpServiceDiscoveryCore.java
b/streampipes-service-discovery/src/main/java/org/apache/streampipes/svcdiscovery/SpServiceDiscoveryCore.java
index 2ad6bed1f..ea3ff85d5 100644
---
a/streampipes-service-discovery/src/main/java/org/apache/streampipes/svcdiscovery/SpServiceDiscoveryCore.java
+++
b/streampipes-service-discovery/src/main/java/org/apache/streampipes/svcdiscovery/SpServiceDiscoveryCore.java
@@ -80,17 +80,19 @@ public class SpServiceDiscoveryCore implements
ISpServiceDiscovery {
private List<SpServiceRegistration> findService(int retryCount) {
var services = serviceStorage.getAll();
- if (services.size() == 0) {
+ if (services.isEmpty()) {
if (retryCount < MAX_RETRIES) {
try {
retryCount++;
- TimeUnit.SECONDS.sleep(10);
+ LOG.info("Could not find any extensions services, retrying ({}/{})",
retryCount, MAX_RETRIES);
+ TimeUnit.MILLISECONDS.sleep(1000);
return findService(retryCount);
} catch (InterruptedException e) {
e.printStackTrace();
return Collections.emptyList();
}
} else {
+ LOG.info("No service found");
return Collections.emptyList();
}
} else {
diff --git
a/streampipes-service-extensions/src/main/java/org/apache/streampipes/service/extensions/security/UnauthenticatedInterfaces.java
b/streampipes-service-extensions/src/main/java/org/apache/streampipes/service/extensions/security/UnauthenticatedInterfaces.java
index bd28e9c29..de1a715db 100644
---
a/streampipes-service-extensions/src/main/java/org/apache/streampipes/service/extensions/security/UnauthenticatedInterfaces.java
+++
b/streampipes-service-extensions/src/main/java/org/apache/streampipes/service/extensions/security/UnauthenticatedInterfaces.java
@@ -27,12 +27,10 @@ public class UnauthenticatedInterfaces {
public static Collection<String> get() {
return Arrays.asList(
- "/svchealth/*",
- "/",
- "/sec/**",
- "/sepa/**",
- "/stream/**",
- "/api/v1/worker/**"
+ "/sec/*/assets/**",
+ "/sepa/*/assets/**",
+ "/stream/*/assets/**",
+ "/api/v1/worker/adapters/*/assets/**"
);
}
}