This is an automated email from the ASF dual-hosted git repository.

tenthe pushed a commit to branch 
4655-optimize-adapter-startstop-api-database-access
in repository https://gitbox.apache.org/repos/asf/streampipes.git


The following commit(s) were added to 
refs/heads/4655-optimize-adapter-startstop-api-database-access by this push:
     new 0a61967ff8 Optimize adapter start and stop lifecycle
0a61967ff8 is described below

commit 0a61967ff8063df23a2444794cd977aa7c34a959
Author: Philipp Zehnder <[email protected]>
AuthorDate: Mon Jun 29 17:36:02 2026 +0200

    Optimize adapter start and stop lifecycle
---
 .../management/AdapterMasterManagement.java        | 27 ++++++++-----
 .../management/AdapterMigrationManager.java        |  2 +-
 .../management/AdapterUpdateManagement.java        |  4 +-
 .../management/management/WorkerRestClient.java    | 47 +++++++++++-----------
 .../export/resolver/AdapterResolver.java           |  2 +-
 .../health/monitoring/AdapterHealthCheck.java      |  2 +-
 .../rest/impl/connect/AdapterResource.java         |  4 +-
 .../rest/impl/connect/CompactAdapterResource.java  |  2 +-
 8 files changed, 48 insertions(+), 42 deletions(-)

