This is an automated email from the ASF dual-hosted git repository.
dominikriemer 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 590b266bec chore: Improve startup logging (#4638)
590b266bec is described below
commit 590b266bec1f7303cc15aac63175cad1c7ebf748
Author: Dominik Riemer <[email protected]>
AuthorDate: Fri Jun 26 13:36:17 2026 +0200
chore: Improve startup logging (#4638)
---
.../health/monitoring/PipelineHealthCheck.java | 3 +--
.../endpoint/ExtensionsServiceEndpointGenerator.java | 2 +-
.../apache/streampipes/service/core/PostStartupTask.java | 8 +++-----
.../service/core/StreamPipesCoreApplication.java | 6 +++---
.../streampipes/service/core/WebSecurityConfig.java | 15 ++++++---------
.../service/core/migrations/MigrationsHandler.java | 2 +-
.../src/main/resources/application.properties | 4 +++-
.../user/management/service/SpUserDetailsService.java | 3 +++
8 files changed, 21 insertions(+), 22 deletions(-)
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 0aebbd81c3..ad1bef456b 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
@@ -147,7 +147,6 @@ public class PipelineHealthCheck {
MAX_FAILED_ATTEMPTS);
} else {
recoveredInstances.add(instanceId);
-
HealthCheckUtils.addSuccessfulRestoreNotification(pipelineNotifications,
pipelineElement);
resetFailedAttempts(instanceId);
LOG.info("Successfully restored pipeline element {} of pipeline
{}",
pipelineElement.getName(), pipeline.getName());
@@ -161,7 +160,7 @@ public class PipelineHealthCheck {
currentPipeline.setHealthStatus(PipelineHealthStatus.FAILURE);
pipelinesStats.failedIncrease();
} else if (!recoveredInstances.isEmpty()) {
-
currentPipeline.setHealthStatus(PipelineHealthStatus.REQUIRES_ATTENTION);
+ currentPipeline.setHealthStatus(PipelineHealthStatus.OK);
pipelinesStats.attentionRequiredIncrease();
}
currentPipeline.setSepas(
diff --git
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/execution/endpoint/ExtensionsServiceEndpointGenerator.java
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/execution/endpoint/ExtensionsServiceEndpointGenerator.java
index 8a8de79d5e..32a0955885 100644
---
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/execution/endpoint/ExtensionsServiceEndpointGenerator.java
+++
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/execution/endpoint/ExtensionsServiceEndpointGenerator.java
@@ -65,7 +65,7 @@ public class ExtensionsServiceEndpointGenerator implements
IExtensionsServiceEnd
}
// If we reach here, no service was found
- LOG.error("Could not find any service endpoints for appId {}, serviceTag
{}", appId,
+ LOG.warn("Could not find any service endpoints for appId {}, serviceTag
{}", appId,
spServiceUrlProvider.getServiceTag(appId).asString());
throw new NoServiceEndpointsAvailableException(
"Could not find any matching service endpoints - are all software
components running?");
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 a6cd6bb66e..dd576ef3eb 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
@@ -133,15 +133,13 @@ public class PostStartupTask implements Runnable {
startPipeline(pipeline, false);
});
- LOG.info("Checking for gracefully shut down pipelines to be restarted...");
-
List<Pipeline> pipelinesToRestart = allPipelines
.stream()
.filter(p -> !(p.isRunning()))
.filter(Pipeline::isRestartOnSystemReboot)
.toList();
- LOG.info("Found {} pipelines that we are attempting to restart...",
pipelinesToRestart.size());
+ LOG.info("Found {} pipelines that will be restarted",
pipelinesToRestart.size());
pipelinesToRestart.forEach(pipeline -> {
startPipeline(pipeline, false);
@@ -162,7 +160,7 @@ public class PostStartupTask implements Runnable {
storeFailedRestartAttempt(pipeline);
int failedAttemptCount = failedPipelines.get(pipeline.getPipelineId());
if (failedAttemptCount <= MAX_PIPELINE_START_RETRIES) {
- LOG.error(
+ LOG.warn(
"Pipeline {} could not be restarted - I'll try again in {} seconds
({}/{} failed attempts)",
pipeline.getName(),
WAIT_TIME_AFTER_FAILURE_IN_SECONDS,
@@ -172,7 +170,7 @@ public class PostStartupTask implements Runnable {
schedulePipelineStart(pipeline, restartOnReboot);
} else {
- LOG.error(
+ LOG.warn(
"Pipeline {} could not be restarted - are all pipeline element
containers running?",
status.getPipelineName()
);
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 2e463403bf..172c281661 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
@@ -223,7 +223,7 @@ public class StreamPipesCoreApplication extends
StreamPipesServiceBase {
)));
var logFetchInterval =
env.getLogFetchIntervalInMillis().getValueOrDefault();
- LOG.info("Extensions logs will be fetched every {} milliseconds",
logFetchInterval);
+ LOG.info("Extensions logs will be fetched every {} seconds",
TimeUnit.MILLISECONDS.toSeconds(logFetchInterval));
logCheckExecutorService.scheduleAtFixedRate(new
ExtensionsServiceLogExecutor(
extensionServiceRequestManager, resourceManager
),
@@ -234,8 +234,8 @@ public class StreamPipesCoreApplication extends
StreamPipesServiceBase {
private void scheduleHealthChecks(int healthCheckIntervalInMillis,
List<Runnable> checks) {
var healthCheckExecutorService =
Executors.newSingleThreadScheduledExecutor();
checks.forEach(check -> {
- LOG.info("Health check {} configured to run every {} {}",
check.getClass().getCanonicalName(),
- healthCheckIntervalInMillis, TimeUnit.MILLISECONDS);
+ LOG.info("Health check {} configured to run every {} seconds",
check.getClass().getSimpleName(),
+ TimeUnit.MILLISECONDS.toSeconds(healthCheckIntervalInMillis));
healthCheckExecutorService.scheduleAtFixedRate(check,
healthCheckIntervalInMillis,
healthCheckIntervalInMillis,
TimeUnit.MILLISECONDS);
diff --git
a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/WebSecurityConfig.java
b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/WebSecurityConfig.java
index eb25ed6f3c..7e93611888 100644
---
a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/WebSecurityConfig.java
+++
b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/WebSecurityConfig.java
@@ -44,10 +44,10 @@ import
org.springframework.security.access.PermissionEvaluator;
import
org.springframework.security.access.expression.method.DefaultMethodSecurityExpressionHandler;
import
org.springframework.security.access.expression.method.MethodSecurityExpressionHandler;
import org.springframework.security.authentication.AuthenticationManager;
+import org.springframework.security.authentication.ProviderManager;
+import
org.springframework.security.authentication.dao.DaoAuthenticationProvider;
import org.springframework.security.config.BeanIds;
import org.springframework.security.config.Customizer;
-import
org.springframework.security.config.annotation.authentication.builders.AuthenticationManagerBuilder;
-import
org.springframework.security.config.annotation.authentication.configuration.AuthenticationConfiguration;
import
org.springframework.security.config.annotation.method.configuration.EnableMethodSecurity;
import
org.springframework.security.config.annotation.web.builders.HttpSecurity;
import
org.springframework.security.config.annotation.web.configuration.EnableWebSecurity;
@@ -106,11 +106,6 @@ public class WebSecurityConfig {
this.permissionStorage = permissionStorage;
}
- @Autowired
- public void configureGlobal(AuthenticationManagerBuilder auth) {
-
auth.userDetailsService(userDetailsService).passwordEncoder(this.passwordEncoder.passwordEncoder());
- }
-
@Bean
MethodSecurityExpressionHandler methodSecurityExpressionHandler(
PermissionEvaluator permissionEvaluator
@@ -167,8 +162,10 @@ public class WebSecurityConfig {
}
@Bean
- public AuthenticationManager
authenticationManager(AuthenticationConfiguration authConfig) throws Exception {
- return authConfig.getAuthenticationManager();
+ public AuthenticationManager authenticationManager() {
+ var authenticationProvider = new
DaoAuthenticationProvider(userDetailsService);
+
authenticationProvider.setPasswordEncoder(this.passwordEncoder.passwordEncoder());
+ return new ProviderManager(authenticationProvider);
}
@Bean
diff --git
a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/MigrationsHandler.java
b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/MigrationsHandler.java
index a3f9b7c0b2..9c02745024 100644
---
a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/MigrationsHandler.java
+++
b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/MigrationsHandler.java
@@ -30,7 +30,7 @@ public class MigrationsHandler {
private static final Logger LOG =
LoggerFactory.getLogger(MigrationsHandler.class);
public void performMigrations(List<Migration> availableMigrations) {
- LOG.info("Running required migrations...");
+ LOG.info("Applying required migrations");
availableMigrations.forEach(migration -> {
if (migration.shouldExecute()) {
LOG.info("Performing migration: {}", migration.getDescription());
diff --git a/streampipes-service-core/src/main/resources/application.properties
b/streampipes-service-core/src/main/resources/application.properties
index 35a668cbbe..11417dd2ff 100644
--- a/streampipes-service-core/src/main/resources/application.properties
+++ b/streampipes-service-core/src/main/resources/application.properties
@@ -29,8 +29,10 @@
streampipes.storage.cache.pipelines.enabled=${SP_STORAGE_CACHE_PIPELINES_ENABLED
streampipes.storage.cache.data-lake-measures.enabled=${SP_STORAGE_CACHE_DATA_LAKE_MEASURES_ENABLED:true}
spring.cache.type=caffeine
spring.cache.cache-names=dataExplorerWidgets,permissions,adapters,dashboards,pipelines,dataLakeMeasures
-spring.cache.caffeine.spec=maximumSize=10000,expireAfterWrite=10m
+spring.cache.caffeine.spec=maximumSize=10000,expireAfterWrite=10m,recordStats
logging.config=classpath:logback.xml
spring.output.ansi.enabled=always
+springdoc.api-docs.enabled=${SP_API_DOCS_ENABLED:true}
springdoc.api-docs.path=/apidocs
+springdoc.swagger-ui.enabled=${SP_SWAGGER_UI_ENABLED:true}
[email protected]@
diff --git
a/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/service/SpUserDetailsService.java
b/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/service/SpUserDetailsService.java
index 5fd1bb74fb..6b90a67deb 100644
---
a/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/service/SpUserDetailsService.java
+++
b/streampipes-user-management/src/main/java/org/apache/streampipes/user/management/service/SpUserDetailsService.java
@@ -40,6 +40,9 @@ public class SpUserDetailsService implements
UserDetailsService {
@Override
public UserDetails loadUserByUsername(String s) throws
UsernameNotFoundException {
Principal user =
StorageDispatcher.INSTANCE.getNoSqlStore().getUserStorageAPI().getUser(s);
+ if (user == null) {
+ throw new UsernameNotFoundException("User not found");
+ }
return user instanceof UserAccount ? new UserAccountDetails((UserAccount)
user, permissionStorage) :
new ServiceAccountDetails((ServiceAccount) user, permissionStorage);
}