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

SvenO3 pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/streampipes.git


The following commit(s) were added to refs/heads/dev by this push:
     new 87ffcff1ae fix(#4618): Health check restarts adapter during stopping 
process (#4621)
87ffcff1ae is described below

commit 87ffcff1ae3343c503133860c266bef44a0d9223
Author: Sven Oehler <[email protected]>
AuthorDate: Mon Jun 29 14:04:05 2026 +0200

    fix(#4618): Health check restarts adapter during stopping process (#4621)
    
    Co-authored-by: Philipp Zehnder <[email protected]>
---
 .../connect/AdapterTransitionRegistry.java         |  80 +++++++++++++
 .../connect/AdapterWorkerManagement.java           |  81 +++++++------
 .../connect/AdapterWorkerRequestManagement.java    |   7 +-
 .../monitoring/HealthCheckManagement.java          |  53 +++++++--
 .../connect/AdapterWorkerManagementTest.java       |  17 ++-
 .../health/monitoring/AdapterHealthCheck.java      |  19 ++-
 .../ExtensionInstanceAvailabilityCheck.java        |   7 +-
 .../health/monitoring/AdapterHealthCheckTest.java  | 129 +++++++++++++++++++++
 ...stanceHealth.java => AdapterInstanceState.java} |   8 +-
 .../model/health/ExtensionInstanceHealth.java      |   3 +-
 10 files changed, 344 insertions(+), 60 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..0882754a57
--- /dev/null
+++ 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterTransitionRegistry.java
@@ -0,0 +1,80 @@
+/*
+ * 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;
+
+/**
+ * Tracks adapter instances that are currently being started or stopped in an 
extension service.
+ *
+ * <p>The health check compares the adapters that should be running, according 
to core storage, with
+ * the adapters reported by each extension service. Starting and stopping are 
short transitional
+ * phases where an adapter can temporarily be absent from the regular 
running-adapter registry, even
+ * though the extension service is already handling the requested lifecycle 
operation.
+ *
+ * <p>This registry was introduced to make those transitional states visible 
to the health check.
+ * Without it, a health check that runs at the same time as an adapter stop 
can interpret the
+ * temporary absence as a crashed adapter and start it again. Reporting the 
adapter as
+ * {@link AdapterInstanceState#STARTING} or {@link 
AdapterInstanceState#STOPPING} prevents this
+ * recovery race while keeping the normal running-adapter registry focused on 
fully running
+ * instances.
+ */
+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..4280723b5d 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,48 +19,79 @@
 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.Map;
+import java.util.Set;
 import java.util.stream.Collectors;
 
 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()
+    return new ExtensionInstanceHealth(
+        getAdapterInstanceStates(),
+        getRunningPipelineElementInstanceIds()
+    );
+  }
+
+  private Map<String, AdapterInstanceState> getAdapterInstanceStates() {
+    var adapterInstanceStates = getRunningAdapterInstanceStates();
+    setTransitioningAdapterInstanceStates(adapterInstanceStates);
+    return adapterInstanceStates;
+  }
+
+  private Map<String, AdapterInstanceState> getRunningAdapterInstanceStates() {
+    return adapterManagement.getAllRunningAdapterInstances()
         .stream()
         .map(NamedStreamPipesEntity::getElementId)
-        .collect(Collectors.toSet());
+        .collect(Collectors.toMap(
+            adapterInstanceId -> adapterInstanceId,
+            adapterInstanceId -> AdapterInstanceState.RUNNING,
+            (existingState, replacementState) -> existingState
+        ));
+  }
+
+  private void setTransitioningAdapterInstanceStates(Map<String, 
AdapterInstanceState> adapterInstanceStates) {
+    
adapterInstanceStates.putAll(adapterTransitionRegistry.getTransitioningAdapterInstanceStates());
+  }
 
-    var runningPipelineElementInstances = 
runningInstances.getRunningInstanceIds()
+  private Set<String> getRunningPipelineElementInstanceIds() {
+    return runningInstances.getRunningInstanceIds()
         .stream()
         .map(InstanceIdExtractor::extractId)
         .collect(Collectors.toSet());
-
-    return new ExtensionInstanceHealth(
-        runningAdapterInstances,
-        runningPipelineElementInstances
-    );
   }
 }
diff --git 
a/streampipes-extensions-management/src/test/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagementTest.java
 
b/streampipes-extensions-management/src/test/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagementTest.java
index 33f49575e3..418162f3f9 100644
--- 
a/streampipes-extensions-management/src/test/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagementTest.java
+++ 
b/streampipes-extensions-management/src/test/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagementTest.java
@@ -20,6 +20,7 @@ package org.apache.streampipes.extensions.management.connect;
 
 import org.apache.streampipes.commons.exceptions.connect.AdapterException;
 import org.apache.streampipes.extensions.management.init.IDeclarersSingleton;
+import org.apache.streampipes.model.health.AdapterInstanceState;
 import org.apache.streampipes.sdk.builder.adapter.AdapterConfigurationBuilder;
 
 import org.junit.jupiter.api.Test;
@@ -27,6 +28,7 @@ import org.junit.jupiter.api.Test;
 import java.util.Optional;
 
 import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.when;