diff --git 
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterMasterManagement.java
 
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterMasterManagement.java
index 53b7c66fad..24f9f828b3 100644
--- 
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterMasterManagement.java
+++ 
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterMasterManagement.java
@@ -120,7 +120,7 @@ public class AdapterMasterManagement {
     try {
       // Stop stream adapter
       try {
-        stopStreamAdapter(elementId, true);
+        stopAdapter(elementId, true);
       } catch (AdapterException e) {
         LOG.info("Could not stop adapter: " + elementId, e);
       }
@@ -144,9 +144,12 @@ public class AdapterMasterManagement {
     return adapterResourceManager.getDb().findAll();
   }
 
-  public void stopStreamAdapter(String elementId, boolean forceStop) throws 
AdapterException {
+  public void stopAdapter(String elementId, boolean forceStop) throws 
AdapterException {
+    stopAdapter(getAdapter(elementId), forceStop);
+  }
+
+  public void stopAdapter(AdapterDescription ad, boolean forceStop) throws 
AdapterException {
     LoadManager.tryLockForAdapter();
-    AdapterDescription ad = 
adapterResourceManager.getDb().getElementById(elementId);
     try {
       try {
         var service = extensionsServiceStorage.findAll().stream()
@@ -154,11 +157,11 @@ public class AdapterMasterManagement {
               if (ad.getSelectedServiceId() != null) {
                 return svc.getSvcId().equals(ad.getSelectedServiceId());
               } else {
-               return svc.getServiceUrl().equals(ad.getSelectedEndpointUrl());
+                return svc.getServiceUrl().equals(ad.getSelectedEndpointUrl());
               }
             })
             .findFirst().orElseThrow(AdapterException::new);
-          workerRestClient.stopStreamAdapter(service, ad);
+        workerRestClient.stopAdapter(service, ad);
       } catch (AdapterException e) {
         if (!forceStop) {
           throw new AdapterException("Could not stop adapter", e);
@@ -168,7 +171,7 @@ public class AdapterMasterManagement {
           adapterResourceManager.getDb().updateElement(ad);
         }
       }
-      ExtensionsLogProvider.INSTANCE.reset(elementId);
+      ExtensionsLogProvider.INSTANCE.reset(ad.getElementId());
 
       // remove the adapter from the metrics manager so that
       // no metrics for this adapter are exposed anymore
@@ -182,11 +185,13 @@ public class AdapterMasterManagement {
     }
   }
 
-  public void startStreamAdapter(String elementId) throws AdapterException {
+  public void startAdapter(String elementId) throws AdapterException {
+    startAdapter(getAdapter(elementId));
+  }
+
+  public void startAdapter(AdapterDescription ad) throws AdapterException {
     LoadManager.tryLockForAdapter();
     try {
-      var ad = adapterResourceManager.getDb().getElementById(elementId);
-
       try {
         // Find endpoint to start adapter on
         var service = new ExtensionsServiceEndpointGenerator()
@@ -199,13 +204,13 @@ public class AdapterMasterManagement {
         adapterResourceManager.getDb().updateElement(ad);
 
         // Invoke adapter instance
-        workerRestClient.invokeStreamAdapter(service, elementId);
+        workerRestClient.invokeStreamAdapter(service, ad);
 
         // register the adapter at the metrics manager so that the 
AdapterHealthCheck
         // can send metrics
         adapterMetrics.register(ad.getElementId(), ad.getName());
 
-        LOG.info("Started adapter " + elementId + " on: " + 
ad.getSelectedServiceId());
+        LOG.info("Started adapter " + ad.getElementId() + " on: " + 
ad.getSelectedServiceId());
       } catch (NoServiceEndpointsAvailableException e) {
         throw new AdapterException("Could not start adapter due to unavailable 
service endpoint",
             e);
diff --git 
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterMigrationManager.java
 
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterMigrationManager.java
index 69f4ece820..0bfcb89a83 100644
--- 
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterMigrationManager.java
+++ 
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterMigrationManager.java
@@ -101,7 +101,7 @@ public class AdapterMigrationManager extends 
AbstractMigrationManager implements
                 migrationResult.element().getElementId()
             );
             try {
-              workerRestClient.stopStreamAdapter(service, adapterDescription);
+              workerRestClient.stopAdapter(service, adapterDescription);
             } catch (AdapterException e) {
               LOG.error("Stopping adapter failed: {}", 
StringUtils.join(e.getStackTrace(), "\n"));
             }
diff --git 
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterUpdateManagement.java
 
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterUpdateManagement.java
index 01811ced18..2c3d9f2504 100644
--- 
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterUpdateManagement.java
+++ 
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/AdapterUpdateManagement.java
@@ -52,7 +52,7 @@ public class AdapterUpdateManagement {
     boolean shouldRestart = ad.isRunning();
 
     if (ad.isRunning()) {
-      this.adapterMasterManagement.stopStreamAdapter(ad.getElementId(), true);
+      this.adapterMasterManagement.stopAdapter(ad.getElementId(), true);
     }
 
     // update data source in database
@@ -61,7 +61,7 @@ public class AdapterUpdateManagement {
     pipelineUpdateCoordinator.updatePipelines(ad);
 
     if (shouldRestart) {
-      this.adapterMasterManagement.startStreamAdapter(ad.getElementId());
+      this.adapterMasterManagement.startAdapter(ad.getElementId());
     }
   }
 
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 27df7db7a1..1416ebf564 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
@@ -42,7 +42,6 @@ import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 import java.io.IOException;
-import java.util.List;
 
 /**
  * This client can be used to interact with the adapter workers executing the 
adapter instances
@@ -61,19 +60,23 @@ public class WorkerRestClient {
 
   public void invokeStreamAdapter(SpServiceRegistration service,
                                   String elementId) throws AdapterException {
-    var adapterStreamDescription = getAndDecryptAdapter(elementId);
+    invokeStreamAdapter(service, 
getAdapterStorage().getElementById(elementId));
+  }
+
+  public void invokeStreamAdapter(SpServiceRegistration service,
+                                  AdapterDescription adapterStreamDescription) 
throws AdapterException {
+    var decryptedAdapterStreamDescription = 
cloneAndDecryptAdapter(adapterStreamDescription);
     var requestTarget = ExtensionServiceRequestTargets.adapterStart(service);
-    startAdapter(requestTarget, adapterStreamDescription);
-    updateStreamAdapterStatus(adapterStreamDescription.getElementId(), true);
+    startAdapter(requestTarget, decryptedAdapterStreamDescription);
+    updateStreamAdapterStatus(decryptedAdapterStreamDescription, true);
   }
 
-  public void stopStreamAdapter(SpServiceRegistration service,
-                                AdapterDescription adapterStreamDescription) 
throws AdapterException {
+  public void stopAdapter(SpServiceRegistration service,
+                          AdapterDescription adapterStreamDescription) throws 
AdapterException {
+    var decryptedAdapterStreamDescription = 
cloneAndDecryptAdapter(adapterStreamDescription);
     var requestTarget = ExtensionServiceRequestTargets.adapterStop(service);
-    var ad = 
getAdapterDescriptionById(resourceManager.manageAdapters().getDb(), 
adapterStreamDescription.getElementId());
-
-    stopAdapter(requestTarget, ad);
-    updateStreamAdapterStatus(adapterStreamDescription.getElementId(), false);
+    stopAdapter(requestTarget, decryptedAdapterStreamDescription);
+    updateStreamAdapterStatus(decryptedAdapterStreamDescription, false);
   }
 
   private void startAdapter(ExtensionServiceRequestTarget requestTarget,
@@ -200,22 +203,14 @@ public class WorkerRestClient {
   }
 
 
-  private AdapterDescription getAdapterDescriptionById(IAdapterStorage 
adapterStorage,
-                                                       String id) {
-    AdapterDescription adapterDescription = null;
-    List<AdapterDescription> allAdapters = adapterStorage.findAll();
-    for (AdapterDescription a : allAdapters) {
-      if (a.getElementId().endsWith(id)) {
-        adapterDescription = a;
-      }
-    }
-
-    return adapterDescription;
-  }
-
   private void updateStreamAdapterStatus(String adapterId,
                                          boolean running) {
     var adapter = getAndDecryptAdapter(adapterId);
+    updateStreamAdapterStatus(adapter, running);
+  }
+
+  private void updateStreamAdapterStatus(AdapterDescription adapter,
+                                         boolean running) {
     adapter.setRunning(running);
     encryptAndUpdateAdapter(adapter);
   }
@@ -232,6 +227,12 @@ public class WorkerRestClient {
     return adapter;
   }
 
+  private AdapterDescription cloneAndDecryptAdapter(AdapterDescription 
adapter) {
+    AdapterDescription decryptedDescription = new 
Cloner().adapterDescription(adapter);
+    SecretProvider.getDecryptionService().apply(decryptedDescription);
+    return decryptedDescription;
+  }
+
   private IAdapterStorage getAdapterStorage() {
     return resourceManager.manageAdapters().getDb();
   }
diff --git 
a/streampipes-data-export/src/main/java/org/apache/streampipes/export/resolver/AdapterResolver.java
 
b/streampipes-data-export/src/main/java/org/apache/streampipes/export/resolver/AdapterResolver.java
index f2ba42bc6e..1de104d39e 100644
--- 
a/streampipes-data-export/src/main/java/org/apache/streampipes/export/resolver/AdapterResolver.java
+++ 
b/streampipes-data-export/src/main/java/org/apache/streampipes/export/resolver/AdapterResolver.java
@@ -104,7 +104,7 @@ public class AdapterResolver extends 
AbstractResolver<AdapterDescription> {
               new WorkerRestClient(extensionServiceRequestManager, 
resourceManager),
               getNoSqlStore().getExtensionsServiceStorage(),
               extensionServiceRequestManager
-          ).stopStreamAdapter(resourceId, true);
+          ).stopAdapter(resourceId, true);
         } catch (AdapterException e) {
           LOG.warn("Error when stopping adapter with id {} and name {}", 
resourceId, existingAdapter.getName());
         }
diff --git 
a/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/AdapterHealthCheck.java
 
b/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/AdapterHealthCheck.java
index aaac997e51..8181647ce5 100644
--- 
a/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/AdapterHealthCheck.java
+++ 
b/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/AdapterHealthCheck.java
@@ -161,7 +161,7 @@ public class AdapterHealthCheck {
       try {
         if (adapterDescription.isRunning()) {
           LOG.debug("Start recovering adapter {} ", 
adapterDescription.getElementId());
-          
this.healthCheckData.resourceProvider().adapterMasterManagement().startStreamAdapter(adapterDescription.getElementId());
+          
this.healthCheckData.resourceProvider().adapterMasterManagement().startAdapter(adapterDescription.getElementId());
           LOG.info("Adapter {} is recovered", 
adapterDescription.getElementId());
         }
       } catch (AdapterException e) {
diff --git 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/AdapterResource.java
 
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/AdapterResource.java
index bc2103895a..84d86461d0 100644
--- 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/AdapterResource.java
+++ 
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/AdapterResource.java
@@ -230,7 +230,7 @@ public class AdapterResource extends 
AbstractAdapterResource<AdapterMasterManage
     try {
       var adapter = getAdapterDescription(elementId);
       if (checkAdapterPermission(adapter, "WRITE")) {
-        managementService.stopStreamAdapter(elementId, forceStop);
+        managementService.stopAdapter(adapter, forceStop);
         return ok(Notifications.success("Adapter stopped"));
       } else {
         return unauthorized();
@@ -247,7 +247,7 @@ public class AdapterResource extends 
AbstractAdapterResource<AdapterMasterManage
     try {
       var adapterDescription = getAdapterDescription(elementId);
       if (checkAdapterPermission(adapterDescription, "WRITE")) {
-        managementService.startStreamAdapter(elementId);
+        managementService.startAdapter(adapterDescription);
         return ok(Notifications.success("Adapter started"));
       } else {
         return unauthorized();
diff --git 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/CompactAdapterResource.java
 
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/CompactAdapterResource.java
index 9d2cc76625..e8f88ae5e9 100644
--- 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/CompactAdapterResource.java
+++ 
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/CompactAdapterResource.java
@@ -144,7 +144,7 @@ public class CompactAdapterResource extends 
AbstractAdapterResource<AdapterMaste
         }
         if (compactAdapter.createOptions()
                           .start()) {
-          managementService.startStreamAdapter(adapterId);
+          managementService.startAdapter(adapterId);
         }
       }
       return ok(Notifications.success(adapterId));

Reply via email to