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