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

temporalio / sdk-java / #393

29 Sep 2026 01:52PM UTC coverage: 68.324% (+0.007%) from 68.317%
#393

push

github

web-flow
Release Java SDK v1.40.0 (#3100)

7893 of 13700 branches covered (57.61%)

Branch coverage included in aggregate %.

32037 of 44742 relevant lines covered (71.6%)

0.72 hits per line

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

90.14
/temporal-sdk/src/main/java/io/temporal/internal/sync/SyncWorkflowContext.java
1
package io.temporal.internal.sync;
2

3
import static io.temporal.client.WorkflowClient.QUERY_TYPE_STACK_TRACE;
4
import static io.temporal.client.WorkflowClient.QUERY_TYPE_WORKFLOW_METADATA;
5
import static io.temporal.internal.common.HeaderUtils.intoPayloadMap;
6
import static io.temporal.internal.common.HeaderUtils.toHeaderGrpc;
7
import static io.temporal.internal.common.RetryOptionsUtils.toRetryPolicy;
8
import static io.temporal.internal.common.WorkflowExecutionUtils.makeUserMetaData;
9
import static io.temporal.internal.sync.WorkflowInternal.DEFAULT_VERSION;
10

11
import com.google.common.base.MoreObjects;
12
import com.google.common.base.Preconditions;
13
import com.uber.m3.tally.Scope;
14
import io.temporal.activity.ActivityOptions;
15
import io.temporal.activity.LocalActivityOptions;
16
import io.temporal.api.command.v1.*;
17
import io.temporal.api.common.v1.ActivityType;
18
import io.temporal.api.common.v1.Memo;
19
import io.temporal.api.common.v1.Payload;
20
import io.temporal.api.common.v1.Payloads;
21
import io.temporal.api.common.v1.SearchAttributes;
22
import io.temporal.api.common.v1.WorkflowExecution;
23
import io.temporal.api.common.v1.WorkflowType;
24
import io.temporal.api.enums.v1.ParentClosePolicy;
25
import io.temporal.api.failure.v1.Failure;
26
import io.temporal.api.history.v1.HistoryEvent;
27
import io.temporal.api.sdk.v1.UserMetadata;
28
import io.temporal.api.sdk.v1.WorkflowDefinition;
29
import io.temporal.api.sdk.v1.WorkflowInteractionDefinition;
30
import io.temporal.api.sdk.v1.WorkflowMetadata;
31
import io.temporal.api.taskqueue.v1.TaskQueue;
32
import io.temporal.api.workflowservice.v1.PollActivityTaskQueueResponse;
33
import io.temporal.client.WorkflowException;
34
import io.temporal.common.RetryOptions;
35
import io.temporal.common.SearchAttributeUpdate;
36
import io.temporal.common.VersioningBehavior;
37
import io.temporal.common.context.ContextPropagator;
38
import io.temporal.common.converter.DataConverter;
39
import io.temporal.common.interceptors.Header;
40
import io.temporal.common.interceptors.WorkflowInboundCallsInterceptor;
41
import io.temporal.common.interceptors.WorkflowOutboundCallsInterceptor;
42
import io.temporal.failure.*;
43
import io.temporal.internal.common.*;
44
import io.temporal.internal.replay.ChildWorkflowTaskFailedException;
45
import io.temporal.internal.replay.ReplayWorkflowContext;
46
import io.temporal.internal.replay.WorkflowContext;
47
import io.temporal.internal.statemachines.*;
48
import io.temporal.payload.context.ActivitySerializationContext;
49
import io.temporal.payload.context.NexusSerializationContext;
50
import io.temporal.payload.context.WorkflowSerializationContext;
51
import io.temporal.worker.WorkflowImplementationOptions;
52
import io.temporal.workflow.*;
53
import io.temporal.workflow.Functions.Func;
54
import java.lang.reflect.Type;
55
import java.time.Duration;
56
import java.time.Instant;
57
import java.util.*;
58
import java.util.concurrent.atomic.AtomicBoolean;
59
import java.util.concurrent.atomic.AtomicReference;
60
import java.util.function.BiPredicate;
61
import java.util.function.Supplier;
62
import javax.annotation.Nonnull;
63
import javax.annotation.Nullable;
64
import org.slf4j.Logger;
65
import org.slf4j.LoggerFactory;
66

67
// TODO separate WorkflowOutboundCallsInterceptor functionality from this class into
68
// RootWorkflowOutboundInterceptor
69

70
/**
71
 * Root, the most top level WorkflowContext that unites all relevant contexts, handlers, options,
72
 * states, etc. It's created when SyncWorkflow which represent a context of a workflow type
73
 * definition on the worker meets ReplayWorkflowContext which contains information from the
74
 * WorkflowTask
75
 */
76
final class SyncWorkflowContext implements WorkflowContext, WorkflowOutboundCallsInterceptor {
77
  private static final Logger log = LoggerFactory.getLogger(SyncWorkflowContext.class);
1 ✔
78

79
  private final String namespace;
80
  private final WorkflowExecution workflowExecution;
81
  private final SyncWorkflowDefinition workflowDefinition;
82
  private final WorkflowImplementationOptions workflowImplementationOptions;
83
  private final DataConverter dataConverter;
84
  // to be used in this class, should not be passed down. Pass the original #dataConverter instead
85
  private final DataConverter dataConverterWithCurrentWorkflowContext;
86
  private final List<ContextPropagator> contextPropagators;
87
  private final SignalDispatcher signalDispatcher;
88
  private final QueryDispatcher queryDispatcher;
89
  private final UpdateDispatcher updateDispatcher;
90

91
  // initialized later when these entities are created
92
  private ReplayWorkflowContext replayContext;
93
  private DeterministicRunner runner;
94

95
  private WorkflowInboundCallsInterceptor headInboundInterceptor;
96
  private WorkflowOutboundCallsInterceptor headOutboundInterceptor;
97

98
  private ActivityOptions defaultActivityOptions = null;
1 ✔
99
  private Map<String, ActivityOptions> activityOptionsMap;
100
  private LocalActivityOptions defaultLocalActivityOptions = null;
1 ✔
101
  private Map<String, LocalActivityOptions> localActivityOptionsMap;
102
  private NexusServiceOptions defaultNexusServiceOptions = null;
1 ✔
103
  private Map<String, NexusServiceOptions> nexusServiceOptionsMap;
104
  private boolean readOnly = false;
1 ✔
105
  private final WorkflowThreadLocal<UpdateInfo> currentUpdateInfo = new WorkflowThreadLocal<>();
1 ✔
106
  @Nullable private String currentDetails;
107

108
  public SyncWorkflowContext(
109
      @Nonnull String namespace,
110
      @Nonnull WorkflowExecution workflowExecution,
111
      @Nullable SyncWorkflowDefinition workflowDefinition,
112
      SignalDispatcher signalDispatcher,
113
      QueryDispatcher queryDispatcher,
114
      UpdateDispatcher updateDispatcher,
115
      @Nullable WorkflowImplementationOptions workflowImplementationOptions,
116
      DataConverter dataConverter,
117
      List<ContextPropagator> contextPropagators) {
1 ✔
118
    this.namespace = namespace;
1 ✔
119
    this.workflowExecution = workflowExecution;
1 ✔
120
    this.workflowDefinition = workflowDefinition;
1 ✔
121
    this.dataConverter = dataConverter;
1 ✔
122
    this.dataConverterWithCurrentWorkflowContext =
1 ✔
123
        dataConverter.withContext(
1 ✔
124
            new WorkflowSerializationContext(namespace, workflowExecution.getWorkflowId()));
1 ✔
125
    this.contextPropagators = contextPropagators;
1 ✔
126
    this.signalDispatcher = signalDispatcher;
1 ✔
127
    this.queryDispatcher = queryDispatcher;
1 ✔
128
    this.updateDispatcher = updateDispatcher;
1 ✔
129
    if (workflowImplementationOptions != null) {
1 ✔
130
      this.defaultActivityOptions = workflowImplementationOptions.getDefaultActivityOptions();
1 ✔
131
      this.activityOptionsMap = new HashMap<>(workflowImplementationOptions.getActivityOptions());
1 ✔
132
      this.defaultLocalActivityOptions =
1 ✔
133
          workflowImplementationOptions.getDefaultLocalActivityOptions();
1 ✔
134
      this.localActivityOptionsMap =
1 ✔
135
          new HashMap<>(workflowImplementationOptions.getLocalActivityOptions());
1 ✔
136
      this.defaultNexusServiceOptions =
1 ✔
137
          workflowImplementationOptions.getDefaultNexusServiceOptions();
1 ✔
138
      this.nexusServiceOptionsMap =
1 ✔
139
          new HashMap<>(workflowImplementationOptions.getNexusServiceOptions());
1 ✔
140
    }
141
    this.workflowImplementationOptions =
1 ✔
142
        workflowImplementationOptions == null
1 ✔
143
            ? WorkflowImplementationOptions.getDefaultInstance()
1 ✔
144
            : workflowImplementationOptions;
1 ✔
145
    // initial values for headInboundInterceptor and headOutboundInterceptor until they initialized
146
    // with actual interceptors through #initHeadInboundCallsInterceptor and
147
    // #initHeadOutboundCallsInterceptor during initialization phase.
148
    // See workflow.initialize() performed inside the workflow root thread inside
149
    // SyncWorkflow#start(HistoryEvent, ReplayWorkflowContext)
150
    this.headInboundInterceptor = new InitialWorkflowInboundCallsInterceptor(this);
1 ✔
151
    this.headOutboundInterceptor = this;
1 ✔
152
  }
1 ✔
153

154
  public void setReplayContext(ReplayWorkflowContext context) {
155
    this.replayContext = context;
1 ✔
156
  }
1 ✔
157

158
  /**
159
   * Using setter, as runner is initialized with this context, so it is not ready during
160
   * construction of this.
161
   */
162
  public void setRunner(DeterministicRunner runner) {
163
    this.runner = runner;
1 ✔
164
  }
1 ✔
165

166
  public DeterministicRunner getRunner() {
167
    return runner;
1 ✔
168
  }
169

170
  public WorkflowOutboundCallsInterceptor getWorkflowOutboundInterceptor() {
171
    return headOutboundInterceptor;
1 ✔
172
  }
173

174
  public WorkflowInboundCallsInterceptor getWorkflowInboundInterceptor() {
175
    return headInboundInterceptor;
1 ✔
176
  }
177

178
  public void initHeadOutboundCallsInterceptor(WorkflowOutboundCallsInterceptor head) {
179
    headOutboundInterceptor = head;
1 ✔
180
  }
1 ✔
181

182
  public void initHeadInboundCallsInterceptor(WorkflowInboundCallsInterceptor head) {
183
    headInboundInterceptor = head;
1 ✔
184
    signalDispatcher.setInboundCallsInterceptor(head);
1 ✔
185
    queryDispatcher.setInboundCallsInterceptor(head);
1 ✔
186
    updateDispatcher.setInboundCallsInterceptor(head);
1 ✔
187
  }
1 ✔
188

189
  public ActivityOptions getDefaultActivityOptions() {
190
    return defaultActivityOptions;
1 ✔
191
  }
192

193
  public @Nonnull Map<String, ActivityOptions> getActivityOptions() {
194
    return activityOptionsMap != null
1 !
195
        ? Collections.unmodifiableMap(activityOptionsMap)
1 ✔
196
        : Collections.emptyMap();
×
197
  }
198

199
  public LocalActivityOptions getDefaultLocalActivityOptions() {
200
    return defaultLocalActivityOptions;
1 ✔
201
  }
202

203
  public @Nonnull Map<String, LocalActivityOptions> getLocalActivityOptions() {
204
    return localActivityOptionsMap != null
1 !
205
        ? Collections.unmodifiableMap(localActivityOptionsMap)
1 ✔
206
        : Collections.emptyMap();
×
207
  }
208

209
  public NexusServiceOptions getDefaultNexusServiceOptions() {
210
    return defaultNexusServiceOptions;
1 ✔
211
  }
212

213
  public @Nonnull Map<String, NexusServiceOptions> getNexusServiceOptions() {
214
    return nexusServiceOptionsMap != null
1 !
215
        ? Collections.unmodifiableMap(nexusServiceOptionsMap)
1 ✔
216
        : Collections.emptyMap();
×
217
  }
218

219
  /**
220
   * Unlike the activity options above, the child workflow options have no runtime mutators, so they
221
   * are served directly from the immutable {@link WorkflowImplementationOptions}.
222
   */
223
  public ChildWorkflowOptions getDefaultChildWorkflowOptions() {
224
    return workflowImplementationOptions.getDefaultChildWorkflowOptions();
1 ✔
225
  }
226

227
  public @Nonnull Map<String, ChildWorkflowOptions> getChildWorkflowOptions() {
228
    return workflowImplementationOptions.getChildWorkflowOptions();
1 ✔
229
  }
230

231
  public void setDefaultActivityOptions(ActivityOptions defaultActivityOptions) {
232
    this.defaultActivityOptions =
1 ✔
233
        (this.defaultActivityOptions == null)
1 !
234
            ? defaultActivityOptions
×
235
            : this.defaultActivityOptions.toBuilder()
1 ✔
236
                .mergeActivityOptions(defaultActivityOptions)
1 ✔
237
                .build();
1 ✔
238
  }
1 ✔
239

240
  public void applyActivityOptions(Map<String, ActivityOptions> activityTypeToOption) {
241
    Objects.requireNonNull(activityTypeToOption);
1 ✔
242
    if (this.activityOptionsMap == null) {
1 !
243
      this.activityOptionsMap = new HashMap<>(activityTypeToOption);
×
244
      return;
×
245
    }
246
    ActivityOptionUtils.mergePredefinedActivityOptions(activityOptionsMap, activityTypeToOption);
1 ✔
247
  }
1 ✔
248

249
  public void setDefaultLocalActivityOptions(LocalActivityOptions defaultLocalActivityOptions) {
250
    this.defaultLocalActivityOptions =
1 ✔
251
        (this.defaultLocalActivityOptions == null)
1 !
252
            ? defaultLocalActivityOptions
×
253
            : this.defaultLocalActivityOptions.toBuilder()
1 ✔
254
                .mergeActivityOptions(defaultLocalActivityOptions)
1 ✔
255
                .build();
1 ✔
256
  }
1 ✔
257

258
  public void applyLocalActivityOptions(Map<String, LocalActivityOptions> activityTypeToOption) {
259
    Objects.requireNonNull(activityTypeToOption);
1 ✔
260
    if (this.localActivityOptionsMap == null) {
1 !
261
      this.localActivityOptionsMap = new HashMap<>(activityTypeToOption);
×
262
      return;
×
263
    }
264
    ActivityOptionUtils.mergePredefinedLocalActivityOptions(
1 ✔
265
        localActivityOptionsMap, activityTypeToOption);
266
  }
1 ✔
267

268
  @Override
269
  public <T> ActivityOutput<T> executeActivity(ActivityInput<T> input) {
270
    ActivitySerializationContext serializationContext =
1 ✔
271
        new ActivitySerializationContext(
272
            replayContext.getNamespace(),
1 ✔
273
            replayContext.getWorkflowId(),
1 ✔
274
            replayContext.getWorkflowType().getName(),
1 ✔
275
            input.getActivityName(),
1 ✔
276
            // input.getOptions().getTaskQueue() may be not specified, workflow task queue is used
277
            // by the Server in this case
278
            MoreObjects.firstNonNull(
1 ✔
279
                input.getOptions().getTaskQueue(), replayContext.getTaskQueue()),
1 ✔
280
            false);
281
    DataConverter dataConverterWithActivityContext =
1 ✔
282
        dataConverter.withContext(serializationContext);
1 ✔
283
    Optional<Payloads> args = dataConverterWithActivityContext.toPayloads(input.getArgs());
1 ✔
284

285
    ActivityOutput<Optional<Payloads>> output =
1 ✔
286
        executeActivityOnce(
1 ✔
287
            input.getActivityName(),
1 ✔
288
            input.getActivityId(),
1 ✔
289
            input.getOptions(),
1 ✔
290
            input.getHeader(),
1 ✔
291
            args);
292

293
    // Avoid passing the input to the output handle as it causes the input to be retained for the
294
    // duration of the operation.
295
    Type resultType = input.getResultType();
1 ✔
296
    Class<T> resultClass = input.getResultClass();
1 ✔
297
    return new ActivityOutput<>(
1 ✔
298
        output.getActivityId(),
1 ✔
299
        output
300
            .getResult()
1 ✔
301
            .handle(
1 ✔
302
                (r, f) -> {
303
                  if (f == null) {
1 ✔
304
                    return resultType != Void.TYPE
1 ✔
305
                        ? dataConverterWithActivityContext.fromPayloads(
1 ✔
306
                            0, r, resultClass, resultType)
307
                        : null;
1 ✔
308
                  } else {
309
                    throw dataConverterWithActivityContext.failureToException(
1 ✔
310
                        ((FailureWrapperException) f).getFailure());
1 ✔
311
                  }
312
                }));
313
  }
314

315
  private ActivityOutput<Optional<Payloads>> executeActivityOnce(
316
      String activityTypeName,
317
      @Nullable String activityId,
318
      ActivityOptions options,
319
      Header header,
320
      Optional<Payloads> input) {
321
    ExecuteActivityParameters params =
1 ✔
322
        constructExecuteActivityParameters(activityTypeName, activityId, options, header, input);
1 ✔
323
    ActivityCallback callback = new ActivityCallback();
1 ✔
324
    ReplayWorkflowContext.ScheduleActivityTaskOutput activityOutput =
1 ✔
325
        replayContext.scheduleActivityTask(params, callback::invoke);
1 ✔
326
    CancellationScope.current()
1 ✔
327
        .getCancellationRequest()
1 ✔
328
        .thenApply(
1 ✔
329
            (reason) -> {
330
              activityOutput.getCancellationHandle().apply(new CanceledFailure(reason));
1 ✔
331
              return null;
1 ✔
332
            });
333
    return new ActivityOutput<>(activityOutput.getActivityId(), callback.result);
1 ✔
334
  }
335

336
  public void handleInterceptedSignal(WorkflowInboundCallsInterceptor.SignalInput input) {
337
    signalDispatcher.handleInterceptedSignal(input);
1 ✔
338
  }
1 ✔
339

340
  public void handleSignal(
341
      String signalName, Optional<Payloads> input, long eventId, Header header) {
342
    signalDispatcher.handleSignal(signalName, input, eventId, header);
1 ✔
343
  }
1 ✔
344

345
  public void handleValidateUpdate(
346
      String updateName, String updateId, Optional<Payloads> input, long eventId, Header header) {
347
    updateDispatcher.handleValidateUpdate(updateName, updateId, input, eventId, header);
1 ✔
348
  }
1 ✔
349

350
  public Optional<Payloads> handleExecuteUpdate(
351
      String updateName, String updateId, Optional<Payloads> input, long eventId, Header header) {
352
    return updateDispatcher.handleExecuteUpdate(updateName, updateId, input, eventId, header);
1 ✔
353
  }
354

355
  public void handleInterceptedValidateUpdate(WorkflowInboundCallsInterceptor.UpdateInput input) {
356
    updateDispatcher.handleInterceptedValidateUpdate(input);
1 ✔
357
  }
1 ✔
358

359
  public WorkflowInboundCallsInterceptor.UpdateOutput handleInterceptedExecuteUpdate(
360
      WorkflowInboundCallsInterceptor.UpdateInput input) {
361
    return updateDispatcher.handleInterceptedExecuteUpdate(input);
1 ✔
362
  }
363

364
  public WorkflowInboundCallsInterceptor.QueryOutput handleInterceptedQuery(
365
      WorkflowInboundCallsInterceptor.QueryInput input) {
366
    return queryDispatcher.handleInterceptedQuery(input);
1 ✔
367
  }
368

369
  public Optional<Payloads> handleQuery(String queryName, Header header, Optional<Payloads> input) {
370
    return queryDispatcher.handleQuery(this, queryName, header, input);
1 ✔
371
  }
372

373
  public boolean isEveryHandlerFinished() {
374
    return updateDispatcher.getRunningUpdateHandlers().isEmpty()
1 ✔
375
        && signalDispatcher.getRunningSignalHandlers().isEmpty();
1 ✔
376
  }
377

378
  public WorkflowMetadata getWorkflowMetadata() {
379
    WorkflowMetadata.Builder workflowMetadata = WorkflowMetadata.newBuilder();
1 ✔
380
    WorkflowDefinition.Builder workflowDefinition = WorkflowDefinition.newBuilder();
1 ✔
381
    // Set the workflow type
382
    if (replayContext.getWorkflowType() != null) {
1 !
383
      workflowDefinition.setType(replayContext.getWorkflowType().getName());
1 ✔
384
    }
385
    // Set built in queries
386
    workflowDefinition.addQueryDefinitions(
1 ✔
387
        WorkflowInteractionDefinition.newBuilder()
1 ✔
388
            .setName(QUERY_TYPE_STACK_TRACE)
1 ✔
389
            .setDescription("Current stack trace")
1 ✔
390
            .build());
1 ✔
391
    workflowDefinition.addQueryDefinitions(
1 ✔
392
        WorkflowInteractionDefinition.newBuilder()
1 ✔
393
            .setName(QUERY_TYPE_WORKFLOW_METADATA)
1 ✔
394
            .setDescription("Metadata about the workflow")
1 ✔
395
            .build());
1 ✔
396
    // Add user defined queries
397
    workflowDefinition.addAllQueryDefinitions(queryDispatcher.getQueryHandlers());
1 ✔
398
    // Add user defined signals
399
    workflowDefinition.addAllSignalDefinitions(signalDispatcher.getSignalHandlers());
1 ✔
400
    // Add user defined update handlers
401
    workflowDefinition.addAllUpdateDefinitions(updateDispatcher.getUpdateHandlers());
1 ✔
402
    // Set the workflow definition
403
    workflowMetadata.setDefinition(workflowDefinition.build());
1 ✔
404
    // Add the current workflow details
405
    if (currentDetails != null) {
1 !
406
      workflowMetadata.setCurrentDetails(currentDetails);
1 ✔
407
    }
408
    return workflowMetadata.build();
1 ✔
409
  }
410

411
  private class ActivityCallback {
1 ✔
412
    private final CompletablePromise<Optional<Payloads>> result = Workflow.newPromise();
1 ✔
413

414
    public void invoke(Optional<Payloads> output, Failure failure) {
415
      if (failure != null) {
1 ✔
416
        runner.executeInWorkflowThread(
1 ✔
417
            "activity failure callback",
418
            () -> result.completeExceptionally(new FailureWrapperException(failure)));
1 ✔
419
      } else {
420
        runner.executeInWorkflowThread(
1 ✔
421
            "activity completion callback", () -> result.complete(output));
1 ✔
422
      }
423
    }
1 ✔
424
  }
425

426
  private class LocalActivityCallbackImpl implements LocalActivityCallback {
1 ✔
427
    private final CompletablePromise<Optional<Payloads>> result = Workflow.newPromise();
1 ✔
428

429
    @Override
430
    public void apply(Optional<Payloads> successOutput, LocalActivityFailedException exception) {
431
      if (exception != null) {
1 ✔
432
        runner.executeInWorkflowThread(
1 ✔
433
            "local activity failure callback", () -> result.completeExceptionally(exception));
1 ✔
434
      } else {
435
        runner.executeInWorkflowThread(
1 ✔
436
            "local activity completion callback", () -> result.complete(successOutput));
1 ✔
437
      }
438
    }
1 ✔
439
  }
440

441
  @Override
442
  public <R> LocalActivityOutput<R> executeLocalActivity(LocalActivityInput<R> input) {
443
    ActivitySerializationContext serializationContext =
1 ✔
444
        new ActivitySerializationContext(
445
            replayContext.getNamespace(),
1 ✔
446
            replayContext.getWorkflowId(),
1 ✔
447
            replayContext.getWorkflowType().getName(),
1 ✔
448
            input.getActivityName(),
1 ✔
449
            replayContext.getTaskQueue(),
1 ✔
450
            true);
451
    DataConverter dataConverterWithActivityContext =
1 ✔
452
        dataConverter.withContext(serializationContext);
1 ✔
453
    Optional<Payloads> payloads = dataConverterWithActivityContext.toPayloads(input.getArgs());
1 ✔
454

455
    long originalScheduledTime = System.currentTimeMillis();
1 ✔
456
    CompletablePromise<Optional<Payloads>> serializedResult =
457
        WorkflowInternal.newCompletablePromise();
1 ✔
458
    executeLocalActivityOverLocalRetryThreshold(
1 ✔
459
        input.getActivityName(),
1 ✔
460
        input.getActivityId(),
1 ✔
461
        input.getOptions(),
1 ✔
462
        input.getHeader(),
1 ✔
463
        payloads,
464
        originalScheduledTime,
465
        1,
466
        null,
467
        serializedResult);
468

469
    // Avoid passing the input to the output handle as it causes the input to be retained for the
470
    // duration of the operation.
471
    Type resultType = input.getResultType();
1 ✔
472
    Class<R> resultClass = input.getResultClass();
1 ✔
473
    Promise<R> result =
1 ✔
474
        serializedResult.handle(
1 ✔
475
            (r, f) -> {
476
              if (f == null) {
1 ✔
477
                return resultClass != Void.TYPE
1 ✔
478
                    ? dataConverterWithActivityContext.fromPayloads(0, r, resultClass, resultType)
1 ✔
479
                    : null;
1 ✔
480
              } else {
481
                throw dataConverterWithActivityContext.failureToException(
1 ✔
482
                    ((LocalActivityCallback.LocalActivityFailedException) f).getFailure());
1 ✔
483
              }
484
            });
485

486
    return new LocalActivityOutput<>(result);
1 ✔
487
  }
488

489
  public void executeLocalActivityOverLocalRetryThreshold(
490
      String activityTypeName,
491
      @Nullable String activityId,
492
      LocalActivityOptions options,
493
      Header header,
494
      Optional<Payloads> input,
495
      long originalScheduledTime,
496
      int attempt,
497
      @Nullable Failure previousExecutionFailure,
498
      CompletablePromise<Optional<Payloads>> result) {
499
    CompletablePromise<Optional<Payloads>> localExecutionResult =
1 ✔
500
        executeLocalActivityLocally(
1 ✔
501
            activityTypeName,
502
            activityId,
503
            options,
504
            header,
505
            input,
506
            originalScheduledTime,
507
            attempt,
508
            previousExecutionFailure);
509

510
    localExecutionResult.handle(
1 ✔
511
        (r, e) -> {
512
          if (e == null) {
1 ✔
513
            result.complete(r);
1 ✔
514
          } else {
515
            if ((e instanceof LocalActivityCallback.LocalActivityFailedException)) {
1 !
516
              LocalActivityCallback.LocalActivityFailedException laException =
1 ✔
517
                  (LocalActivityCallback.LocalActivityFailedException) e;
518
              @Nullable Duration backoff = laException.getBackoff();
1 ✔
519
              if (backoff != null) {
1 ✔
520
                WorkflowInternal.newTimer(backoff)
1 ✔
521
                    .thenApply(
1 ✔
522
                        unused -> {
523
                          executeLocalActivityOverLocalRetryThreshold(
1 ✔
524
                              activityTypeName,
525
                              activityId,
526
                              options,
527
                              header,
528
                              input,
529
                              originalScheduledTime,
530
                              laException.getLastAttempt() + 1,
1 ✔
531
                              // Carry the attempt failure, not the local ActivityFailure wrapper.
532
                              laException.getFailure().getCause(),
1 ✔
533
                              result);
534
                          return null;
1 ✔
535
                        });
536
              } else {
537
                // final failure, report back
538
                result.completeExceptionally(laException);
1 ✔
539
              }
540
            } else {
1 ✔
541
              // Only LocalActivityFailedException is expected
542
              String exceptionMessage =
×
543
                  String.format(
×
544
                      "[BUG] Local Activity State Machine callback for activityType %s returned unexpected exception",
545
                      activityTypeName);
546
              log.warn(exceptionMessage, e);
×
547
              replayContext.failWorkflowTask(new IllegalStateException(exceptionMessage, e));
×
548
            }
549
          }
550
          return null;
1 ✔
551
        });
552
  }
1 ✔
553

554
  private CompletablePromise<Optional<Payloads>> executeLocalActivityLocally(
555
      String activityTypeName,
556
      @Nullable String activityId,
557
      LocalActivityOptions options,
558
      Header header,
559
      Optional<Payloads> input,
560
      long originalScheduledTime,
561
      int attempt,
562
      @Nullable Failure previousExecutionFailure) {
563

564
    LocalActivityCallbackImpl callback = new LocalActivityCallbackImpl();
1 ✔
565
    ExecuteLocalActivityParameters params =
1 ✔
566
        constructExecuteLocalActivityParameters(
1 ✔
567
            activityTypeName,
568
            activityId,
569
            options,
570
            header,
571
            input,
572
            attempt,
573
            originalScheduledTime,
574
            previousExecutionFailure);
575
    Functions.Proc cancellationCallback = replayContext.scheduleLocalActivityTask(params, callback);
1 ✔
576
    CancellationScope.current()
1 ✔
577
        .getCancellationRequest()
1 ✔
578
        .thenApply(
1 ✔
579
            (reason) -> {
580
              cancellationCallback.apply();
×
581
              return null;
×
582
            });
583
    return callback.result;
1 ✔
584
  }
585

586
  @SuppressWarnings("deprecation")
587
  private ExecuteActivityParameters constructExecuteActivityParameters(
588
      String name,
589
      @Nullable String activityId,
590
      ActivityOptions options,
591
      Header header,
592
      Optional<Payloads> input) {
593
    String taskQueue = options.getTaskQueue();
1 ✔
594
    if (taskQueue == null) {
1 ✔
595
      taskQueue = replayContext.getTaskQueue();
1 ✔
596
    }
597
    ScheduleActivityTaskCommandAttributes.Builder attributes =
598
        ScheduleActivityTaskCommandAttributes.newBuilder()
1 ✔
599
            .setActivityType(ActivityType.newBuilder().setName(name))
1 ✔
600
            .setTaskQueue(TaskQueue.newBuilder().setName(taskQueue))
1 ✔
601
            .setScheduleToStartTimeout(
1 ✔
602
                ProtobufTimeUtils.toProtoDuration(options.getScheduleToStartTimeout()))
1 ✔
603
            .setStartToCloseTimeout(
1 ✔
604
                ProtobufTimeUtils.toProtoDuration(options.getStartToCloseTimeout()))
1 ✔
605
            .setScheduleToCloseTimeout(
1 ✔
606
                ProtobufTimeUtils.toProtoDuration(options.getScheduleToCloseTimeout()))
1 ✔
607
            .setHeartbeatTimeout(ProtobufTimeUtils.toProtoDuration(options.getHeartbeatTimeout()))
1 ✔
608
            .setRequestEagerExecution(
1 ✔
609
                !options.isEagerExecutionDisabled()
1 ✔
610
                    && Objects.equals(taskQueue, replayContext.getTaskQueue()));
1 ✔
611

612
    if (activityId != null) {
1 ✔
613
      attributes.setActivityId(activityId);
1 ✔
614
    }
615

616
    input.ifPresent(attributes::setInput);
1 ✔
617
    RetryOptions retryOptions = options.getRetryOptions();
1 ✔
618
    if (retryOptions != null) {
1 ✔
619
      attributes.setRetryPolicy(toRetryPolicy(retryOptions));
1 ✔
620
    }
621

622
    // Set the context value.  Use the context propagators from the ActivityOptions
623
    // if present, otherwise use the ones configured on the WorkflowContext
624
    List<ContextPropagator> propagators = options.getContextPropagators();
1 ✔
625
    if (propagators == null) {
1 !
626
      propagators = this.contextPropagators;
1 ✔
627
    }
628
    io.temporal.api.common.v1.Header grpcHeader =
1 ✔
629
        toHeaderGrpc(header, extractContextsAndConvertToBytes(propagators));
1 ✔
630
    attributes.setHeader(grpcHeader);
1 ✔
631

632
    if (options.getVersioningIntent() != null) {
1 !
633
      attributes.setUseWorkflowBuildId(
1 ✔
634
          options
635
              .getVersioningIntent()
1 ✔
636
              .determineUseCompatibleFlag(
1 ✔
637
                  replayContext.getTaskQueue().equals(options.getTaskQueue())));
1 ✔
638
    }
639

640
    @Nullable
641
    UserMetadata userMetadata =
1 ✔
642
        makeUserMetaData(options.getSummary(), null, dataConverterWithCurrentWorkflowContext);
1 ✔
643

644
    if (options.getPriority() != null) {
1 ✔
645
      attributes.setPriority(ProtoConverters.toProto(options.getPriority()));
1 ✔
646
    }
647

648
    return new ExecuteActivityParameters(attributes, options.getCancellationType(), userMetadata);
1 ✔
649
  }
650

651
  private ExecuteLocalActivityParameters constructExecuteLocalActivityParameters(
652
      String name,
653
      @Nullable String activityId,
654
      LocalActivityOptions options,
655
      Header header,
656
      Optional<Payloads> input,
657
      int attempt,
658
      long originalScheduledTime,
659
      @Nullable Failure previousExecutionFailure) {
660
    options = LocalActivityOptions.newBuilder(options).validateAndBuildWithDefaults();
1 ✔
661

662
    PollActivityTaskQueueResponse.Builder activityTask =
663
        PollActivityTaskQueueResponse.newBuilder()
1 ✔
664
            .setActivityId(
1 ✔
665
                activityId != null ? activityId : this.replayContext.randomUUID().toString())
1 ✔
666
            .setWorkflowNamespace(this.replayContext.getNamespace())
1 ✔
667
            .setWorkflowType(this.replayContext.getWorkflowType())
1 ✔
668
            .setWorkflowExecution(this.replayContext.getWorkflowExecution())
1 ✔
669
            // used to pass scheduled time to the local activity code inside
670
            // ActivityExecutionContext#getInfo
671
            // setCurrentAttemptScheduledTime is called inside LocalActivityWorker before submitting
672
            // into the LA queue
673
            .setScheduledTime(
1 ✔
674
                ProtobufTimeUtils.toProtoTimestamp(Instant.ofEpochMilli(originalScheduledTime)))
1 ✔
675
            .setActivityType(ActivityType.newBuilder().setName(name))
1 ✔
676
            .setAttempt(attempt);
1 ✔
677

678
    Duration scheduleToCloseTimeout = options.getScheduleToCloseTimeout();
1 ✔
679
    if (scheduleToCloseTimeout != null) {
1 ✔
680
      activityTask.setScheduleToCloseTimeout(
1 ✔
681
          ProtobufTimeUtils.toProtoDuration(scheduleToCloseTimeout));
1 ✔
682
    }
683

684
    Duration startToCloseTimeout = options.getStartToCloseTimeout();
1 ✔
685
    if (startToCloseTimeout != null) {
1 ✔
686
      activityTask.setStartToCloseTimeout(ProtobufTimeUtils.toProtoDuration(startToCloseTimeout));
1 ✔
687
    }
688

689
    io.temporal.api.common.v1.Header grpcHeader =
1 ✔
690
        toHeaderGrpc(header, extractContextsAndConvertToBytes(contextPropagators));
1 ✔
691
    activityTask.setHeader(grpcHeader);
1 ✔
692
    input.ifPresent(activityTask::setInput);
1 ✔
693
    RetryOptions retryOptions = options.getRetryOptions();
1 ✔
694
    activityTask.setRetryPolicy(
1 ✔
695
        toRetryPolicy(RetryOptions.newBuilder(retryOptions).validateBuildWithDefaults()));
1 ✔
696
    Duration localRetryThreshold = options.getLocalRetryThreshold();
1 ✔
697
    if (localRetryThreshold == null) {
1 ✔
698
      localRetryThreshold = replayContext.getWorkflowTaskTimeout().multipliedBy(3);
1 ✔
699
    }
700

701
    @Nullable
702
    UserMetadata userMetadata =
1 ✔
703
        makeUserMetaData(options.getSummary(), null, dataConverterWithCurrentWorkflowContext);
1 ✔
704

705
    return new ExecuteLocalActivityParameters(
1 ✔
706
        activityTask,
707
        options.getScheduleToStartTimeout(),
1 ✔
708
        originalScheduledTime,
709
        previousExecutionFailure,
710
        options.isDoNotIncludeArgumentsIntoMarker(),
1 ✔
711
        localRetryThreshold,
712
        userMetadata);
713
  }
714

715
  @Override
716
  public <R> ChildWorkflowOutput<R> executeChildWorkflow(ChildWorkflowInput<R> input) {
717
    if (CancellationScope.current().isCancelRequested()) {
1 ✔
718
      CanceledFailure canceledFailure = new CanceledFailure("execute called from a canceled scope");
1 ✔
719
      return new ChildWorkflowOutput<>(
1 ✔
720
          Workflow.newFailedPromise(canceledFailure), Workflow.newFailedPromise(canceledFailure));
1 ✔
721
    }
722

723
    CompletablePromise<WorkflowExecution> executionPromise = Workflow.newPromise();
1 ✔
724
    CompletablePromise<Optional<Payloads>> resultPromise = Workflow.newPromise();
1 ✔
725

726
    DataConverter dataConverterWithChildWorkflowContext =
1 ✔
727
        dataConverter.withContext(
1 ✔
728
            new WorkflowSerializationContext(replayContext.getNamespace(), input.getWorkflowId()));
1 ✔
729
    Optional<Payloads> payloads = dataConverterWithChildWorkflowContext.toPayloads(input.getArgs());
1 ✔
730

731
    @Nullable
732
    Memo memo =
733
        (input.getOptions().getMemo() != null)
1 ✔
734
            ? Memo.newBuilder()
1 ✔
735
                .putAllFields(
1 ✔
736
                    intoPayloadMap(
1 ✔
737
                        dataConverterWithChildWorkflowContext, input.getOptions().getMemo()))
1 ✔
738
                .build()
1 ✔
739
            : null;
1 ✔
740

741
    @Nullable
742
    UserMetadata userMetadata =
1 ✔
743
        makeUserMetaData(
1 ✔
744
            input.getOptions().getStaticSummary(),
1 ✔
745
            input.getOptions().getStaticDetails(),
1 ✔
746
            dataConverterWithChildWorkflowContext);
747

748
    StartChildWorkflowExecutionParameters parameters =
1 ✔
749
        createChildWorkflowParameters(
1 ✔
750
            input.getWorkflowId(),
1 ✔
751
            input.getWorkflowType(),
1 ✔
752
            input.getOptions(),
1 ✔
753
            input.getHeader(),
1 ✔
754
            payloads,
755
            memo,
756
            userMetadata);
757

758
    Functions.Proc1<Exception> cancellationCallback =
1 ✔
759
        replayContext.startChildWorkflow(
1 ✔
760
            parameters,
761
            (execution, failure) -> {
762
              if (failure != null) {
1 ✔
763
                runner.executeInWorkflowThread(
1 ✔
764
                    "child workflow start failed callback",
765
                    () ->
766
                        executionPromise.completeExceptionally(
1 ✔
767
                            mapChildWorkflowException(
1 ✔
768
                                failure, dataConverterWithChildWorkflowContext)));
769
              } else {
770
                runner.executeInWorkflowThread(
1 ✔
771
                    "child workflow started callback", () -> executionPromise.complete(execution));
1 ✔
772
              }
773
            },
1 ✔
774
            (result, failure) -> {
775
              if (failure != null) {
1 ✔
776
                runner.executeInWorkflowThread(
1 ✔
777
                    "child workflow failure callback",
778
                    () ->
779
                        resultPromise.completeExceptionally(
1 ✔
780
                            mapChildWorkflowException(
1 ✔
781
                                failure, dataConverterWithChildWorkflowContext)));
782
              } else {
783
                runner.executeInWorkflowThread(
1 ✔
784
                    "child workflow completion callback", () -> resultPromise.complete(result));
1 ✔
785
              }
786
            });
1 ✔
787
    AtomicBoolean callbackCalled = new AtomicBoolean();
1 ✔
788
    CancellationScope.current()
1 ✔
789
        .getCancellationRequest()
1 ✔
790
        .thenApply(
1 ✔
791
            (reason) -> {
792
              if (!callbackCalled.getAndSet(true)) {
1 !
793
                cancellationCallback.apply(new CanceledFailure(reason));
1 ✔
794
              }
795
              return null;
1 ✔
796
            });
797

798
    // Avoid passing the input to the output handle as it causes the input to be retained for the
799
    // duration of the operation.
800
    Type resultType = input.getResultType();
1 ✔
801
    Class<R> resultClass = input.getResultClass();
1 ✔
802
    Promise<R> result =
1 ✔
803
        resultPromise.thenApply(
1 ✔
804
            (b) ->
805
                dataConverterWithChildWorkflowContext.fromPayloads(0, b, resultClass, resultType));
1 ✔
806
    return new ChildWorkflowOutput<>(result, executionPromise);
1 ✔
807
  }
808

809
  @Override
810
  public <R> ExecuteNexusOperationOutput<R> executeNexusOperation(
811
      ExecuteNexusOperationInput<R> input) {
812
    Preconditions.checkArgument(
1 ✔
813
        input.getEndpoint() != null && !input.getEndpoint().isEmpty(), "endpoint must be set");
1 !
814
    Preconditions.checkArgument(
1 ✔
815
        input.getService() != null && !input.getService().isEmpty(), "service must be set");
1 !
816

817
    if (CancellationScope.current().isCancelRequested()) {
1 !
818
      CanceledFailure canceledFailure =
×
819
          new CanceledFailure("execute nexus operation called from a canceled scope");
820
      return new ExecuteNexusOperationOutput<>(
×
821
          Workflow.newFailedPromise(canceledFailure), Workflow.newFailedPromise(canceledFailure));
×
822
    }
823

824
    CompletablePromise<NexusOperationExecution> operationPromise = Workflow.newPromise();
1 ✔
825
    CompletablePromise<Optional<Payload>> resultPromise = Workflow.newPromise();
1 ✔
826

827
    // The caller workflow is not available to the operation handler, so Nexus payloads are
828
    // contextualized by the endpoint, service and operation instead. The same converter decodes the
829
    // result and converts failures, so each operation keeps the converter selected for it even when
830
    // several operations are in flight at once.
831
    DataConverter nexusDataConverter =
1 ✔
832
        dataConverter.withContext(
1 ✔
833
            new NexusSerializationContext(
834
                input.getEndpoint(), input.getService(), input.getOperation()));
1 ✔
835

836
    Optional<Payload> payload = nexusDataConverter.toPayload(input.getArg());
1 ✔
837

838
    ScheduleNexusOperationCommandAttributes.Builder attributes =
839
        ScheduleNexusOperationCommandAttributes.newBuilder();
1 ✔
840
    payload.ifPresent(attributes::setInput);
1 ✔
841
    attributes.setOperation(input.getOperation());
1 ✔
842
    attributes.setService(input.getService());
1 ✔
843
    attributes.setEndpoint(input.getEndpoint());
1 ✔
844
    // Ensure that the headers are lowercase
845
    input.getHeaders().forEach((k, v) -> attributes.putNexusHeader(k.toLowerCase(), v));
1 ✔
846
    attributes.setScheduleToCloseTimeout(
1 ✔
847
        ProtobufTimeUtils.toProtoDuration(input.getOptions().getScheduleToCloseTimeout()));
1 ✔
848
    attributes.setScheduleToStartTimeout(
1 ✔
849
        ProtobufTimeUtils.toProtoDuration(input.getOptions().getScheduleToStartTimeout()));
1 ✔
850
    attributes.setStartToCloseTimeout(
1 ✔
851
        ProtobufTimeUtils.toProtoDuration(input.getOptions().getStartToCloseTimeout()));
1 ✔
852

853
    @Nullable
854
    UserMetadata userMetadata =
1 ✔
855
        makeUserMetaData(input.getOptions().getSummary(), null, nexusDataConverter);
1 ✔
856

857
    StartNexusOperationParameters parameters =
1 ✔
858
        new StartNexusOperationParameters(
859
            attributes, input.getOptions().getCancellationType(), userMetadata);
1 ✔
860

861
    Functions.Proc1<Exception> cancellationCallback =
1 ✔
862
        replayContext.startNexusOperation(
1 ✔
863
            parameters,
864
            (operationExec, failure) -> {
865
              if (failure != null) {
1 ✔
866
                runner.executeInWorkflowThread(
1 ✔
867
                    "nexus operation start failed callback",
868
                    () ->
869
                        operationPromise.completeExceptionally(
1 ✔
870
                            nexusDataConverter.failureToException(failure)));
1 ✔
871
              } else {
872
                runner.executeInWorkflowThread(
1 ✔
873
                    "nexus operation started callback",
874
                    () ->
875
                        operationPromise.complete(new NexusOperationExecutionImpl(operationExec)));
1 ✔
876
              }
877
            },
1 ✔
878
            (Optional<Payload> result, Failure failure) -> {
879
              if (failure != null) {
1 ✔
880
                runner.executeInWorkflowThread(
1 ✔
881
                    "nexus operation failure callback",
882
                    () ->
883
                        resultPromise.completeExceptionally(
1 ✔
884
                            nexusDataConverter.failureToException(failure)));
1 ✔
885
              } else {
886
                runner.executeInWorkflowThread(
1 ✔
887
                    "nexus operation completion callback", () -> resultPromise.complete(result));
1 ✔
888
              }
889
            });
1 ✔
890
    AtomicBoolean callbackCalled = new AtomicBoolean();
1 ✔
891
    CancellationScope.current()
1 ✔
892
        .getCancellationRequest()
1 ✔
893
        .thenApply(
1 ✔
894
            (reason) -> {
895
              if (!callbackCalled.getAndSet(true)) {
1 !
896
                cancellationCallback.apply(new CanceledFailure(reason));
1 ✔
897
              }
898
              return null;
1 ✔
899
            });
900
    Promise<R> result =
1 ✔
901
        resultPromise.thenApply(
1 ✔
902
            (b) ->
903
                input.getResultClass() != Void.class
1 ✔
904
                    ? nexusDataConverter.fromPayload(
1 ✔
905
                        b.get(), input.getResultClass(), input.getResultType())
1 ✔
906
                    : null);
1 ✔
907
    // We register an empty handler to make sure that this promise is always "accessed" and never
908
    // leads to a log about it being completed exceptionally and non-accessed.
909
    // The "main" operation promise is the one returned from the execute method and that
910
    // promise will always be logged if not accessed.
911
    operationPromise.handle((ex, failure) -> null);
1 ✔
912
    return new ExecuteNexusOperationOutput<>(result, operationPromise);
1 ✔
913
  }
914

915
  @SuppressWarnings("deprecation")
916
  private StartChildWorkflowExecutionParameters createChildWorkflowParameters(
917
      String workflowId,
918
      String name,
919
      ChildWorkflowOptions options,
920
      Header header,
921
      Optional<Payloads> input,
922
      @Nullable Memo memo,
923
      @Nullable UserMetadata metadata) {
924
    final StartChildWorkflowExecutionCommandAttributes.Builder attributes =
925
        StartChildWorkflowExecutionCommandAttributes.newBuilder()
1 ✔
926
            .setWorkflowType(WorkflowType.newBuilder().setName(name).build());
1 ✔
927
    attributes.setWorkflowId(workflowId);
1 ✔
928
    attributes.setNamespace(OptionsUtils.safeGet(options.getNamespace()));
1 ✔
929
    input.ifPresent(attributes::setInput);
1 ✔
930
    attributes.setWorkflowRunTimeout(
1 ✔
931
        ProtobufTimeUtils.toProtoDuration(options.getWorkflowRunTimeout()));
1 ✔
932
    attributes.setWorkflowExecutionTimeout(
1 ✔
933
        ProtobufTimeUtils.toProtoDuration(options.getWorkflowExecutionTimeout()));
1 ✔
934
    attributes.setWorkflowTaskTimeout(
1 ✔
935
        ProtobufTimeUtils.toProtoDuration(options.getWorkflowTaskTimeout()));
1 ✔
936
    String taskQueue = options.getTaskQueue();
1 ✔
937
    if (taskQueue != null) {
1 ✔
938
      attributes.setTaskQueue(TaskQueue.newBuilder().setName(taskQueue));
1 ✔
939
    }
940
    if (options.getWorkflowIdReusePolicy() != null) {
1 ✔
941
      attributes.setWorkflowIdReusePolicy(options.getWorkflowIdReusePolicy());
1 ✔
942
    }
943
    RetryOptions retryOptions = options.getRetryOptions();
1 ✔
944
    if (retryOptions != null) {
1 ✔
945
      attributes.setRetryPolicy(toRetryPolicy(retryOptions));
1 ✔
946
    }
947
    attributes.setCronSchedule(OptionsUtils.safeGet(options.getCronSchedule()));
1 ✔
948

949
    if (memo != null) {
1 ✔
950
      attributes.setMemo(memo);
1 ✔
951
    }
952

953
    Map<String, Object> searchAttributes = options.getSearchAttributes();
1 ✔
954
    if (searchAttributes != null && !searchAttributes.isEmpty()) {
1 !
955
      if (options.getTypedSearchAttributes() != null) {
1 !
956
        throw new IllegalArgumentException(
×
957
            "Cannot have both typed search attributes and search attributes");
958
      }
959
      attributes.setSearchAttributes(SearchAttributesUtil.encode(searchAttributes));
1 ✔
960
    } else if (options.getTypedSearchAttributes() != null
1 ✔
961
        && options.getTypedSearchAttributes().size() > 0) {
1 !
962
      attributes.setSearchAttributes(
1 ✔
963
          SearchAttributesUtil.encodeTyped(options.getTypedSearchAttributes()));
1 ✔
964
    }
965

966
    List<ContextPropagator> propagators = options.getContextPropagators();
1 ✔
967
    if (propagators == null) {
1 !
968
      propagators = this.contextPropagators;
1 ✔
969
    }
970
    io.temporal.api.common.v1.Header grpcHeader =
1 ✔
971
        toHeaderGrpc(header, extractContextsAndConvertToBytes(propagators));
1 ✔
972
    attributes.setHeader(grpcHeader);
1 ✔
973

974
    ParentClosePolicy parentClosePolicy = options.getParentClosePolicy();
1 ✔
975
    if (parentClosePolicy != null) {
1 ✔
976
      attributes.setParentClosePolicy(parentClosePolicy);
1 ✔
977
    }
978

979
    if (options.getVersioningIntent() != null) {
1 !
980
      attributes.setInheritBuildId(
1 ✔
981
          options
982
              .getVersioningIntent()
1 ✔
983
              .determineUseCompatibleFlag(
1 ✔
984
                  replayContext.getTaskQueue().equals(options.getTaskQueue())));
1 ✔
985
    }
986
    if (options.getPriority() != null) {
1 ✔
987
      attributes.setPriority(ProtoConverters.toProto(options.getPriority()));
1 ✔
988
    }
989
    return new StartChildWorkflowExecutionParameters(
1 ✔
990
        attributes, options.getCancellationType(), metadata);
1 ✔
991
  }
992

993
  private static Header extractContextsAndConvertToBytes(
994
      List<ContextPropagator> contextPropagators) {
995
    if (contextPropagators == null) {
1 ✔
996
      return null;
1 ✔
997
    }
998
    Map<String, Payload> result = new HashMap<>();
1 ✔
999
    for (ContextPropagator propagator : contextPropagators) {
1 ✔
1000
      result.putAll(propagator.serializeContext(propagator.getCurrentContext()));
1 ✔
1001
    }
1 ✔
1002
    return new Header(result);
1 ✔
1003
  }
1004

1005
  private static RuntimeException mapChildWorkflowException(
1006
      Exception failure, DataConverter dataConverterWithChildWorkflowContext) {
1007
    if (failure == null) {
1 !
1008
      return null;
×
1009
    }
1010
    if (failure instanceof TemporalFailure) {
1 ✔
1011
      ((TemporalFailure) failure).setDataConverter(dataConverterWithChildWorkflowContext);
1 ✔
1012
    }
1013
    if (failure instanceof CanceledFailure) {
1 ✔
1014
      return (CanceledFailure) failure;
1 ✔
1015
    }
1016
    if (failure instanceof WorkflowException) {
1 !
1017
      return (RuntimeException) failure;
×
1018
    }
1019
    if (failure instanceof ChildWorkflowFailure) {
1 ✔
1020
      return (ChildWorkflowFailure) failure;
1 ✔
1021
    }
1022
    if (!(failure instanceof ChildWorkflowTaskFailedException)) {
1 !
1023
      return new IllegalArgumentException("Unexpected exception type: ", failure);
×
1024
    }
1025
    ChildWorkflowTaskFailedException taskFailed = (ChildWorkflowTaskFailedException) failure;
1 ✔
1026
    Throwable cause =
1 ✔
1027
        dataConverterWithChildWorkflowContext.failureToException(
1 ✔
1028
            taskFailed.getOriginalCauseFailure());
1 ✔
1029
    ChildWorkflowFailure exception = taskFailed.getException();
1 ✔
1030
    return new ChildWorkflowFailure(
1 ✔
1031
        exception.getInitiatedEventId(),
1 ✔
1032
        exception.getStartedEventId(),
1 ✔
1033
        exception.getWorkflowType(),
1 ✔
1034
        exception.getExecution(),
1 ✔
1035
        exception.getNamespace(),
1 ✔
1036
        exception.getRetryState(),
1 ✔
1037
        cause);
1038
  }
1039

1040
  @Override
1041
  public Promise<Void> newTimer(Duration delay) {
1042
    return newTimer(delay, TimerOptions.newBuilder().build());
1 ✔
1043
  }
1044

1045
  @Override
1046
  public Promise<Void> newTimer(Duration delay, TimerOptions options) {
1047
    CompletablePromise<Void> p = Workflow.newPromise();
1 ✔
1048

1049
    @Nullable
1050
    UserMetadata userMetadata =
1 ✔
1051
        makeUserMetaData(options.getSummary(), null, dataConverterWithCurrentWorkflowContext);
1 ✔
1052

1053
    Functions.Proc1<RuntimeException> cancellationHandler =
1 ✔
1054
        replayContext.newTimer(
1 ✔
1055
            delay,
1056
            userMetadata,
1057
            (e) ->
1058
                runner.executeInWorkflowThread(
1 ✔
1059
                    "timer-callback",
1060
                    () -> {
1061
                      if (e == null) {
1 ✔
1062
                        p.complete(null);
1 ✔
1063
                      } else {
1064
                        p.completeExceptionally(e);
1 ✔
1065
                      }
1066
                    }));
1 ✔
1067
    CancellationScope.current()
1 ✔
1068
        .getCancellationRequest()
1 ✔
1069
        .thenApply(
1 ✔
1070
            (r) -> {
1071
              cancellationHandler.apply(new CanceledFailure(r));
1 ✔
1072
              return r;
1 ✔
1073
            });
1074
    return p;
1 ✔
1075
  }
1076

1077
  @Override
1078
  public <R> R sideEffect(Class<R> resultClass, Type resultType, Func<R> func) {
1079
    return sideEffect(resultClass, resultType, func, SideEffectOptions.newBuilder().build());
1 ✔
1080
  }
1081

1082
  @Override
1083
  public <R> R sideEffect(
1084
      Class<R> resultClass, Type resultType, Func<R> func, SideEffectOptions options) {
1085
    @Nullable
1086
    UserMetadata userMetadata =
1 ✔
1087
        makeUserMetaData(options.getSummary(), null, dataConverterWithCurrentWorkflowContext);
1 ✔
1088
    try {
1089
      CompletablePromise<Optional<Payloads>> result = Workflow.newPromise();
1 ✔
1090
      replayContext.sideEffect(
1 ✔
1091
          () -> {
1092
            try {
1093
              readOnly = true;
1 ✔
1094
              R r = func.apply();
1 ✔
1095
              return dataConverterWithCurrentWorkflowContext.toPayloads(r);
1 ✔
1096
            } finally {
1097
              readOnly = false;
1 ✔
1098
            }
1099
          },
1100
          userMetadata,
1101
          (p) ->
1102
              runner.executeInWorkflowThread(
1 ✔
1103
                  "side-effect-callback", () -> result.complete(Objects.requireNonNull(p))));
1 ✔
1104
      return dataConverterWithCurrentWorkflowContext.fromPayloads(
1 ✔
1105
          0, result.get(), resultClass, resultType);
1 ✔
1106
    } catch (Exception e) {
1 ✔
1107
      // SideEffect cannot throw normal exception as it can lead to non-deterministic behavior. So
1108
      // fail the workflow task by throwing an Error.
1109
      throw new Error(e);
1 ✔
1110
    }
1111
  }
1112

1113
  @Override
1114
  public <R> R mutableSideEffect(
1115
      String id, Class<R> resultClass, Type resultType, BiPredicate<R, R> updated, Func<R> func) {
1116
    return mutableSideEffect(
1 ✔
1117
        id, resultClass, resultType, updated, func, MutableSideEffectOptions.newBuilder().build());
1 ✔
1118
  }
1119

1120
  @Override
1121
  public <R> R mutableSideEffect(
1122
      String id,
1123
      Class<R> resultClass,
1124
      Type resultType,
1125
      BiPredicate<R, R> updated,
1126
      Func<R> func,
1127
      MutableSideEffectOptions options) {
1128
    @Nullable
1129
    UserMetadata userMetadata =
1 ✔
1130
        makeUserMetaData(options.getSummary(), null, dataConverterWithCurrentWorkflowContext);
1 ✔
1131
    try {
1132
      return mutableSideEffectImpl(id, userMetadata, resultClass, resultType, updated, func);
1 ✔
1133
    } catch (Exception e) {
1 ✔
1134
      // MutableSideEffect cannot throw normal exception as it can lead to non-deterministic
1135
      // behavior. So fail the workflow task by throwing an Error.
1136
      throw new Error(e);
1 ✔
1137
    }
1138
  }
1139

1140
  private <R> R mutableSideEffectImpl(
1141
      String id,
1142
      UserMetadata metadata,
1143
      Class<R> resultClass,
1144
      Type resultType,
1145
      BiPredicate<R, R> updated,
1146
      Func<R> func) {
1147
    CompletablePromise<Optional<Payloads>> result = Workflow.newPromise();
1 ✔
1148
    AtomicReference<R> unserializedResult = new AtomicReference<>();
1 ✔
1149
    replayContext.mutableSideEffect(
1 ✔
1150
        id,
1151
        metadata,
1152
        (storedBinary) -> {
1153
          Optional<R> stored =
1 ✔
1154
              storedBinary.map(
1 ✔
1155
                  (b) ->
1156
                      dataConverterWithCurrentWorkflowContext.fromPayloads(
1 ✔
1157
                          0, Optional.of(b), resultClass, resultType));
1 ✔
1158
          try {
1159
            readOnly = true;
1 ✔
1160
            R funcResult =
1 ✔
1161
                Objects.requireNonNull(
1 ✔
1162
                    func.apply(), "mutableSideEffect function " + "returned null");
1 ✔
1163
            if (!stored.isPresent() || updated.test(stored.get(), funcResult)) {
1 ✔
1164
              unserializedResult.set(funcResult);
1 ✔
1165
              return dataConverterWithCurrentWorkflowContext.toPayloads(funcResult);
1 ✔
1166
            }
1167
            return Optional.empty(); // returned only when value doesn't need to be updated
1 ✔
1168
          } finally {
1169
            readOnly = false;
1 ✔
1170
          }
1171
        },
1172
        (p) ->
1173
            runner.executeInWorkflowThread(
1 ✔
1174
                "mutable-side-effect-callback", () -> result.complete(Objects.requireNonNull(p))));
1 ✔
1175

1176
    if (!result.get().isPresent()) {
1 !
1177
      throw new IllegalArgumentException("No value found for mutableSideEffectId=" + id);
×
1178
    }
1179
    // An optimization that avoids unnecessary deserialization of the result.
1180
    R unserialized = unserializedResult.get();
1 ✔
1181
    if (unserialized != null) {
1 ✔
1182
      return unserialized;
1 ✔
1183
    }
1184
    return dataConverterWithCurrentWorkflowContext.fromPayloads(
1 ✔
1185
        0, result.get(), resultClass, resultType);
1 ✔
1186
  }
1187

1188
  @Override
1189
  public int getVersion(String changeId, int minSupported, int maxSupported) {
1190
    CompletablePromise<Integer> result = Workflow.newPromise();
1 ✔
1191
    Integer versionToUse;
1192
    try {
1193
      versionToUse =
1 ✔
1194
          replayContext.getVersion(
1 ✔
1195
              changeId,
1196
              minSupported,
1197
              maxSupported,
1198
              (v, e) ->
1199
                  runner.executeInWorkflowThread(
1 ✔
1200
                      "version-callback",
1201
                      () -> {
1202
                        if (v != null) {
1 !
1203
                          result.complete(v);
1 ✔
1204
                        } else {
1205
                          result.completeExceptionally(e);
×
1206
                        }
1207
                      }));
1 ✔
1208
    } catch (UnsupportedVersion.UnsupportedVersionException ex) {
×
1209
      throw new UnsupportedVersion(ex);
×
1210
    }
1 ✔
1211
    /*
1212
     * If we are replaying a workflow and encounter a getVersion call it is possible that this call did not exist
1213
     * on the original execution. If the call did not exist on the original execution then we cannot block on results
1214
     * because it can lead to non-deterministic scheduling.
1215
     * */
1216
    if (replayContext.isReplaying()
1 ✔
1217
        && versionToUse == null
1218
        && replayContext.tryUseSdkFlag(SdkFlag.SKIP_YIELD_ON_DEFAULT_VERSION)
1 ✔
1219
        && minSupported == DEFAULT_VERSION) {
1220
      return DEFAULT_VERSION;
1 ✔
1221
    }
1222

1223
    /*
1224
     * Previously the SDK would yield on the getVersion call to the scheduler. This is not ideal because it can lead to non-deterministic
1225
     * scheduling if the getVersion call was removed.
1226
     * */
1227
    if (replayContext.tryUseSdkFlag(SdkFlag.SKIP_YIELD_ON_VERSION)) {
1 ✔
1228
      // This can happen if we are replaying a workflow and encounter a getVersion call that did not
1229
      // exist on the original execution and the range does not include the default version.
1230
      if (versionToUse == null) {
1 ✔
1231
        versionToUse = DEFAULT_VERSION;
1 ✔
1232
      }
1233
      if (versionToUse < minSupported || versionToUse > maxSupported) {
1 !
1234
        throw new UnsupportedVersion(
1 ✔
1235
            new UnsupportedVersion.UnsupportedVersionException(
1236
                String.format(
1 ✔
1237
                    "Version %d of changeId %s is not supported. Supported v is between %d and %d.",
1238
                    versionToUse, changeId, minSupported, maxSupported)));
1 ✔
1239
      }
1240
      return versionToUse;
1 ✔
1241
    }
1242
    // Legacy behavior if SKIP_YIELD_ON_VERSION is not set. This means this thread will yield on the
1243
    // getVersion call.
1244
    // while it waits for the result.
1245
    try {
1246
      return result.get();
1 ✔
1247
    } catch (UnsupportedVersion.UnsupportedVersionException ex) {
×
1248
      throw new UnsupportedVersion(ex);
×
1249
    }
1250
  }
1251

1252
  @Override
1253
  public void registerQuery(RegisterQueryInput request) {
1254
    queryDispatcher.registerQueryHandlers(request);
1 ✔
1255
  }
1 ✔
1256

1257
  @Override
1258
  public void registerSignalHandlers(RegisterSignalHandlersInput input) {
1259
    signalDispatcher.registerSignalHandlers(input);
1 ✔
1260
  }
1 ✔
1261

1262
  @Override
1263
  public void registerUpdateHandlers(RegisterUpdateHandlersInput input) {
1264
    updateDispatcher.registerUpdateHandlers(input);
1 ✔
1265
  }
1 ✔
1266

1267
  @Override
1268
  public void registerDynamicSignalHandler(RegisterDynamicSignalHandlerInput input) {
1269
    signalDispatcher.registerDynamicSignalHandler(input);
1 ✔
1270
  }
1 ✔
1271

1272
  @Override
1273
  public void registerDynamicQueryHandler(RegisterDynamicQueryHandlerInput input) {
1274
    queryDispatcher.registerDynamicQueryHandler(input);
1 ✔
1275
  }
1 ✔
1276

1277
  @Override
1278
  public void registerDynamicUpdateHandler(RegisterDynamicUpdateHandlerInput input) {
1279
    updateDispatcher.registerDynamicUpdateHandler(input);
1 ✔
1280
  }
1 ✔
1281

1282
  @Override
1283
  public UUID randomUUID() {
1284
    return replayContext.randomUUID();
1 ✔
1285
  }
1286

1287
  @Override
1288
  public Random newRandom() {
1289
    return replayContext.newRandom();
1 ✔
1290
  }
1291

1292
  public DataConverter getDataConverter() {
1293
    return dataConverter;
×
1294
  }
1295

1296
  public DataConverter getDataConverterWithCurrentWorkflowContext() {
1297
    return dataConverterWithCurrentWorkflowContext;
1 ✔
1298
  }
1299

1300
  boolean isReplaying() {
1301
    return replayContext.isReplaying();
1 ✔
1302
  }
1303

1304
  boolean isReadOnly() {
1305
    return readOnly;
1 ✔
1306
  }
1307

1308
  void setReadOnly(boolean readOnly) {
1309
    this.readOnly = readOnly;
1 ✔
1310
  }
1 ✔
1311

1312
  @Override
1313
  public Map<Long, SignalHandlerInfo> getRunningSignalHandlers() {
1314
    return signalDispatcher.getRunningSignalHandlers();
1 ✔
1315
  }
1316

1317
  @Override
1318
  public Map<String, UpdateHandlerInfo> getRunningUpdateHandlers() {
1319
    return updateDispatcher.getRunningUpdateHandlers();
1 ✔
1320
  }
1321

1322
  @Override
1323
  public ReplayWorkflowContext getReplayContext() {
1324
    return replayContext;
1 ✔
1325
  }
1326

1327
  @Override
1328
  public SignalExternalOutput signalExternalWorkflow(SignalExternalInput input) {
1329
    WorkflowExecution childExecution = input.getExecution();
1 ✔
1330
    DataConverter dataConverterWithChildWorkflowContext =
1 ✔
1331
        dataConverter.withContext(
1 ✔
1332
            new WorkflowSerializationContext(
1333
                replayContext.getNamespace(), childExecution.getWorkflowId()));
1 ✔
1334
    SignalExternalWorkflowExecutionCommandAttributes.Builder attributes =
1335
        SignalExternalWorkflowExecutionCommandAttributes.newBuilder();
1 ✔
1336
    attributes.setSignalName(input.getSignalName());
1 ✔
1337
    attributes.setExecution(childExecution);
1 ✔
1338
    attributes.setHeader(HeaderUtils.toHeaderGrpc(input.getHeader(), null));
1 ✔
1339
    Optional<Payloads> payloads = dataConverterWithChildWorkflowContext.toPayloads(input.getArgs());
1 ✔
1340
    payloads.ifPresent(attributes::setInput);
1 ✔
1341
    CompletablePromise<Void> result = Workflow.newPromise();
1 ✔
1342
    Functions.Proc1<Exception> cancellationCallback =
1 ✔
1343
        replayContext.signalExternalWorkflowExecution(
1 ✔
1344
            attributes,
1345
            (output, failure) -> {
1346
              if (failure != null) {
1 ✔
1347
                runner.executeInWorkflowThread(
1 ✔
1348
                    "child workflow failure callback",
1349
                    () ->
1350
                        result.completeExceptionally(
1 ✔
1351
                            dataConverterWithChildWorkflowContext.failureToException(failure)));
1 ✔
1352
              } else {
1353
                runner.executeInWorkflowThread(
1 ✔
1354
                    "child workflow completion callback", () -> result.complete(output));
1 ✔
1355
              }
1356
            });
1 ✔
1357
    CancellationScope.current()
1 ✔
1358
        .getCancellationRequest()
1 ✔
1359
        .thenApply(
1 ✔
1360
            (reason) -> {
1361
              cancellationCallback.apply(new CanceledFailure(reason));
1 ✔
1362
              return null;
1 ✔
1363
            });
1364
    return new SignalExternalOutput(result);
1 ✔
1365
  }
1366

1367
  @Override
1368
  public void sleep(Duration duration) {
1369
    newTimer(duration).get();
1 ✔
1370
  }
1 ✔
1371

1372
  @Override
1373
  public boolean await(Duration timeout, String reason, Supplier<Boolean> unblockCondition) {
1374
    boolean cancelTimerOnCondition =
1 ✔
1375
        replayContext.tryUseSdkFlag(SdkFlag.CANCEL_AWAIT_TIMER_ON_CONDITION);
1 ✔
1376

1377
    if (cancelTimerOnCondition) {
1 ✔
1378
      // If condition is already satisfied, skip creating timer
1379
      if (unblockCondition.get()) {
1 ✔
1380
        return true;
1 ✔
1381
      }
1382
      // Create timer in a cancellation scope so we can cancel it when condition is satisfied
1383
      CompletablePromise<Void> timer = Workflow.newPromise();
1 ✔
1384
      CancellationScope timerScope =
1 ✔
1385
          Workflow.newCancellationScope(() -> timer.completeFrom(newTimer(timeout)));
1 ✔
1386
      timerScope.run();
1 ✔
1387

1388
      WorkflowThread.await(reason, () -> (timer.isCompleted() || unblockCondition.get()));
1 ✔
1389

1390
      boolean conditionSatisfied = !timer.isCompleted();
1 ✔
1391
      if (conditionSatisfied) {
1 ✔
1392
        timerScope.cancel("await condition resolved");
1 ✔
1393
      }
1394
      return conditionSatisfied;
1 ✔
1395
    } else {
1396
      // Old behavior: timer is not cancelled when condition is satisfied
1397
      Promise<Void> timer = newTimer(timeout);
1 ✔
1398
      WorkflowThread.await(reason, () -> (timer.isCompleted() || unblockCondition.get()));
1 !
1399
      return !timer.isCompleted();
1 !
1400
    }
1401
  }
1402

1403
  @Override
1404
  public void await(String reason, Supplier<Boolean> unblockCondition) {
1405
    WorkflowThread.await(reason, unblockCondition);
1 ✔
1406
  }
1 ✔
1407

1408
  @SuppressWarnings("deprecation")
1409
  @Override
1410
  public void continueAsNew(ContinueAsNewInput input) {
1411
    ContinueAsNewWorkflowExecutionCommandAttributes.Builder attributes =
1412
        ContinueAsNewWorkflowExecutionCommandAttributes.newBuilder();
1 ✔
1413
    String workflowType = input.getWorkflowType();
1 ✔
1414
    if (workflowType != null) {
1 ✔
1415
      attributes.setWorkflowType(WorkflowType.newBuilder().setName(workflowType));
1 ✔
1416
    }
1417
    @Nullable ContinueAsNewOptions options = input.getOptions();
1 ✔
1418
    if (options != null) {
1 ✔
1419
      if (options.getWorkflowRunTimeout() != null) {
1 !
1420
        attributes.setWorkflowRunTimeout(
×
1421
            ProtobufTimeUtils.toProtoDuration(options.getWorkflowRunTimeout()));
×
1422
      }
1423
      if (options.getWorkflowTaskTimeout() != null) {
1 !
1424
        attributes.setWorkflowTaskTimeout(
×
1425
            ProtobufTimeUtils.toProtoDuration(options.getWorkflowTaskTimeout()));
×
1426
      }
1427
      if (options.getBackoffStartInterval() != null) {
1 ✔
1428
        attributes.setBackoffStartInterval(
1 ✔
1429
            ProtobufTimeUtils.toProtoDuration(options.getBackoffStartInterval()));
1 ✔
1430
      }
1431
      if (options.getTaskQueue() != null && !options.getTaskQueue().isEmpty()) {
1 ✔
1432
        attributes.setTaskQueue(TaskQueue.newBuilder().setName(options.getTaskQueue()));
1 ✔
1433
      }
1434
      if (options.getRetryOptions() != null) {
1 ✔
1435
        attributes.setRetryPolicy(toRetryPolicy(options.getRetryOptions()));
1 ✔
1436
      } else if (replayContext.getRetryOptions() != null) {
1 ✔
1437
        attributes.setRetryPolicy(toRetryPolicy(replayContext.getRetryOptions()));
1 ✔
1438
      }
1439
      Map<String, Object> searchAttributes = options.getSearchAttributes();
1 ✔
1440
      if (searchAttributes != null && !searchAttributes.isEmpty()) {
1 !
1441
        if (options.getTypedSearchAttributes() != null) {
×
1442
          throw new IllegalArgumentException(
×
1443
              "Cannot have typed search attributes and search attributes");
1444
        }
1445
        attributes.setSearchAttributes(SearchAttributesUtil.encode(searchAttributes));
×
1446
      } else if (options.getTypedSearchAttributes() != null
1 ✔
1447
          && options.getTypedSearchAttributes().size() > 0) {
1 !
1448
        attributes.setSearchAttributes(
1 ✔
1449
            SearchAttributesUtil.encodeTyped(options.getTypedSearchAttributes()));
1 ✔
1450
      } else if (options.getTypedSearchAttributes() == null && searchAttributes == null) {
1 !
1451
        // Carry over existing search attributes if none are specified.
1452
        SearchAttributes existing = replayContext.getSearchAttributes();
1 ✔
1453
        if (existing != null && !existing.getIndexedFieldsMap().isEmpty()) {
1 !
1454
          attributes.setSearchAttributes(existing);
1 ✔
1455
        }
1456
      }
1457
      Map<String, Object> memo = options.getMemo();
1 ✔
1458
      if (memo != null) {
1 ✔
1459
        attributes.setMemo(
1 ✔
1460
            Memo.newBuilder()
1 ✔
1461
                .putAllFields(intoPayloadMap(dataConverterWithCurrentWorkflowContext, memo)));
1 ✔
1462
      }
1463
      if (options.getVersioningIntent() != null) {
1 !
1464
        attributes.setInheritBuildId(
×
1465
            options
1466
                .getVersioningIntent()
×
1467
                .determineUseCompatibleFlag(
×
1468
                    replayContext.getTaskQueue().equals(options.getTaskQueue())));
×
1469
      }
1470
      if (options.getInitialVersioningBehavior() != null) {
1 !
1471
        switch (options.getInitialVersioningBehavior()) {
×
1472
          case AUTO_UPGRADE:
1473
            attributes.setInitialVersioningBehavior(
×
1474
                io.temporal.api.enums.v1.ContinueAsNewVersioningBehavior
1475
                    .CONTINUE_AS_NEW_VERSIONING_BEHAVIOR_AUTO_UPGRADE);
1476
            break;
×
1477
          case USE_RAMPING_VERSION:
1478
            attributes.setInitialVersioningBehavior(
×
1479
                io.temporal.api.enums.v1.ContinueAsNewVersioningBehavior
1480
                    .CONTINUE_AS_NEW_VERSIONING_BEHAVIOR_USE_RAMPING_VERSION);
1481
            break;
1482
        }
1483
      }
1484
    }
1485

1486
    if (options == null && replayContext.getRetryOptions() != null) {
1 ✔
1487
      // Have to copy certain options as server doesn't copy them.
1488
      attributes.setRetryPolicy(toRetryPolicy(replayContext.getRetryOptions()));
1 ✔
1489
    }
1490

1491
    if (options == null && replayContext.getSearchAttributes() != null) {
1 ✔
1492
      // Carry over existing search attributes if none are specified.
1493
      SearchAttributes existing = replayContext.getSearchAttributes();
1 ✔
1494
      if (existing != null && !existing.getIndexedFieldsMap().isEmpty()) {
1 !
1495
        attributes.setSearchAttributes(existing);
1 ✔
1496
      }
1497
    }
1498

1499
    List<ContextPropagator> propagators =
1500
        options != null && options.getContextPropagators() != null
1 !
1501
            ? options.getContextPropagators()
×
1502
            : this.contextPropagators;
1 ✔
1503
    io.temporal.api.common.v1.Header grpcHeader =
1 ✔
1504
        toHeaderGrpc(input.getHeader(), extractContextsAndConvertToBytes(propagators));
1 ✔
1505
    attributes.setHeader(grpcHeader);
1 ✔
1506

1507
    Optional<Payloads> payloads =
1 ✔
1508
        dataConverterWithCurrentWorkflowContext.toPayloads(input.getArgs());
1 ✔
1509
    payloads.ifPresent(attributes::setInput);
1 ✔
1510

1511
    replayContext.continueAsNewOnCompletion(attributes.build());
1 ✔
1512
    WorkflowThread.exit();
×
1513
  }
×
1514

1515
  @Override
1516
  public CancelWorkflowOutput cancelWorkflow(CancelWorkflowInput input) {
1517
    CompletablePromise<Void> result = Workflow.newPromise();
1 ✔
1518
    replayContext.requestCancelExternalWorkflowExecution(
1 ✔
1519
        input.getExecution(),
1 ✔
1520
        input.getReason(),
1 ✔
1521
        (r, exception) -> {
1522
          if (exception == null) {
1 !
1523
            result.complete(null);
1 ✔
1524
          } else {
1525
            result.completeExceptionally(exception);
×
1526
          }
1527
        });
1 ✔
1528
    return new CancelWorkflowOutput(result);
1 ✔
1529
  }
1530

1531
  @Override
1532
  public Scope getMetricsScope() {
1533
    return replayContext.getMetricsScope();
1 ✔
1534
  }
1535

1536
  public boolean isLoggingEnabledInReplay() {
1537
    return replayContext.getEnableLoggingInReplay();
1 ✔
1538
  }
1539

1540
  @Override
1541
  public void upsertSearchAttributes(Map<String, ?> searchAttributes) {
1542
    Preconditions.checkArgument(searchAttributes != null, "null search attributes");
1 !
1543
    Preconditions.checkArgument(!searchAttributes.isEmpty(), "empty search attributes");
1 ✔
1544
    SearchAttributes attr = SearchAttributesUtil.encode(searchAttributes);
1 ✔
1545
    replayContext.upsertSearchAttributes(attr);
1 ✔
1546
  }
1 ✔
1547

1548
  @Override
1549
  public void upsertTypedSearchAttributes(SearchAttributeUpdate<?>... searchAttributeUpdates) {
1550
    SearchAttributes attr = SearchAttributesUtil.encodeTypedUpdates(searchAttributeUpdates);
1 ✔
1551
    replayContext.upsertSearchAttributes(attr);
1 ✔
1552
  }
1 ✔
1553

1554
  @Override
1555
  public void upsertMemo(Map<String, Object> memo) {
1556
    Preconditions.checkArgument(memo != null, "null memo");
1 !
1557
    Preconditions.checkArgument(!memo.isEmpty(), "empty memo");
1 !
1558
    replayContext.upsertMemo(
1 ✔
1559
        Memo.newBuilder()
1 ✔
1560
            .putAllFields(intoPayloadMap(dataConverterWithCurrentWorkflowContext, memo))
1 ✔
1561
            .build());
1 ✔
1562
  }
1 ✔
1563

1564
  @Nonnull
1565
  public Object newWorkflowMethodThreadIntercepted(Runnable runnable, @Nullable String name) {
1566
    return runner.newWorkflowThread(runnable, false, name);
1 ✔
1567
  }
1568

1569
  @Nonnull
1570
  public Object newWorkflowCallbackThreadIntercepted(Runnable runnable, @Nullable String name) {
1571
    return runner.newCallbackThread(runnable, name);
1 ✔
1572
  }
1573

1574
  @Override
1575
  public Object newChildThread(Runnable runnable, boolean detached, String name) {
1576
    return runner.newWorkflowThread(runnable, detached, name);
1 ✔
1577
  }
1578

1579
  @Override
1580
  public long currentTimeMillis() {
1581
    return replayContext.currentTimeMillis();
1 ✔
1582
  }
1583

1584
  /**
1585
   * This WorkflowInboundCallsInterceptor is used during creation of the initial root workflow
1586
   * thread and should be replaced with another specific implementation during initialization stage
1587
   * {@code workflow.initialize()} performed inside the workflow root thread.
1588
   *
1589
   * @see SyncWorkflow#start(HistoryEvent, ReplayWorkflowContext)
1590
   */
1591
  private static final class InitialWorkflowInboundCallsInterceptor
1592
      extends BaseRootWorkflowInboundCallsInterceptor {
1593

1594
    public InitialWorkflowInboundCallsInterceptor(SyncWorkflowContext workflowContext) {
1595
      super(workflowContext);
1 ✔
1596
    }
1 ✔
1597

1598
    @Override
1599
    public WorkflowOutput execute(WorkflowInput input) {
1600
      throw new UnsupportedOperationException(
×
1601
          "SyncWorkflowContext should be initialized with a non-initial WorkflowInboundCallsInterceptor "
1602
              + "before #execute can be called");
1603
    }
1604
  }
1605

1606
  @Nonnull
1607
  @Override
1608
  public WorkflowImplementationOptions getWorkflowImplementationOptions() {
1609
    return workflowImplementationOptions;
1 ✔
1610
  }
1611

1612
  @Override
1613
  public Failure mapWorkflowExceptionToFailure(Throwable failure) {
1614
    return dataConverterWithCurrentWorkflowContext.exceptionToFailure(failure);
1 ✔
1615
  }
1616

1617
  @Nullable
1618
  @Override
1619
  public <R> R getLastCompletionResult(Class<R> resultClass, Type resultType) {
1620
    return dataConverterWithCurrentWorkflowContext.fromPayloads(
1 ✔
1621
        0, Optional.ofNullable(replayContext.getLastCompletionResult()), resultClass, resultType);
1 ✔
1622
  }
1623

1624
  @Override
1625
  public List<ContextPropagator> getContextPropagators() {
1626
    return contextPropagators;
1 ✔
1627
  }
1628

1629
  @Override
1630
  public Map<String, Object> getPropagatedContexts() {
1631
    if (contextPropagators == null || contextPropagators.isEmpty()) {
1 ✔
1632
      return new HashMap<>();
1 ✔
1633
    }
1634

1635
    Map<String, Payload> headerData = new HashMap<>(replayContext.getHeader());
1 ✔
1636
    Map<String, Object> contextData = new HashMap<>();
1 ✔
1637
    for (ContextPropagator propagator : contextPropagators) {
1 ✔
1638
      contextData.put(propagator.getName(), propagator.deserializeContext(headerData));
1 ✔
1639
    }
1 ✔
1640

1641
    return contextData;
1 ✔
1642
  }
1643

1644
  @Override
1645
  public VersioningBehavior getVersioningBehavior() {
1646
    return workflowDefinition.getVersioningBehavior();
1 ✔
1647
  }
1648

1649
  public void setCurrentUpdateInfo(UpdateInfo updateInfo) {
1650
    currentUpdateInfo.set(updateInfo);
1 ✔
1651
  }
1 ✔
1652

1653
  public void setCurrentDetails(String details) {
1654
    currentDetails = details;
1 ✔
1655
  }
1 ✔
1656

1657
  @Nullable
1658
  public Object getInstance() {
1659
    return workflowDefinition.getInstance();
1 ✔
1660
  }
1661

1662
  @Nullable
1663
  public String getCurrentDetails() {
1664
    return currentDetails;
1 ✔
1665
  }
1666

1667
  public Optional<UpdateInfo> getCurrentUpdateInfo() {
1668
    return Optional.ofNullable(currentUpdateInfo.get());
1 ✔
1669
  }
1670

1671
  /** Simple wrapper over a failure just to allow completing the CompletablePromise as a failure */
1672
  private static class FailureWrapperException extends RuntimeException {
1673
    private final Failure failure;
1674

1675
    public FailureWrapperException(Failure failure) {
1 ✔
1676
      this.failure = failure;
1 ✔
1677
    }
1 ✔
1678

1679
    public Failure getFailure() {
1680
      return failure;
1 ✔
1681
    }
1682
  }
1683
}
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