From d4deac7aececd356efeb44fc716c64d1a538831c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Attila=20M=C3=A9sz=C3=A1ros?= Date: Thu, 17 Sep 2026 09:39:50 +0200 Subject: [PATCH] feat: schedule polling and retries on managed executors instead of Timers PollingEventSource and TimerEventSource each created a java.util.Timer, and PerResourcePollingEventSource a ScheduledThreadPoolExecutor of its own, so every polling event source and every controller cost a thread that the operator neither sized nor shut down. They now schedule on executors managed by ExecutorServiceManager, which builds them from ConfigurationService: - getScheduledExecutorService() backs the polling event sources - getRetryAndRescheduleExecutorService() backs the retry and reschedule timer of every controller, kept separate so that a slow poll cannot delay a retry Both are sized by concurrentScheduledTaskThreads() (4 by default) and concurrentRetryAndRescheduleThreads() (2 by default), overridable programmatically or through the josdk.scheduled-tasks.concurrent-threads and josdk.retry-and-reschedule.concurrent-threads keys. Both default to daemon threads, discard the tasks scheduled for later on shutdown so they don't hold up the termination of the operator, and drop cancelled tasks eagerly. The previous Executors.newScheduledThreadPool(0) was effectively single threaded and created non daemon threads. An event source resolves its executor on every start rather than at creation time, since the manager replaces its pools when the operator is restarted, and only shuts one down if it created it itself. That also fixes PerResourcePollingEventSource shutting down a user supplied executor and not registering its tasks again after a restart. PollingEventSource gains a constructor taking an EventSourceContext, the one without it is deprecated, and both polling configurations accept an executor to run a single event source on a pool of its own. --- .../content/en/docs/documentation/eventing.md | 22 +++ .../documentation/operations/configuration.md | 7 + .../api/config/ConfigurationService.java | 80 +++++++++++ .../config/ConfigurationServiceOverrider.java | 104 ++++++++++++++ .../api/config/ExecutorServiceManager.java | 41 ++++-- .../operator/api/config/Utils.java | 48 +++++++ .../external/PollingDependentResource.java | 1 + .../processing/event/EventSourceManager.java | 4 +- .../processing/event/EventSources.java | 13 +- .../PerResourcePollingConfiguration.java | 12 +- .../PerResourcePollingEventSource.java | 57 +++++++- .../source/polling/PollingConfiguration.java | 18 ++- .../polling/PollingConfigurationBuilder.java | 20 ++- .../source/polling/PollingEventSource.java | 129 ++++++++++++++---- .../event/source/timer/TimerEventSource.java | 91 +++++++++--- .../ConfigurationServiceOverriderTest.java | 59 ++++++++ .../config/ExecutorServiceManagerTest.java | 6 + .../polling/PollingEventSourceTest.java | 77 ++++++++++- .../source/timer/TimerEventSourceTest.java | 57 ++++++++ .../operator/config/loader/ConfigLoader.java | 8 ++ .../config/loader/ConfigLoaderTest.java | 24 ++++ ...edExternalDependentResourceReconciler.java | 1 + 22 files changed, 804 insertions(+), 75 deletions(-) diff --git a/docs/content/en/docs/documentation/eventing.md b/docs/content/en/docs/documentation/eventing.md index e7aea6b065..a739cf5758 100644 --- a/docs/content/en/docs/documentation/eventing.md +++ b/docs/content/en/docs/documentation/eventing.md @@ -222,6 +222,28 @@ is similar to `PerResourcePollingEventSource` except that, contrary to that even doesn't poll a specific API separately per resource, but periodically and independently of actually observed primary resources. +#### Threading of the polling event sources + +Both polling event sources schedule their polls on an executor the operator shares between all of +its polling event sources. A poll therefore only starts once one of that executor's threads is +free, so if your operator registers many polling event sources, or if fetching your external +resources is slow, size the pool accordingly with +`ConfigurationServiceOverrider.withConcurrentScheduledTaskThreads` (4 threads by default): + +```java +Operator operator = new Operator(overrider -> overrider.withConcurrentScheduledTaskThreads(20)); +``` + +Retried and rescheduled reconciliations are triggered on a separate executor, so a slow poll can +never delay them. It is sized with `withConcurrentRetryAndRescheduleThreads` (2 threads by +default); few threads are needed there since triggering a reconciliation only enqueues an event +for one of the reconciliation threads to pick up. + +Use `ConfigurationServiceOverrider.withScheduledExecutorService` to replace the polling executor +altogether, or the `withExecutorService` method of the event source's own configuration builder to +poll a single event source on an executor of its own. An executor provided that way is not managed +by the operator: it is your responsibility to shut it down. + #### Inbound event sources [SimpleInboundEventSource](https://github.com/operator-framework/java-operator-sdk/blob/main/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/inbound/SimpleInboundEventSource.java) diff --git a/docs/content/en/docs/documentation/operations/configuration.md b/docs/content/en/docs/documentation/operations/configuration.md index cdfb1b7fdb..5445be6381 100644 --- a/docs/content/en/docs/documentation/operations/configuration.md +++ b/docs/content/en/docs/documentation/operations/configuration.md @@ -281,6 +281,13 @@ All operator-level keys are prefixed with `josdk.`. |---|---|---| | `josdk.workflow.executor-threads` | `Integer` | Thread pool size for workflow execution | +#### Scheduled Tasks + +| Key | Type | Description | +|---|---|---| +| `josdk.scheduled-tasks.concurrent-threads` | `Integer` | Thread pool size for the operator's scheduled tasks, i.e. the polling event sources | +| `josdk.retry-and-reschedule.concurrent-threads` | `Integer` | Thread pool size for triggering retried and rescheduled reconciliations | + #### Informer | Key | Type | Description | diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/ConfigurationService.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/ConfigurationService.java index 35f46e5019..2dec5aabf2 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/ConfigurationService.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/ConfigurationService.java @@ -21,6 +21,7 @@ import java.util.Set; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; import java.util.function.Consumer; import org.slf4j.Logger; @@ -64,6 +65,20 @@ public interface ConfigurationService { /** The default number of threads used to process dependent workflows */ int DEFAULT_WORKFLOW_EXECUTOR_THREAD_NUMBER = DEFAULT_RECONCILIATION_THREADS_NUMBER; + /** + * The default number of threads used to run the operator's scheduled tasks, i.e. the periodic + * polls of {@link io.javaoperatorsdk.operator.processing.event.source.polling.PollingEventSource} + * and {@link + * io.javaoperatorsdk.operator.processing.event.source.polling.PerResourcePollingEventSource} + */ + int DEFAULT_SCHEDULED_TASK_THREADS_NUMBER = 4; + + /** + * The default number of threads used to trigger the operator's retried and rescheduled + * reconciliations + */ + int DEFAULT_RETRY_AND_RESCHEDULE_THREADS_NUMBER = 2; + /** * Creates a new {@link ConfigurationService} instance used to configure an {@link * io.javaoperatorsdk.operator.Operator} instance, starting from the specified base configuration @@ -219,6 +234,33 @@ default int concurrentWorkflowExecutorThreads() { return DEFAULT_WORKFLOW_EXECUTOR_THREAD_NUMBER; } + /** + * Number of threads the operator can spin out to run its scheduled (i.e. periodic or delayed) + * tasks with the default executor. These threads are shared by all the polling event sources of + * the operator, so this number should be raised when many, or slow, polling event sources are + * registered: a poll only starts once a thread is available, and a slow poll therefore delays the + * polls of the other event sources. + * + * @return the maximum number of concurrent scheduled task threads + * @since 5.6.0 + */ + default int concurrentScheduledTaskThreads() { + return DEFAULT_SCHEDULED_TASK_THREADS_NUMBER; + } + + /** + * Number of threads the operator can spin out to trigger its retried and rescheduled + * reconciliations with the default executor. These threads are shared by all the controllers of + * the operator, but the tasks they run only enqueue an event for the reconciliation to happen on + * a reconciliation thread, so few of them are needed. + * + * @return the maximum number of concurrent retry and reschedule threads + * @since 5.6.0 + */ + default int concurrentRetryAndRescheduleThreads() { + return DEFAULT_RETRY_AND_RESCHEDULE_THREADS_NUMBER; + } + /** * Override to provide a custom {@link Metrics} implementation * @@ -249,6 +291,44 @@ default ExecutorService getWorkflowExecutorService() { return Executors.newFixedThreadPool(concurrentWorkflowExecutorThreads()); } + /** + * Override to provide a custom {@link ScheduledExecutorService} implementation to change how the + * operator's scheduled (i.e. periodic or delayed) tasks are run. This executor is shared by all + * the polling event sources of the operator. The retried and rescheduled reconciliations run on + * an executor of their own, see {@link #getRetryAndRescheduleExecutorService()}, so that a slow + * poll can't delay them. + * + *

