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 6f0300c2ae Update tests
6f0300c2ae is described below

commit 6f0300c2ae064482cca72f1a36a15489b9fc10e2
Author: Sven Oehler <[email protected]>
AuthorDate: Tue Jun 23 15:41:05 2026 +0200

    Update tests
---
 .../connect/AdapterWorkerManagementTest.java       |  17 ++-
 .../health/monitoring/AdapterHealthCheckTest.java  | 129 +++++++++++++++++++++
 2 files changed, 144 insertions(+), 2 deletions(-)

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/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;
+    }
+  }
+}

Reply via email to