From d29c0cbefdb1728e781cfe9a1af924c9b0945fe5 Mon Sep 17 00:00:00 2001 From: AdzerKI Date: Fri, 2 Oct 2026 07:51:13 +0300 Subject: [PATCH 1/2] Wait for workers to terminate when the Spring context closes The root namespace worker factory was shut down by its destroy method without waiting, and the service stubs were closed right after it, so pollers and in-flight activity completions failed with UNAVAILABLE: Channel shutdown invoked. Non-root namespace factories were shut down only when another context closed. Stop every worker factory from a SmartLifecycle that waits for termination within spring.lifecycle.timeout-per-shutdown-phase, before the beans are destroyed. --- CHANGELOG.md | 4 ++ .../NonRootBeanPostProcessor.java | 3 + .../NonRootNamespaceAutoConfiguration.java | 18 +----- .../RootNamespaceAutoConfiguration.java | 7 +++ .../autoconfigure/WorkerFactoryLifecycle.java | 58 +++++++++++++++++++ .../autoconfigure/GracefulShutdownTest.java | 49 ++++++++++++++++ .../gracefulshutdown/SlowActivity.java | 9 +++ .../gracefulshutdown/SlowActivityImpl.java | 35 +++++++++++ .../gracefulshutdown/SlowWorkflow.java | 13 +++++ .../gracefulshutdown/SlowWorkflowImpl.java | 20 +++++++ .../src/test/resources/application.yml | 11 ++++ 11 files changed, 212 insertions(+), 15 deletions(-) create mode 100644 temporal-spring-boot-autoconfigure/src/main/java/io/temporal/spring/boot/autoconfigure/WorkerFactoryLifecycle.java create mode 100644 temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/GracefulShutdownTest.java create mode 100644 temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/gracefulshutdown/SlowActivity.java create mode 100644 temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/gracefulshutdown/SlowActivityImpl.java create mode 100644 temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/gracefulshutdown/SlowWorkflow.java create mode 100644 temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/gracefulshutdown/SlowWorkflowImpl.java diff --git a/CHANGELOG.md b/CHANGELOG.md index 262e716a8d..5017769f2f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -26,6 +26,10 @@ to docs, or any other relevant information. ### Fixed - Test server now honors retry expiration deadlines that fall exactly on a whole second. Previously such deadlines were ignored and retries were scheduled past them instead of failing with `RETRY_STATE_TIMEOUT`. +- Spring Boot: closing the application context now waits for workers to finish in-flight tasks, bounded by + `spring.lifecycle.timeout-per-shutdown-phase`. Previously the root namespace workers were shut down without waiting, + so their pollers and task completions hit the already closed service stubs (`UNAVAILABLE: Channel shutdown invoked`), + and the non-root namespace workers were not shut down when their own context closed. ## Previous releases diff --git a/temporal-spring-boot-autoconfigure/src/main/java/io/temporal/spring/boot/autoconfigure/NonRootBeanPostProcessor.java b/temporal-spring-boot-autoconfigure/src/main/java/io/temporal/spring/boot/autoconfigure/NonRootBeanPostProcessor.java index 2aeef3b18d..06b6558a64 100644 --- a/temporal-spring-boot-autoconfigure/src/main/java/io/temporal/spring/boot/autoconfigure/NonRootBeanPostProcessor.java +++ b/temporal-spring-boot-autoconfigure/src/main/java/io/temporal/spring/boot/autoconfigure/NonRootBeanPostProcessor.java @@ -193,6 +193,9 @@ private void injectBeanByNonRootNamespace(NonRootNamespaceProperties ns) { beanFactory.registerSingleton( beanPrefix + ScheduleClient.class.getSimpleName(), scheduleClient); beanFactory.registerSingleton(beanPrefix + WorkerFactory.class.getSimpleName(), workerFactory); + beanFactory.registerSingleton( + beanPrefix + WorkerFactoryLifecycle.class.getSimpleName(), + new WorkerFactoryLifecycle(workerFactory)); } @Override diff --git a/temporal-spring-boot-autoconfigure/src/main/java/io/temporal/spring/boot/autoconfigure/NonRootNamespaceAutoConfiguration.java b/temporal-spring-boot-autoconfigure/src/main/java/io/temporal/spring/boot/autoconfigure/NonRootNamespaceAutoConfiguration.java index 7a70e2ebaf..0984e393f3 100644 --- a/temporal-spring-boot-autoconfigure/src/main/java/io/temporal/spring/boot/autoconfigure/NonRootNamespaceAutoConfiguration.java +++ b/temporal-spring-boot-autoconfigure/src/main/java/io/temporal/spring/boot/autoconfigure/NonRootNamespaceAutoConfiguration.java @@ -23,7 +23,6 @@ import org.springframework.context.annotation.Conditional; import org.springframework.context.annotation.Lazy; import org.springframework.context.event.ApplicationContextEvent; -import org.springframework.context.event.ContextClosedEvent; import org.springframework.context.event.ContextRefreshedEvent; @AutoConfiguration( @@ -66,12 +65,9 @@ public NonRootNamespaceEventListener( @Override public void onApplicationEvent(ApplicationContextEvent event) { - if (event.getApplicationContext() == this.applicationContext) { - if (event instanceof ContextRefreshedEvent) { - onStart(); - } - } else if (event instanceof ContextClosedEvent) { - onStop(); + if (event.getApplicationContext() == this.applicationContext + && event instanceof ContextRefreshedEvent) { + onStart(); } } @@ -102,14 +98,6 @@ private void onStart() { }); } - private void onStop() { - this.executeByNamespace( - (nonRootNamespaceProperties, workersTemplate) -> { - log.info("shutdown workers for non-root namespace"); - workersTemplate.getWorkerFactory().shutdown(); - }); - } - private void executeByNamespace( BiConsumer consumer) { if (temporalProperties.getNamespaces() == null) { diff --git a/temporal-spring-boot-autoconfigure/src/main/java/io/temporal/spring/boot/autoconfigure/RootNamespaceAutoConfiguration.java b/temporal-spring-boot-autoconfigure/src/main/java/io/temporal/spring/boot/autoconfigure/RootNamespaceAutoConfiguration.java index 9ff3908bb5..006789fb93 100644 --- a/temporal-spring-boot-autoconfigure/src/main/java/io/temporal/spring/boot/autoconfigure/RootNamespaceAutoConfiguration.java +++ b/temporal-spring-boot-autoconfigure/src/main/java/io/temporal/spring/boot/autoconfigure/RootNamespaceAutoConfiguration.java @@ -209,6 +209,13 @@ public WorkerFactory workerFactory( return workersTemplate.getWorkerFactory(); } + @Bean(name = "temporalWorkerFactoryLifecycle") + @Conditional(WorkersPresentCondition.class) + public WorkerFactoryLifecycle workerFactoryLifecycle( + @Qualifier("temporalWorkerFactory") WorkerFactory workerFactory) { + return new WorkerFactoryLifecycle(workerFactory); + } + @Primary @Bean(name = "temporalWorkers") @Conditional(WorkersPresentCondition.class) diff --git a/temporal-spring-boot-autoconfigure/src/main/java/io/temporal/spring/boot/autoconfigure/WorkerFactoryLifecycle.java b/temporal-spring-boot-autoconfigure/src/main/java/io/temporal/spring/boot/autoconfigure/WorkerFactoryLifecycle.java new file mode 100644 index 0000000000..a382691c85 --- /dev/null +++ b/temporal-spring-boot-autoconfigure/src/main/java/io/temporal/spring/boot/autoconfigure/WorkerFactoryLifecycle.java @@ -0,0 +1,58 @@ +package io.temporal.spring.boot.autoconfigure; + +import io.temporal.worker.WorkerFactory; +import java.util.concurrent.TimeUnit; +import org.springframework.context.SmartLifecycle; + +/** + * Shuts down a started {@link WorkerFactory} when the application context stops and waits for its + * workers to terminate, so that in-flight tasks complete before the beans they use, including the + * service stubs, are destroyed. The wait is bounded by {@code + * spring.lifecycle.timeout-per-shutdown-phase}. The factory is started elsewhere. + */ +public class WorkerFactoryLifecycle implements SmartLifecycle { + + private static final String AWAITER_THREAD_NAME = "temporal-worker-factory-stop"; + + private final WorkerFactory workerFactory; + + public WorkerFactoryLifecycle(WorkerFactory workerFactory) { + this.workerFactory = workerFactory; + } + + @Override + public boolean isAutoStartup() { + return false; + } + + @Override + public void start() {} + + @Override + public void stop() { + workerFactory.shutdown(); + workerFactory.awaitTermination(Long.MAX_VALUE, TimeUnit.MILLISECONDS); + } + + @Override + public void stop(Runnable callback) { + workerFactory.shutdown(); + Thread awaiter = + new Thread( + () -> { + try { + workerFactory.awaitTermination(Long.MAX_VALUE, TimeUnit.MILLISECONDS); + } finally { + callback.run(); + } + }, + AWAITER_THREAD_NAME); + awaiter.setDaemon(true); + awaiter.start(); + } + + @Override + public boolean isRunning() { + return workerFactory.isStarted() && !workerFactory.isShutdown(); + } +} diff --git a/temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/GracefulShutdownTest.java b/temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/GracefulShutdownTest.java new file mode 100644 index 0000000000..682ca9107b --- /dev/null +++ b/temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/GracefulShutdownTest.java @@ -0,0 +1,49 @@ +package io.temporal.spring.boot.autoconfigure; + +import static org.junit.jupiter.api.Assertions.assertTrue; + +import io.temporal.client.WorkflowClient; +import io.temporal.client.WorkflowOptions; +import io.temporal.spring.boot.autoconfigure.gracefulshutdown.SlowActivityImpl; +import io.temporal.spring.boot.autoconfigure.gracefulshutdown.SlowWorkflow; +import io.temporal.worker.WorkerFactory; +import java.util.concurrent.TimeUnit; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; +import org.springframework.boot.autoconfigure.EnableAutoConfiguration; +import org.springframework.boot.builder.SpringApplicationBuilder; +import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.context.annotation.Bean; + +public class GracefulShutdownTest { + + @Test + @Timeout(value = 30) + public void testInFlightActivityCompletesBeforeContextIsClosed() throws InterruptedException { + ConfigurableApplicationContext context = + new SpringApplicationBuilder(Configuration.class).profiles("graceful-shutdown").run(); + SlowActivityImpl activity = context.getBean(SlowActivityImpl.class); + WorkerFactory workerFactory = context.getBean(WorkerFactory.class); + SlowWorkflow workflow = + context + .getBean(WorkflowClient.class) + .newWorkflowStub( + SlowWorkflow.class, + WorkflowOptions.newBuilder().setTaskQueue(SlowWorkflow.TASK_QUEUE).build()); + WorkflowClient.start(workflow::execute); + assertTrue(activity.awaitStarted(10, TimeUnit.SECONDS)); + + context.close(); + + assertTrue(activity.isCompleted()); + assertTrue(workerFactory.isTerminated()); + } + + @EnableAutoConfiguration + public static class Configuration { + @Bean + public SlowActivityImpl slowActivity() { + return new SlowActivityImpl(); + } + } +} diff --git a/temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/gracefulshutdown/SlowActivity.java b/temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/gracefulshutdown/SlowActivity.java new file mode 100644 index 0000000000..d069931507 --- /dev/null +++ b/temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/gracefulshutdown/SlowActivity.java @@ -0,0 +1,9 @@ +package io.temporal.spring.boot.autoconfigure.gracefulshutdown; + +import io.temporal.activity.ActivityInterface; + +@ActivityInterface +public interface SlowActivity { + + void run(); +} diff --git a/temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/gracefulshutdown/SlowActivityImpl.java b/temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/gracefulshutdown/SlowActivityImpl.java new file mode 100644 index 0000000000..8934100e8c --- /dev/null +++ b/temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/gracefulshutdown/SlowActivityImpl.java @@ -0,0 +1,35 @@ +package io.temporal.spring.boot.autoconfigure.gracefulshutdown; + +import io.temporal.spring.boot.ActivityImpl; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; + +@ActivityImpl(taskQueues = SlowWorkflow.TASK_QUEUE) +public class SlowActivityImpl implements SlowActivity { + + private static final long DURATION_MILLIS = 2000; + + private final CountDownLatch started = new CountDownLatch(1); + private final AtomicBoolean completed = new AtomicBoolean(); + + @Override + public void run() { + started.countDown(); + try { + Thread.sleep(DURATION_MILLIS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + return; + } + completed.set(true); + } + + public boolean awaitStarted(long timeout, TimeUnit unit) throws InterruptedException { + return started.await(timeout, unit); + } + + public boolean isCompleted() { + return completed.get(); + } +} diff --git a/temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/gracefulshutdown/SlowWorkflow.java b/temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/gracefulshutdown/SlowWorkflow.java new file mode 100644 index 0000000000..12cd7d00c8 --- /dev/null +++ b/temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/gracefulshutdown/SlowWorkflow.java @@ -0,0 +1,13 @@ +package io.temporal.spring.boot.autoconfigure.gracefulshutdown; + +import io.temporal.workflow.WorkflowInterface; +import io.temporal.workflow.WorkflowMethod; + +@WorkflowInterface +public interface SlowWorkflow { + + String TASK_QUEUE = "GracefulShutdown"; + + @WorkflowMethod + void execute(); +} diff --git a/temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/gracefulshutdown/SlowWorkflowImpl.java b/temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/gracefulshutdown/SlowWorkflowImpl.java new file mode 100644 index 0000000000..cabadfcadc --- /dev/null +++ b/temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/gracefulshutdown/SlowWorkflowImpl.java @@ -0,0 +1,20 @@ +package io.temporal.spring.boot.autoconfigure.gracefulshutdown; + +import io.temporal.activity.ActivityOptions; +import io.temporal.spring.boot.WorkflowImpl; +import io.temporal.workflow.Workflow; +import java.time.Duration; + +@WorkflowImpl(taskQueues = SlowWorkflow.TASK_QUEUE) +public class SlowWorkflowImpl implements SlowWorkflow { + + private final SlowActivity activity = + Workflow.newActivityStub( + SlowActivity.class, + ActivityOptions.newBuilder().setStartToCloseTimeout(Duration.ofSeconds(10)).build()); + + @Override + public void execute() { + activity.run(); + } +} diff --git a/temporal-spring-boot-autoconfigure/src/test/resources/application.yml b/temporal-spring-boot-autoconfigure/src/test/resources/application.yml index afc8654c92..be772214d3 100644 --- a/temporal-spring-boot-autoconfigure/src/test/resources/application.yml +++ b/temporal-spring-boot-autoconfigure/src/test/resources/application.yml @@ -209,6 +209,17 @@ spring: - io.temporal.spring.boot.autoconfigure.bytaskqueue start-workers: false +--- +spring: + config: + activate: + on-profile: graceful-shutdown + temporal: + workers-auto-discovery: + workflow-packages: + - io.temporal.spring.boot.autoconfigure.gracefulshutdown + register-activity-beans: true + --- spring: config: From 21d401ccce2b476d22b1c9bafe31e1a19d5e1ac1 Mon Sep 17 00:00:00 2001 From: AdzerKI Date: Sat, 3 Oct 2026 04:39:04 +0300 Subject: [PATCH 2/2] Keep workers running while the Spring context is paused A shut down WorkerFactory cannot be started again. Spring Framework 7 pauses a cached test context when another one is started and restarts it when a test reuses it, so the stopped factory left the reused context without pollers and its workflows never progressed. --- CHANGELOG.md | 3 +- .../autoconfigure/WorkerFactoryLifecycle.java | 9 +++- .../autoconfigure/GracefulShutdownTest.java | 42 +++++++++++++++---- 3 files changed, 44 insertions(+), 10 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 5017769f2f..707fc052d1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -29,7 +29,8 @@ to docs, or any other relevant information. - Spring Boot: closing the application context now waits for workers to finish in-flight tasks, bounded by `spring.lifecycle.timeout-per-shutdown-phase`. Previously the root namespace workers were shut down without waiting, so their pollers and task completions hit the already closed service stubs (`UNAVAILABLE: Channel shutdown invoked`), - and the non-root namespace workers were not shut down when their own context closed. + and the non-root namespace workers were not shut down when their own context closed. A paused context, such as a + cached Spring TestContext one, keeps its workers running: a shut down worker factory cannot be started again. ## Previous releases diff --git a/temporal-spring-boot-autoconfigure/src/main/java/io/temporal/spring/boot/autoconfigure/WorkerFactoryLifecycle.java b/temporal-spring-boot-autoconfigure/src/main/java/io/temporal/spring/boot/autoconfigure/WorkerFactoryLifecycle.java index a382691c85..c62a690b71 100644 --- a/temporal-spring-boot-autoconfigure/src/main/java/io/temporal/spring/boot/autoconfigure/WorkerFactoryLifecycle.java +++ b/temporal-spring-boot-autoconfigure/src/main/java/io/temporal/spring/boot/autoconfigure/WorkerFactoryLifecycle.java @@ -5,10 +5,11 @@ import org.springframework.context.SmartLifecycle; /** - * Shuts down a started {@link WorkerFactory} when the application context stops and waits for its + * Shuts down a started {@link WorkerFactory} when the application context closes and waits for its * workers to terminate, so that in-flight tasks complete before the beans they use, including the * service stubs, are destroyed. The wait is bounded by {@code - * spring.lifecycle.timeout-per-shutdown-phase}. The factory is started elsewhere. + * spring.lifecycle.timeout-per-shutdown-phase}. The factory is started elsewhere. A shut down + * factory cannot be started again, so the factory keeps running while the context is paused. */ public class WorkerFactoryLifecycle implements SmartLifecycle { @@ -25,6 +26,10 @@ public boolean isAutoStartup() { return false; } + public boolean isPauseable() { + return false; + } + @Override public void start() {} diff --git a/temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/GracefulShutdownTest.java b/temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/GracefulShutdownTest.java index 682ca9107b..e0779a19c6 100644 --- a/temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/GracefulShutdownTest.java +++ b/temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/GracefulShutdownTest.java @@ -1,12 +1,15 @@ package io.temporal.spring.boot.autoconfigure; import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.junit.jupiter.api.Assumptions.assumeTrue; import io.temporal.client.WorkflowClient; import io.temporal.client.WorkflowOptions; import io.temporal.spring.boot.autoconfigure.gracefulshutdown.SlowActivityImpl; import io.temporal.spring.boot.autoconfigure.gracefulshutdown.SlowWorkflow; import io.temporal.worker.WorkerFactory; +import java.lang.reflect.Method; +import java.util.Arrays; import java.util.concurrent.TimeUnit; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; @@ -24,13 +27,7 @@ public void testInFlightActivityCompletesBeforeContextIsClosed() throws Interrup new SpringApplicationBuilder(Configuration.class).profiles("graceful-shutdown").run(); SlowActivityImpl activity = context.getBean(SlowActivityImpl.class); WorkerFactory workerFactory = context.getBean(WorkerFactory.class); - SlowWorkflow workflow = - context - .getBean(WorkflowClient.class) - .newWorkflowStub( - SlowWorkflow.class, - WorkflowOptions.newBuilder().setTaskQueue(SlowWorkflow.TASK_QUEUE).build()); - WorkflowClient.start(workflow::execute); + WorkflowClient.start(slowWorkflow(context)::execute); assertTrue(activity.awaitStarted(10, TimeUnit.SECONDS)); context.close(); @@ -39,6 +36,37 @@ public void testInFlightActivityCompletesBeforeContextIsClosed() throws Interrup assertTrue(workerFactory.isTerminated()); } + @Test + @Timeout(value = 30) + public void testWorkersKeepRunningWhenContextIsPausedAndRestarted() throws Exception { + assumeTrue( + Arrays.stream(ConfigurableApplicationContext.class.getMethods()) + .anyMatch(method -> method.getName().equals("pause")), + "Context pausing requires Spring Framework 7"); + Method pause = ConfigurableApplicationContext.class.getMethod("pause"); + Method restart = ConfigurableApplicationContext.class.getMethod("restart"); + ConfigurableApplicationContext context = + new SpringApplicationBuilder(Configuration.class).profiles("graceful-shutdown").run(); + try { + pause.invoke(context); + restart.invoke(context); + + slowWorkflow(context).execute(); + + assertTrue(context.getBean(SlowActivityImpl.class).isCompleted()); + } finally { + context.close(); + } + } + + private static SlowWorkflow slowWorkflow(ConfigurableApplicationContext context) { + return context + .getBean(WorkflowClient.class) + .newWorkflowStub( + SlowWorkflow.class, + WorkflowOptions.newBuilder().setTaskQueue(SlowWorkflow.TASK_QUEUE).build()); + } + @EnableAutoConfiguration public static class Configuration { @Bean