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

temporalio / sdk-java / #400

06 Oct 2026 02:56AM UTC coverage: 68.324% (+0.007%) from 68.317%
#400

push

github

web-flow
Mark dev server APIs as stable (#3112)

7902 of 13718 branches covered (57.6%)

Branch coverage included in aggregate %.

32093 of 44819 relevant lines covered (71.61%)

0.72 hits per line

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

74.06
/temporal-sdk/src/main/java/io/temporal/internal/sync/WorkflowThreadImpl.java
1
package io.temporal.internal.sync;
2

3
import com.google.common.base.Preconditions;
4
import io.temporal.common.context.ContextPropagator;
5
import io.temporal.failure.CanceledFailure;
6
import io.temporal.internal.common.NonIdempotentHandle;
7
import io.temporal.internal.common.SdkFlag;
8
import io.temporal.internal.context.ContextThreadLocal;
9
import io.temporal.internal.logging.LoggerTag;
10
import io.temporal.internal.logging.PrefixedMdc;
11
import io.temporal.internal.replay.ReplayWorkflowContext;
12
import io.temporal.internal.worker.WorkflowExecutorCache;
13
import io.temporal.workflow.Functions;
14
import io.temporal.workflow.Promise;
15
import java.io.PrintWriter;
16
import java.io.StringWriter;
17
import java.util.HashMap;
18
import java.util.List;
19
import java.util.Map;
20
import java.util.Optional;
21
import java.util.concurrent.CompletableFuture;
22
import java.util.concurrent.Future;
23
import java.util.concurrent.RejectedExecutionException;
24
import java.util.function.Supplier;
25
import javax.annotation.Nonnull;
26
import javax.annotation.Nullable;
27
import org.slf4j.Logger;
28
import org.slf4j.LoggerFactory;
29
import org.slf4j.MDC;
30

31
class WorkflowThreadImpl implements WorkflowThread {
32
  /**
33
   * Runnable passed to the thread that wraps a runnable passed to the WorkflowThreadImpl
34
   * constructor.
35
   */
36
  class RunnableWrapper implements Runnable {
37

38
    private final WorkflowThreadContext threadContext;
39
    // TODO: Move MDC injection logic into an interceptor as this context shouldn't be leaked
40
    // to the WorkflowThreadImpl
41
    private final ReplayWorkflowContext replayWorkflowContext;
42
    private String originalName;
43
    private String name;
44
    private final CancellationScopeImpl cancellationScope;
45
    private final List<ContextPropagator> contextPropagators;
46
    private final Map<String, Object> propagatedContexts;
47

48
    RunnableWrapper(
49
        WorkflowThreadContext threadContext,
50
        ReplayWorkflowContext replayWorkflowContext,
51
        String name,
52
        boolean detached,
53
        CancellationScopeImpl parent,
54
        Runnable runnable,
55
        List<ContextPropagator> contextPropagators,
56
        Map<String, Object> propagatedContexts) {
1 ✔
57
      this.threadContext = threadContext;
1 ✔
58
      this.replayWorkflowContext = replayWorkflowContext;
1 ✔
59
      this.name = name;
1 ✔
60
      boolean deterministicCancellationScopeOrder =
1 ✔
61
          replayWorkflowContext.checkSdkFlag(SdkFlag.DETERMINISTIC_CANCELLATION_SCOPE_ORDER);
1 ✔
62
      this.cancellationScope =
1 ✔
63
          new CancellationScopeImpl(
64
              detached, deterministicCancellationScopeOrder, runnable, parent);
65
      Preconditions.checkState(
1 ✔
66
          context.getStatus() == Status.CREATED, "threadContext not in CREATED state");
1 !
67
      this.contextPropagators = contextPropagators;
1 ✔
68
      this.propagatedContexts = propagatedContexts;
1 ✔
69
    }
1 ✔
70

71
    @Override
72
    public void run() {
73
      Thread thread = Thread.currentThread();
1 ✔
74
      originalName = thread.getName();
1 ✔
75
      thread.setName(name);
1 ✔
76

77
      threadContext.initializeCurrentThread(thread);
1 ✔
78
      DeterministicRunnerImpl.setCurrentThreadInternal(WorkflowThreadImpl.this);
1 ✔
79

80
      PrefixedMdc mdc = replayWorkflowContext.getLoggerMdc();
1 ✔
81
      mdc.put(LoggerTag.WORKFLOW_ID, replayWorkflowContext.getWorkflowId());
1 ✔
82
      mdc.put(LoggerTag.WORKFLOW_TYPE, replayWorkflowContext.getWorkflowType().getName());
1 ✔
83
      mdc.put(LoggerTag.RUN_ID, replayWorkflowContext.getRunId());
1 ✔
84
      mdc.put(LoggerTag.TASK_QUEUE, replayWorkflowContext.getTaskQueue());
1 ✔
85
      mdc.put(LoggerTag.NAMESPACE, replayWorkflowContext.getNamespace());
1 ✔
86

87
      // Repopulate the context(s)
88
      ContextThreadLocal.setContextPropagators(this.contextPropagators);
1 ✔
89
      ContextThreadLocal.propagateContextToCurrentThread(this.propagatedContexts);
1 ✔
90
      try {
91
        // initialYield blocks thread until the first runUntilBlocked is called.
92
        // Otherwise, r starts executing without control of the sync.
93
        threadContext.initialYield();
1 ✔
94
        cancellationScope.run();
1 ✔
95
      } catch (DestroyWorkflowThreadError e) {
1 ✔
96
        if (!threadContext.isDestroyRequested()) {
1 ✔
97
          threadContext.setUnhandledException(e);
1 ✔
98
        }
99
      } catch (Error e) {
1 ✔
100
        threadContext.setUnhandledException(e);
1 ✔
101
      } catch (CanceledFailure e) {
×
102
        if (!isCancelRequested()) {
×
103
          threadContext.setUnhandledException(e);
×
104
        }
105
        if (log.isDebugEnabled()) {
×
106
          log.debug(String.format("Workflow thread \"%s\" run canceled", name));
×
107
        }
108
      } catch (Throwable e) {
1 ✔
109
        threadContext.setUnhandledException(e);
1 ✔
110
      } finally {
111
        DeterministicRunnerImpl.setCurrentThreadInternal(null);
1 ✔
112
        threadContext.makeDone();
1 ✔
113
        thread.setName(originalName);
1 ✔
114
        MDC.clear();
1 ✔
115
      }
116
    }
1 ✔
117

118
    public String getName() {
119
      return name;
1 ✔
120
    }
121

122
    StackTraceElement[] getStackTrace() {
123
      @Nullable Thread thread = threadContext.getCurrentThread();
1 ✔
124
      if (thread != null) {
1 ✔
125
        return thread.getStackTrace();
1 ✔
126
      }
127
      return new StackTraceElement[0];
1 ✔
128
    }
129

130
    public void setName(String name) {
131
      this.name = name;
×
132
      @Nullable Thread thread = threadContext.getCurrentThread();
×
133
      if (thread != null) {
×
134
        thread.setName(name);
×
135
      }
136
    }
×
137
  }
138

139
  private static final Logger log = LoggerFactory.getLogger(WorkflowThreadImpl.class);
1 ✔
140

141
  private final WorkflowThreadExecutor workflowThreadExecutor;
142
  private final WorkflowThreadContext context;
143
  private final WorkflowExecutorCache cache;
144
  private final SyncWorkflowContext syncWorkflowContext;
145

146
  private final DeterministicRunnerImpl runner;
147
  private final RunnableWrapper task;
148
  private final int priority;
149
  private Future<?> taskFuture;
150
  private final Map<WorkflowThreadLocalInternal<?>, Object> threadLocalMap = new HashMap<>();
1 ✔
151

152
  WorkflowThreadImpl(
153
      WorkflowThreadExecutor workflowThreadExecutor,
154
      SyncWorkflowContext syncWorkflowContext,
155
      DeterministicRunnerImpl runner,
156
      @Nonnull String name,
157
      int priority,
158
      boolean detached,
159
      CancellationScopeImpl parentCancellationScope,
160
      Runnable runnable,
161
      WorkflowExecutorCache cache,
162
      List<ContextPropagator> contextPropagators,
163
      Map<String, Object> propagatedContexts) {
1 ✔
164
    this.workflowThreadExecutor = workflowThreadExecutor;
1 ✔
165
    this.syncWorkflowContext = Preconditions.checkNotNull(syncWorkflowContext);
1 ✔
166
    this.runner = runner;
1 ✔
167
    this.context = new WorkflowThreadContext(runner.getLock());
1 ✔
168
    this.cache = cache;
1 ✔
169
    this.priority = priority;
1 ✔
170
    this.task =
1 ✔
171
        new RunnableWrapper(
172
            context,
173
            syncWorkflowContext.getReplayContext(),
1 ✔
174
            Preconditions.checkNotNull(name, "Thread name shouldn't be null"),
1 ✔
175
            detached,
176
            parentCancellationScope,
177
            runnable,
178
            contextPropagators,
179
            propagatedContexts);
180
  }
1 ✔
181

182
  @Override
183
  public void run() {
184
    throw new UnsupportedOperationException("not used");
×
185
  }
186

187
  @Override
188
  public boolean isDetached() {
189
    return task.cancellationScope.isDetached();
×
190
  }
191

192
  @Override
193
  public void cancel() {
194
    task.cancellationScope.cancel();
1 ✔
195
  }
1 ✔
196

197
  @Override
198
  public void cancel(String reason) {
199
    task.cancellationScope.cancel(reason);
1 ✔
200
  }
1 ✔
201

202
  @Override
203
  public String getCancellationReason() {
204
    return task.cancellationScope.getCancellationReason();
×
205
  }
206

207
  @Override
208
  public boolean isCancelRequested() {
209
    return task.cancellationScope.isCancelRequested();
×
210
  }
211

212
  @Override
213
  public Promise<String> getCancellationRequest() {
214
    return task.cancellationScope.getCancellationRequest();
×
215
  }
216

217
  @Override
218
  public void start() {
219
    context.verifyAndStart();
1 ✔
220
    while (true) {
221
      try {
222
        taskFuture = workflowThreadExecutor.submit(task);
1 ✔
223
        return;
1 ✔
224
      } catch (RejectedExecutionException e) {
1 ✔
225
        if (cache != null) {
1 ✔
226
          SyncWorkflowContext workflowContext = getWorkflowContext();
1 ✔
227
          ReplayWorkflowContext context = workflowContext.getReplayContext();
1 ✔
228
          boolean evicted =
1 ✔
229
              cache.evictAnyNotInProcessing(
1 ✔
230
                  context.getWorkflowExecution(), workflowContext.getMetricsScope());
1 ✔
231
          if (!evicted) {
1 !
232
            // Note here we need to throw error, not exception. Otherwise it will be
233
            // translated to workflow execution exception and instead of failing the
234
            // workflow task we will be failing the workflow.
235
            throw new WorkflowRejectedExecutionError(e);
×
236
          }
237
        } else {
1 ✔
238
          throw new WorkflowRejectedExecutionError(e);
1 ✔
239
        }
240
      }
1 ✔
241
    }
242
  }
243

244
  @Override
245
  public boolean isStarted() {
246
    return context.getStatus() != Status.CREATED;
×
247
  }
248

249
  @Override
250
  public WorkflowThreadContext getWorkflowThreadContext() {
251
    return context;
1 ✔
252
  }
253

254
  @Override
255
  public DeterministicRunnerImpl getRunner() {
256
    return runner;
1 ✔
257
  }
258

259
  @Override
260
  public SyncWorkflowContext getWorkflowContext() {
261
    return syncWorkflowContext;
1 ✔
262
  }
263

264
  @Override
265
  public void setName(String name) {
266
    task.setName(name);
×
267
  }
×
268

269
  @Override
270
  public String getName() {
271
    return task.getName();
1 ✔
272
  }
273

274
  @Override
275
  public long getId() {
276
    return hashCode();
×
277
  }
278

279
  @Override
280
  public int getPriority() {
281
    return priority;
1 ✔
282
  }
283

284
  @Override
285
  public boolean runUntilBlocked(long deadlockDetectionTimeoutMs) {
286
    if (taskFuture == null) {
1 ✔
287
      start();
1 ✔
288
    }
289
    return context.runUntilBlocked(deadlockDetectionTimeoutMs);
1 ✔
290
  }
291

292
  @Override
293
  public NonIdempotentHandle lockDeadlockDetector() {
294
    return context.lockDeadlockDetector();
1 ✔
295
  }
296

297
  @Override
298
  public boolean isDone() {
299
    return context.isDone();
1 ✔
300
  }
301

302
  @Override
303
  public Throwable getUnhandledException() {
304
    return context.getUnhandledException();
1 ✔
305
  }
306

307
  /**
308
   * Evaluates function in the threadContext of the coroutine without unblocking it. Used to get
309
   * current coroutine status, like stack trace.
310
   *
311
   * @param function Parameter is reason for current goroutine blockage.
312
   */
313
  public void evaluateInCoroutineContext(Functions.Proc1<String> function) {
314
    context.evaluateInCoroutineContext(function);
×
315
  }
×
316

317
  /**
318
   * Interrupt coroutine by throwing DestroyWorkflowThreadError from an await method it is blocked
319
   * on and return underlying Future to be waited on.
320
   */
321
  @Override
322
  public Future<?> stopNow() {
323
    // Cannot call destroy() on itself
324
    @Nullable Thread thread = context.getCurrentThread();
1 ✔
325
    if (Thread.currentThread().equals(thread)) {
1 !
326
      throw new Error("Cannot call destroy on itself: " + thread.getName());
×
327
    }
328
    context.initiateDestroy();
1 ✔
329
    if (taskFuture == null) {
1 ✔
330
      return getCompletedFuture();
1 ✔
331
    }
332
    return taskFuture;
1 ✔
333
  }
334

335
  private Future<?> getCompletedFuture() {
336
    CompletableFuture<String> f = new CompletableFuture<>();
1 ✔
337
    f.complete("done");
1 ✔
338
    return f;
1 ✔
339
  }
340

341
  @Override
342
  public void addStackTrace(StringBuilder result) {
343
    result.append(getName());
1 ✔
344
    @Nullable Thread thread = context.getCurrentThread();
1 ✔
345
    if (thread == null) {
1 !
346
      result.append("(NEW)");
×
347
      return;
×
348
    }
349
    result
1 ✔
350
        .append(": (BLOCKED on ")
1 ✔
351
        .append(getWorkflowThreadContext().getYieldReason())
1 ✔
352
        .append(")\n");
1 ✔
353
    // These numbers might change if implementation changes.
354
    int omitTop = 5;
1 ✔
355
    int omitBottom = 7;
1 ✔
356
    // TODO it's not a good idea to rely on the name to understand the thread type. Instead of that
357
    // we would better
358
    // assign an explicit thread type enum to the threads. This will be especially important when we
359
    // refactor
360
    // root and workflow-method
361
    // thread names into names that will include workflowId
362
    if (DeterministicRunnerImpl.WORKFLOW_ROOT_THREAD_NAME.equals(getName())) {
1 !
363
      // TODO revisit this number
364
      omitBottom = 11;
×
365
    } else if (getName().startsWith(WorkflowMethodThreadNameStrategy.WORKFLOW_MAIN_THREAD_PREFIX)) {
1 ✔
366
      // TODO revisit this number
367
      omitBottom = 11;
1 ✔
368
    }
369
    StackTraceElement[] stackTrace = thread.getStackTrace();
1 ✔
370
    for (int i = omitTop; i < stackTrace.length - omitBottom; i++) {
1 ✔
371
      StackTraceElement e = stackTrace[i];
1 ✔
372
      if (i == omitTop && "await".equals(e.getMethodName())) continue;
1 !
373
      result.append(e);
1 ✔
374
      result.append("\n");
1 ✔
375
    }
376
  }
1 ✔
377

378
  @Override
379
  public void yield(String reason, Supplier<Boolean> unblockCondition) {
380
    context.yield(reason, unblockCondition);
1 ✔
381
  }
1 ✔
382

383
  @Override
384
  public void exitThread() {
385
    runner.exit();
1 ✔
386
    throw new DestroyWorkflowThreadError("exit");
1 ✔
387
  }
388

389
  @Override
390
  public <T> void setThreadLocal(WorkflowThreadLocalInternal<T> key, T value) {
391
    threadLocalMap.put(key, value);
1 ✔
392
  }
1 ✔
393

394
  /**
395
   * Retrieve data from thread locals. Returns 1. not found (an empty Optional) 2. found but null
396
   * (an Optional of an empty Optional) 3. found and non-null (an Optional of an Optional of a
397
   * value). The type nesting is because Java Optionals cannot understand "Some null" vs "None",
398
   * which is exactly what we need here.
399
   *
400
   * @param key
401
   * @return one of three cases
402
   * @param <T>
403
   */
404
  @SuppressWarnings("unchecked")
405
  public <T> Optional<Optional<T>> getThreadLocal(WorkflowThreadLocalInternal<T> key) {
406
    if (!threadLocalMap.containsKey(key)) {
1 ✔
407
      return Optional.empty();
1 ✔
408
    }
409
    return Optional.of(Optional.ofNullable((T) threadLocalMap.get(key)));
1 ✔
410
  }
411

412
  /**
413
   * @return stack trace of the coroutine thread
414
   */
415
  @Override
416
  public String getStackTrace() {
417
    StackTraceElement[] st = task.getStackTrace();
1 ✔
418
    StringWriter sw = new StringWriter();
1 ✔
419
    PrintWriter pw = new PrintWriter(sw);
1 ✔
420
    pw.append(task.getName());
1 ✔
421
    pw.append("\n");
1 ✔
422
    for (StackTraceElement se : st) {
1 ✔
423
      pw.println("\tat " + se);
1 ✔
424
    }
425
    return sw.toString();
1 ✔
426
  }
427

428
  static class YieldWithTimeoutCondition implements Supplier<Boolean> {
429

430
    private final Supplier<Boolean> unblockCondition;
431
    private final long blockedUntil;
432
    private boolean timedOut;
433

434
    YieldWithTimeoutCondition(Supplier<Boolean> unblockCondition, long blockedUntil) {
×
435
      this.unblockCondition = unblockCondition;
×
436
      this.blockedUntil = blockedUntil;
×
437
    }
×
438

439
    boolean isTimedOut() {
440
      return timedOut;
×
441
    }
442

443
    /**
444
     * @return true if condition matched or timed out
445
     */
446
    @Override
447
    public Boolean get() {
448
      boolean result = unblockCondition.get();
×
449
      if (result) {
×
450
        return true;
×
451
      }
452
      long currentTimeMillis = WorkflowInternal.currentTimeMillis();
×
453
      timedOut = currentTimeMillis >= blockedUntil;
×
454
      return timedOut;
×
455
    }
456
  }
457
}
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