Skip to content

Commit e341161

Browse files
authored
feat(trace): add cadenceIsCron tag in the execute workflow span (#1071)
What changed? Add cadenceIsCron tag on execute workflow span. (In Go, we add it in start workflow span) Why? cadence-workflow/cadence-go-client#1518 How did you test it? Unit Test / Integration Test
1 parent 69aa4c6 commit e341161

3 files changed

Lines changed: 140 additions & 0 deletions

File tree

src/main/java/com/uber/cadence/internal/tracing/TracingPropagator.java

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -45,6 +45,7 @@ public class TracingPropagator {
4545
private static final String TAG_WORKFLOW_TYPE = "cadenceWorkflowType";
4646
private static final String TAG_WORKFLOW_RUN_ID = "cadenceRunID";
4747
private static final String TAG_ACTIVITY_TYPE = "cadenceActivityType";
48+
private static final String TAG_IS_CRON = "cadenceIsCron";
4849

4950
private final Tracer tracer;
5051

@@ -70,6 +71,9 @@ public Span spanForExecuteWorkflow(DecisionContext context) {
7071
.withTag(TAG_WORKFLOW_TYPE, context.getWorkflowType().getName())
7172
.withTag(TAG_WORKFLOW_ID, context.getWorkflowId())
7273
.withTag(TAG_WORKFLOW_RUN_ID, context.getRunId())
74+
.withTag(
75+
TAG_IS_CRON,
76+
attributes.getCronSchedule() != null && !attributes.getCronSchedule().isEmpty())
7377
.start();
7478
}
7579

src/test/java/com/uber/cadence/internal/tracing/StartWorkflowTest.java

Lines changed: 99 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -137,6 +137,18 @@ public Integer Double(Integer n) {
137137
}
138138
}
139139

140+
public interface CronWorkflow {
141+
@WorkflowMethod(executionStartToCloseTimeoutSeconds = 120, taskList = TASK_LIST)
142+
String execute();
143+
}
144+
145+
public static class CronWorkflowImpl implements CronWorkflow {
146+
@Override
147+
public String execute() {
148+
return "done";
149+
}
150+
}
151+
140152
private static final boolean useDockerService = TestEnvironment.isUseDockerService();
141153
private static final Logger logger = LoggerFactory.getLogger(StartWorkflowTest.class);
142154
private static final String DOMAIN = "test-domain";
@@ -259,6 +271,93 @@ public void testSignalStartWorkflowGRPCNoPropagation() {
259271
testSignalWithStartWorkflowHelper(service, mockTracer, false);
260272
}
261273

274+
@Test
275+
public void testCronWorkflowSetsIsCronSpanTagGRPC() {
276+
Assume.assumeTrue(useDockerService);
277+
MockTracer mockTracer = new MockTracer();
278+
IWorkflowService service =
279+
new WorkflowServiceGrpc(
280+
ClientOptions.newBuilder().setTracer(mockTracer).setPort(7833).build());
281+
testCronWorkflowHelper(service, mockTracer);
282+
}
283+
284+
private void testCronWorkflowHelper(IWorkflowService service, MockTracer mockTracer) {
285+
try {
286+
service.RegisterDomain(new RegisterDomainRequest().setName(DOMAIN));
287+
} catch (DomainAlreadyExistsError e) {
288+
logger.info("domain already registered");
289+
} catch (Exception e) {
290+
fail("fail to register domain: " + e);
291+
}
292+
293+
WorkflowClient client =
294+
WorkflowClient.newInstance(
295+
service, WorkflowClientOptions.newBuilder().setDomain(DOMAIN).build());
296+
297+
WorkerFactory workerFactory =
298+
WorkerFactory.newInstance(client, WorkerFactoryOptions.newBuilder().build());
299+
Worker worker = workerFactory.newWorker(TASK_LIST, WorkerOptions.newBuilder().build());
300+
worker.registerWorkflowImplementationTypes(CronWorkflowImpl.class);
301+
workerFactory.start();
302+
303+
Span rootSpan = mockTracer.buildSpan("Test Started").start();
304+
mockTracer.activateSpan(rootSpan);
305+
306+
WorkflowStub wf =
307+
client.newUntypedWorkflowStub(
308+
"CronWorkflow::execute",
309+
new WorkflowOptions.Builder()
310+
.setExecutionStartToCloseTimeout(Duration.ofMinutes(2))
311+
.setTaskList(TASK_LIST)
312+
.setCronSchedule("* * * * *")
313+
.build());
314+
try {
315+
wf.start();
316+
317+
// Cron's minimum interval is 1 minute, so wait for the first scheduled run to execute and
318+
// produce a finished cadence-ExecuteWorkflow span.
319+
MockSpan executeWorkflowSpan = awaitExecuteWorkflowSpan(mockTracer, Duration.ofSeconds(150));
320+
assertNotNull("cadence-ExecuteWorkflow span not found", executeWorkflowSpan);
321+
assertEquals(
322+
"cadenceIsCron tag should be true for cron workflows",
323+
Boolean.TRUE,
324+
executeWorkflowSpan.tags().get("cadenceIsCron"));
325+
} catch (Exception e) {
326+
fail("workflow failure: " + e);
327+
} finally {
328+
try {
329+
wf.cancel();
330+
} catch (Exception ignored) {
331+
// best effort: stop further cron runs
332+
}
333+
rootSpan.finish();
334+
workerFactory.shutdown();
335+
}
336+
}
337+
338+
private MockSpan awaitExecuteWorkflowSpan(MockTracer mockTracer, Duration timeout) {
339+
long deadline = System.currentTimeMillis() + timeout.toMillis();
340+
while (System.currentTimeMillis() < deadline) {
341+
MockSpan span =
342+
mockTracer
343+
.finishedSpans()
344+
.stream()
345+
.filter(s -> s.operationName().equals("cadence-ExecuteWorkflow"))
346+
.findFirst()
347+
.orElse(null);
348+
if (span != null) {
349+
return span;
350+
}
351+
try {
352+
Thread.sleep(1000);
353+
} catch (InterruptedException e) {
354+
Thread.currentThread().interrupt();
355+
break;
356+
}
357+
}
358+
return null;
359+
}
360+
262361
private void testStartWorkflowHelper(
263362
IWorkflowService service, MockTracer mockTracer, boolean shouldPropagate) {
264363
try {

src/test/java/com/uber/cadence/internal/tracing/TracingPropagatorTest.java

Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,10 +19,15 @@
1919

2020
import static org.junit.Assert.assertEquals;
2121
import static org.junit.Assert.assertFalse;
22+
import static org.mockito.Mockito.mock;
23+
import static org.mockito.Mockito.when;
2224

2325
import com.google.common.collect.ImmutableMap;
2426
import com.uber.cadence.Header;
2527
import com.uber.cadence.PollForActivityTaskResponse;
28+
import com.uber.cadence.WorkflowExecutionStartedEventAttributes;
29+
import com.uber.cadence.WorkflowType;
30+
import com.uber.cadence.internal.replay.DecisionContext;
2631
import io.opentracing.Span;
2732
import io.opentracing.mock.MockSpan;
2833
import io.opentracing.mock.MockTracer;
@@ -31,6 +36,8 @@
3136

3237
public class TracingPropagatorTest {
3338

39+
private static final String CADENCE_IS_CRON = "cadenceIsCron";
40+
3441
private final MockTracer mockTracer = new MockTracer();
3542
private final TracingPropagator propagator = new TracingPropagator(mockTracer);
3643

@@ -58,4 +65,34 @@ public void testSpanForExecuteActivity_allowReusingHeaders() {
5865
assertEquals("follows_from", from.getReferenceType());
5966
}
6067
}
68+
69+
@Test
70+
public void testSpanForExecuteWorkflow_withCronSchedule_setsIsCronTagTrue() {
71+
Span span = propagator.spanForExecuteWorkflow(newDecisionContext("0 * * * *"));
72+
span.finish();
73+
74+
MockSpan mockSpan = mockTracer.finishedSpans().get(0);
75+
assertEquals(Boolean.TRUE, mockSpan.tags().get(CADENCE_IS_CRON));
76+
}
77+
78+
@Test
79+
public void testSpanForExecuteWorkflow_withoutCronSchedule_setsIsCronTagFalse() {
80+
Span span = propagator.spanForExecuteWorkflow(newDecisionContext(""));
81+
span.finish();
82+
83+
MockSpan mockSpan = mockTracer.finishedSpans().get(0);
84+
assertEquals(Boolean.FALSE, mockSpan.tags().get(CADENCE_IS_CRON));
85+
}
86+
87+
private static DecisionContext newDecisionContext(String cronSchedule) {
88+
WorkflowExecutionStartedEventAttributes attributes =
89+
new WorkflowExecutionStartedEventAttributes().setCronSchedule(cronSchedule);
90+
91+
DecisionContext context = mock(DecisionContext.class);
92+
when(context.getWorkflowExecutionStartedEventAttributes()).thenReturn(attributes);
93+
when(context.getWorkflowType()).thenReturn(new WorkflowType().setName("TestWorkflow"));
94+
when(context.getWorkflowId()).thenReturn("workflow-id");
95+
when(context.getRunId()).thenReturn("run-id");
96+
return context;
97+
}
6198
}

0 commit comments

Comments
 (0)