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

SvenO3 pushed a commit to branch 
4618-health-check-restarts-adapter-during-stopping-process
in repository https://gitbox.apache.org/repos/asf/streampipes.git


The following commit(s) were added to 
refs/heads/4618-health-check-restarts-adapter-during-stopping-process by this 
push:
     new 4d1a973d20 Add adapter transition registry
4d1a973d20 is described below

commit 4d1a973d20a74a0dfa86a8e95829bfd353781a58
Author: Sven Oehler <[email protected]>
AuthorDate: Tue Jun 23 15:17:09 2026 +0200

    Add adapter transition registry
---
 .../connect/AdapterTransitionRegistry.java         | 65 +++++++++++++++++
 .../connect/AdapterWorkerManagement.java           | 81 +++++++++++++---------
 .../connect/AdapterWorkerRequestManagement.java    |  7 +-
 .../monitoring/HealthCheckManagement.java          | 28 ++++++--
 .../health/monitoring/AdapterHealthCheck.java      | 19 +++--
 .../ExtensionInstanceAvailabilityCheck.java        |  7 +-
 ...stanceHealth.java => AdapterInstanceState.java} |  8 +--
 .../model/health/ExtensionInstanceHealth.java      |  3 +-
 8 files changed, 165 insertions(+), 53 deletions(-)

diff --git 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterTransitionRegistry.java
 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterTransitionRegistry.java
