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

temporalio / sdk-java / #356

13 Aug 2026 04:59PM UTC coverage: 68.023% (+0.04%) from 67.982%
#356

push

github

web-flow
💥 Fix unbounded timeout failure chain in local activity  (#3006)

Fix unbounded timeout failure chain

7404 of 12970 branches covered (57.09%)

Branch coverage included in aggregate %.

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

5 existing lines in 2 files now uncovered.

30504 of 42758 relevant lines covered (71.34%)

0.71 hits per line

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

76.42
/temporal-sdk/src/main/java/io/temporal/internal/worker/AsyncWorkflowPollTask.java
1
package io.temporal.internal.worker;
2

3
import static io.temporal.serviceclient.MetricsTag.METRICS_TAGS_CALL_OPTIONS_KEY;
4

5
import com.uber.m3.tally.Scope;
6
import com.uber.m3.util.ImmutableMap;
7
import io.grpc.Context;
8
import io.temporal.api.common.v1.WorkerVersionCapabilities;
9
import io.temporal.api.enums.v1.TaskQueueKind;
10
import io.temporal.api.taskqueue.v1.TaskQueue;
11
import io.temporal.api.workflowservice.v1.*;
12
import io.temporal.internal.common.GrpcUtils;
13
import io.temporal.internal.common.ProtobufTimeUtils;
14
import io.temporal.serviceclient.MetricsTag;
15
import io.temporal.serviceclient.WorkflowServiceStubs;
16
import io.temporal.worker.MetricsType;
17
import io.temporal.worker.PollerTypeMetricsTag;
18
import io.temporal.worker.tuning.SlotPermit;
19
import io.temporal.worker.tuning.WorkflowSlotInfo;
20
import java.util.Objects;
21
import java.util.concurrent.CompletableFuture;
22
import java.util.concurrent.atomic.AtomicBoolean;
23
import java.util.function.Supplier;
24
import javax.annotation.Nonnull;
25
import javax.annotation.Nullable;
26
import org.slf4j.Logger;
27
import org.slf4j.LoggerFactory;
28

29
public class AsyncWorkflowPollTask
30
    implements AsyncPoller.PollTaskAsync<WorkflowTask>, DisableNormalPolling {
31
  private static final Logger log = LoggerFactory.getLogger(AsyncWorkflowPollTask.class);
1✔
32
  private final TrackingSlotSupplier<WorkflowSlotInfo> slotSupplier;
33
  private final WorkflowServiceStubs service;
34
  private final Scope metricsScope;
35
  private final Scope pollerMetricScope;
36
  private final PollWorkflowTaskQueueRequest pollRequest;
37
  private final MetricsTag.TagValue taskQueueTagValue;
38
  private final boolean stickyPoller;
39
  private final Context.CancellableContext grpcContext = Context.ROOT.withCancellation();
1✔
40
  private final AtomicBoolean shutdown = new AtomicBoolean(false);
1✔
41
  private final PollerTracker pollerTracker;
42

43
  @Override
44
  public String toString() {
45
    return "AsyncWorkflowPollTask{" + "stickyPoller=" + stickyPoller + '}';
×
46
  }
47

48
  @SuppressWarnings("deprecation")
49
  public AsyncWorkflowPollTask(
50
      @Nonnull WorkflowServiceStubs service,
51
      @Nonnull String namespace,
52
      @Nonnull String taskQueue,
53
      @Nullable String stickyTaskQueue,
54
      @Nonnull String identity,
55
      @Nonnull String workerInstanceKey,
56
      @Nonnull WorkerVersioningOptions versioningOptions,
57
      @Nonnull TrackingSlotSupplier<WorkflowSlotInfo> slotSupplier,
58
      @Nonnull Scope metricsScope,
59
      @Nonnull Supplier<GetSystemInfoResponse.Capabilities> serverCapabilities,
60
      @Nonnull PollerTracker pollerTracker,
61
      String workerControlTaskQueue) {
1✔
62
    this.service = service;
1✔
63
    this.slotSupplier = slotSupplier;
1✔
64
    this.metricsScope = metricsScope;
1✔
65
    this.pollerTracker = Objects.requireNonNull(pollerTracker);
1✔
66

67
    PollWorkflowTaskQueueRequest.Builder pollRequestBuilder =
68
        PollWorkflowTaskQueueRequest.newBuilder()
1✔
69
            .setNamespace(Objects.requireNonNull(namespace))
1✔
70
            .setIdentity(Objects.requireNonNull(identity));
1✔
71

72
    pollRequestBuilder.setWorkerInstanceKey(workerInstanceKey);
1✔
73
    if (workerControlTaskQueue != null) {
1!
74
      pollRequestBuilder.setWorkerControlTaskQueue(workerControlTaskQueue);
×
75
    }
76

77
    if (versioningOptions.getWorkerDeploymentOptions() != null) {
1!
78
      pollRequestBuilder.setDeploymentOptions(
×
79
          WorkerVersioningProtoUtils.deploymentOptionsToProto(
×
80
              versioningOptions.getWorkerDeploymentOptions()));
×
81
    } else if (serverCapabilities.get().getBuildIdBasedVersioning()) {
1!
82
      pollRequestBuilder.setWorkerVersionCapabilities(
×
83
          WorkerVersionCapabilities.newBuilder()
×
84
              .setBuildId(versioningOptions.getBuildId())
×
85
              .setUseVersioning(versioningOptions.isUsingVersioning())
×
86
              .build());
×
87
    } else {
88
      pollRequestBuilder.setBinaryChecksum(versioningOptions.getBuildId());
1✔
89
    }
90
    stickyPoller = stickyTaskQueue != null && !stickyTaskQueue.isEmpty();
1!
91
    if (!stickyPoller) {
1✔
92
      taskQueueTagValue = PollerTypeMetricsTag.PollerType.WORKFLOW_TASK;
1✔
93
      this.pollRequest =
1✔
94
          pollRequestBuilder
95
              .setTaskQueue(
1✔
96
                  TaskQueue.newBuilder()
1✔
97
                      .setName(taskQueue)
1✔
98
                      .setKind(TaskQueueKind.TASK_QUEUE_KIND_NORMAL)
1✔
99
                      .build())
1✔
100
              .build();
1✔
101
      this.pollerMetricScope =
1✔
102
          metricsScope.tagged(
1✔
103
              new ImmutableMap.Builder<String, String>(1)
104
                  .put(MetricsTag.TASK_QUEUE, String.format("%s:%s", taskQueue, "sticky"))
1✔
105
                  .build());
1✔
106
    } else {
107
      taskQueueTagValue = PollerTypeMetricsTag.PollerType.WORKFLOW_STICKY_TASK;
1✔
108
      this.pollRequest =
1✔
109
          pollRequestBuilder
110
              .setTaskQueue(
1✔
111
                  TaskQueue.newBuilder()
1✔
112
                      .setName(stickyTaskQueue)
1✔
113
                      .setKind(TaskQueueKind.TASK_QUEUE_KIND_STICKY)
1✔
114
                      .setNormalName(taskQueue)
1✔
115
                      .build())
1✔
116
              .build();
1✔
117
      this.pollerMetricScope = metricsScope;
1✔
118
    }
119
  }
1✔
120

121
  @Override
122
  public CompletableFuture<WorkflowTask> poll(SlotPermit permit)
123
      throws AsyncPoller.PollTaskAsyncAbort {
124
    if (shutdown.get()) {
1✔
125
      throw new AsyncPoller.PollTaskAsyncAbort("Normal poller is disabled");
1✔
126
    }
127
    if (log.isTraceEnabled()) {
1!
128
      log.trace("poll request begin: " + pollRequest);
×
129
    }
130

131
    MetricsTag.tagged(metricsScope, taskQueueTagValue)
1✔
132
        .gauge(MetricsType.NUM_POLLERS)
1✔
133
        .update(pollerTracker.pollStarted());
1✔
134

135
    CompletableFuture<PollWorkflowTaskQueueResponse> response = null;
1✔
136
    try {
137
      response =
1✔
138
          grpcContext.call(
1✔
139
              () ->
140
                  GrpcUtils.toCompletableFuture(
1✔
141
                      service
142
                          .futureStub()
1✔
143
                          .withOption(METRICS_TAGS_CALL_OPTIONS_KEY, metricsScope)
1✔
144
                          .pollWorkflowTaskQueue(pollRequest)));
1✔
145
    } catch (Exception e) {
×
146
      MetricsTag.tagged(metricsScope, taskQueueTagValue)
×
147
          .gauge(MetricsType.NUM_POLLERS)
×
148
          .update(pollerTracker.pollCompleted());
×
149
      throw new RuntimeException(e);
×
150
    }
1✔
151

152
    return response
1✔
153
        .thenApply(
1✔
154
            r -> {
155
              if (r == null || r.getTaskToken().isEmpty()) {
1!
UNCOV
156
                pollerMetricScope
×
UNCOV
157
                    .counter(MetricsType.WORKFLOW_TASK_QUEUE_POLL_EMPTY_COUNTER)
×
UNCOV
158
                    .inc(1);
×
UNCOV
159
                return null;
×
160
              }
161
              pollerTracker.pollSucceeded();
1✔
162
              slotSupplier.markSlotUsed(new WorkflowSlotInfo(r, pollRequest), permit);
1✔
163
              pollerMetricScope
1✔
164
                  .counter(MetricsType.WORKFLOW_TASK_QUEUE_POLL_SUCCEED_COUNTER)
1✔
165
                  .inc(1);
1✔
166
              pollerMetricScope
1✔
167
                  .timer(MetricsType.WORKFLOW_TASK_SCHEDULE_TO_START_LATENCY)
1✔
168
                  .record(ProtobufTimeUtils.toM3Duration(r.getStartedTime(), r.getScheduledTime()));
1✔
169
              return new WorkflowTask(r, (reason) -> slotSupplier.releaseSlot(reason, permit));
1✔
170
            })
171
        .whenComplete(
1✔
172
            (r, e) -> {
173
              MetricsTag.tagged(metricsScope, taskQueueTagValue)
1✔
174
                  .gauge(MetricsType.NUM_POLLERS)
1✔
175
                  .update(pollerTracker.pollCompleted());
1✔
176
            });
1✔
177
  }
178

179
  @Override
180
  public void cancel(Throwable cause) {
181
    grpcContext.cancel(cause);
1✔
182
  }
1✔
183

184
  @Override
185
  public void disableNormalPoll() {
186
    if (stickyPoller) {
1!
187
      throw new IllegalStateException("Cannot disable normal poll for sticky poller");
×
188
    }
189
    shutdown.set(true);
1✔
190
  }
1✔
191

192
  @Override
193
  public String getLabel() {
194
    return stickyPoller ? "StickyWorkflowPollTask" : "NormalWorkflowPollTask";
1✔
195
  }
196
}
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