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) {
}