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

temporalio / sdk-java / #349

03 Aug 2026 05:42AM UTC coverage: 68.221% (+0.08%) from 68.143%
#349

push

github

web-flow
Add Standalone Activities to Temporal Nexus Operation Handler (#2918)

* Enable Nexus activity operations without regressing workflow updates

Compose standalone activity support with the workflow-update Nexus model already present on master. Shared token, callback, link, client, and cancellation paths retain both operation families.

Constraint: Preserve the workflow-update token and API contracts from master
Rejected: Choose one conflict side | each side would drop a supported Nexus operation family
Confidence: high
Scope-risk: moderate
Reversibility: clean
Directive: Keep activity execution at token type 2 and workflow update at token type 3
Tested: Focused temporal-sdk token, link, invoker, client, async activity, and cancellation tests
Not-tested: Full repository suite and real-server-only standalone activity cases

* Make sure we test attaching

* Make sure to handle ActivityAlreadyStartedException

* Add test with SANO

* Bump test server

7220 of 12610 branches covered (57.26%)

Branch coverage included in aggregate %.

182 of 283 new or added lines in 15 files covered. (64.31%)

13 existing lines in 3 files now uncovered.

29918 of 41828 relevant lines covered (71.53%)

0.72 hits per line

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

25.9
/temporal-sdk/src/main/java/io/temporal/internal/client/RootActivityClientInvoker.java
1
package io.temporal.internal.client;
2

3
import static io.temporal.internal.common.RetryOptionsUtils.toRetryPolicy;
4
import static io.temporal.internal.common.WorkflowExecutionUtils.makeUserMetaData;
5

6
import com.google.common.base.Strings;
7
import com.google.common.collect.Iterators;
8
import io.grpc.Deadline;
9
import io.grpc.Status;
10
import io.grpc.StatusRuntimeException;
11
import io.temporal.api.activity.v1.ActivityExecutionOutcome;
12
import io.temporal.api.common.v1.ActivityType;
13
import io.temporal.api.common.v1.Callback;
14
import io.temporal.api.common.v1.Link;
15
import io.temporal.api.common.v1.Payloads;
16
import io.temporal.api.errordetails.v1.ActivityExecutionAlreadyStartedFailure;
17
import io.temporal.api.sdk.v1.UserMetadata;
18
import io.temporal.api.taskqueue.v1.TaskQueue;
19
import io.temporal.api.workflowservice.v1.*;
20
import io.temporal.client.*;
21
import io.temporal.common.converter.DataConverter;
22
import io.temporal.common.interceptors.ActivityClientCallsInterceptor;
23
import io.temporal.internal.client.external.GenericWorkflowClient;
24
import io.temporal.internal.common.HeaderUtils;
25
import io.temporal.internal.common.InternalUtils;
26
import io.temporal.internal.common.ProtoConverters;
27
import io.temporal.internal.common.ProtobufTimeUtils;
28
import io.temporal.internal.common.SearchAttributesUtil;
29
import io.temporal.internal.nexus.CurrentNexusOperationContext;
30
import io.temporal.internal.nexus.InternalNexusOperationContext;
31
import io.temporal.internal.nexus.NexusOperationMetadata;
32
import io.temporal.serviceclient.StatusUtils;
33
import java.lang.reflect.Type;
34
import java.util.*;
35
import java.util.concurrent.CompletableFuture;
36
import java.util.concurrent.CompletionException;
37
import java.util.concurrent.TimeoutException;
38
import java.util.stream.StreamSupport;
39

40
/**
41
 * Terminus of the activity interceptor chain. Implements all activity RPCs against the Temporal
42
 * service.
43
 */
44
public class RootActivityClientInvoker implements ActivityClientCallsInterceptor {
45

46
  private final GenericWorkflowClient genericClient;
47
  private final ActivityClientOptions clientOptions;
48

49
  public RootActivityClientInvoker(
50
      GenericWorkflowClient genericClient, ActivityClientOptions clientOptions) {
1✔
51
    this.genericClient = genericClient;
1✔
52
    this.clientOptions = clientOptions;
1✔
53
  }
1✔
54

55
  @Override
56
  public StartActivityOutput startActivity(StartActivityInput input) {
57
    StartActivityOptions options = input.getOptions();
1✔
58
    DataConverter dc = clientOptions.getDataConverter();
1✔
59
    InternalNexusOperationContext nexusContext =
60
        CurrentNexusOperationContext.isNexusContext() ? CurrentNexusOperationContext.get() : null;
1!
61
    NexusOperationMetadata nexusOperationMetadata =
62
        nexusContext == null ? null : nexusContext.getNexusOperationMetadata();
1!
63

64
    StartActivityExecutionRequest.Builder request =
65
        StartActivityExecutionRequest.newBuilder()
1✔
66
            .setNamespace(clientOptions.getNamespace())
1✔
67
            .setIdentity(clientOptions.getIdentity())
1✔
68
            .setRequestId(
1✔
69
                nexusOperationMetadata == null
1✔
70
                    ? UUID.randomUUID().toString()
1✔
71
                    : nexusOperationMetadata.requestId)
1✔
72
            .setActivityId(options.getId())
1✔
73
            .setActivityType(ActivityType.newBuilder().setName(input.getActivityType()).build())
1✔
74
            .setTaskQueue(TaskQueue.newBuilder().setName(options.getTaskQueue()).build())
1✔
75
            .setIdReusePolicy(options.getIdReusePolicy())
1✔
76
            .setIdConflictPolicy(options.getIdConflictPolicy());
1✔
77

78
    Optional<Payloads> activityInput = dc.toPayloads(input.getArgs().toArray());
1✔
79
    activityInput.ifPresent(request::setInput);
1✔
80

81
    if (options.getScheduleToCloseTimeout() != null) {
1!
82
      request.setScheduleToCloseTimeout(
×
83
          ProtobufTimeUtils.toProtoDuration(options.getScheduleToCloseTimeout()));
×
84
    }
85
    if (options.getScheduleToStartTimeout() != null) {
1!
86
      request.setScheduleToStartTimeout(
×
87
          ProtobufTimeUtils.toProtoDuration(options.getScheduleToStartTimeout()));
×
88
    }
89
    if (options.getStartToCloseTimeout() != null) {
1!
90
      request.setStartToCloseTimeout(
1✔
91
          ProtobufTimeUtils.toProtoDuration(options.getStartToCloseTimeout()));
1✔
92
    }
93
    if (options.getHeartbeatTimeout() != null) {
1!
94
      request.setHeartbeatTimeout(ProtobufTimeUtils.toProtoDuration(options.getHeartbeatTimeout()));
×
95
    }
96
    if (options.getRetryOptions() != null) {
1!
97
      request.setRetryPolicy(toRetryPolicy(options.getRetryOptions()));
×
98
    }
99
    if (options.getTypedSearchAttributes() != null
1!
100
        && options.getTypedSearchAttributes().size() > 0) {
×
101
      request.setSearchAttributes(
×
102
          SearchAttributesUtil.encodeTyped(options.getTypedSearchAttributes()));
×
103
    }
104
    if (options.getStaticSummary() != null || options.getStaticDetails() != null) {
1!
105
      UserMetadata userMetadata =
×
106
          makeUserMetaData(options.getStaticSummary(), options.getStaticDetails(), dc);
×
107
      if (userMetadata != null) {
×
108
        request.setUserMetadata(userMetadata);
×
109
      }
110
    }
111
    if (options.getPriority() != null) {
1!
112
      request.setPriority(ProtoConverters.toProto(options.getPriority()));
×
113
    }
114
    if (options.getStartDelay() != null) {
1!
115
      request.setStartDelay(ProtobufTimeUtils.toProtoDuration(options.getStartDelay()));
×
116
    }
117

118
    io.temporal.api.common.v1.Header grpcHeader = HeaderUtils.toHeaderGrpc(input.getHeader(), null);
1✔
119
    request.setHeader(grpcHeader);
1✔
120

121
    if (nexusOperationMetadata != null) {
1✔
122
      List<Link> protoLinks = nexusContext.getRequestLinks();
1✔
123
      request.addAllLinks(protoLinks);
1✔
124
      request.setOnConflictOptions(
1✔
125
          io.temporal.api.common.v1.OnConflictOptions.newBuilder()
1✔
126
              .setAttachRequestId(true)
1✔
127
              .setAttachLinks(true)
1✔
128
              .setAttachCompletionCallbacks(true));
1✔
129
      // Generate the operation token from the user-supplied activity ID and namespace so the
130
      // dual OPERATION_ID + OPERATION_TOKEN headers can be injected before the start RPC fires.
131
      try {
132
        nexusOperationMetadata.operationToken =
1✔
133
            io.temporal.internal.nexus.OperationTokenUtil.generateActivityExecutionOperationToken(
1✔
134
                options.getId(), clientOptions.getNamespace());
1✔
NEW
135
      } catch (com.fasterxml.jackson.core.JsonProcessingException e) {
×
NEW
136
        throw new io.nexusrpc.handler.HandlerException(
×
137
            io.nexusrpc.handler.HandlerException.ErrorType.BAD_REQUEST,
138
            "failed to generate activity operation token",
139
            e);
140
      }
1✔
141
      if (!Strings.isNullOrEmpty(nexusOperationMetadata.callbackUrl)) {
1✔
142
        Callback cb =
1✔
143
            InternalUtils.buildNexusCallback(
1✔
144
                nexusOperationMetadata.callbackUrl,
145
                nexusOperationMetadata.callbackHeaders,
146
                nexusOperationMetadata.operationToken,
147
                protoLinks);
148
        request.addCompletionCallbacks(cb);
1✔
149
      }
150
    }
151

152
    StartActivityExecutionResponse response;
153
    try {
154
      response = genericClient.startActivity(request.build());
1✔
155
    } catch (StatusRuntimeException e) {
×
156
      if (e.getStatus().getCode() == Status.Code.ALREADY_EXISTS) {
×
157
        ActivityExecutionAlreadyStartedFailure detail =
×
158
            StatusUtils.getFailure(e, ActivityExecutionAlreadyStartedFailure.class);
×
159
        if (detail != null) {
×
160
          String runId = detail.getRunId().isEmpty() ? null : detail.getRunId();
×
161
          throw new ActivityAlreadyStartedException(
×
162
              options.getId(), input.getActivityType(), runId, e);
×
163
        }
164
      }
165
      throw e;
×
166
    }
1✔
167

168
    if (nexusOperationMetadata != null && response.hasLink()) {
1!
169
      nexusContext.addResponseLink(response.getLink());
1✔
170
    }
171

172
    String runId = response.getRunId().isEmpty() ? null : response.getRunId();
1!
173
    return new StartActivityOutput(options.getId(), runId);
1✔
174
  }
175

176
  @Override
177
  public <R> GetActivityResultOutput<R> getActivityResult(GetActivityResultInput<R> input)
178
      throws TimeoutException {
179
    String namespace = clientOptions.getNamespace();
×
180
    DataConverter dc = clientOptions.getDataConverter();
×
181
    Deadline deadline = Deadline.after(input.getTimeout(), input.getTimeoutUnit());
×
182

183
    while (true) {
184
      PollActivityExecutionRequest.Builder pollRequest =
185
          PollActivityExecutionRequest.newBuilder()
×
186
              .setNamespace(namespace)
×
187
              .setActivityId(input.getActivityId());
×
188
      if (input.getRunId() != null) {
×
189
        pollRequest.setRunId(input.getRunId());
×
190
      }
191

192
      PollActivityExecutionResponse pollResponse;
193
      try {
194
        pollResponse = genericClient.pollActivity(pollRequest.build(), deadline);
×
195
      } catch (StatusRuntimeException e) {
×
196
        if (deadline.isExpired() && Status.Code.DEADLINE_EXCEEDED.equals(e.getStatus().getCode())) {
×
197
          throw new TimeoutException(
×
198
              "Activity did not complete within timeout: activityId='"
199
                  + input.getActivityId()
×
200
                  + "'");
201
        }
202
        throw e;
×
203
      }
×
204

205
      if (!pollResponse.hasOutcome()) {
×
206
        if (Thread.currentThread().isInterrupted()) {
×
207
          throw new ActivityFailedException(
×
208
              "Interrupted while waiting for activity result for activityId='"
209
                  + input.getActivityId()
×
210
                  + "'",
211
              input.getActivityId(),
×
212
              input.getRunId(),
×
213
              new InterruptedException());
214
        }
215
        continue;
216
      }
217

218
      ActivityExecutionOutcome outcome = pollResponse.getOutcome();
×
219
      switch (outcome.getValueCase()) {
×
220
        case RESULT:
221
          Type resultType =
222
              input.getResultType() != null ? input.getResultType() : input.getResultClass();
×
223
          @SuppressWarnings("unchecked")
224
          R result =
×
225
              (R)
226
                  dc.fromPayloads(
×
227
                      0,
228
                      outcome.hasResult() ? Optional.of(outcome.getResult()) : Optional.empty(),
×
229
                      input.getResultClass(),
×
230
                      resultType);
231
          return new GetActivityResultOutput<>(result);
×
232
        case FAILURE:
233
          throw new ActivityFailedException(
×
234
              "Activity failed: activityId='" + input.getActivityId() + "'",
×
235
              input.getActivityId(),
×
236
              input.getRunId(),
×
237
              dc.failureToException(outcome.getFailure()));
×
238
        default:
239
          throw new ActivityFailedException(
×
240
              "Activity completed with unexpected outcome '"
241
                  + outcome.getValueCase()
×
242
                  + "' for activityId='"
243
                  + input.getActivityId()
×
244
                  + "'",
245
              input.getActivityId(),
×
246
              input.getRunId(),
×
247
              null);
248
      }
249
    }
250
  }
251

252
  @Override
253
  public <R> CompletableFuture<GetActivityResultOutput<R>> getActivityResultAsync(
254
      GetActivityResultInput<R> input) {
255
    DataConverter dc = clientOptions.getDataConverter();
×
256
    Deadline deadline = Deadline.after(input.getTimeout(), input.getTimeoutUnit());
×
257
    return pollActivityUntilOutcome(input, deadline)
×
258
        .handle(
×
259
            (outcome, e) -> {
260
              if (e == null) {
×
261
                return decodeOutcome(outcome, input, dc);
×
262
              }
263
              throw handleAsyncException(e, deadline, input.getActivityId());
×
264
            });
265
  }
266

267
  private CompletableFuture<ActivityExecutionOutcome> pollActivityUntilOutcome(
268
      GetActivityResultInput<?> input, Deadline deadline) {
269
    PollActivityExecutionRequest.Builder pollRequest =
270
        PollActivityExecutionRequest.newBuilder()
×
271
            .setNamespace(clientOptions.getNamespace())
×
272
            .setActivityId(input.getActivityId());
×
273
    if (input.getRunId() != null) {
×
274
      pollRequest.setRunId(input.getRunId());
×
275
    }
276
    return genericClient
×
277
        .pollActivityAsync(pollRequest.build(), deadline)
×
278
        .thenComposeAsync(
×
279
            response -> {
280
              if (!response.hasOutcome()) {
×
281
                return pollActivityUntilOutcome(input, deadline);
×
282
              }
283
              return CompletableFuture.completedFuture(response.getOutcome());
×
284
            });
285
  }
286

287
  private static CompletionException handleAsyncException(
288
      Throwable e, Deadline deadline, String activityId) {
289
    Throwable cause = e instanceof CompletionException ? e.getCause() : e;
×
290
    if (deadline.isExpired()
×
291
        && cause instanceof StatusRuntimeException
292
        && Status.Code.DEADLINE_EXCEEDED.equals(
×
293
            ((StatusRuntimeException) cause).getStatus().getCode())) {
×
294
      return new CompletionException(
×
295
          new TimeoutException(
296
              "Activity did not complete within timeout: activityId='" + activityId + "'"));
297
    }
298
    return e instanceof CompletionException ? (CompletionException) e : new CompletionException(e);
×
299
  }
300

301
  private <R> GetActivityResultOutput<R> decodeOutcome(
302
      ActivityExecutionOutcome outcome, GetActivityResultInput<R> input, DataConverter dc) {
303
    switch (outcome.getValueCase()) {
×
304
      case RESULT:
305
        Type resultType =
306
            input.getResultType() != null ? input.getResultType() : input.getResultClass();
×
307
        @SuppressWarnings("unchecked")
308
        R result =
×
309
            (R)
310
                dc.fromPayloads(
×
311
                    0,
312
                    outcome.hasResult() ? Optional.of(outcome.getResult()) : Optional.empty(),
×
313
                    input.getResultClass(),
×
314
                    resultType);
315
        return new GetActivityResultOutput<>(result);
×
316
      case FAILURE:
317
        throw new java.util.concurrent.CompletionException(
×
318
            new ActivityFailedException(
319
                "Activity failed: activityId='" + input.getActivityId() + "'",
×
320
                input.getActivityId(),
×
321
                input.getRunId(),
×
322
                dc.failureToException(outcome.getFailure())));
×
323
      default:
324
        throw new java.util.concurrent.CompletionException(
×
325
            new ActivityFailedException(
326
                "Activity completed with unexpected outcome '"
327
                    + outcome.getValueCase()
×
328
                    + "' for activityId='"
329
                    + input.getActivityId()
×
330
                    + "'",
331
                input.getActivityId(),
×
332
                input.getRunId(),
×
333
                null));
334
    }
335
  }
336

337
  @Override
338
  public DescribeActivityOutput describeActivity(DescribeActivityInput input) {
339
    DescribeActivityExecutionRequest.Builder req =
340
        DescribeActivityExecutionRequest.newBuilder()
×
341
            .setNamespace(clientOptions.getNamespace())
×
342
            .setActivityId(input.getId());
×
343
    if (input.getRunId() != null) {
×
344
      req.setRunId(input.getRunId());
×
345
    }
346
    DescribeActivityExecutionResponse response = genericClient.describeActivity(req.build());
×
347
    return new DescribeActivityOutput(
×
348
        new ActivityExecutionDescription(
349
            response.getInfo(), clientOptions.getDataConverter(), clientOptions.getNamespace()));
×
350
  }
351

352
  @Override
353
  public CancelActivityOutput cancelActivity(CancelActivityInput input) {
354
    RequestCancelActivityExecutionRequest.Builder req =
355
        RequestCancelActivityExecutionRequest.newBuilder()
×
356
            .setNamespace(clientOptions.getNamespace())
×
357
            .setIdentity(clientOptions.getIdentity())
×
358
            .setRequestId(UUID.randomUUID().toString())
×
359
            .setActivityId(input.getId());
×
360
    if (input.getRunId() != null) {
×
361
      req.setRunId(input.getRunId());
×
362
    }
363
    if (input.getReason() != null) {
×
364
      req.setReason(input.getReason());
×
365
    }
366
    genericClient.cancelActivity(req.build());
×
367
    return new CancelActivityOutput();
×
368
  }
369

370
  @Override
371
  public TerminateActivityOutput terminateActivity(TerminateActivityInput input) {
372
    TerminateActivityExecutionRequest.Builder req =
373
        TerminateActivityExecutionRequest.newBuilder()
×
374
            .setNamespace(clientOptions.getNamespace())
×
375
            .setIdentity(clientOptions.getIdentity())
×
376
            .setRequestId(UUID.randomUUID().toString())
×
377
            .setActivityId(input.getId());
×
378
    if (input.getRunId() != null) {
×
379
      req.setRunId(input.getRunId());
×
380
    }
381
    if (input.getReason() != null) {
×
382
      req.setReason(input.getReason());
×
383
    }
384
    genericClient.terminateActivity(req.build());
×
385
    return new TerminateActivityOutput();
×
386
  }
387

388
  @Override
389
  public ListActivitiesOutput listActivities(ListActivitiesInput input) {
390
    ListActivityExecutionIterator iterator =
×
391
        new ListActivityExecutionIterator(
392
            input.getQuery(), clientOptions.getNamespace(), genericClient);
×
393
    iterator.init();
×
394
    Iterator<ActivityExecutionMetadata> wrappedIterator =
×
395
        Iterators.transform(iterator, ActivityExecutionMetadata::fromListInfo);
×
396

397
    final int CHARACTERISTICS = Spliterator.ORDERED | Spliterator.NONNULL | Spliterator.IMMUTABLE;
×
398
    return new ListActivitiesOutput(
×
399
        StreamSupport.stream(
×
400
            Spliterators.spliteratorUnknownSize(wrappedIterator, CHARACTERISTICS), false));
×
401
  }
402

403
  @Override
404
  public CountActivitiesOutput countActivities(CountActivitiesInput input) {
405
    CountActivityExecutionsRequest.Builder req =
406
        CountActivityExecutionsRequest.newBuilder().setNamespace(clientOptions.getNamespace());
×
407
    if (input.getQuery() != null) {
×
408
      req.setQuery(input.getQuery());
×
409
    }
410
    CountActivityExecutionsResponse resp = genericClient.countActivities(req.build());
×
411
    return new CountActivitiesOutput(new ActivityExecutionCount(resp));
×
412
  }
413
}
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