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