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();
   }
 

Reply via email to