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/**"
     );
   }
 }


Reply via email to