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