new file mode 100644
index 0000000000..e994378174
--- /dev/null
+++ 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterTransitionRegistry.java
@@ -0,0 +1,65 @@
+/*
+ * 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.extensions.management.connect;
+
+import org.apache.streampipes.model.health.AdapterInstanceState;
+
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+
+public class AdapterTransitionRegistry {
+
+  public static final AdapterTransitionRegistry INSTANCE = new 
AdapterTransitionRegistry();
+
+  private final Map<String, AdapterInstanceState> 
transitioningAdapterInstanceStates = new ConcurrentHashMap<>();
+
+  public void registerStarting(String adapterInstanceId) {
+    register(adapterInstanceId, AdapterInstanceState.STARTING);
+  }
+
+  public void deregisterStarting(String adapterInstanceId) {
+    deregister(adapterInstanceId, AdapterInstanceState.STARTING);
+  }
+
+  public void registerStopping(String adapterInstanceId) {
+    register(adapterInstanceId, AdapterInstanceState.STOPPING);
+  }
+
+  public void deregisterStopping(String adapterInstanceId) {
+    deregister(adapterInstanceId, AdapterInstanceState.STOPPING);
+  }
+
+  public Map<String, AdapterInstanceState> 
getTransitioningAdapterInstanceStates() {
+    return Map.copyOf(transitioningAdapterInstanceStates);
+  }
+
+  private void register(String adapterInstanceId,
+                        AdapterInstanceState adapterInstanceState) {
+    if (adapterInstanceId != null) {
+      transitioningAdapterInstanceStates.put(adapterInstanceId, 
adapterInstanceState);
+    }
+  }
+
+  private void deregister(String adapterInstanceId,
+                          AdapterInstanceState adapterInstanceState) {
+    if (adapterInstanceId != null) {
+      transitioningAdapterInstanceStates.remove(adapterInstanceId, 
adapterInstanceState);
+    }
+  }
+}
diff --git 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagement.java
 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagement.java
index 3b495a3cb8..4df5066ff1 100644
--- 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagement.java
+++ 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagement.java
@@ -41,11 +41,14 @@ public class AdapterWorkerManagement {
 
   private final RunningAdapterInstances runningAdapterInstances;
   private final IDeclarersSingleton declarers;
+  private final AdapterTransitionRegistry adapterTransitionRegistry;
 
   public AdapterWorkerManagement(RunningAdapterInstances 
runningAdapterInstances,
-                                 IDeclarersSingleton declarers) {
+                                 IDeclarersSingleton declarers,
+                                 AdapterTransitionRegistry 
adapterTransitionRegistry) {
     this.runningAdapterInstances = runningAdapterInstances;
     this.declarers = declarers;
+    this.adapterTransitionRegistry = adapterTransitionRegistry;
   }
 
   public Collection<AdapterDescription> getAllRunningAdapterInstances() {
@@ -53,31 +56,36 @@ public class AdapterWorkerManagement {
   }
 
   public void invokeAdapter(AdapterDescription adapterDescription) throws 
AdapterException {
-    var adapter = declarers
-        .getAdapter(adapterDescription.getAppId());
-
-    if (adapter.isPresent()) {
-      var newAdapterInstance = 
adapter.get().declareConfig().getSupplier().get();
-      runningAdapterInstances.addAdapter(
-          adapterDescription.getElementId(),
-          newAdapterInstance,
-          adapterDescription);
-
-      // This method allows adapters to modify the adapter description prior 
to invocation.
-      // It is particularly useful for adapters like FileReplayAdapter that 
need to manipulate timestamp values
-      // internally, bypassing the adapter preprocessing pipeline.
-      newAdapterInstance.preprocessAdapterDescription(adapterDescription);
-
-      var registeredParsers = 
newAdapterInstance.declareConfig().getSupportedParsers();
-      var extractor = AdapterParameterExtractor.from(adapterDescription, 
registeredParsers);
-      var runtimeContext = 
makeRuntimeContext(adapterDescription.getElementId());
-      var eventCollector = EventCollector.from(adapterDescription, 
runtimeContext);
-
-      newAdapterInstance.onAdapterStarted(extractor, eventCollector, 
runtimeContext);
-    } else {
-      var errorMessage = "Adapter with id %s could not be 
found".formatted(adapterDescription.getAppId());
-      LOG.error(errorMessage);
-      throw new AdapterException(errorMessage);
+    
adapterTransitionRegistry.registerStarting(adapterDescription.getElementId());
+    try {
+      var adapter = declarers
+          .getAdapter(adapterDescription.getAppId());
+
+      if (adapter.isPresent()) {
+        var newAdapterInstance = 
adapter.get().declareConfig().getSupplier().get();
+        runningAdapterInstances.addAdapter(
+            adapterDescription.getElementId(),
+            newAdapterInstance,
+            adapterDescription);
+
+        // This method allows adapters to modify the adapter description prior 
to invocation.
+        // It is particularly useful for adapters like FileReplayAdapter that 
need to manipulate timestamp values
+        // internally, bypassing the adapter preprocessing pipeline.
+        newAdapterInstance.preprocessAdapterDescription(adapterDescription);
+
+        var registeredParsers = 
newAdapterInstance.declareConfig().getSupportedParsers();
+        var extractor = AdapterParameterExtractor.from(adapterDescription, 
registeredParsers);
+        var runtimeContext = 
makeRuntimeContext(adapterDescription.getElementId());
+        var eventCollector = EventCollector.from(adapterDescription, 
runtimeContext);
+
+        newAdapterInstance.onAdapterStarted(extractor, eventCollector, 
runtimeContext);
+      } else {
+        var errorMessage = "Adapter with id %s could not be 
found".formatted(adapterDescription.getAppId());
+        LOG.error(errorMessage);
+        throw new AdapterException(errorMessage);
+      }
+    } finally {
+      
adapterTransitionRegistry.deregisterStarting(adapterDescription.getElementId());
     }
   }
 
@@ -85,17 +93,22 @@ public class AdapterWorkerManagement {
 
     String elementId = adapterDescription.getElementId();
 
-    StreamPipesAdapter adapter = 
RunningAdapterInstances.INSTANCE.removeAdapter(elementId);
+    adapterTransitionRegistry.registerStopping(elementId);
+    try {
+      StreamPipesAdapter adapter = 
RunningAdapterInstances.INSTANCE.removeAdapter(elementId);
 
-    if (adapter != null) {
+      if (adapter != null) {
 
-      var registeredParsers = adapter.declareConfig().getSupportedParsers();
-      var extractor = AdapterParameterExtractor.from(adapterDescription, 
registeredParsers);
-      var runtimeContext = makeRuntimeContext(elementId);
-      adapter.onAdapterStopped(extractor, runtimeContext);
-    }
+        var registeredParsers = adapter.declareConfig().getSupportedParsers();
+        var extractor = AdapterParameterExtractor.from(adapterDescription, 
registeredParsers);
+        var runtimeContext = makeRuntimeContext(elementId);
+        adapter.onAdapterStopped(extractor, runtimeContext);
+      }
 
-    resetMonitoring(elementId);
+      resetMonitoring(elementId);
+    } finally {
+      adapterTransitionRegistry.deregisterStopping(elementId);
+    }
   }
 
   private IAdapterRuntimeContext makeRuntimeContext(String adapterInstanceId) {
diff --git 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerRequestManagement.java
 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerRequestManagement.java
index f6a5dc4053..a5c099c7f2 100644
--- 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerRequestManagement.java
+++ 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerRequestManagement.java
@@ -37,9 +37,14 @@ public class AdapterWorkerRequestManagement {
   private final AdapterWorkerManagement adapterManagement;
 
   public AdapterWorkerRequestManagement() {
+    this(AdapterTransitionRegistry.INSTANCE);
+  }
+
+  public AdapterWorkerRequestManagement(AdapterTransitionRegistry 
adapterTransitionRegistry) {
     this(new AdapterWorkerManagement(
         RunningAdapterInstances.INSTANCE,
-        DeclarersSingleton.getInstance()
+        DeclarersSingleton.getInstance(),
+        adapterTransitionRegistry
     ));
   }
 
diff --git 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/monitoring/HealthCheckManagement.java
 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/monitoring/HealthCheckManagement.java
index 1512fd854d..6130e1cf74 100644
--- 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/monitoring/HealthCheckManagement.java
+++ 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/monitoring/HealthCheckManagement.java
@@ -19,11 +19,13 @@
 package org.apache.streampipes.extensions.management.monitoring;
 
 import org.apache.streampipes.commons.constants.InstanceIdExtractor;
+import 
org.apache.streampipes.extensions.management.connect.AdapterTransitionRegistry;
 import 
org.apache.streampipes.extensions.management.connect.AdapterWorkerManagement;
 import org.apache.streampipes.extensions.management.init.DeclarersSingleton;
 import 
org.apache.streampipes.extensions.management.init.RunningAdapterInstances;
 import org.apache.streampipes.extensions.management.init.RunningInstances;
 import org.apache.streampipes.model.base.NamedStreamPipesEntity;
+import org.apache.streampipes.model.health.AdapterInstanceState;
 import org.apache.streampipes.model.health.ExtensionInstanceHealth;
 
 import java.util.stream.Collectors;
@@ -32,26 +34,40 @@ public class HealthCheckManagement {
 
   private final AdapterWorkerManagement adapterManagement;
   private final RunningInstances runningInstances;
+  private final AdapterTransitionRegistry adapterTransitionRegistry;
 
   public HealthCheckManagement() {
+    this(AdapterTransitionRegistry.INSTANCE);
+  }
+
+  public HealthCheckManagement(AdapterTransitionRegistry 
adapterTransitionRegistry) {
     this(new AdapterWorkerManagement(
              RunningAdapterInstances.INSTANCE,
-             DeclarersSingleton.getInstance()
+             DeclarersSingleton.getInstance(),
+             adapterTransitionRegistry
          ),
-         RunningInstances.INSTANCE);
+         RunningInstances.INSTANCE,
+         adapterTransitionRegistry);
   }
 
   public HealthCheckManagement(AdapterWorkerManagement adapterManagement,
-                               RunningInstances runningInstances) {
+                               RunningInstances runningInstances,
+                               AdapterTransitionRegistry 
adapterTransitionRegistry) {
     this.adapterManagement = adapterManagement;
     this.runningInstances = runningInstances;
+    this.adapterTransitionRegistry = adapterTransitionRegistry;
   }
 
   public ExtensionInstanceHealth getExtensionInstanceHealth() {
-    var runningAdapterInstances = 
adapterManagement.getAllRunningAdapterInstances()
+    var adapterInstanceStates = 
adapterManagement.getAllRunningAdapterInstances()
         .stream()
         .map(NamedStreamPipesEntity::getElementId)
-        .collect(Collectors.toSet());
+        .collect(Collectors.toMap(
+            adapterInstanceId -> adapterInstanceId,
+            adapterInstanceId -> AdapterInstanceState.RUNNING,
+            (existingState, replacementState) -> existingState
+        ));
+    
adapterInstanceStates.putAll(adapterTransitionRegistry.getTransitioningAdapterInstanceStates());
 
     var runningPipelineElementInstances = 
runningInstances.getRunningInstanceIds()
         .stream()
@@ -59,7 +75,7 @@ public class HealthCheckManagement {
         .collect(Collectors.toSet());
 
     return new ExtensionInstanceHealth(
-        runningAdapterInstances,
+        adapterInstanceStates,
         runningPipelineElementInstances
     );
   }
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 991577d5da..aaac997e51 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
@@ -23,11 +23,13 @@ import 
org.apache.streampipes.commons.prometheus.adapter.AdapterMetricsManager;
 import org.apache.streampipes.health.monitoring.model.HealthCheckData;
 import org.apache.streampipes.loadbalance.pipeline.ExtensionsLogProvider;
 import org.apache.streampipes.model.connect.adapter.AdapterDescription;
+import org.apache.streampipes.model.health.AdapterInstanceState;
 
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 import java.util.List;
+import java.util.Map;
 import java.util.NoSuchElementException;
 import java.util.Objects;
 import java.util.stream.Collectors;
@@ -128,22 +130,31 @@ public class AdapterHealthCheck {
    */
   public List<AdapterDescription> getAdaptersToRecover() {
 
-    var runningAdapterIds =
+    var adapterInstanceStates =
         healthCheckData.activeExtensionInstances()
             .values()
             .stream()
-            .flatMap(h -> h.runningAdapterInstanceIds().stream())
-            .collect(Collectors.toSet());
+            .flatMap(h -> 
adapterInstanceStates(h.adapterInstanceStates()).entrySet().stream())
+            .collect(Collectors.toMap(
+                Map.Entry::getKey,
+                Map.Entry::getValue,
+                (existingState, replacementState) -> existingState
+            ));
 
     return healthCheckData.activeResources()
         .runningAdapters()
         .stream()
         .filter(Objects::nonNull)
         .filter(a -> a.getElementId() != null)
-        .filter(a -> !runningAdapterIds.contains(a.getElementId()))
+        .filter(a -> !adapterInstanceStates.containsKey(a.getElementId()))
         .toList();
   }
 