@@ -38,14 +40,25 @@ public class AdapterWorkerManagementTest {
     var adapterDescription = AdapterConfigurationBuilder
         .create("id", 0,  null)
         .build();
+    adapterDescription.setElementId("adapter-id");
+    var adapterTransitionRegistry = new AdapterTransitionRegistry();
 
     var declarerSingleton = mock(IDeclarersSingleton.class);
-    when(declarerSingleton.getAdapter(any())).thenReturn(Optional.empty());
+    when(declarerSingleton.getAdapter(any())).thenAnswer(invocation -> {
+      assertTrue(adapterTransitionRegistry
+          .getTransitioningAdapterInstanceStates()
+          .containsKey("adapter-id"));
+      assertTrue(adapterTransitionRegistry
+          .getTransitioningAdapterInstanceStates()
+          .containsValue(AdapterInstanceState.STARTING));
+      return Optional.empty();
+    });
 
     var adapterWorkerManagement = new AdapterWorkerManagement(
-        null, declarerSingleton);
+        null, declarerSingleton, adapterTransitionRegistry);
 
     assertThrows(AdapterException.class, () -> 
adapterWorkerManagement.invokeAdapter(adapterDescription));
+    
assertTrue(adapterTransitionRegistry.getTransitioningAdapterInstanceStates().isEmpty());
 
   }
 }
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-health-monitoring/src/test/java/org/apache/streampipes/health/monitoring/AdapterHealthCheckTest.java
 
b/streampipes-health-monitoring/src/test/java/org/apache/streampipes/health/monitoring/AdapterHealthCheckTest.java
new file mode 100644
index 0000000000..47207d027d
--- /dev/null
+++ 
b/streampipes-health-monitoring/src/test/java/org/apache/streampipes/health/monitoring/AdapterHealthCheckTest.java
@@ -0,0 +1,129 @@
+/*
+ * 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.health.monitoring;
+
+import org.apache.streampipes.health.monitoring.model.ActiveResources;
+import org.apache.streampipes.health.monitoring.model.HealthCheckData;
+import org.apache.streampipes.model.connect.adapter.AdapterDescription;
+import org.apache.streampipes.model.health.AdapterInstanceState;
+import org.apache.streampipes.model.health.ExtensionInstanceHealth;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+public class AdapterHealthCheckTest {
+
+  @Test
+  public void getAdaptersToRecoverIgnoresTransitioningAdapters() {
+    var adapterDescription = new AdapterDescription();
+    adapterDescription.setElementId("adapter-id");
+    adapterDescription.setRunning(true);
+
+    var activeResources = new ActiveResources(
+        List.of(),
+        List.of(),
+        List.of(adapterDescription),
+        List.of(adapterDescription)
+    );
+    var extensionInstanceHealth = new ExtensionInstanceHealth(
+        Map.of("adapter-id", AdapterInstanceState.STOPPING),
+        Set.of()
+    );
+    var healthCheckData = new HealthCheckData(
+        null,
+        activeResources,
+        Map.of(),
+        Map.of("service-id", extensionInstanceHealth)
+    );
+
+    var adaptersToRecover = new 
AdapterHealthCheck(healthCheckData).getAdaptersToRecover();
+
+    assertTrue(adaptersToRecover.isEmpty());
+  }
+
+  @Test
+  public void runCheckDoesNotRestartStartingAdapters() {
+    var adapterHealthCheck = new TestAdapterHealthCheck(
+        healthCheckData(AdapterInstanceState.STARTING)
+    );
+
+    adapterHealthCheck.runCheck();
+
+    assertTrue(adapterHealthCheck.recoveredAdapters.isEmpty());
+  }
+
+  @Test
+  public void runCheckDoesNotRestartStoppingAdapters() {
+    var adapterHealthCheck = new TestAdapterHealthCheck(
+        healthCheckData(AdapterInstanceState.STOPPING)
+    );
+
+    adapterHealthCheck.runCheck();
+
+    assertTrue(adapterHealthCheck.recoveredAdapters.isEmpty());
+  }
+
+  private HealthCheckData healthCheckData(AdapterInstanceState 
adapterInstanceState) {
+    var adapterDescription = new AdapterDescription();
+    adapterDescription.setElementId("adapter-id");
+    adapterDescription.setRunning(true);
+
+    var activeResources = new ActiveResources(
+        List.of(),
+        List.of(),
+        List.of(adapterDescription),
+        List.of(adapterDescription)
+    );
+    var extensionInstanceHealth = new ExtensionInstanceHealth(
+        Map.of("adapter-id", adapterInstanceState),
+        Set.of()
+    );
+
+    return new HealthCheckData(
+        null,
+        activeResources,
+        Map.of(),
+        Map.of("service-id", extensionInstanceHealth)
+    );
+  }
+
+  private static class TestAdapterHealthCheck extends AdapterHealthCheck {
+
+    private List<AdapterDescription> recoveredAdapters = List.of();
+
+    TestAdapterHealthCheck(HealthCheckData healthCheckData) {
+      super(healthCheckData);
+    }
+
+    @Override
+    protected void updateMonitoringMetrics(List<AdapterDescription> 
runningAdapterDescriptions) {
+
+    }
+
+    @Override
+    public void recoverAdapters(List<AdapterDescription> adaptersToRecover) {
+      this.recoveredAdapters = adaptersToRecover;
+    }
+  }
+}
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