Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down Expand Up @@ -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();
}
}

Expand Down Expand Up @@ -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<NonRootNamespaceProperties, WorkersTemplate> consumer) {
if (temporalProperties.getNamespaces() == null) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
@@ -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();
}
}
Original file line number Diff line number Diff line change
@@ -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();
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
package io.temporal.spring.boot.autoconfigure.gracefulshutdown;

import io.temporal.activity.ActivityInterface;

@ActivityInterface
public interface SlowActivity {

void run();
}
Original file line number Diff line number Diff line change
@@ -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();
}
}
Original file line number Diff line number Diff line change
@@ -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();
}
Original file line number Diff line number Diff line change
@@ -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();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
Loading