+  private Map<String, AdapterInstanceState> adapterInstanceStates(
+      Map<String, AdapterInstanceState> adapterInstanceStates) {
+    return adapterInstanceStates == null ? Map.of() : adapterInstanceStates;
+  }
+
   public void recoverAdapters(List<AdapterDescription> adaptersToRecover) {
     for (AdapterDescription adapterDescription : adaptersToRecover) {
       // Invoke all adapters that were running when the adapter container was 
stopped
diff --git 
a/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/ExtensionInstanceAvailabilityCheck.java
 
b/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/ExtensionInstanceAvailabilityCheck.java
index 5966344be6..69124efd72 100644
--- 
a/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/ExtensionInstanceAvailabilityCheck.java
+++ 
b/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/ExtensionInstanceAvailabilityCheck.java
@@ -32,6 +32,7 @@ import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 import java.io.IOException;
+import java.util.Map;
 import java.util.Set;
 
 public class ExtensionInstanceAvailabilityCheck {
@@ -60,21 +61,21 @@ public class ExtensionInstanceAvailabilityCheck {
           .findFirst();
 
       if (service.isEmpty()) {
-        return new ExtensionInstanceHealth(Set.of(), Set.of());
+        return new ExtensionInstanceHealth(Map.of(), Set.of());
       } else {
         var response = extensionRequestManager.request(
             ExtensionServiceRequests
                 .extensionInstanceHealth(makeRequestTarget(service.get()), 
resourceManager)
         );
         if (response.statusCode() != 200) {
-          return new ExtensionInstanceHealth(Set.of(), Set.of());
+          return new ExtensionInstanceHealth(Map.of(), Set.of());
         }
         return deserialize(response.responseBody());
       }
 
     } catch (IOException e) {
       LOG.error("Extension service {} is unavailable", serviceId);
-      return new ExtensionInstanceHealth(Set.of(), Set.of());
+      return new ExtensionInstanceHealth(Map.of(), Set.of());
     }
   }
 
diff --git 
a/streampipes-model/src/main/java/org/apache/streampipes/model/health/ExtensionInstanceHealth.java
 
b/streampipes-model/src/main/java/org/apache/streampipes/model/health/AdapterInstanceState.java
similarity index 82%
copy from 
streampipes-model/src/main/java/org/apache/streampipes/model/health/ExtensionInstanceHealth.java
copy to 
streampipes-model/src/main/java/org/apache/streampipes/model/health/AdapterInstanceState.java
index 2adf8a0839..97e0f1c367 100644
--- 
a/streampipes-model/src/main/java/org/apache/streampipes/model/health/ExtensionInstanceHealth.java
+++ 
b/streampipes-model/src/main/java/org/apache/streampipes/model/health/AdapterInstanceState.java
@@ -18,8 +18,8 @@
 
 package org.apache.streampipes.model.health;
 
-import java.util.Set;
-
-public record ExtensionInstanceHealth(Set<String> runningAdapterInstanceIds,
-                                      Set<String> 
runningPipelineElementInstanceIds) {
+public enum AdapterInstanceState {
+  RUNNING,
+  STARTING,
+  STOPPING
 }
diff --git 
a/streampipes-model/src/main/java/org/apache/streampipes/model/health/ExtensionInstanceHealth.java
 
b/streampipes-model/src/main/java/org/apache/streampipes/model/health/ExtensionInstanceHealth.java
index 2adf8a0839..859a644144 100644
--- 
a/streampipes-model/src/main/java/org/apache/streampipes/model/health/ExtensionInstanceHealth.java
+++ 
b/streampipes-model/src/main/java/org/apache/streampipes/model/health/ExtensionInstanceHealth.java
@@ -18,8 +18,9 @@
 
 package org.apache.streampipes.model.health;
 
+import java.util.Map;
 import java.util.Set;
 
-public record ExtensionInstanceHealth(Set<String> runningAdapterInstanceIds,
+public record ExtensionInstanceHealth(Map<String, AdapterInstanceState> 
adapterInstanceStates,
                                       Set<String> 
runningPipelineElementInstanceIds) {
 }

Reply via email to