This is an automated email from the ASF dual-hosted git repository. dominikriemer pushed a commit to branch improve-startup-behaviour in repository https://gitbox.apache.org/repos/asf/streampipes.git
commit 3574ac2e0001dc7088522cc7d8dee6731fcaaf51 Author: Dominik Riemer <[email protected]> AuthorDate: Thu Jun 25 23:12:58 2026 +0200 chore: Improve startup logging --- .../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); }
