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

temporalio / sdk-java / #387

15 Sep 2026 10:00PM UTC coverage: 68.408% (+0.04%) from 68.372%
#387

push

github

web-flow
Account for dynamic workflows when setting WorkflowImplementationOptions (#3042)

7844 of 13600 branches covered (57.68%)

Branch coverage included in aggregate %.

6 of 6 new or added lines in 1 file covered. (100.0%)

3 existing lines in 3 files now uncovered.

31842 of 44414 relevant lines covered (71.69%)

0.72 hits per line

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

90.76
/temporal-sdk/src/main/java/io/temporal/internal/sync/SyncWorkflow.java
1
package io.temporal.internal.sync;
2

3
import io.temporal.api.common.v1.*;
4
import io.temporal.api.enums.v1.EventType;
5
import io.temporal.api.history.v1.HistoryEvent;
6
import io.temporal.api.history.v1.WorkflowExecutionStartedEventAttributes;
7
import io.temporal.api.query.v1.WorkflowQuery;
8
import io.temporal.client.WorkflowClient;
9
import io.temporal.common.context.ContextPropagator;
10
import io.temporal.common.converter.DataConverter;
11
import io.temporal.common.converter.DefaultDataConverter;
12
import io.temporal.common.converter.RawValue;
13
import io.temporal.internal.logging.LoggerTag;
14
import io.temporal.internal.replay.ReplayWorkflow;
15
import io.temporal.internal.replay.ReplayWorkflowContext;
16
import io.temporal.internal.replay.WorkflowContext;
17
import io.temporal.internal.statemachines.UpdateProtocolCallback;
18
import io.temporal.internal.worker.WorkflowExecutionException;
19
import io.temporal.internal.worker.WorkflowExecutorCache;
20
import io.temporal.payload.context.WorkflowSerializationContext;
21
import io.temporal.worker.WorkflowImplementationOptions;
22
import io.temporal.workflow.UpdateInfo;
23
import java.util.List;
24
import java.util.Objects;
25
import java.util.Optional;
26
import javax.annotation.Nonnull;
27
import javax.annotation.Nullable;
28
import org.slf4j.Logger;
29
import org.slf4j.LoggerFactory;
30
import org.slf4j.MDC;
31

32
/**
33
 * SyncWorkflow supports workflows that use synchronous blocking code. An instance is created per
34
 * cached workflow run.
35
 */
36
class SyncWorkflow implements ReplayWorkflow {
37

38
  private static final Logger log = LoggerFactory.getLogger(SyncWorkflow.class);
1✔
39

40
  private final WorkflowThreadExecutor workflowThreadExecutor;
41
  private final SyncWorkflowDefinition workflow;
42
  @Nonnull private final WorkflowImplementationOptions workflowImplementationOptions;
43
  private final WorkflowExecutorCache cache;
44
  private final long defaultDeadlockDetectionTimeout;
45
  private final WorkflowMethodThreadNameStrategy workflowMethodThreadNameStrategy =
1✔
46
      ExecutionInfoStrategy.INSTANCE;
47
  private final SyncWorkflowContext workflowContext;
48
  private WorkflowExecutionHandler workflowProc;
49
  private DeterministicRunner runner;
50
  private DataConverter dataConverter;
51
  private DataConverter dataConverterWithWorkflowContext;
52

53
  public SyncWorkflow(
54
      String namespace,
55
      WorkflowExecution workflowExecution,
56
      SyncWorkflowDefinition workflow,
57
      SignalDispatcher signalDispatcher,
58
      QueryDispatcher queryDispatcher,
59
      UpdateDispatcher updateDispatcher,
60
      @Nullable WorkflowImplementationOptions workflowImplementationOptions,
61
      DataConverter dataConverter,
62
      WorkflowThreadExecutor workflowThreadExecutor,
63
      WorkflowExecutorCache cache,
64
      List<ContextPropagator> contextPropagators,
65
      long defaultDeadlockDetectionTimeout) {
1✔
66
    this.workflow = Objects.requireNonNull(workflow);
1✔
67
    this.workflowImplementationOptions =
1✔
68
        workflowImplementationOptions == null
1!
UNCOV
69
            ? WorkflowImplementationOptions.getDefaultInstance()
×
70
            : workflowImplementationOptions;
1✔
71
    this.workflowThreadExecutor = Objects.requireNonNull(workflowThreadExecutor);
1✔
72
    this.cache = cache;
1✔
73
    this.defaultDeadlockDetectionTimeout = defaultDeadlockDetectionTimeout;
1✔
74
    this.dataConverter = dataConverter;
1✔
75
    this.dataConverterWithWorkflowContext =
1✔
76
        dataConverter.withContext(
1✔
77
            new WorkflowSerializationContext(namespace, workflowExecution.getWorkflowId()));
1✔
78
    this.workflowContext =
1✔
79
        new SyncWorkflowContext(
80
            namespace,
81
            workflowExecution,
82
            workflow,
83
            signalDispatcher,
84
            queryDispatcher,
85
            updateDispatcher,
86
            workflowImplementationOptions,
87
            dataConverter,
88
            contextPropagators);
89
  }
1✔
90

91
  @Override
92
  public void start(HistoryEvent event, ReplayWorkflowContext context) {
93
    if (event.getEventType() != EventType.EVENT_TYPE_WORKFLOW_EXECUTION_STARTED
1!
94
        || !event.hasWorkflowExecutionStartedEventAttributes()) {
1!
95
      throw new IllegalArgumentException(
×
96
          "first event is not WorkflowExecutionStarted, but " + event.getEventType());
×
97
    }
98

99
    WorkflowExecutionStartedEventAttributes startEvent =
1✔
100
        event.getWorkflowExecutionStartedEventAttributes();
1✔
101
    WorkflowType workflowType = startEvent.getWorkflowType();
1✔
102
    if (workflow == null) {
1!
103
      throw new IllegalArgumentException("Unknown workflow type: " + workflowType);
×
104
    }
105

106
    this.workflowContext.setReplayContext(context);
1✔
107

108
    workflowProc =
1✔
109
        new WorkflowExecutionHandler(
110
            workflowContext, workflow, startEvent, workflowImplementationOptions);
111
    // The following order is ensured by this code and DeterministicRunner implementation:
112
    // 1. workflow.initialize
113
    // 2. signal handler (if signalWithStart was called)
114
    // 3. main workflow method
115
    runner =
1✔
116
        DeterministicRunner.newRunner(
1✔
117
            workflowThreadExecutor,
118
            workflowContext,
119
            () -> {
120
              workflowProc.runConstructor();
1✔
121
              WorkflowInternal.newWorkflowMethodThread(
1✔
122
                      () -> workflowProc.runWorkflowMethod(),
1✔
123
                      workflowMethodThreadNameStrategy.createThreadName(
1✔
124
                          context.getWorkflowExecution()))
1✔
125
                  .start();
1✔
126
            },
1✔
127
            cache);
128
  }
1✔
129

130
  @Override
131
  public void handleSignal(
132
      String signalName, Optional<Payloads> input, long eventId, Header header) {
133
    // Signals can trigger completion
134
    runner.executeInWorkflowThread(
1✔
135
        "signal " + signalName,
136
        () -> {
137
          workflowProc.handleSignal(signalName, input, eventId, header);
1✔
138
        });
1✔
139
  }
1✔
140

141
  @Override
142
  public void handleUpdate(
143
      String updateName,
144
      String updateId,
145
      Optional<Payloads> input,
146
      long eventId,
147
      Header header,
148
      UpdateProtocolCallback callbacks) {
149
    final UpdateInfo updateInfo = new UpdateInfoImpl(updateName, updateId);
1✔
150
    runner.executeInWorkflowThread(
1✔
151
        "update " + updateName,
152
        () -> {
153
          try {
154
            workflowContext.setCurrentUpdateInfo(updateInfo);
1✔
155
            MDC.put(LoggerTag.UPDATE_ID, updateInfo.getUpdateId());
1✔
156
            MDC.put(LoggerTag.UPDATE_NAME, updateInfo.getUpdateName());
1✔
157
            // Skip validator on replay
158
            if (!callbacks.isReplaying()) {
1✔
159
              try {
160
                workflowContext.setReadOnly(true);
1✔
161
                workflowProc.handleValidateUpdate(updateName, updateId, input, eventId, header);
1✔
162
              } catch (ReadOnlyException r) {
1✔
163
                // Rethrow instead on rejecting the update to fail the WFT
164
                throw r;
1✔
165
              } catch (WorkflowExecutionException e) {
1✔
166
                callbacks.reject(e.getFailure());
1✔
167
                return;
1✔
168
              } catch (Exception e) {
1✔
169
                callbacks.reject(
1✔
170
                    workflowContext
171
                        .getDataConverterWithCurrentWorkflowContext()
1✔
172
                        .exceptionToFailure(e));
1✔
173
                return;
1✔
174
              } finally {
175
                workflowContext.setReadOnly(false);
1✔
176
              }
177
            }
178
            callbacks.accept();
1✔
179
            try {
180
              Optional<Payloads> result =
1✔
181
                  workflowProc.handleExecuteUpdate(updateName, updateId, input, eventId, header);
1✔
182
              callbacks.complete(result, null);
1✔
183
            } catch (WorkflowExecutionException e) {
1✔
184
              callbacks.complete(Optional.empty(), e.getFailure());
1✔
185
            }
1✔
186
          } finally {
187
            workflowContext.setCurrentUpdateInfo(null);
1✔
188
          }
189
        });
1✔
190
  }
1✔
191

192
  @Override
193
  public boolean eventLoop() {
194
    if (runner == null) {
1!
195
      return false;
×
196
    }
197
    runner.runUntilAllBlocked(defaultDeadlockDetectionTimeout);
1✔
198
    return runner.isDone() || workflowProc.isDone(); // Do not wait for all other threads.
1✔
199
  }
200

201
  @Override
202
  public Optional<Payloads> getOutput() {
203
    return workflowProc.getOutput();
1✔
204
  }
205

206
  @Override
207
  public void cancel(String reason) {
208
    runner.cancel(reason);
1✔
209
  }
1✔
210

211
  @Override
212
  public void close() {
213
    if (runner != null) {
1!
214
      runner.close();
1✔
215
    }
216
  }
1✔
217

218
  @Override
219
  public Optional<Payloads> query(WorkflowQuery query) {
220
    if (WorkflowClient.QUERY_TYPE_REPLAY_ONLY.equals(query.getQueryType())) {
1✔
221
      return Optional.empty();
1✔
222
    }
223
    if (WorkflowClient.QUERY_TYPE_STACK_TRACE.equals(query.getQueryType())) {
1✔
224
      // stack trace query result should be readable for UI even if user specifies a custom data
225
      // converter
226
      return DefaultDataConverter.STANDARD_INSTANCE.toPayloads(runner.stackTrace());
1✔
227
    }
228
    if (WorkflowClient.QUERY_TYPE_WORKFLOW_METADATA.equals(query.getQueryType())) {
1✔
229
      // metadata should be readable independent of user DataConverter settings
230
      Payload payload =
1✔
231
          DefaultDataConverter.STANDARD_INSTANCE
232
              .toPayload(workflowContext.getWorkflowMetadata())
1✔
233
              .orElseThrow(() -> new IllegalStateException("Failed to serialize metadata"));
1✔
234
      return dataConverterWithWorkflowContext.toPayloads(new RawValue(payload));
1✔
235
    }
236
    Optional<Payloads> args =
237
        query.hasQueryArgs() ? Optional.of(query.getQueryArgs()) : Optional.empty();
1✔
238
    return workflowProc.handleQuery(query.getQueryType(), query.getHeader(), args);
1✔
239
  }
240

241
  @Override
242
  public WorkflowContext getWorkflowContext() {
243
    return workflowContext;
1✔
244
  }
245
}
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