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

commit eb552f802114d6cff3c0cd9dcc8b6363d74c754c
Author: Sven Oehler <[email protected]>
AuthorDate: Tue Jun 30 11:01:29 2026 +0200

    Add health check registration
---
 .../health/monitoring/ExtensionHealthCheck.java    | 18 +++++++++++++++-
 .../monitoring/RegisteredExtensionHealthCheck.java | 24 ++++++++++++++++++++++
 .../streampipes/service/core/PostStartupTask.java  |  7 +++++--
 .../service/core/StreamPipesCoreApplication.java   | 11 ++++++++--
 4 files changed, 55 insertions(+), 5 deletions(-)

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..9f02ea7e62 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
@@ -28,6 +28,7 @@ import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 import java.util.HashMap;
+import java.util.List;
 
 public class ExtensionHealthCheck implements Runnable {
 
@@ -37,15 +38,18 @@ public class ExtensionHealthCheck implements Runnable {
   private final ExtensionServiceRequestManager extensionRequestManager;
   private final IExtensionsServiceStorage extensionsServiceStorage;
   private final SpResourceManager resourceManager;
+  private final List<RegisteredExtensionHealthCheck> registeredHealthChecks;
 
   public ExtensionHealthCheck(ResourceProvider resourceProvider,
                               IExtensionsServiceStorage 
extensionsServiceStorage,
                               ExtensionServiceRequestManager 
extensionRequestManager,
-                              SpResourceManager resourceManager) {
+                              SpResourceManager resourceManager,
+                              List<RegisteredExtensionHealthCheck> 
registeredHealthChecks) {
     this.resourceProvider = resourceProvider;
     this.extensionsServiceStorage = extensionsServiceStorage;
     this.extensionRequestManager = extensionRequestManager;
     this.resourceManager = resourceManager;
+    this.registeredHealthChecks = registeredHealthChecks;
   }
 
   @Override
@@ -66,8 +70,20 @@ public class ExtensionHealthCheck implements Runnable {
       var healthCheckData = new HealthCheckData(resourceProvider, 
activeResources, activeCoreInstances, activeExtensionInstances);
       new PipelineHealthCheck(healthCheckData, extensionRequestManager, 
resourceProvider, resourceManager).runCheck();
       new AdapterHealthCheck(healthCheckData).runCheck();
+      runRegisteredHealthChecks();
     } catch (Exception e) {
       LOG.warn("An unhandled error occurred while running health check.", e);
     }
   }
+
+  private void runRegisteredHealthChecks() {
+    registeredHealthChecks.forEach(healthCheck -> {
+      try {
+        healthCheck.runCheck();
+      } catch (Exception e) {
+        LOG.warn("An unhandled error occurred while running registered health 
check {}.",
+            healthCheck.getClass().getSimpleName(), e);
+      }
+    });
+  }
 }
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/RegisteredExtensionHealthCheck.java
new file mode 100644
index 0000000000..564555afa0
--- /dev/null
+++ 
b/streampipes-health-monitoring/src/main/java/org/apache/streampipes/health/monitoring/RegisteredExtensionHealthCheck.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 RegisteredExtensionHealthCheck {
+
+  void runCheck();
+}
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..2601af9ad4 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,6 +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.ResourceProvider;
 import org.apache.streampipes.health.monitoring.ServiceHealthCheck;
 import 
org.apache.streampipes.manager.api.extensions.ExtensionServiceRequestManager;
@@ -66,7 +67,8 @@ public class PostStartupTask implements Runnable {
   public PostStartupTask(IPipelineStorage pipelineStorage,
                          ExtensionServiceRequestManager 
extensionServiceRequestManager,
                          WorkerRestClient workerRestClient,
-                         SpResourceManager resourceManager) {
+                         SpResourceManager resourceManager,
+                         List<RegisteredExtensionHealthCheck> 
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..09767a2100 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.RegisteredExtensionHealthCheck;
 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<RegisteredExtensionHealthCheck> 
getRegisteredExtensionHealthChecks() {
+    return List.of();
+  }
+
   private boolean isConfigured() {
     return new UserStorage().existsDatabase();
   }

Reply via email to