Note that the default implementation lets the executor discard the tasks that were scheduled + * for later when it is shut down, so that they don't delay the termination of the operator, and + * that it creates daemon threads so that a never stopped operator doesn't keep the JVM alive. + * Custom implementations are advised to do the same. + * + * @return the {@link ScheduledExecutorService} implementation to use to run scheduled tasks + * @since 5.6.0 + */ + default ScheduledExecutorService getScheduledExecutorService() { + return Utils.daemonScheduledThreadPool( + concurrentScheduledTaskThreads(), "josdk-scheduled-task"); + } + + /** + * Override to provide a custom {@link ScheduledExecutorService} implementation to change how the + * operator's retried and rescheduled reconciliations are triggered. This executor is kept + * separate from the one the polling event sources use, see {@link + * #getScheduledExecutorService()}, so that a slow poll can't delay a retry. + * + *

The same notes as for {@link #getScheduledExecutorService()} apply to custom + * implementations. + * + * @return the {@link ScheduledExecutorService} implementation to use to trigger retried and + * rescheduled reconciliations + * @since 5.6.0 + */ + default ScheduledExecutorService getRetryAndRescheduleExecutorService() { + return Utils.daemonScheduledThreadPool( + concurrentRetryAndRescheduleThreads(), "josdk-retry-reschedule"); + } + /** * Determines whether the associated Kubernetes client should be closed when the associated {@link * io.javaoperatorsdk.operator.Operator} is stopped. diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/ConfigurationServiceOverrider.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/ConfigurationServiceOverrider.java index 2cf6540af0..9ac8668b20 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/ConfigurationServiceOverrider.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/ConfigurationServiceOverrider.java @@ -21,6 +21,7 @@ import java.util.Optional; import java.util.Set; import java.util.concurrent.ExecutorService; +import java.util.concurrent.ScheduledExecutorService; import java.util.function.Function; import org.slf4j.Logger; @@ -45,11 +46,15 @@ public class ConfigurationServiceOverrider { private Boolean checkCR; private Integer concurrentReconciliationThreads; private Integer concurrentWorkflowExecutorThreads; + private Integer concurrentScheduledTaskThreads; + private Integer concurrentRetryAndRescheduleThreads; private Cloner cloner; private Boolean closeClientOnStop; private KubernetesClient client; private ExecutorService executorService; private ExecutorService workflowExecutorService; + private ScheduledExecutorService scheduledExecutorService; + private ScheduledExecutorService retryAndRescheduleExecutorService; private LeaderElectionConfiguration leaderElectionConfiguration; private String clusterScopedEventNamespace; private EventRecorder eventRecorder; @@ -86,6 +91,34 @@ public ConfigurationServiceOverrider withConcurrentWorkflowExecutorThreads(int t return this; } + /** + * Sets the number of threads used to run the operator's scheduled (i.e. periodic or delayed) + * tasks, which are shared by all its polling event sources. + * + * @param threadNumber the maximum number of concurrent scheduled task threads + * @return this {@link ConfigurationServiceOverrider} for chained customization + * @see ConfigurationService#concurrentScheduledTaskThreads() + * @since 5.6.0 + */ + public ConfigurationServiceOverrider withConcurrentScheduledTaskThreads(int threadNumber) { + this.concurrentScheduledTaskThreads = threadNumber; + return this; + } + + /** + * Sets the number of threads used to trigger the operator's retried and rescheduled + * reconciliations, which are shared by all its controllers. + * + * @param threadNumber the maximum number of concurrent retry and reschedule threads + * @return this {@link ConfigurationServiceOverrider} for chained customization + * @see ConfigurationService#concurrentRetryAndRescheduleThreads() + * @since 5.6.0 + */ + public ConfigurationServiceOverrider withConcurrentRetryAndRescheduleThreads(int threadNumber) { + this.concurrentRetryAndRescheduleThreads = threadNumber; + return this; + } + @SuppressWarnings("rawtypes") public ConfigurationServiceOverrider withDependentResourceFactory( DependentResourceFactory dependentResourceFactory) { @@ -119,6 +152,37 @@ public ConfigurationServiceOverrider withWorkflowExecutorService( return this; } + /** + * Replaces the executor used to run the operator's scheduled (i.e. periodic or delayed) tasks, + * which are shared by all its polling event sources. + * + * @param scheduledExecutorService the executor to run scheduled tasks on + * @return this {@link ConfigurationServiceOverrider} for chained customization + * @see ConfigurationService#getScheduledExecutorService() + * @since 5.6.0 + */ + public ConfigurationServiceOverrider withScheduledExecutorService( + ScheduledExecutorService scheduledExecutorService) { + this.scheduledExecutorService = scheduledExecutorService; + return this; + } + + /** + * Replaces the executor used to trigger the operator's retried and rescheduled reconciliations, + * which is shared by all its controllers. + * + * @param retryAndRescheduleExecutorService the executor to trigger retried and rescheduled + * reconciliations on + * @return this {@link ConfigurationServiceOverrider} for chained customization + * @see ConfigurationService#getRetryAndRescheduleExecutorService() + * @since 5.6.0 + */ + public ConfigurationServiceOverrider withRetryAndRescheduleExecutorService( + ScheduledExecutorService retryAndRescheduleExecutorService) { + this.retryAndRescheduleExecutorService = retryAndRescheduleExecutorService; + return this; + } + /** * Replaces the default {@link KubernetesClient} instance by the specified one. This is the * preferred mechanism to configure which client will be used to access the cluster. @@ -312,6 +376,28 @@ public int concurrentWorkflowExecutorThreads() { original.concurrentWorkflowExecutorThreads()); } + @Override + public int concurrentScheduledTaskThreads() { + return Utils.ensureValid( + overriddenValueOrDefault( + concurrentScheduledTaskThreads, + ConfigurationService::concurrentScheduledTaskThreads), + "maximum scheduled task threads", + 1, + original.concurrentScheduledTaskThreads()); + } + + @Override + public int concurrentRetryAndRescheduleThreads() { + return Utils.ensureValid( + overriddenValueOrDefault( + concurrentRetryAndRescheduleThreads, + ConfigurationService::concurrentRetryAndRescheduleThreads), + "maximum retry and reschedule threads", + 1, + original.concurrentRetryAndRescheduleThreads()); + } + @Override public Metrics getMetrics() { return overriddenValueOrDefault(metrics, ConfigurationService::getMetrics); @@ -340,6 +426,24 @@ public ExecutorService getWorkflowExecutorService() { } } + @Override + public ScheduledExecutorService getScheduledExecutorService() { + if (scheduledExecutorService != null) { + return scheduledExecutorService; + } else { + return super.getScheduledExecutorService(); + } + } + + @Override + public ScheduledExecutorService getRetryAndRescheduleExecutorService() { + if (retryAndRescheduleExecutorService != null) { + return retryAndRescheduleExecutorService; + } else { + return super.getRetryAndRescheduleExecutorService(); + } + } + @Override public Optional getLeaderElectionConfiguration() { return leaderElectionConfiguration != null diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/ExecutorServiceManager.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/ExecutorServiceManager.java index cdcafcaa46..437bbfeae9 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/ExecutorServiceManager.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/ExecutorServiceManager.java @@ -42,6 +42,7 @@ public class ExecutorServiceManager { private ExecutorService workflowExecutor; private ExecutorService cachingExecutorService; private ScheduledExecutorService scheduledExecutorService; + private ScheduledExecutorService retryAndRescheduleExecutorService; private boolean started; private ConfigurationService configurationService; @@ -128,30 +129,54 @@ public ExecutorService cachingExecutorService() { return cachingExecutorService; } + /** + * The executor the operator runs its scheduled (i.e. periodic or delayed) tasks on, shared by its + * polling event sources. Note that it is only valid while the manager is started: it is shut down + * by {@link #stop(Duration)} and replaced by a fresh one on the next {@link + * #start(ConfigurationService)}, so callers should retrieve it when they start rather than hold + * on to it. + * + * @return the executor to run scheduled tasks on + */ public ScheduledExecutorService scheduledExecutorService() { return scheduledExecutorService; } + /** + * The executor the operator triggers its retried and rescheduled reconciliations on, kept + * separate from {@link #scheduledExecutorService()} so that a slow poll can't delay a retry. The + * same lifecycle caveat as for {@link #scheduledExecutorService()} applies. + * + * @return the executor to trigger retried and rescheduled reconciliations on + */ + public ScheduledExecutorService retryAndRescheduleExecutorService() { + return retryAndRescheduleExecutorService; + } + public synchronized void start(ConfigurationService configurationService) { if (!started) { this.configurationService = configurationService; // used to lazy init workflow executor this.cachingExecutorService = Executors.newCachedThreadPool(); - this.scheduledExecutorService = Executors.newScheduledThreadPool(0); + this.scheduledExecutorService = configurationService.getScheduledExecutorService(); + this.retryAndRescheduleExecutorService = + configurationService.getRetryAndRescheduleExecutorService(); this.executor = new InstrumentedExecutorService(configurationService.getExecutorService()); started = true; } } public synchronized void stop(Duration gracefulShutdownTimeout) { - var parallelExec = Executors.newFixedThreadPool(4); + var shutdowns = + List.of( + shutdown(executor, gracefulShutdownTimeout), + shutdown(workflowExecutor, gracefulShutdownTimeout), + shutdown(cachingExecutorService, gracefulShutdownTimeout), + shutdown(scheduledExecutorService, gracefulShutdownTimeout), + shutdown(retryAndRescheduleExecutorService, gracefulShutdownTimeout)); + var parallelExec = Executors.newFixedThreadPool(shutdowns.size()); try { log.debug("Closing executor"); - parallelExec.invokeAll( - List.of( - shutdown(executor, gracefulShutdownTimeout), - shutdown(workflowExecutor, gracefulShutdownTimeout), - shutdown(cachingExecutorService, gracefulShutdownTimeout), - shutdown(scheduledExecutorService, gracefulShutdownTimeout))); + parallelExec.invokeAll(shutdowns); } catch (InterruptedException e) { log.debug("Exception closing executor: {}", e.getLocalizedMessage()); Thread.currentThread().interrupt(); diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/Utils.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/Utils.java index 6ad4928c86..9e33469900 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/Utils.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/config/Utils.java @@ -26,6 +26,11 @@ import java.util.Date; import java.util.Optional; import java.util.Properties; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledThreadPoolExecutor; +import java.util.concurrent.ThreadFactory; +import java.util.concurrent.atomic.AtomicLong; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -102,6 +107,49 @@ public static int ensureValid(int value, String description, int minValue, int d return value; } + /** + * Creates a {@link ThreadFactory} producing daemon threads named after the specified prefix. + * Daemon threads don't keep the JVM alive if the {@link io.javaoperatorsdk.operator.Operator} is + * never stopped, and naming them makes the pool they belong to identifiable in thread dumps. + * + * @param namePrefix the prefix the created threads are named after + * @return a {@link ThreadFactory} creating named daemon threads + * @since 5.6.0 + */ + public static ThreadFactory daemonThreadFactory(String namePrefix) { + final var defaultThreadFactory = Executors.defaultThreadFactory(); + final var counter = new AtomicLong(); + return runnable -> { + final var thread = defaultThreadFactory.newThread(runnable); + thread.setName(namePrefix + "-" + counter.incrementAndGet()); + thread.setDaemon(true); + return thread; + }; + } + + /** + * Creates the kind of {@link ScheduledExecutorService} the operator runs its scheduled (i.e. + * periodic or delayed) tasks on: one whose threads are daemon threads named after the specified + * prefix, and which doesn't let the tasks that were scheduled for later delay its shutdown. + * + * @param corePoolSize the number of threads to keep in the pool + * @param threadNamePrefix the prefix the threads of the pool are named after + * @return a {@link ScheduledExecutorService} to run scheduled tasks on + * @since 5.6.0 + */ + public static ScheduledExecutorService daemonScheduledThreadPool( + int corePoolSize, String threadNamePrefix) { + final var executor = + new ScheduledThreadPoolExecutor(corePoolSize, daemonThreadFactory(threadNamePrefix)); + // tasks that are scheduled far out (a reconciliation rescheduled in an hour, say) would + // otherwise keep the pool from terminating until the graceful shutdown timeout expires + executor.setExecuteExistingDelayedTasksAfterShutdownPolicy(false); + // cancelled tasks are frequent (every retry that is superseded by a new event cancels one) and + // would otherwise be retained until their delay elapses + executor.setRemoveOnCancelPolicy(true); + return executor; + } + @SuppressWarnings("unused") // this is used in the Quarkus extension public static boolean isValidateCustomResourcesEnvVarSet() { diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/dependent/external/PollingDependentResource.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/dependent/external/PollingDependentResource.java index 894e359d57..a524fc2ec1 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/dependent/external/PollingDependentResource.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/dependent/external/PollingDependentResource.java @@ -55,6 +55,7 @@ protected ExternalResourceCachingEventSource createEventSource( EventSourceContext

context) { return new PollingEventSource<>( resourceType(), + context, new PollingConfiguration<>(name(), this, getPollingPeriod(), resourceIDMapper)); } } diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/EventSourceManager.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/EventSourceManager.java index d553d14cf9..bd5222fbbb 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/EventSourceManager.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/EventSourceManager.java @@ -55,7 +55,9 @@ public class EventSourceManager

public EventSourceManager(Controller

controller) { this( controller, - new EventSources<>(controller.getConfiguration().triggerReconcilerOnAllEvents())); + new EventSources<>( + controller.getConfiguration().triggerReconcilerOnAllEvents(), + () -> controller.getExecutorServiceManager().retryAndRescheduleExecutorService())); } EventSourceManager(Controller

controller, EventSources

eventSources) { diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/EventSources.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/EventSources.java index b482d26dca..2e8a523de0 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/EventSources.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/EventSources.java @@ -23,6 +23,8 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentNavigableMap; import java.util.concurrent.ConcurrentSkipListMap; +import java.util.concurrent.ScheduledExecutorService; +import java.util.function.Supplier; import java.util.stream.Stream; import io.fabric8.kubernetes.api.model.HasMetadata; @@ -41,8 +43,17 @@ class EventSources

{ private ControllerEventSource

controllerEventSource; public EventSources(boolean triggerReconcilerOnAllEvents) { + this(triggerReconcilerOnAllEvents, null); + } + + public EventSources( + boolean triggerReconcilerOnAllEvents, + Supplier scheduledExecutorServiceSupplier) { retryAndRescheduleTimerEventSource = - new TimerEventSource<>("RetryAndRescheduleTimerEventSource", triggerReconcilerOnAllEvents); + new TimerEventSource<>( + "RetryAndRescheduleTimerEventSource", + triggerReconcilerOnAllEvents, + scheduledExecutorServiceSupplier); } EventSources() { diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PerResourcePollingConfiguration.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PerResourcePollingConfiguration.java index 599647ff29..5c84a800a7 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PerResourcePollingConfiguration.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PerResourcePollingConfiguration.java @@ -18,12 +18,15 @@ import java.time.Duration; import java.util.Objects; import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.ScheduledThreadPoolExecutor; import java.util.function.Predicate; import io.fabric8.kubernetes.api.model.HasMetadata; import io.javaoperatorsdk.operator.processing.ResourceIDMapper; +/** + * @param executorService the executor to run the polls on, {@code null} (the default) to run them + * on the executor the operator shares between all its scheduled tasks + */ public record PerResourcePollingConfiguration( String name, ScheduledExecutorService executorService, @@ -32,8 +35,6 @@ public record PerResourcePollingConfiguration( Predicate

registerPredicate, Duration defaultPollingPeriod) { - public static final int DEFAULT_EXECUTOR_THREAD_NUMBER = 1; - public PerResourcePollingConfiguration( String name, ScheduledExecutorService executorService, @@ -42,10 +43,7 @@ public PerResourcePollingConfiguration( Predicate

registerPredicate, Duration defaultPollingPeriod) { this.name = name; - this.executorService = - executorService == null - ? new ScheduledThreadPoolExecutor(DEFAULT_EXECUTOR_THREAD_NUMBER) - : executorService; + this.executorService = executorService; this.resourceIDMapper = resourceIDMapper == null ? ResourceIDMapper.resourceIdProviderMapper() : resourceIDMapper; this.resourceFetcher = Objects.requireNonNull(resourceFetcher); diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PerResourcePollingEventSource.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PerResourcePollingEventSource.java index 1ab750d8f0..6303912c5a 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PerResourcePollingEventSource.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PerResourcePollingEventSource.java @@ -21,6 +21,7 @@ import java.util.Optional; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; @@ -31,6 +32,8 @@ import io.fabric8.kubernetes.api.model.HasMetadata; import io.javaoperatorsdk.operator.OperatorException; +import io.javaoperatorsdk.operator.api.config.ConfigurationService; +import io.javaoperatorsdk.operator.api.config.Utils; import io.javaoperatorsdk.operator.api.reconciler.EventSourceContext; import io.javaoperatorsdk.operator.processing.event.ResourceID; import io.javaoperatorsdk.operator.processing.event.source.Cache; @@ -56,11 +59,15 @@ public class PerResourcePollingEventSource private final Cache

primaryResourceCache; private final Set fetchedForPrimaries = ConcurrentHashMap.newKeySet(); - private final ScheduledExecutorService executorService; + private final ScheduledExecutorService configuredExecutorService; + private final ConfigurationService configurationService; + private final boolean ownsExecutorService; private final ResourceFetcher resourceFetcher; private final Predicate

registerPredicate; private final Duration period; + private volatile ScheduledExecutorService executorService; + public PerResourcePollingEventSource( Class resourceClass, EventSourceContext

context, @@ -69,10 +76,48 @@ public PerResourcePollingEventSource( this.primaryResourceCache = context.getPrimaryCache(); this.resourceFetcher = config.resourceFetcher(); this.registerPredicate = config.registerPredicate(); - this.executorService = config.executorService(); + this.configuredExecutorService = config.executorService(); + this.configurationService = configurationServiceOf(context); + // when neither the configuration nor the operator provides one, the event source has to create + // an executor of its own, and is then the one responsible for shutting it down + this.ownsExecutorService = configuredExecutorService == null && configurationService == null; this.period = config.defaultPollingPeriod(); } + /** + * The configuration of the operator the event source belongs to, or {@code null} if it doesn't + * belong to one, which only happens when the event source is used standalone, outside an + * operator. + */ + private static ConfigurationService configurationServiceOf(EventSourceContext context) { + final var controllerConfiguration = context.getControllerConfiguration(); + return controllerConfiguration == null + ? null + : controllerConfiguration.getConfigurationService(); + } + + @Override + public void start() throws OperatorException { + executorService = resolveExecutorService(); + super.start(); + } + + /** + * Resolves the executor to poll on. Note that this happens on every start, and not once at + * creation time, since the operator shuts its executors down when it is stopped and creates new + * ones if it is started again. + */ + private ScheduledExecutorService resolveExecutorService() { + if (configuredExecutorService != null) { + return configuredExecutorService; + } + if (configurationService != null) { + return configurationService.getExecutorServiceManager().scheduledExecutorService(); + } + return Executors.newSingleThreadScheduledExecutor( + Utils.daemonThreadFactory("josdk-polling-" + name())); + } + private Set getAndCacheResource(P primary, boolean fromGetter) { var values = resourceFetcher.fetchResources(primary); var primaryID = ResourceID.fromResource(primary); @@ -202,6 +247,12 @@ default Optional fetchDelay(Set lastFetchedResource, P primary) { @Override public void stop() throws OperatorException { super.stop(); - executorService.shutdownNow(); + // the tasks have to be cancelled explicitly now that the executor can be shared with the rest + // of the operator, and the map cleared so that they are registered again on a restart + scheduledFutures.values().forEach(future -> future.cancel(true)); + scheduledFutures.clear(); + if (ownsExecutorService && executorService != null) { + executorService.shutdownNow(); + } } } diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PollingConfiguration.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PollingConfiguration.java index 9ac1b8cc96..0836058cf4 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PollingConfiguration.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PollingConfiguration.java @@ -17,24 +17,40 @@ import java.time.Duration; import java.util.Objects; +import java.util.concurrent.ScheduledExecutorService; import io.javaoperatorsdk.operator.processing.ResourceIDMapper; +/** + * @param executorService the executor to run the polls on, {@code null} (the default) to run them + * on the executor the operator shares between all its scheduled tasks + */ public record PollingConfiguration( String name, PollingEventSource.GenericResourceFetcher genericResourceFetcher, Duration period, - ResourceIDMapper resourceIDMapper) { + ResourceIDMapper resourceIDMapper, + ScheduledExecutorService executorService) { public PollingConfiguration( String name, PollingEventSource.GenericResourceFetcher genericResourceFetcher, Duration period, ResourceIDMapper resourceIDMapper) { + this(name, genericResourceFetcher, period, resourceIDMapper, null); + } + + public PollingConfiguration( + String name, + PollingEventSource.GenericResourceFetcher genericResourceFetcher, + Duration period, + ResourceIDMapper resourceIDMapper, + ScheduledExecutorService executorService) { this.name = name; this.genericResourceFetcher = Objects.requireNonNull(genericResourceFetcher); this.period = period; this.resourceIDMapper = resourceIDMapper == null ? ResourceIDMapper.resourceIdProviderMapper() : resourceIDMapper; + this.executorService = executorService; } } diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PollingConfigurationBuilder.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PollingConfigurationBuilder.java index 0e68876b60..5a12fc4d50 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PollingConfigurationBuilder.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PollingConfigurationBuilder.java @@ -16,6 +16,7 @@ package io.javaoperatorsdk.operator.processing.event.source.polling; import java.time.Duration; +import java.util.concurrent.ScheduledExecutorService; import io.javaoperatorsdk.operator.processing.ResourceIDMapper; @@ -24,6 +25,7 @@ public final class PollingConfigurationBuilder { private final PollingEventSource.GenericResourceFetcher genericResourceFetcher; private ResourceIDMapper resourceIDMapper; private String name; + private ScheduledExecutorService executorService; public PollingConfigurationBuilder( PollingEventSource.GenericResourceFetcher fetcher, Duration period) { @@ -42,7 +44,23 @@ public PollingConfigurationBuilder withName(String name) { return this; } + /** + * Runs the polls on the specified executor instead of the one the operator shares between all its + * scheduled tasks. Note that an explicitly provided executor is not managed by the operator: it + * is the caller's responsibility to shut it down. + * + * @param executorService the executor to run the polls on + * @return this builder for chained customization + * @since 5.6.0 + */ + public PollingConfigurationBuilder withExecutorService( + ScheduledExecutorService executorService) { + this.executorService = executorService; + return this; + } + public PollingConfiguration build() { - return new PollingConfiguration<>(name, genericResourceFetcher, period, resourceIDMapper); + return new PollingConfiguration<>( + name, genericResourceFetcher, period, resourceIDMapper, executorService); } } diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PollingEventSource.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PollingEventSource.java index 5c85c24f90..8e615a6059 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PollingEventSource.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/polling/PollingEventSource.java @@ -18,8 +18,10 @@ import java.time.Duration; import java.util.Map; import java.util.Set; -import java.util.Timer; -import java.util.TimerTask; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import org.slf4j.Logger; @@ -27,6 +29,9 @@ import io.fabric8.kubernetes.api.model.HasMetadata; import io.javaoperatorsdk.operator.OperatorException; +import io.javaoperatorsdk.operator.api.config.ConfigurationService; +import io.javaoperatorsdk.operator.api.config.Utils; +import io.javaoperatorsdk.operator.api.reconciler.EventSourceContext; import io.javaoperatorsdk.operator.health.Status; import io.javaoperatorsdk.operator.processing.event.ResourceID; import io.javaoperatorsdk.operator.processing.event.source.ExternalResourceCachingEventSource; @@ -62,45 +67,109 @@ public class PollingEventSource private static final Logger log = LoggerFactory.getLogger(PollingEventSource.class); - private Timer timer; private final GenericResourceFetcher genericResourceFetcher; private final Duration period; private final AtomicBoolean healthy = new AtomicBoolean(true); + private final ScheduledExecutorService configuredExecutorService; + private final ConfigurationService configurationService; + private final boolean ownsExecutorService; + private volatile ScheduledExecutorService executorService; + private volatile ScheduledFuture pollingTask; + + /** + * Creates an event source polling on the executor the operator shares between all its scheduled + * tasks, unless the configuration provides one of its own. + * + * @param resourceClass the type of the polled resource + * @param context the context this event source is created for + * @param config the configuration of the polling + * @since 5.6.0 + */ + public PollingEventSource( + Class resourceClass, EventSourceContext

context, PollingConfiguration config) { + this(resourceClass, config, configurationServiceOf(context)); + } + + /** + * @deprecated use {@link #PollingEventSource(Class, EventSourceContext, PollingConfiguration)} + * instead: without a context, and unless the configuration provides an executor, this event + * source has to create a thread of its own to poll on instead of using the one the operator + * shares between all its scheduled tasks. + */ + @Deprecated(since = "5.6.0") public PollingEventSource(Class resourceClass, PollingConfiguration config) { + this(resourceClass, config, null); + } + + private PollingEventSource( + Class resourceClass, + PollingConfiguration config, + ConfigurationService configurationService) { super(config.name(), resourceClass, config.resourceIDMapper()); this.genericResourceFetcher = config.genericResourceFetcher(); this.period = config.period(); + this.configuredExecutorService = config.executorService(); + this.configurationService = configurationService; + // when neither the configuration nor the operator provides one, the event source has to create + // an executor of its own, and is then the one responsible for shutting it down + this.ownsExecutorService = configuredExecutorService == null && configurationService == null; + } + + /** + * The configuration of the operator the event source belongs to, or {@code null} if it doesn't + * belong to one, which only happens when the event source is used standalone, outside an + * operator. + */ + private static ConfigurationService configurationServiceOf(EventSourceContext context) { + final var controllerConfiguration = context.getControllerConfiguration(); + return controllerConfiguration == null + ? null + : controllerConfiguration.getConfigurationService(); } @Override public void start() throws OperatorException { - if (timer != null) { + if (pollingTask != null) { return; } super.start(); + executorService = resolveExecutorService(); getStateAndFillCache(); - timer = new Timer(true); - timer.schedule( - new TimerTask() { - @Override - public void run() { - try { - if (!isRunning()) { - log.debug("Event source not yet started. Will not run."); - return; - } - getStateAndFillCache(); - healthy.set(true); - } catch (Exception e) { - // Exception is required because of Kotlin - healthy.set(false); - log.error("Error during polling.", e); - } - } - }, - period.toMillis(), - period.toMillis()); + pollingTask = + executorService.scheduleWithFixedDelay( + this::poll, period.toMillis(), period.toMillis(), TimeUnit.MILLISECONDS); + } + + private void poll() { + try { + if (!isRunning()) { + log.debug("Event source not yet started. Will not run."); + return; + } + getStateAndFillCache(); + healthy.set(true); + } catch (Exception e) { + // Exception is required because of Kotlin + healthy.set(false); + log.error("Error during polling.", e); + } + } + + /** + * Resolves the executor to poll on. Note that this happens on every start, and not once at + * creation time, since the operator shuts its executors down when it is stopped and creates new + * ones if it is started again. + */ + private ScheduledExecutorService resolveExecutorService() { + if (configuredExecutorService != null) { + return configuredExecutorService; + } + if (configurationService != null) { + return configurationService.getExecutorServiceManager().scheduledExecutorService(); + } + return Executors.newSingleThreadScheduledExecutor( + Utils.daemonThreadFactory("josdk-polling-" + name())); } protected synchronized void getStateAndFillCache() { @@ -115,9 +184,13 @@ public interface GenericResourceFetcher { @Override public void stop() throws OperatorException { super.stop(); - if (timer != null) { - timer.cancel(); - timer = null; + if (pollingTask != null) { + // as with the java.util.Timer this replaces, an ongoing poll is left to finish + pollingTask.cancel(false); + pollingTask = null; + } + if (ownsExecutorService && executorService != null) { + executorService.shutdownNow(); } } diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/timer/TimerEventSource.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/timer/TimerEventSource.java index eae9663fe6..b4d7aa36a3 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/timer/TimerEventSource.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/timer/TimerEventSource.java @@ -17,14 +17,18 @@ import java.util.Map; import java.util.Set; -import java.util.Timer; -import java.util.TimerTask; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.TimeUnit; +import java.util.function.Supplier; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import io.fabric8.kubernetes.api.model.HasMetadata; +import io.javaoperatorsdk.operator.api.config.Utils; import io.javaoperatorsdk.operator.api.reconciler.BaseControl; import io.javaoperatorsdk.operator.health.Status; import io.javaoperatorsdk.operator.processing.event.Event; @@ -36,17 +40,47 @@ public class TimerEventSource extends AbstractEventSource implements ResourceEventAware { private static final Logger log = LoggerFactory.getLogger(TimerEventSource.class); - private Timer timer; - private final Map onceTasks = new ConcurrentHashMap<>(); + private final Map> onceTasks = new ConcurrentHashMap<>(); + private final Supplier executorServiceSupplier; + private final boolean ownsExecutorService; private boolean triggerReconcilerOnAllEvents; + private volatile ScheduledExecutorService executorService; public TimerEventSource() { - super(Void.class); + this((Supplier) null); } public TimerEventSource(String name, boolean triggerReconcilerOnAllEvents) { + this(name, triggerReconcilerOnAllEvents, null); + } + + /** + * Creates an event source scheduling on the executor provided by the specified supplier. The + * supplier is called on every start, and not once at creation time, since the operator shuts its + * executors down when it is stopped and creates new ones if it is started again. + * + * @param executorServiceSupplier supplies the executor to schedule on, {@code null} to have the + * event source create, and shut down, an executor of its own + * @since 5.6.0 + */ + public TimerEventSource(Supplier executorServiceSupplier) { + super(Void.class); + this.executorServiceSupplier = executorServiceSupplier; + this.ownsExecutorService = executorServiceSupplier == null; + } + + /** + * @see #TimerEventSource(Supplier) + * @since 5.6.0 + */ + public TimerEventSource( + String name, + boolean triggerReconcilerOnAllEvents, + Supplier executorServiceSupplier) { super(Void.class, name); this.triggerReconcilerOnAllEvents = triggerReconcilerOnAllEvents; + this.executorServiceSupplier = executorServiceSupplier; + this.ownsExecutorService = executorServiceSupplier == null; } @SuppressWarnings("unused") @@ -55,20 +89,25 @@ public void scheduleOnce(R resource, long delay) { } public void scheduleOnce(ResourceID resourceID, long delay) { - if (!isRunning()) { + final var executor = executorService; + if (!isRunning() || executor == null) { throw new IllegalStateException("The TimerEventSource is not running"); } - if (onceTasks.containsKey(resourceID)) { - cancelOnceSchedule(resourceID); - } - EventProducerTimeTask task = new EventProducerTimeTask(resourceID); if (delay == BaseControl.INSTANT_RESCHEDULE) { - task.run(); - } else { - onceTasks.put(resourceID, task); - timer.schedule(task, delay); + cancelOnceSchedule(resourceID); + new EventProducerTimeTask(resourceID).run(); + return; } + + onceTasks.compute( + resourceID, + (id, alreadyScheduled) -> { + if (alreadyScheduled != null) { + alreadyScheduled.cancel(false); + } + return executor.schedule(new EventProducerTimeTask(id), delay, TimeUnit.MILLISECONDS); + }); } @Override @@ -81,9 +120,11 @@ public void onResourceDeleted(R resource) { } public void cancelOnceSchedule(ResourceID customResourceUid) { - TimerTask timerTask = onceTasks.remove(customResourceUid); - if (timerTask != null) { - timerTask.cancel(); + var scheduled = onceTasks.remove(customResourceUid); + if (scheduled != null) { + // as with the java.util.TimerTask this replaces, a task that is already running is left to + // finish + scheduled.cancel(false); } } @@ -91,16 +132,26 @@ public void cancelOnceSchedule(ResourceID customResourceUid) { public void start() { if (!isRunning()) { super.start(); - timer = new Timer(true); + executorService = resolveExecutorService(); } } + private ScheduledExecutorService resolveExecutorService() { + if (executorServiceSupplier != null) { + return executorServiceSupplier.get(); + } + return Executors.newSingleThreadScheduledExecutor( + Utils.daemonThreadFactory("josdk-timer-" + name())); + } + @Override public void stop() { if (isRunning()) { onceTasks.keySet().forEach(this::cancelOnceSchedule); - timer.cancel(); super.stop(); + if (ownsExecutorService && executorService != null) { + executorService.shutdownNow(); + } } } @@ -114,7 +165,7 @@ public Set getSecondaryResources(HasMetadata primary) { return Set.of(); } - public class EventProducerTimeTask extends TimerTask { + public class EventProducerTimeTask implements Runnable { protected final ResourceID customResourceUid; diff --git a/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/api/config/ConfigurationServiceOverriderTest.java b/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/api/config/ConfigurationServiceOverriderTest.java index aec8381135..aebd7bc9b0 100644 --- a/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/api/config/ConfigurationServiceOverriderTest.java +++ b/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/api/config/ConfigurationServiceOverriderTest.java @@ -20,6 +20,7 @@ import java.util.Optional; import java.util.Set; import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledThreadPoolExecutor; import java.util.concurrent.ThreadPoolExecutor; import org.junit.jupiter.api.Test; @@ -117,6 +118,64 @@ public R clone(R object) { config.reconciliationTerminationTimeout(), overridden.reconciliationTerminationTimeout()); } + @Test + void scheduledExecutorCanBeOverridden() { + final var scheduledExecutorService = Executors.newScheduledThreadPool(1); + final var retryExecutorService = Executors.newScheduledThreadPool(1); + try { + final var overridden = + new ConfigurationServiceOverrider(config) + .withConcurrentScheduledTaskThreads(7) + .withScheduledExecutorService(scheduledExecutorService) + .withConcurrentRetryAndRescheduleThreads(5) + .withRetryAndRescheduleExecutorService(retryExecutorService) + .build(); + + assertThat(overridden.concurrentScheduledTaskThreads()).isEqualTo(7); + assertThat(overridden.getScheduledExecutorService()).isSameAs(scheduledExecutorService); + assertThat(overridden.concurrentRetryAndRescheduleThreads()).isEqualTo(5); + assertThat(overridden.getRetryAndRescheduleExecutorService()).isSameAs(retryExecutorService); + } finally { + scheduledExecutorService.shutdownNow(); + retryExecutorService.shutdownNow(); + } + } + + @Test + void scheduledExecutorDefaultsToADaemonPoolOfTheConfiguredSize() { + final var overridden = + (ScheduledThreadPoolExecutor) + new ConfigurationServiceOverrider(config) + .withConcurrentScheduledTaskThreads(3) + .build() + .getScheduledExecutorService(); + + try { + assertThat(overridden.getCorePoolSize()).isEqualTo(3); + // scheduled tasks must not delay the termination of the operator, nor keep the JVM alive + assertThat(overridden.getExecuteExistingDelayedTasksAfterShutdownPolicy()).isFalse(); + assertThat(overridden.getThreadFactory().newThread(() -> {}).isDaemon()).isTrue(); + } finally { + overridden.shutdownNow(); + } + } + + @Test + void retryAndRescheduleExecutorIsSeparateFromTheScheduledTaskOne() { + final var scheduled = config.getScheduledExecutorService(); + final var retryAndReschedule = config.getRetryAndRescheduleExecutorService(); + + try { + // a slow poll on the scheduled task executor must not be able to delay a retry + assertThat(retryAndReschedule).isNotSameAs(scheduled); + assertThat(((ScheduledThreadPoolExecutor) retryAndReschedule).getCorePoolSize()) + .isEqualTo(ConfigurationService.DEFAULT_RETRY_AND_RESCHEDULE_THREADS_NUMBER); + } finally { + scheduled.shutdownNow(); + retryAndReschedule.shutdownNow(); + } + } + @Test void eventRecorderIsNotConfiguredByDefaultAndCanBeOverridden() { final var eventRecorder = diff --git a/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/api/config/ExecutorServiceManagerTest.java b/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/api/config/ExecutorServiceManagerTest.java index 40ffded241..995a6e374f 100644 --- a/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/api/config/ExecutorServiceManagerTest.java +++ b/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/api/config/ExecutorServiceManagerTest.java @@ -30,14 +30,19 @@ void stopShutsDownTheScheduledExecutorService() { ConfigurationService configurationService = new BaseConfigurationService(); var manager = configurationService.getExecutorServiceManager(); var scheduled = manager.scheduledExecutorService(); + var retryAndReschedule = manager.retryAndRescheduleExecutorService(); try { assertThat(scheduled.isShutdown()).isFalse(); + assertThat(retryAndReschedule.isShutdown()).isFalse(); + // retries are triggered on an executor of their own so that a slow poll can't delay them + assertThat(retryAndReschedule).isNotSameAs(scheduled); } finally { manager.stop(SHUTDOWN_TIMEOUT); } assertThat(scheduled.isShutdown()).isTrue(); + assertThat(retryAndReschedule.isShutdown()).isTrue(); } @Test @@ -53,6 +58,7 @@ void canBeRestartedAfterStop() { assertThat(manager.reconcileExecutorService().isShutdown()).isFalse(); assertThat(manager.cachingExecutorService().isShutdown()).isFalse(); assertThat(manager.scheduledExecutorService().isShutdown()).isFalse(); + assertThat(manager.retryAndRescheduleExecutorService().isShutdown()).isFalse(); manager.stop(SHUTDOWN_TIMEOUT); } diff --git a/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/processing/event/source/polling/PollingEventSourceTest.java b/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/processing/event/source/polling/PollingEventSourceTest.java index 92400df4de..43b20d368f 100644 --- a/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/processing/event/source/polling/PollingEventSourceTest.java +++ b/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/processing/event/source/polling/PollingEventSourceTest.java @@ -19,11 +19,16 @@ import java.util.HashMap; import java.util.Map; import java.util.Set; +import java.util.concurrent.CopyOnWriteArrayList; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import io.fabric8.kubernetes.api.model.HasMetadata; +import io.javaoperatorsdk.operator.api.config.BaseConfigurationService; +import io.javaoperatorsdk.operator.api.config.ConfigurationService; +import io.javaoperatorsdk.operator.api.config.ControllerConfiguration; +import io.javaoperatorsdk.operator.api.reconciler.EventSourceContext; import io.javaoperatorsdk.operator.health.Status; import io.javaoperatorsdk.operator.processing.event.EventHandler; import io.javaoperatorsdk.operator.processing.event.ResourceID; @@ -46,6 +51,7 @@ class PollingEventSourceTest private final PollingEventSource.GenericResourceFetcher resourceFetcher = mock(PollingEventSource.GenericResourceFetcher.class); + @SuppressWarnings("deprecation") private final PollingEventSource pollingEventSource = new PollingEventSource<>( SampleExternalResource.class, @@ -67,7 +73,7 @@ void canBeRestartedAfterStop() throws InterruptedException { Thread.sleep(DEFAULT_WAIT_PERIOD); pollingEventSource.stop(); - // a cancelled java.util.Timer cannot be reused, a new one has to be created on start + // the polling task is cancelled on stop, a new one has to be scheduled on start pollingEventSource.start(); Thread.sleep(DEFAULT_WAIT_PERIOD); @@ -75,7 +81,20 @@ void canBeRestartedAfterStop() throws InterruptedException { } @Test - void timerThreadIsADaemonSoItDoesNotKeepTheJvmAlive() throws InterruptedException { + void stopCancelsThePollingTask() throws InterruptedException { + when(resourceFetcher.fetchResources()).thenReturn(testResponseWithTwoValues()); + pollingEventSource.start(); + Thread.sleep(DEFAULT_WAIT_PERIOD); + pollingEventSource.stop(); + clearInvocations(resourceFetcher); + + Thread.sleep(DEFAULT_WAIT_PERIOD); + + verify(resourceFetcher, never()).fetchResources(); + } + + @Test + void pollingThreadIsADaemonSoItDoesNotKeepTheJvmAlive() throws InterruptedException { when(resourceFetcher.fetchResources()).thenReturn(testResponseWithTwoValues()); var threadsBeforeStart = Thread.getAllStackTraces().keySet(); @@ -83,13 +102,13 @@ void timerThreadIsADaemonSoItDoesNotKeepTheJvmAlive() throws InterruptedExceptio pollingEventSource.start(); Thread.sleep(DEFAULT_WAIT_PERIOD); - var newTimerThreads = + var newPollingThreads = Thread.getAllStackTraces().keySet().stream() - .filter(t -> t.getName().startsWith("Timer-")) + .filter(t -> t.getName().startsWith("josdk-polling-")) .filter(t -> !threadsBeforeStart.contains(t)) .toList(); - assertThat(newTimerThreads).isNotEmpty().allMatch(Thread::isDaemon); + assertThat(newPollingThreads).isNotEmpty().allMatch(Thread::isDaemon); } @Test @@ -154,6 +173,54 @@ void updatesHealthIndicatorBasedOnExceptionsInFetcher() { .untilAsserted(() -> assertThat(pollingEventSource.getStatus()).isEqualTo(Status.HEALTHY)); } + @Test + void pollsOnTheOperatorsSharedSchedulerWhenCreatedWithAContext() { + var pollingThreadNames = new CopyOnWriteArrayList(); + when(resourceFetcher.fetchResources()) + .thenAnswer( + invocation -> { + pollingThreadNames.add(Thread.currentThread().getName()); + return testResponseWithOneValue(); + }); + + var configurationService = new BaseConfigurationService(); + var executorServiceManager = configurationService.getExecutorServiceManager(); + var eventSource = + new PollingEventSource( + SampleExternalResource.class, + contextFor(configurationService), + new PollingConfiguration<>(null, resourceFetcher, POLL_PERIOD, null)); + eventSource.setEventHandler(mock(EventHandler.class)); + + try { + eventSource.start(); + + // the initial fetch happens on the calling thread, the scheduled ones on the shared pool + await() + .untilAsserted( + () -> + assertThat(pollingThreadNames) + .anyMatch(name -> name.startsWith("josdk-scheduled-task-"))); + + eventSource.stop(); + + // the shared executor belongs to the operator, stopping an event source must not shut it down + assertThat(executorServiceManager.scheduledExecutorService().isShutdown()).isFalse(); + } finally { + executorServiceManager.stop(Duration.ofMillis(100)); + } + } + + @SuppressWarnings("unchecked") + private static EventSourceContext contextFor( + ConfigurationService configurationService) { + var controllerConfiguration = mock(ControllerConfiguration.class); + when(controllerConfiguration.getConfigurationService()).thenReturn(configurationService); + var context = mock(EventSourceContext.class); + when(context.getControllerConfiguration()).thenReturn(controllerConfiguration); + return context; + } + private Map> testResponseWithTwoValueForSameId() { Map> res = new HashMap<>(); res.put(primaryID1(), Set.of(testResource1(), testResource2())); diff --git a/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/processing/event/source/timer/TimerEventSourceTest.java b/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/processing/event/source/timer/TimerEventSourceTest.java index 3a4e1cb80d..e8d36c0034 100644 --- a/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/processing/event/source/timer/TimerEventSourceTest.java +++ b/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/processing/event/source/timer/TimerEventSourceTest.java @@ -17,6 +17,7 @@ import java.util.List; import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; import org.awaitility.Awaitility; @@ -26,6 +27,7 @@ import org.junit.jupiter.api.Test; import io.javaoperatorsdk.operator.TestUtils; +import io.javaoperatorsdk.operator.api.config.Utils; import io.javaoperatorsdk.operator.api.reconciler.BaseControl; import io.javaoperatorsdk.operator.health.Status; import io.javaoperatorsdk.operator.processing.event.Event; @@ -125,6 +127,59 @@ public void handlesInstanceReschedule() { assertThat(eventHandler.events).hasSize(1); } + @Test + public void schedulesOnTheProvidedExecutorAndLeavesItRunningOnStop() { + var providedExecutor = + Executors.newSingleThreadScheduledExecutor( + Utils.daemonThreadFactory("provided-timer-executor")); + var handler = new CapturingEventHandler(); + var eventSource = new TimerEventSource(() -> providedExecutor); + eventSource.setEventHandler(handler); + + try { + eventSource.start(); + eventSource.scheduleOnce(ResourceID.fromResource(TestUtils.testCustomResource()), PERIOD); + + untilAsserted( + () -> + assertThat(handler.eventProducingThreadNames) + .containsExactly("provided-timer-executor-1")); + + eventSource.stop(); + + // the executor is not the event source's to shut down + assertThat(providedExecutor.isShutdown()).isFalse(); + } finally { + providedExecutor.shutdownNow(); + } + } + + @Test + public void shutsDownTheExecutorItCreatedItselfAndCreatesANewOneOnRestart() { + var eventSource = new TimerEventSource(); + var handler = new CapturingEventHandler(); + eventSource.setEventHandler(handler); + + eventSource.start(); + eventSource.scheduleOnce(ResourceID.fromResource(TestUtils.testCustomResource()), PERIOD); + untilAsserted(() -> assertThat(handler.events).hasSize(1)); + var firstThreadName = handler.eventProducingThreadNames.get(0); + assertThat(firstThreadName).startsWith("josdk-timer-"); + + eventSource.stop(); + Awaitility.await() + .untilAsserted( + () -> + assertThat(Thread.getAllStackTraces().keySet()) + .noneMatch(t -> t.getName().equals(firstThreadName))); + + eventSource.start(); + eventSource.scheduleOnce(ResourceID.fromResource(TestUtils.testCustomResource()), PERIOD); + untilAsserted(() -> assertThat(handler.events).hasSize(2)); + + eventSource.stop(); + } + private void untilAsserted(ThrowingRunnable assertion) { untilAsserted(INITIAL_DELAY, PERIOD, assertion); } @@ -150,10 +205,12 @@ private void untilAsserted(long initialDelay, long interval, ThrowingRunnable as public static class CapturingEventHandler implements EventHandler { private final List events = new CopyOnWriteArrayList<>(); + private final List eventProducingThreadNames = new CopyOnWriteArrayList<>(); @Override public void handleEvent(Event event) { events.add(event); + eventProducingThreadNames.add(Thread.currentThread().getName()); } } } diff --git a/operator-framework/src/main/java/io/javaoperatorsdk/operator/config/loader/ConfigLoader.java b/operator-framework/src/main/java/io/javaoperatorsdk/operator/config/loader/ConfigLoader.java index c8daf89724..8aa92ebce4 100644 --- a/operator-framework/src/main/java/io/javaoperatorsdk/operator/config/loader/ConfigLoader.java +++ b/operator-framework/src/main/java/io/javaoperatorsdk/operator/config/loader/ConfigLoader.java @@ -78,6 +78,14 @@ public static ConfigLoader getDefault() { "workflow.executor-threads", Integer.class, ConfigurationServiceOverrider::withConcurrentWorkflowExecutorThreads), + new ConfigBinding<>( + "scheduled-tasks.concurrent-threads", + Integer.class, + ConfigurationServiceOverrider::withConcurrentScheduledTaskThreads), + new ConfigBinding<>( + "retry-and-reschedule.concurrent-threads", + Integer.class, + ConfigurationServiceOverrider::withConcurrentRetryAndRescheduleThreads), new ConfigBinding<>( "close-client-on-stop", Boolean.class, diff --git a/operator-framework/src/test/java/io/javaoperatorsdk/operator/config/loader/ConfigLoaderTest.java b/operator-framework/src/test/java/io/javaoperatorsdk/operator/config/loader/ConfigLoaderTest.java index 44fac32b7d..08bcdfba91 100644 --- a/operator-framework/src/test/java/io/javaoperatorsdk/operator/config/loader/ConfigLoaderTest.java +++ b/operator-framework/src/test/java/io/javaoperatorsdk/operator/config/loader/ConfigLoaderTest.java @@ -84,6 +84,30 @@ void applyConfigsAppliesConcurrentWorkflowExecutorThreads() { assertThat(result.concurrentWorkflowExecutorThreads()).isEqualTo(3); } + @Test + void applyConfigsAppliesConcurrentScheduledTaskThreads() { + var loader = + new ConfigLoader(mapProvider(Map.of("josdk.scheduled-tasks.concurrent-threads", 8))); + + var base = new BaseConfigurationService(null); + var result = + ConfigurationService.newOverriddenConfigurationService(base, loader.applyConfigs()); + + assertThat(result.concurrentScheduledTaskThreads()).isEqualTo(8); + } + + @Test + void applyConfigsAppliesConcurrentRetryAndRescheduleThreads() { + var loader = + new ConfigLoader(mapProvider(Map.of("josdk.retry-and-reschedule.concurrent-threads", 6))); + + var base = new BaseConfigurationService(null); + var result = + ConfigurationService.newOverriddenConfigurationService(base, loader.applyConfigs()); + + assertThat(result.concurrentRetryAndRescheduleThreads()).isEqualTo(6); + } + @Test void applyConfigsAppliesBooleanFlags() { var values = new HashMap(); diff --git a/operator-framework/src/test/java/io/javaoperatorsdk/operator/dependent/multiplemanagedexternaldependenttype/MultipleManagedExternalDependentResourceReconciler.java b/operator-framework/src/test/java/io/javaoperatorsdk/operator/dependent/multiplemanagedexternaldependenttype/MultipleManagedExternalDependentResourceReconciler.java index dec9422486..33b7f965f3 100644 --- a/operator-framework/src/test/java/io/javaoperatorsdk/operator/dependent/multiplemanagedexternaldependenttype/MultipleManagedExternalDependentResourceReconciler.java +++ b/operator-framework/src/test/java/io/javaoperatorsdk/operator/dependent/multiplemanagedexternaldependenttype/MultipleManagedExternalDependentResourceReconciler.java @@ -96,6 +96,7 @@ public int getNumberOfExecutions() { pollingEventSource = new PollingEventSource<>( ExternalResource.class, + context, new PollingConfigurationBuilder( fetcher, Duration.ofMillis(1000L)) .withName(EVENT_SOURCE_NAME)