• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

temporalio / sdk-java / #397

02 Oct 2026 03:49PM UTC coverage: 68.315% (-0.009%) from 68.324%
#397

push

github

web-flow
Configurable MDC tag prefix (#3098)

7897 of 13704 branches covered (57.63%)

Branch coverage included in aggregate %.

97 of 99 new or added lines in 16 files covered. (97.98%)

18 existing lines in 6 files now uncovered.

32058 of 44782 relevant lines covered (71.59%)

0.72 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

3.33
/temporal-sdk/src/main/java/io/temporal/internal/worker/WorkerCommandTaskHandler.java
1
package io.temporal.internal.worker;
2

3
import com.google.protobuf.InvalidProtocolBufferException;
4
import com.uber.m3.tally.Scope;
5
import io.temporal.api.common.v1.Payload;
6
import io.temporal.api.nexus.v1.Response;
7
import io.temporal.api.nexus.v1.StartOperationRequest;
8
import io.temporal.api.nexus.v1.StartOperationResponse;
9
import io.temporal.api.nexusservices.workerservice.v1.ExecuteCommandsRequest;
10
import io.temporal.api.nexusservices.workerservice.v1.ExecuteCommandsResponse;
11
import io.temporal.api.worker.v1.CancelActivityResult;
12
import io.temporal.api.worker.v1.WorkerCommand;
13
import io.temporal.api.worker.v1.WorkerCommandResult;
14
import io.temporal.common.converter.DataConverter;
15
import io.temporal.common.converter.GlobalDataConverter;
16
import io.temporal.internal.logging.PrefixedMdc;
17
import io.temporal.serviceclient.Version;
18
import io.temporal.serviceclient.WorkflowServiceStubs;
19
import io.temporal.worker.tuning.FixedSizeSlotSupplier;
20
import io.temporal.worker.tuning.NexusSlotInfo;
21
import io.temporal.worker.tuning.PollerBehaviorSimpleMaximum;
22
import java.util.Objects;
23
import java.util.concurrent.TimeoutException;
24
import java.util.function.Function;
25
import javax.annotation.Nonnull;
26
import org.slf4j.Logger;
27
import org.slf4j.LoggerFactory;
28

29
/** Handles server-to-worker commands delivered on the worker command Nexus task queue. */
30
public final class WorkerCommandTaskHandler implements NexusTaskHandler {
31
  private static final Logger log = LoggerFactory.getLogger(WorkerCommandTaskHandler.class);
1 ✔
32
  private static final String TASK_QUEUE_PREFIX = "temporal-sys/worker-commands";
33

34
  private final Function<byte[], Boolean> activityCancelCallback;
35

36
  public WorkerCommandTaskHandler(Function<byte[], Boolean> activityCancelCallback) {
×
37
    this.activityCancelCallback = Objects.requireNonNull(activityCancelCallback);
×
38
  }
×
39

40
  public static String workerControlTaskQueue(String namespace, String workerGroupingKey) {
41
    return String.format("%s/%s/%s", TASK_QUEUE_PREFIX, namespace, workerGroupingKey);
1 ✔
42
  }
43

44
  public static SuspendableWorker newWorkerCommandWorker(
45
      @Nonnull WorkflowServiceStubs service,
46
      @Nonnull String namespace,
47
      @Nonnull String identity,
48
      @Nonnull String workerGroupingKey,
49
      @Nonnull Function<byte[], Boolean> activityCancelCallback,
50
      @Nonnull Scope metricsScope,
51
      @Nonnull NamespaceCapabilities namespaceCapabilities,
52
      @Nonnull PrefixedMdc mdc) {
53
    String taskQueue = workerControlTaskQueue(namespace, workerGroupingKey);
×
54
    DataConverter dataConverter = GlobalDataConverter.get();
×
55
    SingleWorkerOptions options =
56
        SingleWorkerOptions.newBuilder()
×
57
            .setIdentity(identity)
×
58
            .setBuildId(Version.LIBRARY_VERSION)
×
59
            .setWorkerInstanceKey(workerGroupingKey)
×
60
            .setDataConverter(dataConverter)
×
61
            .setMetricsScope(metricsScope)
×
62
            .setPollerOptions(
×
63
                PollerOptions.newBuilder()
×
64
                    .setPollerBehavior(new PollerBehaviorSimpleMaximum(1))
×
65
                    .setPollThreadNamePrefix("WorkerCommandNexusPoller")
×
66
                    .build())
×
NEW
67
            .setLoggerMdc(mdc)
×
68
            .build();
×
69
    return new NexusWorker(
×
70
        service,
71
        namespace,
72
        taskQueue,
73
        options,
74
        new WorkerCommandTaskHandler(activityCancelCallback),
75
        dataConverter,
76
        new FixedSizeSlotSupplier<NexusSlotInfo>(5),
77
        namespaceCapabilities,
78
        true);
79
  }
80

81
  @Override
82
  public boolean start() {
83
    return true;
×
84
  }
85

86
  @Override
87
  public Result handle(NexusTask task, Scope metricsScope) throws TimeoutException {
88
    ExecuteCommandsRequest request = decodeRequest(task);
×
89
    ExecuteCommandsResponse.Builder response = ExecuteCommandsResponse.newBuilder();
×
90
    for (WorkerCommand command : request.getCommandsList()) {
×
91
      response.addResults(handleCommand(command));
×
92
    }
×
93
    return new Result(
×
94
        Response.newBuilder()
×
95
            .setStartOperation(
×
96
                StartOperationResponse.newBuilder()
×
97
                    .setSyncSuccess(
×
98
                        StartOperationResponse.Sync.newBuilder()
×
99
                            .setPayload(
×
100
                                Payload.newBuilder().setData(response.build().toByteString()))))
×
101
            .build());
×
102
  }
103

104
  private ExecuteCommandsRequest decodeRequest(NexusTask task) {
105
    StartOperationRequest request = task.getResponse().getRequest().getStartOperation();
×
106
    if (!request.hasPayload()) {
×
107
      throw new IllegalArgumentException(
×
108
          "Worker command Nexus task missing ExecuteCommands payload");
109
    }
110
    try {
111
      return ExecuteCommandsRequest.parseFrom(request.getPayload().getData());
×
112
    } catch (InvalidProtocolBufferException e) {
×
113
      throw new IllegalArgumentException("Failed to decode ExecuteCommandsRequest", e);
×
114
    }
115
  }
116

117
  private WorkerCommandResult handleCommand(WorkerCommand command) {
118
    WorkerCommandResult.Builder result = WorkerCommandResult.newBuilder();
×
119
    if (command.hasCancelActivity()) {
×
120
      byte[] taskToken = command.getCancelActivity().getTaskToken().toByteArray();
×
121
      Boolean found = activityCancelCallback.apply(taskToken);
×
122
      if (!Boolean.TRUE.equals(found)) {
×
123
        log.debug("Activity task token from worker command was not found");
×
124
      }
125
      result.setCancelActivity(CancelActivityResult.newBuilder());
×
126
    } else {
×
127
      log.warn("Unsupported worker command");
×
128
    }
129
    return result.build();
×
130
  }
131
}
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc