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