From 3459cb476869ca33ccb0963b2f9f4a71c789ccf8 Mon Sep 17 00:00:00 2001 From: Varshi Bachu Date: Fri, 7 Aug 2026 09:04:19 -0700 Subject: [PATCH 1/7] 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/7] 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/7] 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/7] 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/7] 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/7] 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); } From ad9da817f7049bd50f98b2f52aa63a1c6dd16bbd Mon Sep 17 00:00:00 2001 From: Varshi Bachu Date: Tue, 25 Aug 2026 19:40:07 -0700 Subject: [PATCH 7/7] addressed pr feedback --- .../durabletask/GrpcDurableEntityClient.java | 36 ++- .../durabletask/TaskEntityExecutor.java | 224 +++++++++++++----- .../microsoft/durabletask/TracingHelper.java | 101 +------- .../GrpcDurableEntityClientTracingTest.java | 114 +++++++++ .../TaskEntityExecutorTracingTest.java | 76 ++++-- .../TaskOrchestrationEntityTracingTest.java | 22 -- .../durabletask/TracingHelperTest.java | 97 -------- 7 files changed, 365 insertions(+), 305 deletions(-) create mode 100644 client/src/test/java/com/microsoft/durabletask/GrpcDurableEntityClientTracingTest.java diff --git a/client/src/main/java/com/microsoft/durabletask/GrpcDurableEntityClient.java b/client/src/main/java/com/microsoft/durabletask/GrpcDurableEntityClient.java index 343e3856..ece6dfdc 100644 --- a/client/src/main/java/com/microsoft/durabletask/GrpcDurableEntityClient.java +++ b/client/src/main/java/com/microsoft/durabletask/GrpcDurableEntityClient.java @@ -7,6 +7,9 @@ import com.microsoft.durabletask.implementation.protobuf.OrchestratorService.*; import com.microsoft.durabletask.implementation.protobuf.TaskHubSidecarServiceGrpc.TaskHubSidecarServiceBlockingStub; +import io.opentelemetry.api.trace.Span; +import io.opentelemetry.context.Scope; + import javax.annotation.Nullable; import java.time.Instant; import java.util.ArrayList; @@ -39,10 +42,12 @@ public void signalEntity( Helpers.throwIfArgumentNull(entityId, "entityId"); Helpers.throwIfArgumentNull(operationName, "operationName"); + Instant requestTime = Instant.now(); SignalEntityRequest.Builder builder = SignalEntityRequest.newBuilder() .setInstanceId(entityId.toString()) .setName(operationName) - .setRequestId(UUID.randomUUID().toString()); + .setRequestId(UUID.randomUUID().toString()) + .setRequestTime(DataConverter.getTimestampFromInstant(requestTime)); if (input != null) { String serializedInput = this.dataConverter.serialize(input); @@ -58,11 +63,34 @@ public void signalEntity( // Capture and propagate distributed trace context (matching .NET SDK pattern) TraceContext traceContext = TracingHelper.getCurrentTraceContext(); - if (traceContext != null) { - builder.setParentTraceContext(traceContext); + String scheduledTime = options != null && options.getScheduledTime() != null + ? options.getScheduledTime().toString() : null; + Span producerSpan = TracingHelper.startEntitySignalProducerSpan( + entityId.getName(), + operationName, + entityId.toString(), + null, + traceContext, + requestTime, + scheduledTime); + TraceContext producerTraceContext = TracingHelper.getCurrentTraceContext(producerSpan); + TraceContext propagatedTraceContext = producerTraceContext != null + ? producerTraceContext : traceContext; + if (propagatedTraceContext != null) { + builder.setParentTraceContext(propagatedTraceContext); } - this.sidecarClient.signalEntity(builder.build()); + Scope producerScope = producerTraceContext != null ? producerSpan.makeCurrent() : null; + try { + this.sidecarClient.signalEntity(builder.build()); + } finally { + if (producerScope != null) { + producerScope.close(); + } + if (producerSpan != null) { + producerSpan.end(); + } + } } @Override diff --git a/client/src/main/java/com/microsoft/durabletask/TaskEntityExecutor.java b/client/src/main/java/com/microsoft/durabletask/TaskEntityExecutor.java index 7176f6ca..f22e3834 100644 --- a/client/src/main/java/com/microsoft/durabletask/TaskEntityExecutor.java +++ b/client/src/main/java/com/microsoft/durabletask/TaskEntityExecutor.java @@ -7,6 +7,7 @@ import com.microsoft.durabletask.implementation.protobuf.OrchestratorService.*; import io.opentelemetry.api.trace.Span; +import io.opentelemetry.context.Scope; import javax.annotation.Nonnull; import javax.annotation.Nullable; @@ -132,22 +133,11 @@ EntityBatchResult execute(@Nonnull EntityBatchRequest request) { Instant startTime = Instant.now(); - // 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, - 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)); + // The dispatcher owns the processing span and propagates its context with the operation. + // Child spans and actions created by user code use that context without re-emitting it. + TraceContext operationTraceContext = opRequest.hasTraceContext() + ? opRequest.getTraceContext() : null; + context.setCurrentOperationTraceContext(operationTraceContext); try { // Build the operation @@ -155,7 +145,15 @@ EntityBatchResult execute(@Nonnull EntityBatchRequest request) { operationName, serializedInput, context, entityState, this.dataConverter); // Execute - Object result = entity.run(operation); + Object result; + Scope processingScope = TracingHelper.makeTraceContextCurrent(operationTraceContext); + try { + result = entity.run(operation); + } finally { + if (processingScope != null) { + processingScope.close(); + } + } Instant endTime = Instant.now(); @@ -179,8 +177,6 @@ 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}); @@ -212,8 +208,6 @@ EntityBatchResult execute(@Nonnull EntityBatchRequest request) { // Rollback state and actions on failure entityState.rollback(); context.rollback(); - - TracingHelper.endEntityProcessingSpan(processingSpan, e.getMessage()); } } @@ -276,10 +270,11 @@ public void signalEntity( Objects.requireNonNull(targetEntityId, "targetEntityId must not be null"); Objects.requireNonNull(operationName, "operationName must not be null"); + Instant requestTime = Instant.now(); SendSignalAction.Builder signalBuilder = SendSignalAction.newBuilder() .setInstanceId(targetEntityId.toString()) .setName(operationName) - .setRequestTime(toTimestamp(Instant.now())); + .setRequestTime(toTimestamp(requestTime)); if (input != null) { String serializedInput = this.dataConverter.serialize(input); @@ -300,23 +295,17 @@ 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)); + String signalScheduledTime = (options != null && options.getScheduledTime() != null) + ? options.getScheduledTime().toString() : null; + this.pendingActions.add(PendingAction.forSignal( + signalBuilder, + targetEntityId.getName(), + operationName, + targetEntityId.toString(), + this.entityId.toString(), + this.currentOperationTraceContext, + requestTime, + signalScheduledTime)); } @Nonnull @@ -331,10 +320,11 @@ public String startNewOrchestration( ? options.getInstanceId() : UUID.randomUUID().toString(); + Instant requestTime = Instant.now(); StartNewOrchestrationAction.Builder orchBuilder = StartNewOrchestrationAction.newBuilder() .setInstanceId(instanceId) .setName(name) - .setRequestTime(toTimestamp(Instant.now())); + .setRequestTime(toTimestamp(requestTime)); if (input != null) { String serializedInput = this.dataConverter.serialize(input); @@ -360,23 +350,16 @@ 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())); + String orchScheduledTime = (options != null && options.getStartTime() != null) + ? options.getStartTime().toString() : null; + this.pendingActions.add(PendingAction.forOrchestration( + orchBuilder, + this.entityId.getName(), + this.entityId.toString(), + instanceId, + this.currentOperationTraceContext, + requestTime, + orchScheduledTime)); return instanceId; } @@ -385,6 +368,9 @@ public String startNewOrchestration( * Marks the current set of pending actions as committed (snapshot for rollback). */ void commit() { + for (int i = this.committedActionCount; i < this.pendingActions.size(); i++) { + this.pendingActions.get(i).commit(this.emitTraceSpans); + } this.committedActionCount = this.pendingActions.size(); } @@ -411,9 +397,9 @@ List getCommittedActions(int startId) { OperationAction.Builder actionBuilder = OperationAction.newBuilder() .setId(id++); if (pending.type == PendingAction.Type.SEND_SIGNAL) { - actionBuilder.setSendSignal(pending.sendSignal); + actionBuilder.setSendSignal(pending.sendSignal.build()); } else { - actionBuilder.setStartNewOrchestration(pending.startNewOrchestration); + actionBuilder.setStartNewOrchestration(pending.startNewOrchestration.build()); } actions.add(actionBuilder.build()); } @@ -427,13 +413,125 @@ private static class PendingAction { enum Type { SEND_SIGNAL, START_NEW_ORCHESTRATION } final Type type; - final SendSignalAction sendSignal; - final StartNewOrchestrationAction startNewOrchestration; - - PendingAction(Type type, SendSignalAction sendSignal, StartNewOrchestrationAction startNewOrchestration) { + final SendSignalAction.Builder sendSignal; + final StartNewOrchestrationAction.Builder startNewOrchestration; + final String sourceEntityName; + final String sourceEntityInstanceId; + final String targetName; + final String operationName; + final String targetInstanceId; + final TraceContext parentTraceContext; + final Instant requestTime; + final String scheduledTime; + + private PendingAction( + Type type, + SendSignalAction.Builder sendSignal, + StartNewOrchestrationAction.Builder startNewOrchestration, + String sourceEntityName, + String sourceEntityInstanceId, + String targetName, + String operationName, + String targetInstanceId, + TraceContext parentTraceContext, + Instant requestTime, + String scheduledTime) { this.type = type; this.sendSignal = sendSignal; this.startNewOrchestration = startNewOrchestration; + this.sourceEntityName = sourceEntityName; + this.sourceEntityInstanceId = sourceEntityInstanceId; + this.targetName = targetName; + this.operationName = operationName; + this.targetInstanceId = targetInstanceId; + this.parentTraceContext = parentTraceContext; + this.requestTime = requestTime; + this.scheduledTime = scheduledTime; + } + + static PendingAction forSignal( + SendSignalAction.Builder signal, + String targetEntityName, + String operationName, + String targetEntityInstanceId, + String sourceEntityInstanceId, + TraceContext parentTraceContext, + Instant requestTime, + String scheduledTime) { + return new PendingAction( + Type.SEND_SIGNAL, + signal, + null, + null, + sourceEntityInstanceId, + targetEntityName, + operationName, + targetEntityInstanceId, + parentTraceContext, + requestTime, + scheduledTime); + } + + static PendingAction forOrchestration( + StartNewOrchestrationAction.Builder orchestration, + String sourceEntityName, + String sourceEntityInstanceId, + String targetOrchestrationInstanceId, + TraceContext parentTraceContext, + Instant requestTime, + String scheduledTime) { + return new PendingAction( + Type.START_NEW_ORCHESTRATION, + null, + orchestration, + sourceEntityName, + sourceEntityInstanceId, + null, + null, + targetOrchestrationInstanceId, + parentTraceContext, + requestTime, + scheduledTime); + } + + void commit(boolean emitTraceSpans) { + if (!emitTraceSpans || this.parentTraceContext == null) { + return; + } + + Span producerSpan; + if (this.type == Type.SEND_SIGNAL) { + producerSpan = TracingHelper.startEntitySignalProducerSpan( + this.targetName, + this.operationName, + this.targetInstanceId, + this.sourceEntityInstanceId, + this.parentTraceContext, + this.requestTime, + this.scheduledTime); + } else { + producerSpan = TracingHelper.startEntityStartOrchestrationSpan( + this.sourceEntityName, + this.sourceEntityInstanceId, + this.targetInstanceId, + this.parentTraceContext, + this.requestTime, + this.scheduledTime); + } + + if (producerSpan == null) { + return; + } + + TraceContext producerTraceContext = TracingHelper.getCurrentTraceContext(producerSpan); + if (producerTraceContext != null) { + if (this.type == Type.SEND_SIGNAL) { + this.sendSignal.setParentTraceContext(producerTraceContext); + } else { + this.startNewOrchestration.setParentTraceContext(producerTraceContext); + } + } + producerSpan.end(); } } } diff --git a/client/src/main/java/com/microsoft/durabletask/TracingHelper.java b/client/src/main/java/com/microsoft/durabletask/TracingHelper.java index f2fec008..9f774109 100644 --- a/client/src/main/java/com/microsoft/durabletask/TracingHelper.java +++ b/client/src/main/java/com/microsoft/durabletask/TracingHelper.java @@ -16,6 +16,7 @@ import io.opentelemetry.api.trace.Tracer; import io.opentelemetry.api.GlobalOpenTelemetry; import io.opentelemetry.context.Context; +import io.opentelemetry.context.Scope; import javax.annotation.Nullable; import java.lang.reflect.Field; @@ -79,11 +80,9 @@ final class TracingHelper { // Entity span constants matching .NET SDK schema. Entity spans deliberately do NOT set // 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"; 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 @@ -242,6 +241,13 @@ static Context extractTraceContext(@Nullable TraceContext protoCtx) { return Context.current().with(Span.wrap(remoteContext)); } + /** Makes a propagated trace context current until the returned scope is closed. */ + @Nullable + static Scope makeTraceContextCurrent(@Nullable TraceContext traceContext) { + Context context = extractTraceContext(traceContext); + return context != null ? context.makeCurrent() : null; + } + /** * Starts a new span as a child of the given trace context. * @@ -498,97 +504,6 @@ 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, 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( - 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. - */ - 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. 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; {@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(endTime); - } 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 diff --git a/client/src/test/java/com/microsoft/durabletask/GrpcDurableEntityClientTracingTest.java b/client/src/test/java/com/microsoft/durabletask/GrpcDurableEntityClientTracingTest.java new file mode 100644 index 00000000..7ec68e79 --- /dev/null +++ b/client/src/test/java/com/microsoft/durabletask/GrpcDurableEntityClientTracingTest.java @@ -0,0 +1,114 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. +// Licensed under the MIT License. +package com.microsoft.durabletask; + +import com.microsoft.durabletask.implementation.protobuf.OrchestratorService.SignalEntityRequest; +import com.microsoft.durabletask.implementation.protobuf.OrchestratorService.SignalEntityResponse; +import com.microsoft.durabletask.implementation.protobuf.TaskHubSidecarServiceGrpc; + +import io.grpc.ManagedChannel; +import io.grpc.Server; +import io.grpc.inprocess.InProcessChannelBuilder; +import io.grpc.inprocess.InProcessServerBuilder; +import io.grpc.stub.StreamObserver; +import io.opentelemetry.api.GlobalOpenTelemetry; +import io.opentelemetry.api.trace.Span; +import io.opentelemetry.api.trace.SpanKind; +import io.opentelemetry.context.Scope; +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.concurrent.atomic.AtomicReference; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +public class GrpcDurableEntityClientTracingTest { + + private final AtomicReference capturedRequest = new AtomicReference<>(); + private InMemorySpanExporter spanExporter; + private OpenTelemetrySdk openTelemetry; + private Server inProcessServer; + private ManagedChannel inProcessChannel; + private DurableTaskClient client; + + @BeforeEach + void setUp() throws Exception { + GlobalOpenTelemetry.resetForTest(); + this.spanExporter = InMemorySpanExporter.create(); + SdkTracerProvider tracerProvider = SdkTracerProvider.builder() + .addSpanProcessor(SimpleSpanProcessor.create(this.spanExporter)) + .build(); + this.openTelemetry = OpenTelemetrySdk.builder() + .setTracerProvider(tracerProvider) + .buildAndRegisterGlobal(); + + String serverName = InProcessServerBuilder.generateName(); + this.inProcessServer = InProcessServerBuilder.forName(serverName) + .directExecutor() + .addService(new TaskHubSidecarServiceGrpc.TaskHubSidecarServiceImplBase() { + @Override + public void signalEntity( + SignalEntityRequest request, + StreamObserver responseObserver) { + capturedRequest.set(request); + responseObserver.onNext(SignalEntityResponse.getDefaultInstance()); + responseObserver.onCompleted(); + } + }) + .build() + .start(); + this.inProcessChannel = InProcessChannelBuilder.forName(serverName).directExecutor().build(); + this.client = new DurableTaskGrpcClientBuilder().grpcChannel(this.inProcessChannel).build(); + } + + @AfterEach + void tearDown() { + if (this.inProcessChannel != null) { + this.inProcessChannel.shutdownNow(); + } + if (this.inProcessServer != null) { + this.inProcessServer.shutdownNow(); + } + if (this.openTelemetry != null) { + this.openTelemetry.close(); + } + GlobalOpenTelemetry.resetForTest(); + } + + @Test + void signalEntity_emitsProducerSpanAndPropagatesItsContext() { + Span parentSpan = GlobalOpenTelemetry.getTracer("test").spanBuilder("parent").startSpan(); + try (Scope ignored = parentSpan.makeCurrent()) { + this.client.getEntities().signalEntity( + new EntityInstanceId("Counter", "c1"), + "add", + 5); + } finally { + parentSpan.end(); + } + + SignalEntityRequest request = this.capturedRequest.get(); + assertNotNull(request); + assertTrue(request.hasRequestTime()); + assertTrue(request.getRequestTime().getSeconds() > 0); + + SpanData producer = this.spanExporter.getFinishedSpanItems().stream() + .filter(span -> span.getKind() == SpanKind.PRODUCER) + .findFirst() + .orElse(null); + assertNotNull(producer, "expected external entity signal PRODUCER span"); + assertEquals("entity:counter:add", producer.getName()); + assertEquals(parentSpan.getSpanContext().getSpanId(), producer.getParentSpanId()); + assertEquals(producer.getSpanId(), + request.getParentTraceContext().getTraceParent().split("-")[2]); + } +} \ No newline at end of file diff --git a/client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTracingTest.java b/client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTracingTest.java index 9d4040f4..505417bc 100644 --- a/client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTracingTest.java +++ b/client/src/test/java/com/microsoft/durabletask/TaskEntityExecutorTracingTest.java @@ -10,7 +10,6 @@ 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; @@ -30,8 +29,7 @@ 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. + * Verifies entity action spans and propagation from the dispatcher-owned processing context. */ public class TaskEntityExecutorTracingTest { @@ -71,6 +69,19 @@ public void signalOther(int amount) { this.context.signalEntity(new EntityInstanceId("counter", "c2"), "add", amount); } + public void signalThenThrow(int amount) { + this.context.signalEntity(new EntityInstanceId("counter", "c2"), "add", amount); + throw new IllegalStateException("operation failed"); + } + + public void createNestedSpan(int amount) { + io.opentelemetry.api.GlobalOpenTelemetry.getTracer("test") + .spanBuilder("entity-user-code") + .setAttribute("amount", amount) + .startSpan() + .end(); + } + public void startOrch() { this.context.startNewOrchestration("DownstreamOrch", null); } @@ -119,25 +130,12 @@ private static TraceContext parentTraceContext() { } @Test - void execute_emitsEntityProcessingServerSpanUnderParent() { + void execute_doesNotDuplicateDispatcherProcessingSpan() { 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()); + assertTrue(spanExporter.getFinishedSpanItems().isEmpty()); } @Test @@ -161,39 +159,65 @@ void execute_noParentTraceContext_emitsNoSpan() { } @Test - void execute_entitySignalsEntity_emitsProducerSpanNestedUnderProcessingSpan() { + void execute_entitySignalsEntity_emitsProducerSpanUnderProcessingContext() { 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()); + assertEquals(PARENT_SPAN_ID, producer.getParentSpanId()); + assertEquals(producer.getSpanId(), result.getActions(0).getSendSignal() + .getParentTraceContext().getTraceParent().split("-")[2]); } @Test - void execute_entityStartsOrchestration_emitsProducerSpanNestedUnderProcessingSpan() { + void execute_entityStartsOrchestration_emitsProducerSpanUnderProcessingContext() { 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()); + assertEquals(PARENT_SPAN_ID, producer.getParentSpanId()); + assertEquals(producer.getSpanId(), result.getActions(0).getStartNewOrchestration() + .getParentTraceContext().getTraceParent().split("-")[2]); + } + + @Test + void execute_makesProcessingSpanCurrentForUserCode() { + TaskEntityExecutor executor = createExecutor(true); + + EntityBatchResult result = executor.execute(requestWithOp("createNestedSpan", 3, parentTraceContext())); + + assertTrue(result.getResults(0).hasSuccess()); + List spans = spanExporter.getFinishedSpanItems(); + SpanData nested = spans.stream().filter(s -> s.getName().equals("entity-user-code")).findFirst().orElse(null); + assertNotNull(nested, "expected span created by entity user code"); + assertEquals(TRACE_ID, nested.getTraceId()); + assertEquals(PARENT_SPAN_ID, nested.getParentSpanId()); + } + + @Test + void execute_entityActionRollsBack_emitsNoProducerSpan() { + TaskEntityExecutor executor = createExecutor(true); + + EntityBatchResult result = executor.execute(requestWithOp("signalThenThrow", 3, parentTraceContext())); + + assertTrue(result.getResults(0).hasFailure()); + assertEquals(0, result.getActionsCount()); + assertTrue(spanExporter.getFinishedSpanItems().stream() + .noneMatch(span -> span.getKind() == SpanKind.PRODUCER)); } } diff --git a/client/src/test/java/com/microsoft/durabletask/TaskOrchestrationEntityTracingTest.java b/client/src/test/java/com/microsoft/durabletask/TaskOrchestrationEntityTracingTest.java index 47c06c33..2f31fdb3 100644 --- a/client/src/test/java/com/microsoft/durabletask/TaskOrchestrationEntityTracingTest.java +++ b/client/src/test/java/com/microsoft/durabletask/TaskOrchestrationEntityTracingTest.java @@ -163,26 +163,4 @@ void signalEntity_emitTraceSpansDisabled_suppressesProducerSpan() { "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"); - } } diff --git a/client/src/test/java/com/microsoft/durabletask/TracingHelperTest.java b/client/src/test/java/com/microsoft/durabletask/TracingHelperTest.java index cfa5a90b..79aa9f7d 100644 --- a/client/src/test/java/com/microsoft/durabletask/TracingHelperTest.java +++ b/client/src/test/java/com/microsoft/durabletask/TracingHelperTest.java @@ -385,103 +385,6 @@ void createEntityStartOrchestrationSpanName_isInverted() { 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(