diff --git a/CHANGELOG.md b/CHANGELOG.md index 262e716a8d..707fc052d1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -26,6 +26,11 @@ 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. 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/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..c62a690b71 --- /dev/null +++ b/temporal-spring-boot-autoconfigure/src/main/java/io/temporal/spring/boot/autoconfigure/WorkerFactoryLifecycle.java @@ -0,0 +1,63 @@ +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 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. A shut down + * factory cannot be started again, so the factory keeps running while the context is paused. + */ +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; + } + + public boolean isPauseable() { + 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..e0779a19c6 --- /dev/null +++ b/temporal-spring-boot-autoconfigure/src/test/java/io/temporal/spring/boot/autoconfigure/GracefulShutdownTest.java @@ -0,0 +1,77 @@ +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; +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); + WorkflowClient.start(slowWorkflow(context)::execute); + assertTrue(activity.awaitStarted(10, TimeUnit.SECONDS)); + + context.close(); + + assertTrue(activity.isCompleted()); + 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 + 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: