This is an automated email from the ASF dual-hosted git repository.
SvenO3 pushed a commit to branch allow-registering-new-health-checks
in repository https://gitbox.apache.org/repos/asf/streampipes.git
The following commit(s) were added to
refs/heads/allow-registering-new-health-checks by this push:
new 79c4a935bf Unify health checks
79c4a935bf is described below
commit 79c4a935bfacf54d267f5bbd3cb1e9f8fb2a9cab
Author: Sven Oehler <[email protected]>
AuthorDate: Tue Jun 30 11:28:58 2026 +0200
Unify health checks
---
.../health/monitoring/AdapterHealthCheck.java | 3 +-
.../health/monitoring/ExtensionHealthCheck.java | 34 +++++++++++++++++-----
...dExtensionHealthCheck.java => HealthCheck.java} | 2 +-
.../health/monitoring/PipelineHealthCheck.java | 3 +-
.../streampipes/service/core/PostStartupTask.java | 4 +--
.../service/core/StreamPipesCoreApplication.java | 4 +--
6 files changed, 35 insertions(+), 15 deletions(-)
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 aaac997e51..cd8ceae390 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
@@ -34,7 +34,7 @@ import java.util.NoSuchElementException;
import java.util.Objects;
import java.util.stream.Collectors;
-public class AdapterHealthCheck {
+public class AdapterHealthCheck implements HealthCheck {
private static final Logger LOG =
LoggerFactory.getLogger(AdapterHealthCheck.class);
@@ -51,6 +51,7 @@ public class AdapterHealthCheck {
* running adapters (in line with
* {@link PipelineHealthCheck}).
*/
+ @Override
public void runCheck() {
LOG.debug("Adapter health check started");
diff --git
a/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/ExtensionHealthCheck.java
b/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/ExtensionHealthCheck.java
index 9f02ea7e62..265cca67ea 100644
---
a/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/ExtensionHealthCheck.java
+++
b/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/ExtensionHealthCheck.java
@@ -27,6 +27,7 @@ import
org.apache.streampipes.storage.api.system.IExtensionsServiceStorage;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
@@ -38,13 +39,13 @@ public class ExtensionHealthCheck implements Runnable {
private final ExtensionServiceRequestManager extensionRequestManager;
private final IExtensionsServiceStorage extensionsServiceStorage;
private final SpResourceManager resourceManager;
- private final List<RegisteredExtensionHealthCheck> registeredHealthChecks;
+ private final List<HealthCheck> registeredHealthChecks;
public ExtensionHealthCheck(ResourceProvider resourceProvider,
IExtensionsServiceStorage
extensionsServiceStorage,
ExtensionServiceRequestManager
extensionRequestManager,
SpResourceManager resourceManager,
- List<RegisteredExtensionHealthCheck>
registeredHealthChecks) {
+ List<HealthCheck> registeredHealthChecks) {
this.resourceProvider = resourceProvider;
this.extensionsServiceStorage = extensionsServiceStorage;
this.extensionRequestManager = extensionRequestManager;
@@ -67,17 +68,21 @@ public class ExtensionHealthCheck implements Runnable {
).checkRunningInstances());
});
- var healthCheckData = new HealthCheckData(resourceProvider,
activeResources, activeCoreInstances, activeExtensionInstances);
- new PipelineHealthCheck(healthCheckData, extensionRequestManager,
resourceProvider, resourceManager).runCheck();
- new AdapterHealthCheck(healthCheckData).runCheck();
- runRegisteredHealthChecks();
+ var healthCheckData = new HealthCheckData(
+ resourceProvider,
+ activeResources,
+ activeCoreInstances,
+ activeExtensionInstances
+ );
+ var healthChecks = getHealthChecks(healthCheckData);
+ runHealthChecks(healthChecks);
} catch (Exception e) {
LOG.warn("An unhandled error occurred while running health check.", e);
}
}
- private void runRegisteredHealthChecks() {
- registeredHealthChecks.forEach(healthCheck -> {
+ private void runHealthChecks(List<HealthCheck> healthChecks) {
+ healthChecks.forEach(healthCheck -> {
try {
healthCheck.runCheck();
} catch (Exception e) {
@@ -86,4 +91,17 @@ public class ExtensionHealthCheck implements Runnable {
}
});
}
+
+ private List<HealthCheck> getHealthChecks(HealthCheckData healthCheckData) {
+ var healthChecks = new
ArrayList<>(getBuiltInHealthChecks(healthCheckData));
+ healthChecks.addAll(registeredHealthChecks);
+ return healthChecks;
+ }
+
+ protected List<HealthCheck> getBuiltInHealthChecks(HealthCheckData
healthCheckData) {
+ return List.of(
+ new PipelineHealthCheck(healthCheckData, extensionRequestManager,
resourceProvider, resourceManager),
+ new AdapterHealthCheck(healthCheckData)
+ );
+ }
}
diff --git
a/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/RegisteredExtensionHealthCheck.java
b/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/HealthCheck.java
similarity index 94%
rename from
streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/RegisteredExtensionHealthCheck.java
rename to
streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/HealthCheck.java
index 564555afa0..9cf5cf74e8 100644
---
a/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/RegisteredExtensionHealthCheck.java
+++
b/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/HealthCheck.java
@@ -18,7 +18,7 @@
package org.apache.streampipes.health.monitoring;
-public interface RegisteredExtensionHealthCheck {
+public interface HealthCheck {
void runCheck();
}
diff --git
a/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/PipelineHealthCheck.java
b/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/PipelineHealthCheck.java
index ad1bef456b..f3d88fad16 100644
---
a/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/PipelineHealthCheck.java
+++
b/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/PipelineHealthCheck.java
@@ -48,7 +48,7 @@ import java.util.Objects;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.stream.Stream;
-public class PipelineHealthCheck {
+public class PipelineHealthCheck implements HealthCheck {
private static final Logger LOG =
LoggerFactory.getLogger(PipelineHealthCheck.class);
private static final int MAX_FAILED_ATTEMPTS = 10;
@@ -71,6 +71,7 @@ public class PipelineHealthCheck {
this.resourceManager = resourceManager;
}
+ @Override
public void runCheck() {
try {
initPipelineMetrics();
diff --git
a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/PostStartupTask.java
b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/PostStartupTask.java
index 2601af9ad4..bf1078f941 100644
---
a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/PostStartupTask.java
+++
b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/PostStartupTask.java
@@ -24,7 +24,7 @@ import
org.apache.streampipes.connect.management.management.WorkerAdministration
import org.apache.streampipes.connect.management.management.WorkerRestClient;
import org.apache.streampipes.health.monitoring.ExtensionHealthCheck;
import org.apache.streampipes.health.monitoring.PostStartupRecovery;
-import org.apache.streampipes.health.monitoring.RegisteredExtensionHealthCheck;
+import org.apache.streampipes.health.monitoring.HealthCheck;
import org.apache.streampipes.health.monitoring.ResourceProvider;
import org.apache.streampipes.health.monitoring.ServiceHealthCheck;
import
org.apache.streampipes.manager.api.extensions.ExtensionServiceRequestManager;
@@ -68,7 +68,7 @@ public class PostStartupTask implements Runnable {
ExtensionServiceRequestManager
extensionServiceRequestManager,
WorkerRestClient workerRestClient,
SpResourceManager resourceManager,
- List<RegisteredExtensionHealthCheck>
registeredHealthChecks) {
+ List<HealthCheck> registeredHealthChecks) {
this.pipelineStorage = pipelineStorage;
this.extensionServiceRequestManager = extensionServiceRequestManager;
this.executorService = Executors.newSingleThreadScheduledExecutor();
diff --git
a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/StreamPipesCoreApplication.java
b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/StreamPipesCoreApplication.java
index 09767a2100..42bd160b7c 100644
---
a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/StreamPipesCoreApplication.java
+++
b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/StreamPipesCoreApplication.java
@@ -25,7 +25,7 @@ import
org.apache.streampipes.connect.transformer.api.TransformationEngine;
import org.apache.streampipes.connect.transformer.api.TransformationEngines;
import org.apache.streampipes.connect.transformer.js.GraalJsScriptEngine;
import org.apache.streampipes.health.monitoring.ExtensionHealthCheck;
-import org.apache.streampipes.health.monitoring.RegisteredExtensionHealthCheck;
+import org.apache.streampipes.health.monitoring.HealthCheck;
import org.apache.streampipes.health.monitoring.ResourceProvider;
import org.apache.streampipes.health.monitoring.ServiceHealthCheck;
import org.apache.streampipes.loadbalance.LoadManager;
@@ -251,7 +251,7 @@ public class StreamPipesCoreApplication extends
StreamPipesServiceBase {
return new AvailableMigrations(resourceManager).getAvailableMigrations();
}
- protected List<RegisteredExtensionHealthCheck>
getRegisteredExtensionHealthChecks() {
+ protected List<HealthCheck> getRegisteredExtensionHealthChecks() {
return List.of();
}