This is an automated email from the ASF dual-hosted git repository.

SvenO3 pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/streampipes.git


The following commit(s) were added to refs/heads/dev by this push:
     new 8dd5b58b2d feat: Allow registration of new health checks (#4667)
8dd5b58b2d is described below

commit 8dd5b58b2d5e3815686cf216c32d2829a4ea0ae1
Author: Sven Oehler <[email protected]>
AuthorDate: Tue Jun 30 16:12:42 2026 +0200

    feat: Allow registration of new health checks (#4667)
---
 .../health/monitoring/AdapterHealthCheck.java      |  3 +-
 .../health/monitoring/ExtensionHealthCheck.java    | 42 +++++++++++++++++++---
 .../streampipes/health/monitoring/HealthCheck.java | 24 +++++++++++++
 .../health/monitoring/PipelineHealthCheck.java     |  3 +-
 .../streampipes/service/core/PostStartupTask.java  |  7 ++--
 .../service/core/StreamPipesCoreApplication.java   | 11 ++++--
 6 files changed, 80 insertions(+), 10 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 9c5ea287fb..4fb8cb5bb2 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
@@ -35,7 +35,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);
 
@@ -52,6 +52,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 77cce0a29a..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,7 +27,9 @@ 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;
 
 public class ExtensionHealthCheck implements Runnable {
 
@@ -37,15 +39,18 @@ public class ExtensionHealthCheck implements Runnable {
   private final ExtensionServiceRequestManager extensionRequestManager;
   private final IExtensionsServiceStorage extensionsServiceStorage;
   private final SpResourceManager resourceManager;
+  private final List<HealthCheck> registeredHealthChecks;
 
   public ExtensionHealthCheck(ResourceProvider resourceProvider,
                               IExtensionsServiceStorage 
extensionsServiceStorage,
                               ExtensionServiceRequestManager 
extensionRequestManager,
-                              SpResourceManager resourceManager) {
+                              SpResourceManager resourceManager,
+                              List<HealthCheck> registeredHealthChecks) {
     this.resourceProvider = resourceProvider;
     this.extensionsServiceStorage = extensionsServiceStorage;
     this.extensionRequestManager = extensionRequestManager;
     this.resourceManager = resourceManager;
+    this.registeredHealthChecks = registeredHealthChecks;
   }
 
   @Override
@@ -63,11 +68,40 @@ 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();
+      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 runHealthChecks(List<HealthCheck> healthChecks) {
+    healthChecks.forEach(healthCheck -> {
+      try {
+        healthCheck.runCheck();
+      } catch (Exception e) {
+        LOG.warn("An unhandled error occurred while running registered health 
check {}.",
+            healthCheck.getClass().getSimpleName(), e);
+      }
+    });
+  }
+
+  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/HealthCheck.java
 
b/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/HealthCheck.java
new file mode 100644
index 0000000000..9cf5cf74e8
--- /dev/null
+++ 
b/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/HealthCheck.java
@@ -0,0 +1,24 @@
+/*
+ * 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;
+
+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 dd576ef3eb..ccd7ee61a2 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
@@ -23,6 +23,7 @@ import 
org.apache.streampipes.connect.management.management.AdapterMasterManagem
 import 
org.apache.streampipes.connect.management.management.WorkerAdministrationManagement;
 import org.apache.streampipes.connect.management.management.WorkerRestClient;
 import org.apache.streampipes.health.monitoring.ExtensionHealthCheck;
+import org.apache.streampipes.health.monitoring.HealthCheck;
 import org.apache.streampipes.health.monitoring.PostStartupRecovery;
 import org.apache.streampipes.health.monitoring.ResourceProvider;
 import org.apache.streampipes.health.monitoring.ServiceHealthCheck;
@@ -66,7 +67,8 @@ public class PostStartupTask implements Runnable {
   public PostStartupTask(IPipelineStorage pipelineStorage,
                          ExtensionServiceRequestManager 
extensionServiceRequestManager,
                          WorkerRestClient workerRestClient,
-                         SpResourceManager resourceManager) {
+                         SpResourceManager resourceManager,
+                         List<HealthCheck> registeredHealthChecks) {
     this.pipelineStorage = pipelineStorage;
     this.extensionServiceRequestManager = extensionServiceRequestManager;
     this.executorService = Executors.newSingleThreadScheduledExecutor();
@@ -90,7 +92,8 @@ public class PostStartupTask implements Runnable {
             ),
             
StorageDispatcher.INSTANCE.getNoSqlStore().getExtensionsServiceStorage(),
             extensionServiceRequestManager,
-            resourceManager
+            resourceManager,
+            registeredHealthChecks
         )
     );
   }
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 ab3570a0ea..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,6 +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.HealthCheck;
 import org.apache.streampipes.health.monitoring.ResourceProvider;
 import org.apache.streampipes.health.monitoring.ServiceHealthCheck;
 import org.apache.streampipes.loadbalance.LoadManager;
@@ -199,7 +200,8 @@ public class StreamPipesCoreApplication extends 
StreamPipesServiceBase {
             getPipelineStorage(),
             extensionServiceRequestManager,
             workerRestClient,
-            resourceManager),
+            resourceManager,
+            getRegisteredExtensionHealthChecks()),
         env.getInitialHealthCheckDelayInMillis().getValueOrDefault(),
         TimeUnit.MILLISECONDS);
 
@@ -221,7 +223,8 @@ public class StreamPipesCoreApplication extends 
StreamPipesServiceBase {
                     )),
                 
StorageDispatcher.INSTANCE.getNoSqlStore().getExtensionsServiceStorage(),
                 extensionServiceRequestManager,
-                resourceManager
+                resourceManager,
+                getRegisteredExtensionHealthChecks()
             )));
 
     var logFetchInterval = 
env.getLogFetchIntervalInMillis().getValueOrDefault();
@@ -248,6 +251,10 @@ public class StreamPipesCoreApplication extends 
StreamPipesServiceBase {
     return new AvailableMigrations(resourceManager).getAvailableMigrations();
   }
 
+  protected List<HealthCheck> getRegisteredExtensionHealthChecks() {
+    return List.of();
+  }
+
   private boolean isConfigured() {
     return new UserStorage().existsDatabase();
   }

Reply via email to