Skip to content

Commit cd01312

Browse files
committed
Fix shutting down status and add shutdown integration test
1 parent 4d6a9fc commit cd01312

3 files changed

Lines changed: 12 additions & 0 deletions

File tree

.github/workflows/ci.yml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -110,6 +110,7 @@ jobs:
110110
--dynamic-config-value history.enableRequestIdRefLinks=true \
111111
--dynamic-config-value frontend.WorkerHeartbeatsEnabled=true \
112112
--dynamic-config-value frontend.ListWorkersEnabled=true \
113+
--dynamic-config-value frontend.enableCancelWorkerPollsOnShutdown=true \
113114
--dynamic-config-value 'component.callbacks.allowedAddresses=[{"Pattern":"localhost:7243","AllowInsecure":true}]' \
114115
--dynamic-config-value frontend.activityAPIsEnabled=true \
115116
--dynamic-config-value activity.enableStandalone=true \

temporal-sdk/src/main/java/io/temporal/worker/Worker.java

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -482,6 +482,7 @@ void start() {
482482
}
483483

484484
CompletableFuture<Void> shutdown(ShutdownManager shutdownManager, boolean interruptUserTasks) {
485+
shuttingDown = true;
485486
ShutdownWorkerRequest.Builder requestBuilder =
486487
ShutdownWorkerRequest.newBuilder()
487488
.setNamespace(namespace)

temporal-sdk/src/test/java/io/temporal/worker/WorkerShutdownTest.java

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,8 @@
1313
import io.temporal.activity.ActivityInterface;
1414
import io.temporal.activity.ActivityMethod;
1515
import io.temporal.api.enums.v1.TaskQueueType;
16+
import io.temporal.api.enums.v1.WorkerStatus;
17+
import io.temporal.api.worker.v1.WorkerHeartbeat;
1618
import io.temporal.api.workflowservice.v1.GetSystemInfoResponse;
1719
import io.temporal.api.workflowservice.v1.ShutdownWorkerRequest;
1820
import io.temporal.api.workflowservice.v1.ShutdownWorkerResponse;
@@ -31,6 +33,7 @@
3133
import java.util.Collections;
3234
import java.util.List;
3335
import java.util.concurrent.TimeUnit;
36+
import java.util.function.Supplier;
3437
import org.junit.Test;
3538
import org.mockito.ArgumentCaptor;
3639

@@ -121,6 +124,9 @@ public void activeTaskQueueTypesEvaluatedAtShutdownTime() throws Exception {
121124
worker.registerWorkflowImplementationTypes(TestWorkflowImpl.class);
122125
worker.registerActivitiesImplementations(new TestActivityImpl());
123126
worker.registerNexusServiceImplementation(new TestNexusServiceImpl());
127+
Supplier<WorkerHeartbeat> heartbeatSupplier =
128+
() -> WorkerHeartbeat.newBuilder().setStatus(WorkerStatus.WORKER_STATUS_RUNNING).build();
129+
worker.setHeartbeatSupplier(heartbeatSupplier);
124130

125131
worker.shutdown(new ShutdownManager(), true).get(5, TimeUnit.SECONDS);
126132

@@ -137,5 +143,9 @@ public void activeTaskQueueTypesEvaluatedAtShutdownTime() throws Exception {
137143
assertTrue(
138144
"ShutdownWorkerRequest should include NEXUS type registered after construction",
139145
shutdownTypes.contains(TaskQueueType.TASK_QUEUE_TYPE_NEXUS));
146+
assertEquals(
147+
"ShutdownWorkerRequest heartbeat should report SHUTTING_DOWN",
148+
WorkerStatus.WORKER_STATUS_SHUTTING_DOWN,
149+
captor.getValue().getWorkerHeartbeat().getStatus());
140150
}
141151
}

0 commit comments

Comments
 (0)