From 3459cb476869ca33ccb0963b2f9f4a71c789ccf8 Mon Sep 17 00:00:00 2001 From: Varshi Bachu Date: Fri, 7 Aug 2026 09:04:19 -0700 Subject: [PATCH 1/6] initial commit --- .../durabletask/TaskEntityExecutor.java | 18 ++ .../microsoft/durabletask/TracingHelper.java | 190 ++++++++++++++++++ .../durabletask/TaskEntityExecutorTest.java | 63 ++++++ .../durabletask/TracingHelperTest.java | 175 ++++++++++++++++ 4 files changed, 446 insertions(+) diff --git a/client/src/main/java/com/microsoft/durabletask/TaskEntityExecutor.java b/client/src/main/java/com/microsoft/durabletask/TaskEntityExecutor.java index e2352da2..c2380f8b 100644 --- a/client/src/main/java/com/microsoft/durabletask/TaskEntityExecutor.java +++ b/client/src/main/java/com/microsoft/durabletask/TaskEntityExecutor.java @@ -121,6 +121,10 @@ EntityBatchResult execute(@Nonnull EntityBatchRequest request) { logger.log(Level.FINE, "Executing operation '{0}' (requestId={1}) on entity '{2}'.", new Object[]{operationName, requestId, instanceId}); + // Parent context for any signals/orchestrations this operation produces, so the host can link them. + context.setCurrentOperationTraceContext( + opRequest.hasTraceContext() ? opRequest.getTraceContext() : null); + // Snapshot state and actions before each operation (for rollback on failure) entityState.commit(); context.commit(); @@ -221,12 +225,18 @@ private static class TaskEntityContextImpl extends TaskEntityContext { private final DataConverter dataConverter; private final List pendingActions = new ArrayList<>(); private int committedActionCount = 0; + @Nullable + private TraceContext currentOperationTraceContext; TaskEntityContextImpl(EntityInstanceId entityId, DataConverter dataConverter) { this.entityId = entityId; this.dataConverter = dataConverter; } + void setCurrentOperationTraceContext(@Nullable TraceContext traceContext) { + this.currentOperationTraceContext = traceContext; + } + @Nonnull @Override public EntityInstanceId getId() { @@ -261,6 +271,10 @@ public void signalEntity( .build()); } + if (this.currentOperationTraceContext != null) { + signalBuilder.setParentTraceContext(this.currentOperationTraceContext); + } + this.pendingActions.add(new PendingAction(PendingAction.Type.SEND_SIGNAL, signalBuilder.build(), null)); } @@ -300,6 +314,10 @@ public String startNewOrchestration( } } + if (this.currentOperationTraceContext != null) { + orchBuilder.setParentTraceContext(this.currentOperationTraceContext); + } + this.pendingActions.add(new PendingAction( PendingAction.Type.START_NEW_ORCHESTRATION, null, orchBuilder.build())); diff --git a/client/src/main/java/com/microsoft/durabletask/TracingHelper.java b/client/src/main/java/com/microsoft/durabletask/TracingHelper.java index 79f9dcd7..9e610308 100644 --- a/client/src/main/java/com/microsoft/durabletask/TracingHelper.java +++ b/client/src/main/java/com/microsoft/durabletask/TracingHelper.java @@ -76,6 +76,16 @@ final class TracingHelper { static final String ATTR_FIRE_AT = "durabletask.fire_at"; static final String ATTR_EVENT_TARGET_INSTANCE_ID = "durabletask.event.target_instance_id"; + // Entity span constants matching .NET SDK schema. Entity spans deliberately do NOT set + // durabletask.task.name/version/task_id: no .NET entity trace helper sets them, and the + // entity name already appears in the span name. + static final String TYPE_ENTITY = "entity"; + static final String OP_CALL_ENTITY = "call_entity"; + static final String OP_SIGNAL_ENTITY = "signal_entity"; + static final String ATTR_OPERATION = "durabletask.task.operation"; + static final String ATTR_SCHEDULED_TIME = "durabletask.task.scheduled_time"; + static final String ATTR_ENTITY_ERROR_MESSAGE = "durabletask.entity.error_message"; + private TracingHelper() { // Static utility class } @@ -476,4 +486,184 @@ static void emitEventSpan( spanBuilder.startSpan().end(); } + + // region Entity spans + + /** Builds an entity span name: {@code entity::}. */ + static String createEntitySpanName(String entityName, String operation) { + return TYPE_ENTITY + ":" + entityName + ":" + operation; + } + + /** Builds an entity-starts-orchestration span name: {@code :create_orchestration}. */ + static String createEntityStartOrchestrationSpanName(String entityName) { + return entityName + ":" + TYPE_CREATE_ORCHESTRATION; + } + + /** + * Starts a processing span for an entity operation: {@link SpanKind#SERVER} for a call, + * {@link SpanKind#CONSUMER} for a signal. Returns {@code null} when the parent context is absent + * or invalid, matching the .NET guard that avoids attaching entity spans to an unrelated ambient + * trace. The caller makes the span current and later calls {@link #endEntityProcessingSpan}. + */ + @Nullable + static Span startEntityProcessingSpan( + String entityName, + String operation, + boolean signal, + String entityInstanceId, + @Nullable TraceContext parentContext) { + Context parentCtx = extractTraceContext(parentContext); + if (parentCtx == null) { + return null; + } + Tracer tracer = GlobalOpenTelemetry.getTracer(TRACER_NAME); + return tracer.spanBuilder(createEntitySpanName(entityName, operation)) + .setSpanKind(signal ? SpanKind.CONSUMER : SpanKind.SERVER) + .setParent(parentCtx) + .setAttribute(ATTR_TYPE, TYPE_ENTITY) + .setAttribute(ATTR_OPERATION, signal ? OP_SIGNAL_ENTITY : OP_CALL_ENTITY) + .setAttribute(ATTR_INSTANCE_ID, entityInstanceId) + .startSpan(); + } + + /** + * Ends a processing span with {@code OK}/{@code Completed} on success or {@code ERROR} plus + * {@code durabletask.entity.error_message} on failure, matching .NET's + * {@code EndActivitiesForProcessingEntityInvocation}. + */ + static void endEntityProcessingSpan(@Nullable Span span, @Nullable String errorMessage) { + if (span == null) { + return; + } + if (errorMessage != null) { + span.setAttribute(ATTR_ENTITY_ERROR_MESSAGE, errorMessage); + span.setStatus(StatusCode.ERROR, errorMessage); + } else { + span.setStatus(StatusCode.OK, "Completed"); + } + span.end(); + } + + /** + * Emits the retroactive {@link SpanKind#CLIENT} span for a call to an entity, covering the + * request-to-response interval and sharing {@code syntheticSpanId} with the SERVER processing + * span. The status is left unset for normal completion or entity-processing failure (matching + * .NET); {@code errorDescription} is set only for timeout/cancellation closure boundaries. + * Does nothing when the parent context is absent or invalid. + */ + static void emitEntityCallClientSpan( + String entityName, + String operation, + String targetEntityInstanceId, + @Nullable TraceContext parentContext, + @Nullable java.time.Instant startTime, + @Nullable java.time.Instant endTime, + @Nullable String syntheticSpanId, + @Nullable String scheduledTime, + @Nullable String errorDescription) { + Context parentCtx = extractTraceContext(parentContext); + if (parentCtx == null) { + return; + } + Tracer tracer = GlobalOpenTelemetry.getTracer(TRACER_NAME); + SpanBuilder spanBuilder = tracer.spanBuilder(createEntitySpanName(entityName, operation)) + .setSpanKind(SpanKind.CLIENT) + .setParent(parentCtx) + .setAttribute(ATTR_TYPE, TYPE_ENTITY) + .setAttribute(ATTR_OPERATION, OP_CALL_ENTITY) + .setAttribute(ATTR_EVENT_TARGET_INSTANCE_ID, targetEntityInstanceId); + if (scheduledTime != null) { + spanBuilder.setAttribute(ATTR_SCHEDULED_TIME, scheduledTime); + } + if (startTime != null) { + spanBuilder.setStartTimestamp(startTime); + } + Span span = spanBuilder.startSpan(); + setSpanId(span, syntheticSpanId); + if (errorDescription != null) { + span.setStatus(StatusCode.ERROR, errorDescription); + } + if (endTime != null) { + span.end(toEpochNanos(endTime), java.util.concurrent.TimeUnit.NANOSECONDS); + } else { + span.end(); + } + } + + /** + * Starts a {@link SpanKind#PRODUCER} span for signaling an entity. Used by orchestration signals, + * entity-to-entity signals, and the external client signal path. Returns {@code null} when the + * parent context is absent or invalid. Short-lived callers end the span immediately; the client + * path ends it in a {@code finally} block after the gRPC call. + */ + @Nullable + static Span startEntitySignalProducerSpan( + String targetEntityName, + String operation, + String targetEntityInstanceId, + @Nullable String sourceEntityInstanceId, + @Nullable TraceContext parentContext, + @Nullable java.time.Instant startTime, + @Nullable String scheduledTime) { + Context parentCtx = extractTraceContext(parentContext); + if (parentCtx == null) { + return null; + } + Tracer tracer = GlobalOpenTelemetry.getTracer(TRACER_NAME); + SpanBuilder spanBuilder = tracer.spanBuilder(createEntitySpanName(targetEntityName, operation)) + .setSpanKind(SpanKind.PRODUCER) + .setParent(parentCtx) + .setAttribute(ATTR_TYPE, TYPE_ENTITY) + .setAttribute(ATTR_OPERATION, OP_SIGNAL_ENTITY) + .setAttribute(ATTR_EVENT_TARGET_INSTANCE_ID, targetEntityInstanceId); + if (sourceEntityInstanceId != null) { + spanBuilder.setAttribute(ATTR_INSTANCE_ID, sourceEntityInstanceId); + } + if (scheduledTime != null) { + spanBuilder.setAttribute(ATTR_SCHEDULED_TIME, scheduledTime); + } + if (startTime != null) { + spanBuilder.setStartTimestamp(startTime); + } + return spanBuilder.startSpan(); + } + + /** + * Starts a {@link SpanKind#PRODUCER} span for an entity starting an orchestration. The span name + * is {@code :create_orchestration}. Returns {@code null} when the parent context is + * absent or invalid. + */ + @Nullable + static Span startEntityStartOrchestrationSpan( + String sourceEntityName, + String sourceEntityInstanceId, + String targetOrchestrationInstanceId, + @Nullable TraceContext parentContext, + @Nullable java.time.Instant startTime, + @Nullable String scheduledTime) { + Context parentCtx = extractTraceContext(parentContext); + if (parentCtx == null) { + return null; + } + Tracer tracer = GlobalOpenTelemetry.getTracer(TRACER_NAME); + SpanBuilder spanBuilder = tracer.spanBuilder(createEntityStartOrchestrationSpanName(sourceEntityName)) + .setSpanKind(SpanKind.PRODUCER) + .setParent(parentCtx) + .setAttribute(ATTR_TYPE, TYPE_ENTITY) + .setAttribute(ATTR_EVENT_TARGET_INSTANCE_ID, targetOrchestrationInstanceId) + .setAttribute(ATTR_INSTANCE_ID, sourceEntityInstanceId); + if (scheduledTime != null) { + spanBuilder.setAttribute(ATTR_SCHEDULED_TIME, scheduledTime); + } + if (startTime != null) { + spanBuilder.setStartTimestamp(startTime); + } + return spanBuilder.startSpan(); + } + + private static long toEpochNanos(java.time.Instant instant) { + return instant.getEpochSecond() * 1_000_000_000L + instant.getNano(); + } + + // endregion } diff --git a/client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTest.java b/client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTest.java index 7b678df8..30de8732 100644 --- a/client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTest.java +++ b/client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTest.java @@ -380,6 +380,69 @@ void execute_entityStartsOrchestration_actionsIncluded() { assertEquals("MyOrchestration", orchAction.getName()); } + @Test + void execute_entitySignalsOther_propagatesOperationTraceContext() { + TaskEntityExecutor executor = createExecutor("Signaler", SignalingEntity::new); + + TraceContext opTraceContext = TraceContext.newBuilder() + .setTraceParent("00-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331-01") + .setTraceState(StringValue.of("rojo=00f067aa0ba902b7")) + .build(); + OperationRequest op = OperationRequest.newBuilder() + .setOperation("signalOther") + .setRequestId("req-signalOther") + .setTraceContext(opTraceContext) + .build(); + + EntityBatchResult result = executor.execute(buildBatchRequest("Signaler", "s1", null, op)); + + assertEquals(1, result.getActionsCount()); + assertTrue(result.getActions(0).hasSendSignal()); + SendSignalAction signalAction = result.getActions(0).getSendSignal(); + assertTrue(signalAction.hasParentTraceContext(), + "SendSignalAction should carry the operation's trace context"); + assertEquals("00-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331-01", + signalAction.getParentTraceContext().getTraceParent()); + assertEquals("rojo=00f067aa0ba902b7", signalAction.getParentTraceContext().getTraceState().getValue()); + } + + @Test + void execute_entityStartsOrchestration_propagatesOperationTraceContext() { + TaskEntityExecutor executor = createExecutor("OrchStarter", OrchestrationStartingEntity::new); + + TraceContext opTraceContext = TraceContext.newBuilder() + .setTraceParent("00-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331-01") + .build(); + OperationRequest op = OperationRequest.newBuilder() + .setOperation("startOrch") + .setRequestId("req-startOrch") + .setTraceContext(opTraceContext) + .build(); + + EntityBatchResult result = executor.execute(buildBatchRequest("OrchStarter", "o1", null, op)); + + assertEquals(1, result.getActionsCount()); + assertTrue(result.getActions(0).hasStartNewOrchestration()); + StartNewOrchestrationAction orchAction = result.getActions(0).getStartNewOrchestration(); + assertTrue(orchAction.hasParentTraceContext(), + "StartNewOrchestrationAction should carry the operation's trace context"); + assertEquals("00-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331-01", + orchAction.getParentTraceContext().getTraceParent()); + } + + @Test + void execute_entitySignalsOther_noTraceContext_noParentTraceContext() { + TaskEntityExecutor executor = createExecutor("Signaler", SignalingEntity::new); + + EntityBatchResult result = executor.execute( + buildBatchRequest("Signaler", "s1", null, buildOperationRequest("signalOther"))); + + assertEquals(1, result.getActionsCount()); + assertTrue(result.getActions(0).hasSendSignal()); + assertFalse(result.getActions(0).getSendSignal().hasParentTraceContext(), + "SendSignalAction should not carry a trace context when the operation had none"); + } + @Test void execute_failedOperationRollsBackActions() { // Create an entity that signals another entity then fails diff --git a/client/src/test/java/com/microsoft/durabletask/TracingHelperTest.java b/client/src/test/java/com/microsoft/durabletask/TracingHelperTest.java index c7e0a035..cfa5a90b 100644 --- a/client/src/test/java/com/microsoft/durabletask/TracingHelperTest.java +++ b/client/src/test/java/com/microsoft/durabletask/TracingHelperTest.java @@ -358,4 +358,179 @@ void emitEventSpan_fromClient_createsProducerSpan() { assertEquals("event", sd.getAttributes().get(io.opentelemetry.api.common.AttributeKey.stringKey("durabletask.type"))); assertEquals("target-orch-1", sd.getAttributes().get(io.opentelemetry.api.common.AttributeKey.stringKey("durabletask.event.target_instance_id"))); } + + // region Entity spans + + private static final String TRACE_ID = "0af7651916cd43dd8448eb211c80319c"; + private static final String PARENT_SPAN_ID = "b7ad6b7169203331"; + + private static TraceContext parentCtx() { + return TraceContext.newBuilder() + .setTraceParent("00-" + TRACE_ID + "-" + PARENT_SPAN_ID + "-01") + .build(); + } + + private static String attr(SpanData sd, String key) { + return sd.getAttributes().get(io.opentelemetry.api.common.AttributeKey.stringKey(key)); + } + + @Test + void createEntitySpanName_usesEntityAndOperation() { + assertEquals("entity:Counter:add", TracingHelper.createEntitySpanName("Counter", "add")); + } + + @Test + void createEntityStartOrchestrationSpanName_isInverted() { + assertEquals("Counter:create_orchestration", + TracingHelper.createEntityStartOrchestrationSpanName("Counter")); + } + + @Test + void startEntityProcessingSpan_call_createsServerSpanUnderParent() { + Span span = TracingHelper.startEntityProcessingSpan("Counter", "add", false, "@counter@c1", parentCtx()); + assertNotNull(span); + TracingHelper.endEntityProcessingSpan(span, null); + + SpanData sd = spanExporter.getFinishedSpanItems().get(0); + assertEquals("entity:Counter:add", sd.getName()); + assertEquals(SpanKind.SERVER, sd.getKind()); + assertEquals(TRACE_ID, sd.getTraceId()); + assertEquals(PARENT_SPAN_ID, sd.getParentSpanId()); + assertEquals("entity", attr(sd, "durabletask.type")); + assertEquals("call_entity", attr(sd, "durabletask.task.operation")); + assertEquals("@counter@c1", attr(sd, "durabletask.task.instance_id")); + assertEquals(io.opentelemetry.api.trace.StatusCode.OK, sd.getStatus().getStatusCode()); + } + + @Test + void startEntityProcessingSpan_signal_createsConsumerSpan() { + Span span = TracingHelper.startEntityProcessingSpan("Counter", "add", true, "@counter@c1", parentCtx()); + assertNotNull(span); + TracingHelper.endEntityProcessingSpan(span, null); + + SpanData sd = spanExporter.getFinishedSpanItems().get(0); + assertEquals(SpanKind.CONSUMER, sd.getKind()); + assertEquals("signal_entity", attr(sd, "durabletask.task.operation")); + } + + @Test + void startEntityProcessingSpan_missingParent_returnsNullAndEmitsNothing() { + assertNull(TracingHelper.startEntityProcessingSpan("Counter", "add", false, "@counter@c1", null)); + assertTrue(spanExporter.getFinishedSpanItems().isEmpty()); + } + + @Test + void endEntityProcessingSpan_failure_setsErrorAndMessage() { + Span span = TracingHelper.startEntityProcessingSpan("Counter", "add", false, "@counter@c1", parentCtx()); + TracingHelper.endEntityProcessingSpan(span, "boom"); + + SpanData sd = spanExporter.getFinishedSpanItems().get(0); + assertEquals(io.opentelemetry.api.trace.StatusCode.ERROR, sd.getStatus().getStatusCode()); + assertEquals("boom", attr(sd, "durabletask.entity.error_message")); + } + + @Test + void endEntityProcessingSpan_nullSpan_doesNotThrow() { + assertDoesNotThrow(() -> TracingHelper.endEntityProcessingSpan(null, null)); + } + + @Test + void entityProcessingSpan_omitsTaskNameVersionAndTaskId() { + Span span = TracingHelper.startEntityProcessingSpan("Counter", "add", false, "@counter@c1", parentCtx()); + TracingHelper.endEntityProcessingSpan(span, null); + + SpanData sd = spanExporter.getFinishedSpanItems().get(0); + assertNull(attr(sd, "durabletask.task.name")); + assertNull(attr(sd, "durabletask.task.version")); + assertNull(attr(sd, "durabletask.task.task_id")); + } + + @Test + void emitEntityCallClientSpan_createsClientSpanWithSyntheticIdAndTimestamps() { + java.time.Instant start = java.time.Instant.parse("2026-01-01T00:00:00Z"); + java.time.Instant end = java.time.Instant.parse("2026-01-01T00:00:05Z"); + String syntheticId = "abcdef1234567890"; + + TracingHelper.emitEntityCallClientSpan( + "Counter", "add", "@counter@c1", parentCtx(), start, end, syntheticId, null, null); + + SpanData sd = spanExporter.getFinishedSpanItems().get(0); + assertEquals("entity:Counter:add", sd.getName()); + assertEquals(SpanKind.CLIENT, sd.getKind()); + assertEquals(syntheticId, sd.getSpanContext().getSpanId()); + assertEquals(PARENT_SPAN_ID, sd.getParentSpanId()); + assertEquals("call_entity", attr(sd, "durabletask.task.operation")); + assertEquals("@counter@c1", attr(sd, "durabletask.event.target_instance_id")); + assertEquals(io.opentelemetry.api.trace.StatusCode.UNSET, sd.getStatus().getStatusCode()); + assertEquals(start.getEpochSecond() * 1_000_000_000L + start.getNano(), sd.getStartEpochNanos()); + assertEquals(end.getEpochSecond() * 1_000_000_000L + end.getNano(), sd.getEndEpochNanos()); + } + + @Test + void emitEntityCallClientSpan_withErrorDescription_setsError() { + TracingHelper.emitEntityCallClientSpan( + "Counter", "add", "@counter@c1", parentCtx(), null, null, null, null, "call timed out"); + + SpanData sd = spanExporter.getFinishedSpanItems().get(0); + assertEquals(io.opentelemetry.api.trace.StatusCode.ERROR, sd.getStatus().getStatusCode()); + } + + @Test + void emitEntityCallClientSpan_missingParent_emitsNothing() { + TracingHelper.emitEntityCallClientSpan( + "Counter", "add", "@counter@c1", null, null, null, null, null, null); + assertTrue(spanExporter.getFinishedSpanItems().isEmpty()); + } + + @Test + void startEntitySignalProducerSpan_setsTargetAndSource() { + Span span = TracingHelper.startEntitySignalProducerSpan( + "Audit", "record", "@audit@a1", "@counter@c1", parentCtx(), null, null); + assertNotNull(span); + span.end(); + + SpanData sd = spanExporter.getFinishedSpanItems().get(0); + assertEquals("entity:Audit:record", sd.getName()); + assertEquals(SpanKind.PRODUCER, sd.getKind()); + assertEquals("signal_entity", attr(sd, "durabletask.task.operation")); + assertEquals("@audit@a1", attr(sd, "durabletask.event.target_instance_id")); + assertEquals("@counter@c1", attr(sd, "durabletask.task.instance_id")); + } + + @Test + void startEntitySignalProducerSpan_scheduledTime_setsAttribute() { + Span span = TracingHelper.startEntitySignalProducerSpan( + "Audit", "record", "@audit@a1", null, parentCtx(), null, "2026-01-01T00:00:00Z"); + span.end(); + + SpanData sd = spanExporter.getFinishedSpanItems().get(0); + assertEquals("2026-01-01T00:00:00Z", attr(sd, "durabletask.task.scheduled_time")); + assertNull(attr(sd, "durabletask.task.instance_id")); + } + + @Test + void startEntityStartOrchestrationSpan_setsInvertedNameAndAttributes() { + Span span = TracingHelper.startEntityStartOrchestrationSpan( + "Counter", "@counter@c1", "orch-2", parentCtx(), null, null); + assertNotNull(span); + span.end(); + + SpanData sd = spanExporter.getFinishedSpanItems().get(0); + assertEquals("Counter:create_orchestration", sd.getName()); + assertEquals(SpanKind.PRODUCER, sd.getKind()); + assertEquals("entity", attr(sd, "durabletask.type")); + assertEquals("orch-2", attr(sd, "durabletask.event.target_instance_id")); + assertEquals("@counter@c1", attr(sd, "durabletask.task.instance_id")); + } + + @Test + void entityProducerSpans_missingParent_returnNull() { + assertNull(TracingHelper.startEntitySignalProducerSpan( + "Audit", "record", "@audit@a1", null, null, null, null)); + assertNull(TracingHelper.startEntityStartOrchestrationSpan( + "Counter", "@counter@c1", "orch-2", null, null, null)); + assertTrue(spanExporter.getFinishedSpanItems().isEmpty()); + } + + // endregion } From e1d64302e6302d03b4bd05ba13e4ae084223f694 Mon Sep 17 00:00:00 2001 From: Varshi Bachu Date: Wed, 12 Aug 2026 16:06:12 -0700 Subject: [PATCH 2/6] address copilot comment --- .../main/java/com/microsoft/durabletask/TracingHelper.java | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/client/src/main/java/com/microsoft/durabletask/TracingHelper.java b/client/src/main/java/com/microsoft/durabletask/TracingHelper.java index 9e610308..84cfd097 100644 --- a/client/src/main/java/com/microsoft/durabletask/TracingHelper.java +++ b/client/src/main/java/com/microsoft/durabletask/TracingHelper.java @@ -546,9 +546,10 @@ static void endEntityProcessingSpan(@Nullable Span span, @Nullable String errorM /** * Emits the retroactive {@link SpanKind#CLIENT} span for a call to an entity, covering the - * request-to-response interval and sharing {@code syntheticSpanId} with the SERVER processing - * span. The status is left unset for normal completion or entity-processing failure (matching - * .NET); {@code errorDescription} is set only for timeout/cancellation closure boundaries. + * request-to-response interval. The CLIENT span's ID is set to {@code syntheticSpanId}, which + * the SERVER processing span uses as its parent span ID. The status is left unset for normal + * completion or entity-processing failure (matching .NET); {@code errorDescription} is set only + * for timeout/cancellation closure boundaries. * Does nothing when the parent context is absent or invalid. */ static void emitEntityCallClientSpan( From 5d9c0934a07501397450cd6720a1799b463dade6 Mon Sep 17 00:00:00 2001 From: Varshi Bachu Date: Fri, 14 Aug 2026 09:34:09 -0700 Subject: [PATCH 3/6] wire entity producer spans and trace-context propagation --- .../durabletask/DurableTaskGrpcWorker.java | 8 +- .../DurableTaskGrpcWorkerBuilder.java | 17 ++ .../microsoft/durabletask/EntityRunner.java | 5 +- .../durabletask/TaskEntityExecutor.java | 71 ++++++- .../TaskOrchestrationExecutor.java | 59 ++++++ .../durabletask/TaskEntityExecutorTest.java | 2 +- .../TaskEntityExecutorTracingTest.java | 199 ++++++++++++++++++ .../TaskOrchestrationEntityEventTest.java | 33 +++ .../TaskOrchestrationEntityTracingTest.java | 165 +++++++++++++++ 9 files changed, 548 insertions(+), 11 deletions(-) create mode 100644 client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTracingTest.java create mode 100644 client/src/test/java/com/microsoft/durabletask/TaskOrchestrationEntityTracingTest.java diff --git a/client/src/main/java/com/microsoft/durabletask/DurableTaskGrpcWorker.java b/client/src/main/java/com/microsoft/durabletask/DurableTaskGrpcWorker.java index 086b1cfe..d640f3b2 100644 --- a/client/src/main/java/com/microsoft/durabletask/DurableTaskGrpcWorker.java +++ b/client/src/main/java/com/microsoft/durabletask/DurableTaskGrpcWorker.java @@ -59,6 +59,7 @@ public final class DurableTaskGrpcWorker implements AutoCloseable { private final boolean supportsLargePayloads; private final int maxChunkSizeBytes; private final int largePayloadThresholdBytes; + private final boolean emitTraceSpans; DurableTaskGrpcWorker(DurableTaskGrpcWorkerBuilder builder, WorkItemFilter workItemFilter) { this.orchestrationFactories.putAll(builder.orchestrationFactories); @@ -95,6 +96,7 @@ public final class DurableTaskGrpcWorker implements AutoCloseable { this.supportsLargePayloads = builder.supportsLargePayloads; this.maxChunkSizeBytes = builder.maxChunkSizeBytes; this.largePayloadThresholdBytes = builder.largePayloadThresholdBytes; + this.emitTraceSpans = builder.emitTraceSpans; this.dataConverter = builder.dataConverter != null ? builder.dataConverter : new JacksonDataConverter(); this.maximumTimerInterval = builder.maximumTimerInterval != null ? builder.maximumTimerInterval : DEFAULT_MAXIMUM_TIMER_INTERVAL; this.versioningOptions = builder.versioningOptions; @@ -175,7 +177,8 @@ public void startAndBlock() { logger, this.versioningOptions, true, - this.exceptionPropertiesProvider); + this.exceptionPropertiesProvider, + this.emitTraceSpans); TaskActivityExecutor taskActivityExecutor = new TaskActivityExecutor( this.activityFactories, this.dataConverter, @@ -183,7 +186,8 @@ public void startAndBlock() { TaskEntityExecutor taskEntityExecutor = new TaskEntityExecutor( this.entityFactories, this.dataConverter, - logger); + logger, + this.emitTraceSpans); // TODO: How do we interrupt manually? while (true) { diff --git a/client/src/main/java/com/microsoft/durabletask/DurableTaskGrpcWorkerBuilder.java b/client/src/main/java/com/microsoft/durabletask/DurableTaskGrpcWorkerBuilder.java index 3df0a289..9a6ea803 100644 --- a/client/src/main/java/com/microsoft/durabletask/DurableTaskGrpcWorkerBuilder.java +++ b/client/src/main/java/com/microsoft/durabletask/DurableTaskGrpcWorkerBuilder.java @@ -28,6 +28,7 @@ public final class DurableTaskGrpcWorkerBuilder { ExceptionPropertiesProvider exceptionPropertiesProvider; int maxConcurrentEntityWorkItems = 1; int maxWorkItemThreads; + boolean emitTraceSpans = true; private WorkItemFilter workItemFilter; private boolean autoGenerateWorkItemFilters; final List interceptors = new ArrayList<>(); @@ -452,6 +453,22 @@ public DurableTaskGrpcWorkerBuilder setMaxChunkSizeBytes(int maxChunkSizeBytes) return this; } + /** + * Sets whether this worker emits its own OpenTelemetry spans for orchestrations, activities, and + * entities. Defaults to {@code true}. + *

+ * Set to {@code false} when running under a host that already emits Durable Task spans (for + * example, the Azure Functions Durable extension, which emits {@code DurableTask.Core} spans), to + * avoid a duplicate worker-side span layer. Trace-context propagation is unaffected either way. + * + * @param emitTraceSpans whether the worker emits its own spans + * @return this builder object + */ + public DurableTaskGrpcWorkerBuilder setEmitTraceSpans(boolean emitTraceSpans) { + this.emitTraceSpans = emitTraceSpans; + return this; + } + /** * Initializes a new {@link DurableTaskGrpcWorker} object with the settings specified in the current builder object. * @return a new {@link DurableTaskGrpcWorker} object diff --git a/client/src/main/java/com/microsoft/durabletask/EntityRunner.java b/client/src/main/java/com/microsoft/durabletask/EntityRunner.java index 680cced2..1630e01a 100644 --- a/client/src/main/java/com/microsoft/durabletask/EntityRunner.java +++ b/client/src/main/java/com/microsoft/durabletask/EntityRunner.java @@ -90,7 +90,10 @@ public static byte[] loadAndRun(byte[] entityRequestBytes, TaskEntityFactory ent TaskEntityExecutor executor = new TaskEntityExecutor( factories, new JacksonDataConverter(), - logger); + logger, + // EntityRunner is the Azure Functions entry point; the Durable extension host already + // emits DurableTask.Core entity spans, so the worker suppresses its own to avoid duplicates. + false); EntityBatchResult result = executor.execute(request); return result.toByteArray(); diff --git a/client/src/main/java/com/microsoft/durabletask/TaskEntityExecutor.java b/client/src/main/java/com/microsoft/durabletask/TaskEntityExecutor.java index c2380f8b..ee0034a9 100644 --- a/client/src/main/java/com/microsoft/durabletask/TaskEntityExecutor.java +++ b/client/src/main/java/com/microsoft/durabletask/TaskEntityExecutor.java @@ -6,6 +6,8 @@ import com.google.protobuf.Timestamp; import com.microsoft.durabletask.implementation.protobuf.OrchestratorService.*; +import io.opentelemetry.api.trace.Span; + import javax.annotation.Nonnull; import javax.annotation.Nullable; import java.time.Instant; @@ -24,14 +26,17 @@ final class TaskEntityExecutor { private final HashMap entityFactories; private final DataConverter dataConverter; private final Logger logger; + private final boolean emitTraceSpans; TaskEntityExecutor( HashMap entityFactories, DataConverter dataConverter, - Logger logger) { + Logger logger, + boolean emitTraceSpans) { this.entityFactories = entityFactories; this.dataConverter = dataConverter; this.logger = logger; + this.emitTraceSpans = emitTraceSpans; } /** @@ -80,7 +85,7 @@ EntityBatchResult execute(@Nonnull EntityBatchRequest request) { TaskEntityState entityState = new TaskEntityState(this.dataConverter, initialState); // Create the concrete context that collects actions - TaskEntityContextImpl context = new TaskEntityContextImpl(entityId, this.dataConverter); + TaskEntityContextImpl context = new TaskEntityContextImpl(entityId, this.dataConverter, this.emitTraceSpans); // Process each operation List results = new ArrayList<>(); @@ -121,16 +126,31 @@ EntityBatchResult execute(@Nonnull EntityBatchRequest request) { logger.log(Level.FINE, "Executing operation '{0}' (requestId={1}) on entity '{2}'.", new Object[]{operationName, requestId, instanceId}); - // Parent context for any signals/orchestrations this operation produces, so the host can link them. - context.setCurrentOperationTraceContext( - opRequest.hasTraceContext() ? opRequest.getTraceContext() : null); - // Snapshot state and actions before each operation (for rollback on failure) entityState.commit(); context.commit(); Instant startTime = Instant.now(); + // Entity processing (SERVER) span. All operations use call_entity semantics: the + // OperationRequest carries no call-vs-signal bit, so we cannot emit CONSUMER for signals + // the way the .NET engine does. Emission is suppressed when emitTraceSpans is false (e.g. + // under Azure Functions, where the host already emits DurableTask.Core entity spans). + Span processingSpan = this.emitTraceSpans + ? TracingHelper.startEntityProcessingSpan( + entityName, + operationName, + false, + instanceId, + opRequest.hasTraceContext() ? opRequest.getTraceContext() : null) + : null; + + // Signals/orchestrations this operation produces nest under the processing span (or the + // raw incoming context when spans are suppressed), so the host can link them downstream. + context.setCurrentOperationTraceContext(processingSpan != null + ? TracingHelper.getCurrentTraceContext(processingSpan) + : (opRequest.hasTraceContext() ? opRequest.getTraceContext() : null)); + try { // Build the operation TaskEntityOperation operation = new TaskEntityOperation( @@ -161,6 +181,8 @@ EntityBatchResult execute(@Nonnull EntityBatchRequest request) { entityState.commit(); context.commit(); + TracingHelper.endEntityProcessingSpan(processingSpan, null); + logger.log(Level.FINE, "Operation '{0}' on entity '{1}' completed successfully.", new Object[]{operationName, instanceId}); @@ -192,6 +214,8 @@ EntityBatchResult execute(@Nonnull EntityBatchRequest request) { // Rollback state and actions on failure entityState.rollback(); context.rollback(); + + TracingHelper.endEntityProcessingSpan(processingSpan, e.getMessage()); } } @@ -223,14 +247,16 @@ private static Timestamp toTimestamp(Instant instant) { private static class TaskEntityContextImpl extends TaskEntityContext { private final EntityInstanceId entityId; private final DataConverter dataConverter; + private final boolean emitTraceSpans; private final List pendingActions = new ArrayList<>(); private int committedActionCount = 0; @Nullable private TraceContext currentOperationTraceContext; - TaskEntityContextImpl(EntityInstanceId entityId, DataConverter dataConverter) { + TaskEntityContextImpl(EntityInstanceId entityId, DataConverter dataConverter, boolean emitTraceSpans) { this.entityId = entityId; this.dataConverter = dataConverter; + this.emitTraceSpans = emitTraceSpans; } void setCurrentOperationTraceContext(@Nullable TraceContext traceContext) { @@ -275,6 +301,22 @@ public void signalEntity( signalBuilder.setParentTraceContext(this.currentOperationTraceContext); } + if (this.emitTraceSpans && this.currentOperationTraceContext != null) { + String signalScheduledTime = (options != null && options.getScheduledTime() != null) + ? options.getScheduledTime().toString() : null; + Span producerSpan = TracingHelper.startEntitySignalProducerSpan( + targetEntityId.getName(), + operationName, + targetEntityId.toString(), + this.entityId.toString(), + this.currentOperationTraceContext, + null, + signalScheduledTime); + if (producerSpan != null) { + producerSpan.end(); + } + } + this.pendingActions.add(new PendingAction(PendingAction.Type.SEND_SIGNAL, signalBuilder.build(), null)); } @@ -318,6 +360,21 @@ public String startNewOrchestration( orchBuilder.setParentTraceContext(this.currentOperationTraceContext); } + if (this.emitTraceSpans && this.currentOperationTraceContext != null) { + String orchScheduledTime = (options != null && options.getStartTime() != null) + ? options.getStartTime().toString() : null; + Span producerSpan = TracingHelper.startEntityStartOrchestrationSpan( + this.entityId.getName(), + this.entityId.toString(), + instanceId, + this.currentOperationTraceContext, + null, + orchScheduledTime); + if (producerSpan != null) { + producerSpan.end(); + } + } + this.pendingActions.add(new PendingAction( PendingAction.Type.START_NEW_ORCHESTRATION, null, orchBuilder.build())); diff --git a/client/src/main/java/com/microsoft/durabletask/TaskOrchestrationExecutor.java b/client/src/main/java/com/microsoft/durabletask/TaskOrchestrationExecutor.java index 5bdbad2e..82c7b621 100644 --- a/client/src/main/java/com/microsoft/durabletask/TaskOrchestrationExecutor.java +++ b/client/src/main/java/com/microsoft/durabletask/TaskOrchestrationExecutor.java @@ -14,6 +14,8 @@ import com.microsoft.durabletask.implementation.protobuf.OrchestratorService.ScheduleTaskAction.Builder; import com.microsoft.durabletask.util.UUIDGenerator; +import io.opentelemetry.api.trace.Span; + import javax.annotation.Nullable; import java.time.Duration; import java.time.Instant; @@ -39,6 +41,7 @@ final class TaskOrchestrationExecutor { private final DurableTaskGrpcWorkerVersioningOptions versioningOptions; private final ExceptionPropertiesProvider exceptionPropertiesProvider; private final boolean useNativeEntityActions; + private final boolean emitTraceSpans; public TaskOrchestrationExecutor( HashMap orchestrationFactories, @@ -67,6 +70,19 @@ public TaskOrchestrationExecutor( DurableTaskGrpcWorkerVersioningOptions versioningOptions, boolean useNativeEntityActions, ExceptionPropertiesProvider exceptionPropertiesProvider) { + this(orchestrationFactories, dataConverter, maximumTimerInterval, logger, versioningOptions, + useNativeEntityActions, exceptionPropertiesProvider, false); + } + + public TaskOrchestrationExecutor( + HashMap orchestrationFactories, + DataConverter dataConverter, + Duration maximumTimerInterval, + Logger logger, + DurableTaskGrpcWorkerVersioningOptions versioningOptions, + boolean useNativeEntityActions, + ExceptionPropertiesProvider exceptionPropertiesProvider, + boolean emitTraceSpans) { this.orchestrationFactories = orchestrationFactories; this.dataConverter = dataConverter; this.maximumTimerInterval = maximumTimerInterval; @@ -74,6 +90,7 @@ public TaskOrchestrationExecutor( this.versioningOptions = versioningOptions; this.useNativeEntityActions = useNativeEntityActions; this.exceptionPropertiesProvider = exceptionPropertiesProvider; + this.emitTraceSpans = emitTraceSpans; } public TaskOrchestratorResult execute( @@ -446,6 +463,24 @@ public UUID newUUID() { // region Entity integration methods (Phase 4) + // Writes the orchestration trace context into the legacy DTFx RequestMessage JSON as + // parentTraceContext (DistributedTraceContext, PascalCase members) so the Azure Functions + // host can link its entity spans. The orchestration context is deterministic (from history). + private void addLegacyEntityParentTraceContext(ObjectNode requestMessage) { + TraceContext propagatedCtx = this.orchestrationSpanContext != null + ? this.orchestrationSpanContext : this.parentTraceContext; + if (propagatedCtx == null || propagatedCtx.getTraceParent() == null + || propagatedCtx.getTraceParent().isEmpty()) { + return; + } + ObjectNode ptc = requestMessage.putObject("parentTraceContext"); + ptc.put("TraceParent", propagatedCtx.getTraceParent()); + if (propagatedCtx.hasTraceState() && propagatedCtx.getTraceState().getValue() != null + && !propagatedCtx.getTraceState().getValue().isEmpty()) { + ptc.put("TraceState", propagatedCtx.getTraceState().getValue()); + } + } + @Override public void signalEntity(EntityInstanceId entityId, String operationName, Object input, SignalEntityOptions options) { Helpers.throwIfOrchestratorComplete(this.isComplete); @@ -490,6 +525,7 @@ public void signalEntity(EntityInstanceId entityId, String operationName, Object requestMessage.put("due", scheduledTimeStr); eventName = "op@" + scheduledTimeStr; } + this.addLegacyEntityParentTraceContext(requestMessage); this.pendingActions.put(id, OrchestratorAction.newBuilder() .setId(id) .setSendEvent(SendEventAction.newBuilder() @@ -500,6 +536,28 @@ public void signalEntity(EntityInstanceId entityId, String operationName, Object .build()); } + // PRODUCER span for the signal so standalone/DTS workers record the client side. + // Suppressed under Azure Functions, where the host emits it. + if (TaskOrchestrationExecutor.this.emitTraceSpans && !this.isReplaying) { + TraceContext signalParentCtx = this.orchestrationSpanContext != null + ? this.orchestrationSpanContext : this.parentTraceContext; + if (signalParentCtx != null) { + String signalScheduledTime = (options != null && options.getScheduledTime() != null) + ? options.getScheduledTime().toString() : null; + Span signalSpan = TracingHelper.startEntitySignalProducerSpan( + entityId.getName(), + operationName, + entityId.toString(), + this.instanceId, + signalParentCtx, + null, + signalScheduledTime); + if (signalSpan != null) { + signalSpan.end(); + } + } + } + if (!this.isReplaying) { this.logger.fine(() -> String.format( "%s: signaling entity '%s' operation '%s' (#%d)", @@ -567,6 +625,7 @@ public Task callEntity(EntityInstanceId entityId, String operationName, O if (this.executionId != null) { requestMessage.put("parentExecution", this.executionId); } + this.addLegacyEntityParentTraceContext(requestMessage); this.pendingActions.put(id, OrchestratorAction.newBuilder() .setId(id) .setSendEvent(SendEventAction.newBuilder() diff --git a/client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTest.java b/client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTest.java index 30de8732..7d369be0 100644 --- a/client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTest.java +++ b/client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTest.java @@ -116,7 +116,7 @@ public Object run(TaskEntityOperation operation) throws Exception { private TaskEntityExecutor createExecutor(String entityName, TaskEntityFactory factory) { HashMap factories = new HashMap<>(); factories.put(entityName.toLowerCase(java.util.Locale.ROOT), factory); - return new TaskEntityExecutor(factories, dataConverter, logger); + return new TaskEntityExecutor(factories, dataConverter, logger, true); } private OperationRequest buildOperationRequest(String operationName, Object input, String requestId) { diff --git a/client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTracingTest.java b/client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTracingTest.java new file mode 100644 index 00000000..7fabc7ca --- /dev/null +++ b/client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTracingTest.java @@ -0,0 +1,199 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. +// Licensed under the MIT License. +package com.microsoft.durabletask; + +import com.google.protobuf.StringValue; +import com.microsoft.durabletask.implementation.protobuf.OrchestratorService.EntityBatchRequest; +import com.microsoft.durabletask.implementation.protobuf.OrchestratorService.EntityBatchResult; +import com.microsoft.durabletask.implementation.protobuf.OrchestratorService.OperationRequest; +import com.microsoft.durabletask.implementation.protobuf.OrchestratorService.TraceContext; + +import io.opentelemetry.api.common.AttributeKey; +import io.opentelemetry.api.trace.SpanKind; +import io.opentelemetry.api.trace.StatusCode; +import io.opentelemetry.sdk.OpenTelemetrySdk; +import io.opentelemetry.sdk.testing.exporter.InMemorySpanExporter; +import io.opentelemetry.sdk.trace.SdkTracerProvider; +import io.opentelemetry.sdk.trace.data.SpanData; +import io.opentelemetry.sdk.trace.export.SimpleSpanProcessor; + +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.util.HashMap; +import java.util.List; +import java.util.logging.Logger; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Verifies that {@link TaskEntityExecutor} emits the entity processing (SERVER) span when + * {@code emitTraceSpans} is enabled and suppresses it otherwise. + */ +public class TaskEntityExecutorTracingTest { + + private static final Logger logger = Logger.getLogger(TaskEntityExecutorTracingTest.class.getName()); + private static final DataConverter dataConverter = new JacksonDataConverter(); + private static final String TRACE_ID = "0af7651916cd43dd8448eb211c80319c"; + private static final String PARENT_SPAN_ID = "b7ad6b7169203331"; + + private InMemorySpanExporter spanExporter; + private OpenTelemetrySdk openTelemetry; + + @BeforeEach + void setUp() { + io.opentelemetry.api.GlobalOpenTelemetry.resetForTest(); + spanExporter = InMemorySpanExporter.create(); + SdkTracerProvider tracerProvider = SdkTracerProvider.builder() + .addSpanProcessor(SimpleSpanProcessor.create(spanExporter)) + .build(); + openTelemetry = OpenTelemetrySdk.builder() + .setTracerProvider(tracerProvider) + .buildAndRegisterGlobal(); + } + + @AfterEach + void tearDown() { + openTelemetry.close(); + io.opentelemetry.api.GlobalOpenTelemetry.resetForTest(); + } + + /** Minimal counter entity used for the batch under test. */ + static class CounterEntity extends AbstractTaskEntity { + public void add(int amount) { + this.state += amount; + } + + public void signalOther(int amount) { + this.context.signalEntity(new EntityInstanceId("counter", "c2"), "add", amount); + } + + public void startOrch(int amount) { + this.context.startNewOrchestration("DownstreamOrch", null); + } + + @Override + protected Integer initializeState(TaskEntityOperation operation) { + return 0; + } + + @Override + protected Class getStateType() { + return Integer.class; + } + } + + private TaskEntityExecutor createExecutor(boolean emitTraceSpans) { + HashMap factories = new HashMap<>(); + factories.put("counter", CounterEntity::new); + return new TaskEntityExecutor(factories, dataConverter, logger, emitTraceSpans); + } + + private EntityBatchRequest requestWith(@javax.annotation.Nullable TraceContext traceContext) { + return requestWithOp("add", 5, traceContext); + } + + private EntityBatchRequest requestWithOp( + String operation, int input, @javax.annotation.Nullable TraceContext traceContext) { + OperationRequest.Builder op = OperationRequest.newBuilder() + .setOperation(operation) + .setRequestId("req-1") + .setInput(StringValue.of(dataConverter.serialize(input))); + if (traceContext != null) { + op.setTraceContext(traceContext); + } + return EntityBatchRequest.newBuilder() + .setInstanceId("@counter@c1") + .setEntityState(StringValue.of(dataConverter.serialize(10))) + .addOperations(op.build()) + .build(); + } + + private static TraceContext parentTraceContext() { + return TraceContext.newBuilder() + .setTraceParent("00-" + TRACE_ID + "-" + PARENT_SPAN_ID + "-01") + .build(); + } + + @Test + void execute_emitsEntityProcessingServerSpanUnderParent() { + TaskEntityExecutor executor = createExecutor(true); + + EntityBatchResult result = executor.execute(requestWith(parentTraceContext())); + assertTrue(result.getResults(0).hasSuccess()); + + List spans = spanExporter.getFinishedSpanItems(); + assertEquals(1, spans.size()); + SpanData span = spans.get(0); + assertEquals("entity:counter:add", span.getName()); + assertEquals(SpanKind.SERVER, span.getKind()); + assertEquals(TRACE_ID, span.getTraceId()); + assertEquals(PARENT_SPAN_ID, span.getParentSpanId()); + assertEquals("entity", span.getAttributes().get(AttributeKey.stringKey("durabletask.type"))); + assertEquals("call_entity", + span.getAttributes().get(AttributeKey.stringKey("durabletask.task.operation"))); + assertEquals("@counter@c1", + span.getAttributes().get(AttributeKey.stringKey("durabletask.task.instance_id"))); + assertEquals(StatusCode.OK, span.getStatus().getStatusCode()); + } + + @Test + void execute_emitTraceSpansDisabled_suppressesSpan() { + TaskEntityExecutor executor = createExecutor(false); + + EntityBatchResult result = executor.execute(requestWith(parentTraceContext())); + assertTrue(result.getResults(0).hasSuccess()); + + assertTrue(spanExporter.getFinishedSpanItems().isEmpty()); + } + + @Test + void execute_noParentTraceContext_emitsNoSpan() { + TaskEntityExecutor executor = createExecutor(true); + + EntityBatchResult result = executor.execute(requestWith(null)); + assertTrue(result.getResults(0).hasSuccess()); + + assertTrue(spanExporter.getFinishedSpanItems().isEmpty()); + } + + @Test + void execute_entitySignalsEntity_emitsProducerSpanNestedUnderProcessingSpan() { + TaskEntityExecutor executor = createExecutor(true); + + EntityBatchResult result = executor.execute(requestWithOp("signalOther", 3, parentTraceContext())); + assertTrue(result.getResults(0).hasSuccess()); + + List spans = spanExporter.getFinishedSpanItems(); + SpanData server = spans.stream().filter(s -> s.getKind() == SpanKind.SERVER).findFirst().orElse(null); + SpanData producer = spans.stream().filter(s -> s.getKind() == SpanKind.PRODUCER).findFirst().orElse(null); + assertNotNull(server, "expected SERVER processing span"); + assertNotNull(producer, "expected PRODUCER signal span"); + assertEquals("entity:counter:add", producer.getName()); + assertEquals("signal_entity", + producer.getAttributes().get(AttributeKey.stringKey("durabletask.task.operation"))); + assertEquals(TRACE_ID, producer.getTraceId()); + assertEquals(server.getSpanId(), producer.getParentSpanId()); + } + + @Test + void execute_entityStartsOrchestration_emitsProducerSpanNestedUnderProcessingSpan() { + TaskEntityExecutor executor = createExecutor(true); + + EntityBatchResult result = executor.execute(requestWithOp("startOrch", 0, parentTraceContext())); + assertTrue(result.getResults(0).hasSuccess()); + + List spans = spanExporter.getFinishedSpanItems(); + SpanData server = spans.stream().filter(s -> s.getKind() == SpanKind.SERVER).findFirst().orElse(null); + SpanData producer = spans.stream().filter(s -> s.getKind() == SpanKind.PRODUCER).findFirst().orElse(null); + assertNotNull(server, "expected SERVER processing span"); + assertNotNull(producer, "expected PRODUCER create_orchestration span"); + assertEquals("counter:create_orchestration", producer.getName()); + assertEquals("entity", producer.getAttributes().get(AttributeKey.stringKey("durabletask.type"))); + assertEquals(TRACE_ID, producer.getTraceId()); + assertEquals(server.getSpanId(), producer.getParentSpanId()); + } +} diff --git a/client/src/test/java/com/microsoft/durabletask/TaskOrchestrationEntityEventTest.java b/client/src/test/java/com/microsoft/durabletask/TaskOrchestrationEntityEventTest.java index b07ac430..a83033eb 100644 --- a/client/src/test/java/com/microsoft/durabletask/TaskOrchestrationEntityEventTest.java +++ b/client/src/test/java/com/microsoft/durabletask/TaskOrchestrationEntityEventTest.java @@ -309,6 +309,39 @@ private boolean hasLockRequestAction(Collection actions) thr // region signalEntity tests + @Test + void signalEntity_legacyPath_propagatesParentTraceContextInJson() throws Exception { + final String orchestratorName = "SignalEntityTraceOrchestration"; + EntityInstanceId entityId = new EntityInstanceId("Counter", "c1"); + + TaskOrchestrationExecutor executor = createExecutor(orchestratorName, ctx -> { + ctx.signalEntity(entityId, "add", 5); + ctx.complete("done"); + }); + + String traceParent = "00-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331-01"; + TraceContext orchCtx = TraceContext.newBuilder().setTraceParent(traceParent).build(); + + List pastEvents = Arrays.asList( + orchestratorStarted(), + executionStarted(orchestratorName, "null")); + List newEvents = Collections.singletonList(orchestratorCompleted()); + + TaskOrchestratorResult result = executor.execute(pastEvents, newEvents, orchCtx); + + boolean found = false; + for (OrchestratorAction action : result.getActions()) { + if (action.hasSendEvent() + && action.getSendEvent().getInstance().getInstanceId().contains("@counter@c1")) { + JsonNode json = JSON_MAPPER.readTree(action.getSendEvent().getData().getValue()); + assertTrue(json.has("parentTraceContext"), "expected parentTraceContext in signal JSON"); + assertEquals(traceParent, json.get("parentTraceContext").get("TraceParent").asText()); + found = true; + } + } + assertTrue(found, "expected a SendEvent signal action carrying parentTraceContext"); + } + @Test void signalEntity_producesSendEventAction() throws Exception { final String orchestratorName = "SignalEntityOrchestration"; diff --git a/client/src/test/java/com/microsoft/durabletask/TaskOrchestrationEntityTracingTest.java b/client/src/test/java/com/microsoft/durabletask/TaskOrchestrationEntityTracingTest.java new file mode 100644 index 00000000..c381184b --- /dev/null +++ b/client/src/test/java/com/microsoft/durabletask/TaskOrchestrationEntityTracingTest.java @@ -0,0 +1,165 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. +// Licensed under the MIT License. +package com.microsoft.durabletask; + +import com.google.protobuf.StringValue; +import com.google.protobuf.Timestamp; +import com.microsoft.durabletask.implementation.protobuf.OrchestratorService.ExecutionStartedEvent; +import com.microsoft.durabletask.implementation.protobuf.OrchestratorService.HistoryEvent; +import com.microsoft.durabletask.implementation.protobuf.OrchestratorService.OrchestrationInstance; +import com.microsoft.durabletask.implementation.protobuf.OrchestratorService.OrchestratorStartedEvent; +import com.microsoft.durabletask.implementation.protobuf.OrchestratorService.TraceContext; + +import io.opentelemetry.api.common.AttributeKey; +import io.opentelemetry.api.trace.SpanKind; +import io.opentelemetry.sdk.OpenTelemetrySdk; +import io.opentelemetry.sdk.testing.exporter.InMemorySpanExporter; +import io.opentelemetry.sdk.trace.SdkTracerProvider; +import io.opentelemetry.sdk.trace.data.SpanData; +import io.opentelemetry.sdk.trace.export.SimpleSpanProcessor; + +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.time.Duration; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.logging.Logger; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Verifies that {@link TaskOrchestrationExecutor} emits the orchestration-initiated entity signal + * PRODUCER span when {@code emitTraceSpans} is enabled (standalone/DTS worker) and suppresses it + * otherwise (Azure Functions, where the host emits the entity spans). + */ +public class TaskOrchestrationEntityTracingTest { + + private static final Logger logger = Logger.getLogger(TaskOrchestrationEntityTracingTest.class.getName()); + private static final String TRACE_ID = "0af7651916cd43dd8448eb211c80319c"; + private static final String ORCH_SPAN_ID = "b7ad6b7169203331"; + + private InMemorySpanExporter spanExporter; + private OpenTelemetrySdk openTelemetry; + + @BeforeEach + void setUp() { + io.opentelemetry.api.GlobalOpenTelemetry.resetForTest(); + spanExporter = InMemorySpanExporter.create(); + SdkTracerProvider tracerProvider = SdkTracerProvider.builder() + .addSpanProcessor(SimpleSpanProcessor.create(spanExporter)) + .build(); + openTelemetry = OpenTelemetrySdk.builder() + .setTracerProvider(tracerProvider) + .buildAndRegisterGlobal(); + } + + @AfterEach + void tearDown() { + openTelemetry.close(); + io.opentelemetry.api.GlobalOpenTelemetry.resetForTest(); + } + + private TaskOrchestrationExecutor createExecutor( + String orchestratorName, TaskOrchestration orchestration, boolean emitTraceSpans) { + HashMap factories = new HashMap<>(); + factories.put(orchestratorName, new TaskOrchestrationFactory() { + @Override + public String getName() { + return orchestratorName; + } + + @Override + public TaskOrchestration create() { + return orchestration; + } + }); + return new TaskOrchestrationExecutor( + factories, + new JacksonDataConverter(), + Duration.ofDays(1), + logger, + null, + true, + null, + emitTraceSpans); + } + + private HistoryEvent orchestratorStarted() { + return HistoryEvent.newBuilder() + .setEventId(-1) + .setTimestamp(Timestamp.getDefaultInstance()) + .setOrchestratorStarted(OrchestratorStartedEvent.getDefaultInstance()) + .build(); + } + + private HistoryEvent executionStarted(String name) { + return HistoryEvent.newBuilder() + .setEventId(-1) + .setTimestamp(Timestamp.getDefaultInstance()) + .setExecutionStarted(ExecutionStartedEvent.newBuilder() + .setName(name) + .setVersion(StringValue.of("")) + .setInput(StringValue.of("null")) + .setOrchestrationInstance(OrchestrationInstance.newBuilder() + .setInstanceId("test-instance-id") + .build()) + .build()) + .build(); + } + + private static TraceContext orchestrationContext() { + return TraceContext.newBuilder() + .setTraceParent("00-" + TRACE_ID + "-" + ORCH_SPAN_ID + "-01") + .build(); + } + + @Test + void signalEntity_emitsProducerSpanUnderOrchestrationContext() { + String orchestratorName = "SignalOrch"; + EntityInstanceId entityId = new EntityInstanceId("Counter", "c1"); + TaskOrchestrationExecutor executor = createExecutor(orchestratorName, ctx -> { + ctx.signalEntity(entityId, "add", 5); + ctx.complete("done"); + }, true); + + executor.execute( + Collections.emptyList(), + Arrays.asList(orchestratorStarted(), executionStarted(orchestratorName)), + orchestrationContext()); + + List spans = spanExporter.getFinishedSpanItems(); + SpanData producer = spans.stream() + .filter(s -> s.getKind() == SpanKind.PRODUCER).findFirst().orElse(null); + assertNotNull(producer, "expected PRODUCER signal span"); + assertEquals("entity:counter:add", producer.getName()); + assertEquals("signal_entity", + producer.getAttributes().get(AttributeKey.stringKey("durabletask.task.operation"))); + assertEquals(TRACE_ID, producer.getTraceId()); + assertEquals(ORCH_SPAN_ID, producer.getParentSpanId()); + } + + @Test + void signalEntity_emitTraceSpansDisabled_suppressesProducerSpan() { + String orchestratorName = "SignalOrchDisabled"; + EntityInstanceId entityId = new EntityInstanceId("Counter", "c1"); + TaskOrchestrationExecutor executor = createExecutor(orchestratorName, ctx -> { + ctx.signalEntity(entityId, "add", 5); + ctx.complete("done"); + }, false); + + executor.execute( + Collections.emptyList(), + Arrays.asList(orchestratorStarted(), executionStarted(orchestratorName)), + orchestrationContext()); + + List spans = spanExporter.getFinishedSpanItems(); + assertTrue(spans.stream().noneMatch(s -> s.getKind() == SpanKind.PRODUCER), + "expected no PRODUCER span when emitTraceSpans is disabled"); + } +} From ac4f9ac2f82e06768d081979562aa291dc0f6fc3 Mon Sep 17 00:00:00 2001 From: Varshi Bachu Date: Fri, 14 Aug 2026 14:17:20 -0700 Subject: [PATCH 4/6] Set requestTime on entity actions and edit comments --- .../durabletask/DurableTaskGrpcWorker.java | 2 +- .../durabletask/TaskEntityExecutor.java | 12 +++++----- .../microsoft/durabletask/TracingHelper.java | 14 +++++------ .../durabletask/TaskEntityExecutorTest.java | 6 +++++ .../TaskOrchestrationEntityTracingTest.java | 23 +++++++++++++++++++ 5 files changed, 42 insertions(+), 15 deletions(-) diff --git a/client/src/main/java/com/microsoft/durabletask/DurableTaskGrpcWorker.java b/client/src/main/java/com/microsoft/durabletask/DurableTaskGrpcWorker.java index d640f3b2..0bd158cf 100644 --- a/client/src/main/java/com/microsoft/durabletask/DurableTaskGrpcWorker.java +++ b/client/src/main/java/com/microsoft/durabletask/DurableTaskGrpcWorker.java @@ -454,7 +454,7 @@ public void startAndBlock() { EntityRequest entityRequestV2 = workItem.getEntityRequestV2(); this.workItemExecutor.submit(() -> { try { - // Convert V2 (history-based) format to V1 (flat) format + // Convert V2 (history-based) format to V1 (flat) format. EntityBatchRequest.Builder batchBuilder = EntityBatchRequest.newBuilder() .setInstanceId(entityRequestV2.getInstanceId()); if (entityRequestV2.hasEntityState()) { diff --git a/client/src/main/java/com/microsoft/durabletask/TaskEntityExecutor.java b/client/src/main/java/com/microsoft/durabletask/TaskEntityExecutor.java index ee0034a9..7176f6ca 100644 --- a/client/src/main/java/com/microsoft/durabletask/TaskEntityExecutor.java +++ b/client/src/main/java/com/microsoft/durabletask/TaskEntityExecutor.java @@ -132,10 +132,8 @@ EntityBatchResult execute(@Nonnull EntityBatchRequest request) { Instant startTime = Instant.now(); - // Entity processing (SERVER) span. All operations use call_entity semantics: the - // OperationRequest carries no call-vs-signal bit, so we cannot emit CONSUMER for signals - // the way the .NET engine does. Emission is suppressed when emitTraceSpans is false (e.g. - // under Azure Functions, where the host already emits DurableTask.Core entity spans). + // Entity processing span parented on the incoming operation trace context; + // suppressed when the host emits its own (emitTraceSpans is false). Span processingSpan = this.emitTraceSpans ? TracingHelper.startEntityProcessingSpan( entityName, @@ -280,7 +278,8 @@ public void signalEntity( SendSignalAction.Builder signalBuilder = SendSignalAction.newBuilder() .setInstanceId(targetEntityId.toString()) - .setName(operationName); + .setName(operationName) + .setRequestTime(toTimestamp(Instant.now())); if (input != null) { String serializedInput = this.dataConverter.serialize(input); @@ -334,7 +333,8 @@ public String startNewOrchestration( StartNewOrchestrationAction.Builder orchBuilder = StartNewOrchestrationAction.newBuilder() .setInstanceId(instanceId) - .setName(name); + .setName(name) + .setRequestTime(toTimestamp(Instant.now())); if (input != null) { String serializedInput = this.dataConverter.serialize(input); diff --git a/client/src/main/java/com/microsoft/durabletask/TracingHelper.java b/client/src/main/java/com/microsoft/durabletask/TracingHelper.java index 84cfd097..10a36ae6 100644 --- a/client/src/main/java/com/microsoft/durabletask/TracingHelper.java +++ b/client/src/main/java/com/microsoft/durabletask/TracingHelper.java @@ -77,8 +77,7 @@ final class TracingHelper { static final String ATTR_EVENT_TARGET_INSTANCE_ID = "durabletask.event.target_instance_id"; // Entity span constants matching .NET SDK schema. Entity spans deliberately do NOT set - // durabletask.task.name/version/task_id: no .NET entity trace helper sets them, and the - // entity name already appears in the span name. + // durabletask.task.name/version/task_id; the entity name already appears in the span name. static final String TYPE_ENTITY = "entity"; static final String OP_CALL_ENTITY = "call_entity"; static final String OP_SIGNAL_ENTITY = "signal_entity"; @@ -502,8 +501,8 @@ static String createEntityStartOrchestrationSpanName(String entityName) { /** * Starts a processing span for an entity operation: {@link SpanKind#SERVER} for a call, * {@link SpanKind#CONSUMER} for a signal. Returns {@code null} when the parent context is absent - * or invalid, matching the .NET guard that avoids attaching entity spans to an unrelated ambient - * trace. The caller makes the span current and later calls {@link #endEntityProcessingSpan}. + * or invalid, which avoids attaching entity spans to an unrelated ambient trace. The caller makes + * the span current and later calls {@link #endEntityProcessingSpan}. */ @Nullable static Span startEntityProcessingSpan( @@ -528,8 +527,7 @@ static Span startEntityProcessingSpan( /** * Ends a processing span with {@code OK}/{@code Completed} on success or {@code ERROR} plus - * {@code durabletask.entity.error_message} on failure, matching .NET's - * {@code EndActivitiesForProcessingEntityInvocation}. + * {@code durabletask.entity.error_message} on failure. */ static void endEntityProcessingSpan(@Nullable Span span, @Nullable String errorMessage) { if (span == null) { @@ -548,8 +546,8 @@ static void endEntityProcessingSpan(@Nullable Span span, @Nullable String errorM * Emits the retroactive {@link SpanKind#CLIENT} span for a call to an entity, covering the * request-to-response interval. The CLIENT span's ID is set to {@code syntheticSpanId}, which * the SERVER processing span uses as its parent span ID. The status is left unset for normal - * completion or entity-processing failure (matching .NET); {@code errorDescription} is set only - * for timeout/cancellation closure boundaries. + * completion or entity-processing failure; {@code errorDescription} is set only for + * timeout/cancellation closure boundaries. * Does nothing when the parent context is absent or invalid. */ static void emitEntityCallClientSpan( diff --git a/client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTest.java b/client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTest.java index 7d369be0..df3d1374 100644 --- a/client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTest.java +++ b/client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTest.java @@ -358,6 +358,9 @@ void execute_entitySignalsOther_actionsIncluded() { SendSignalAction signalAction = result.getActions(0).getSendSignal(); assertEquals("@counter@target1", signalAction.getInstanceId()); assertEquals("add", signalAction.getName()); + assertTrue(signalAction.hasRequestTime(), "SendSignalAction should set requestTime"); + assertTrue(signalAction.getRequestTime().getSeconds() > 0, + "requestTime should be a real timestamp, not the unset Unix epoch"); } @Test @@ -378,6 +381,9 @@ void execute_entityStartsOrchestration_actionsIncluded() { StartNewOrchestrationAction orchAction = result.getActions(0).getStartNewOrchestration(); assertEquals("MyOrchestration", orchAction.getName()); + assertTrue(orchAction.hasRequestTime(), "StartNewOrchestrationAction should set requestTime"); + assertTrue(orchAction.getRequestTime().getSeconds() > 0, + "requestTime should be a real timestamp, not the unset Unix epoch"); } @Test diff --git a/client/src/test/java/com/microsoft/durabletask/TaskOrchestrationEntityTracingTest.java b/client/src/test/java/com/microsoft/durabletask/TaskOrchestrationEntityTracingTest.java index c381184b..47c06c33 100644 --- a/client/src/test/java/com/microsoft/durabletask/TaskOrchestrationEntityTracingTest.java +++ b/client/src/test/java/com/microsoft/durabletask/TaskOrchestrationEntityTracingTest.java @@ -162,4 +162,27 @@ void signalEntity_emitTraceSpansDisabled_suppressesProducerSpan() { assertTrue(spans.stream().noneMatch(s -> s.getKind() == SpanKind.PRODUCER), "expected no PRODUCER span when emitTraceSpans is disabled"); } + + @Test + void callAndSignal_signalEmitsProducerSpan_callClientSpanDeferred() { + String orchestratorName = "CallAndSignalOrch"; + EntityInstanceId signalTarget = new EntityInstanceId("Counter", "c1"); + EntityInstanceId callTarget = new EntityInstanceId("Counter", "c2"); + TaskOrchestrationExecutor executor = createExecutor(orchestratorName, ctx -> { + ctx.signalEntity(signalTarget, "add", 1); + ctx.callEntity(callTarget, "get", null, Integer.class); + ctx.complete("done"); + }, true); + + executor.execute( + Collections.emptyList(), + Arrays.asList(orchestratorStarted(), executionStarted(orchestratorName)), + orchestrationContext()); + + List spans = spanExporter.getFinishedSpanItems(); + long producers = spans.stream().filter(s -> s.getKind() == SpanKind.PRODUCER).count(); + long clients = spans.stream().filter(s -> s.getKind() == SpanKind.CLIENT).count(); + assertEquals(1L, producers, "signalEntity should emit exactly one PRODUCER span"); + assertEquals(0L, clients, "callEntity CLIENT span is deferred pending protocol support"); + } } From 8ac0fc61b9092735d77f6901a5956bd022cc5423 Mon Sep 17 00:00:00 2001 From: Varshi Bachu Date: Fri, 14 Aug 2026 14:20:32 -0700 Subject: [PATCH 5/6] use span.end(endTime) --- .../main/java/com/microsoft/durabletask/TracingHelper.java | 6 +----- 1 file changed, 1 insertion(+), 5 deletions(-) diff --git a/client/src/main/java/com/microsoft/durabletask/TracingHelper.java b/client/src/main/java/com/microsoft/durabletask/TracingHelper.java index 10a36ae6..f2fec008 100644 --- a/client/src/main/java/com/microsoft/durabletask/TracingHelper.java +++ b/client/src/main/java/com/microsoft/durabletask/TracingHelper.java @@ -583,7 +583,7 @@ static void emitEntityCallClientSpan( span.setStatus(StatusCode.ERROR, errorDescription); } if (endTime != null) { - span.end(toEpochNanos(endTime), java.util.concurrent.TimeUnit.NANOSECONDS); + span.end(endTime); } else { span.end(); } @@ -660,9 +660,5 @@ static Span startEntityStartOrchestrationSpan( return spanBuilder.startSpan(); } - private static long toEpochNanos(java.time.Instant instant) { - return instant.getEpochSecond() * 1_000_000_000L + instant.getNano(); - } - // endregion } From 27f6ea2a483ca1c1792af1e3500e254decf4bab8 Mon Sep 17 00:00:00 2001 From: Varshitha Bachu Date: Fri, 14 Aug 2026 14:26:35 -0700 Subject: [PATCH 6/6] Potential fix for pull request finding 'Useless parameter' Co-authored-by: Copilot Autofix powered by AI <223894421+github-code-quality[bot]@users.noreply.github.com> --- .../microsoft/durabletask/TaskEntityExecutorTracingTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTracingTest.java b/client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTracingTest.java index 7fabc7ca..9d4040f4 100644 --- a/client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTracingTest.java +++ b/client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTracingTest.java @@ -71,7 +71,7 @@ public void signalOther(int amount) { this.context.signalEntity(new EntityInstanceId("counter", "c2"), "add", amount); } - public void startOrch(int amount) { + public void startOrch() { this.context.startNewOrchestration("DownstreamOrch", null); }