From 1ecaa96363685543b73d5cdad08902c74ba69fec Mon Sep 17 00:00:00 2001 From: jcstorms1 Date: Wed, 23 Sep 2026 11:27:25 -0400 Subject: [PATCH 01/10] Add Azure Durable Functions worker instrumentation --- .../azure-functions-worker-2.7/build.gradle | 26 ++ .../gradle.lockfile | 130 +++++++ .../AzureFunctionsWorkerInstrumentation.java | 106 ++++++ .../worker/DurableFunctionsDecorator.java | 46 +++ .../worker/DurableFunctionsUtils.java | 165 ++++++++ .../worker/TraceContextExtractAdapter.java | 87 +++++ .../groovy/AzureFunctionsWorkerTest.groovy | 357 ++++++++++++++++++ .../chain/FunctionExecutionMiddleware.java | 10 + .../OrchestratorBlockedException.java | 3 + .../ContinueAsNewInterruption.java | 3 + .../OrchestratorBlockedException.java | 3 + .../core/propagation/ExtractedContext.java | 11 + .../propagation/ExtractedContextTest.java | 40 ++ .../instrumentation/api/AgentSpanContext.java | 8 + .../instrumentation/api/TagContext.java | 6 +- settings.gradle.kts | 1 + 16 files changed, 1001 insertions(+), 1 deletion(-) create mode 100644 dd-java-agent/instrumentation/azure-functions-worker-2.7/build.gradle create mode 100644 dd-java-agent/instrumentation/azure-functions-worker-2.7/gradle.lockfile create mode 100644 dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java create mode 100644 dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsDecorator.java create mode 100644 dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java create mode 100644 dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/TraceContextExtractAdapter.java create mode 100644 dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/groovy/AzureFunctionsWorkerTest.groovy create mode 100644 dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/com/microsoft/azure/functions/worker/chain/FunctionExecutionMiddleware.java create mode 100644 dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/com/microsoft/durabletask/OrchestratorBlockedException.java create mode 100644 dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/com/microsoft/durabletask/interruption/ContinueAsNewInterruption.java create mode 100644 dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/com/microsoft/durabletask/interruption/OrchestratorBlockedException.java create mode 100644 dd-trace-core/src/test/java/datadog/trace/core/propagation/ExtractedContextTest.java diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/build.gradle b/dd-java-agent/instrumentation/azure-functions-worker-2.7/build.gradle new file mode 100644 index 00000000000..930115252af --- /dev/null +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/build.gradle @@ -0,0 +1,26 @@ +plugins { + id 'dd-trace-java.module.instrumentation' +} + +muzzle { + pass { + group = 'com.microsoft.azure.functions' + module = 'azure-functions-java-spi' + versions = '[1.0.0,)' + extraDependency 'com.microsoft.azure.functions:azure-functions-java-core-library:1.2.0' + } +} + +addTestSuiteForDir('latestDepTest', 'test') + +dependencies { + compileOnly group: 'com.microsoft.azure.functions', name: 'azure-functions-java-core-library', version: '1.2.0' + compileOnly group: 'com.microsoft.azure.functions', name: 'azure-functions-java-spi', version: '1.0.0' + + testImplementation group: 'com.microsoft.azure.functions', name: 'azure-functions-java-core-library', version: '1.2.0' + testImplementation group: 'com.microsoft.azure.functions', name: 'azure-functions-java-spi', version: '1.0.0' + testImplementation libs.bundles.mockito + + latestDepTestImplementation group: 'com.microsoft.azure.functions', name: 'azure-functions-java-core-library', version: '+' + latestDepTestImplementation group: 'com.microsoft.azure.functions', name: 'azure-functions-java-spi', version: '+' +} diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/gradle.lockfile b/dd-java-agent/instrumentation/azure-functions-worker-2.7/gradle.lockfile new file mode 100644 index 00000000000..a8ef06fa2bb --- /dev/null +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/gradle.lockfile @@ -0,0 +1,130 @@ +# This is a Gradle generated file for dependency locking. +# Manual edits can break the build and are not advised. +# This file is expected to be part of source control. +# To regenerate this file, run: ./gradlew :dd-java-agent:instrumentation:azure-functions-worker-2.7:dependencies --write-locks +cafe.cryptography:curve25519-elisabeth:0.1.0=latestDepTestRuntimeClasspath,testRuntimeClasspath +cafe.cryptography:ed25519-elisabeth:0.1.0=latestDepTestRuntimeClasspath,testRuntimeClasspath +ch.qos.logback:logback-classic:1.2.13=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +ch.qos.logback:logback-core:1.2.13=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +com.blogspot.mydailyjava:weak-lock-free:0.17=buildTimeInstrumentationPlugin,compileClasspath,latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,muzzleTooling,runtimeClasspath,testCompileClasspath,testRuntimeClasspath +com.datadoghq.okhttp3:okhttp:3.12.15=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +com.datadoghq.okio:okio:1.17.6=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +com.datadoghq:dd-instrument-java:0.0.5=buildTimeInstrumentationPlugin,compileClasspath,latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,muzzleBootstrap,muzzleTooling,runtimeClasspath,testCompileClasspath,testRuntimeClasspath +com.datadoghq:dd-javac-plugin-client:0.2.2=buildTimeInstrumentationPlugin,compileClasspath,latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,muzzleBootstrap,muzzleTooling,runtimeClasspath,testCompileClasspath,testRuntimeClasspath +com.datadoghq:java-dogstatsd-client:4.4.5=latestDepTestRuntimeClasspath,testRuntimeClasspath +com.datadoghq:sketches-java:0.8.3=latestDepTestRuntimeClasspath,testRuntimeClasspath +com.github.javaparser:javaparser-core:3.25.6=codenarc +com.github.jnr:jffi:1.3.15=latestDepTestRuntimeClasspath,testRuntimeClasspath +com.github.jnr:jnr-a64asm:1.0.0=latestDepTestRuntimeClasspath,testRuntimeClasspath +com.github.jnr:jnr-constants:0.10.4=latestDepTestRuntimeClasspath,testRuntimeClasspath +com.github.jnr:jnr-enxio:0.32.20=latestDepTestRuntimeClasspath,testRuntimeClasspath +com.github.jnr:jnr-ffi:2.2.19=latestDepTestRuntimeClasspath,testRuntimeClasspath +com.github.jnr:jnr-posix:3.1.22=latestDepTestRuntimeClasspath,testRuntimeClasspath +com.github.jnr:jnr-unixsocket:0.38.25=latestDepTestRuntimeClasspath,testRuntimeClasspath +com.github.jnr:jnr-x86asm:1.0.2=latestDepTestRuntimeClasspath,testRuntimeClasspath +com.github.spotbugs:spotbugs-annotations:4.10.3=compileClasspath,spotbugs +com.github.spotbugs:spotbugs:4.10.3=spotbugs +com.github.stephenc.jcip:jcip-annotations:1.0-1=spotbugs +com.google.auto.service:auto-service-annotations:1.1.1=annotationProcessor,compileClasspath,latestDepTestAnnotationProcessor,latestDepTestCompileClasspath,testAnnotationProcessor,testCompileClasspath +com.google.auto.service:auto-service:1.1.1=annotationProcessor,latestDepTestAnnotationProcessor,testAnnotationProcessor +com.google.auto:auto-common:1.2.1=annotationProcessor,latestDepTestAnnotationProcessor,testAnnotationProcessor +com.google.code.findbugs:jsr305:3.0.2=annotationProcessor,compileClasspath,latestDepTestAnnotationProcessor,latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,spotbugs,testAnnotationProcessor,testCompileClasspath,testRuntimeClasspath +com.google.code.gson:gson:2.14.0=spotbugs +com.google.errorprone:error_prone_annotations:2.18.0=annotationProcessor,latestDepTestAnnotationProcessor,testAnnotationProcessor +com.google.errorprone:error_prone_annotations:2.47.0=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +com.google.errorprone:error_prone_annotations:2.48.0=spotbugs +com.google.guava:failureaccess:1.0.1=annotationProcessor,latestDepTestAnnotationProcessor,testAnnotationProcessor +com.google.guava:failureaccess:1.0.3=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +com.google.guava:guava:32.0.1-jre=annotationProcessor,latestDepTestAnnotationProcessor,testAnnotationProcessor +com.google.guava:guava:33.6.0-jre=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +com.google.guava:listenablefuture:9999.0-empty-to-avoid-conflict-with-guava=annotationProcessor,latestDepTestAnnotationProcessor,latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testAnnotationProcessor,testCompileClasspath,testRuntimeClasspath +com.google.j2objc:j2objc-annotations:2.8=annotationProcessor,latestDepTestAnnotationProcessor,testAnnotationProcessor +com.google.j2objc:j2objc-annotations:3.1=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +com.google.re2j:re2j:1.8=latestDepTestRuntimeClasspath,testRuntimeClasspath +com.microsoft.azure.functions:azure-functions-java-core-library:1.2.0=compileClasspath,testCompileClasspath,testRuntimeClasspath +com.microsoft.azure.functions:azure-functions-java-core-library:1.3.0=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath +com.microsoft.azure.functions:azure-functions-java-spi:1.0.0=compileClasspath,testCompileClasspath,testRuntimeClasspath +com.microsoft.azure.functions:azure-functions-java-spi:1.1.0=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath +com.squareup.moshi:moshi:1.11.0=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +com.squareup.okhttp3:logging-interceptor:3.12.12=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +com.squareup.okhttp3:okhttp:3.12.12=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +com.squareup.okio:okio:1.17.5=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +com.thoughtworks.qdox:qdox:1.12.1=codenarc +commons-fileupload:commons-fileupload:1.5=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +commons-io:commons-io:2.11.0=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +commons-io:commons-io:2.21.0=spotbugs +de.thetaphi:forbiddenapis:3.10=compileClasspath,latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +io.leangen.geantyref:geantyref:1.3.16=latestDepTestRuntimeClasspath,testRuntimeClasspath +io.sqreen:libsqreen:17.5.0=latestDepTestRuntimeClasspath,testRuntimeClasspath +javax.servlet:javax.servlet-api:3.1.0=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +jaxen:jaxen:2.0.6=spotbugs +junit:junit:4.13.2=latestDepTestRuntimeClasspath,testRuntimeClasspath +net.bytebuddy:byte-buddy-agent:1.18.12=buildTimeInstrumentationPlugin,compileClasspath,latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,muzzleTooling,runtimeClasspath,testCompileClasspath,testRuntimeClasspath +net.bytebuddy:byte-buddy:1.18.12=buildTimeInstrumentationPlugin,compileClasspath,latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,muzzleTooling,runtimeClasspath,testCompileClasspath,testRuntimeClasspath +net.java.dev.jna:jna-platform:5.8.0=latestDepTestRuntimeClasspath,testRuntimeClasspath +net.java.dev.jna:jna:5.8.0=latestDepTestRuntimeClasspath,testRuntimeClasspath +net.sf.saxon:Saxon-HE:12.10=spotbugs +org.apache.ant:ant-antlr:1.10.14=codenarc +org.apache.ant:ant-junit:1.10.14=codenarc +org.apache.bcel:bcel:6.12.0=spotbugs +org.apache.commons:commons-lang3:3.20.0=spotbugs +org.apache.commons:commons-text:1.15.0=spotbugs +org.apache.logging.log4j:log4j-api:2.26.1=spotbugs +org.apache.logging.log4j:log4j-core:2.26.1=spotbugs +org.apiguardian:apiguardian-api:1.1.2=latestDepTestCompileClasspath,testCompileClasspath +org.checkerframework:checker-qual:3.33.0=annotationProcessor,latestDepTestAnnotationProcessor,testAnnotationProcessor +org.codehaus.groovy:groovy-ant:3.0.23=codenarc +org.codehaus.groovy:groovy-docgenerator:3.0.23=codenarc +org.codehaus.groovy:groovy-groovydoc:3.0.23=codenarc +org.codehaus.groovy:groovy-json:3.0.23=codenarc +org.codehaus.groovy:groovy-json:3.0.25=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +org.codehaus.groovy:groovy-templates:3.0.23=codenarc +org.codehaus.groovy:groovy-xml:3.0.23=codenarc +org.codehaus.groovy:groovy:3.0.23=codenarc +org.codehaus.groovy:groovy:3.0.25=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +org.codenarc:CodeNarc:3.7.0=codenarc +org.dom4j:dom4j:2.2.0=spotbugs +org.gmetrics:GMetrics:2.1.0=codenarc +org.hamcrest:hamcrest-core:1.3=latestDepTestRuntimeClasspath,testRuntimeClasspath +org.hamcrest:hamcrest:3.0=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +org.jctools:jctools-core-jdk11:4.0.6=latestDepTestRuntimeClasspath,testRuntimeClasspath +org.jctools:jctools-core:4.0.6=latestDepTestRuntimeClasspath,testRuntimeClasspath +org.jspecify:jspecify:1.0.0=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +org.junit.jupiter:junit-jupiter-api:5.14.1=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +org.junit.jupiter:junit-jupiter-engine:5.14.1=latestDepTestRuntimeClasspath,testRuntimeClasspath +org.junit.jupiter:junit-jupiter-params:5.14.1=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +org.junit.jupiter:junit-jupiter:5.14.1=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +org.junit.platform:junit-platform-commons:1.14.1=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +org.junit.platform:junit-platform-engine:1.14.1=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +org.junit.platform:junit-platform-launcher:1.14.1=latestDepTestRuntimeClasspath,testRuntimeClasspath +org.junit.platform:junit-platform-runner:1.14.1=latestDepTestRuntimeClasspath,testRuntimeClasspath +org.junit.platform:junit-platform-suite-api:1.14.1=latestDepTestRuntimeClasspath,testRuntimeClasspath +org.junit.platform:junit-platform-suite-commons:1.14.1=latestDepTestRuntimeClasspath,testRuntimeClasspath +org.junit:junit-bom:5.14.1=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +org.junit:junit-bom:6.1.2=spotbugs +org.mockito:mockito-core:4.4.0=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +org.mockito:mockito-junit-jupiter:4.4.0=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +org.objenesis:objenesis:3.3=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +org.opentest4j:opentest4j:1.3.0=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +org.ow2.asm:asm-analysis:9.10.1=spotbugs +org.ow2.asm:asm-analysis:9.7.1=latestDepTestRuntimeClasspath,testRuntimeClasspath +org.ow2.asm:asm-commons:9.10.1=latestDepTestRuntimeClasspath,spotbugs,testRuntimeClasspath +org.ow2.asm:asm-tree:9.10.1=latestDepTestRuntimeClasspath,spotbugs,testRuntimeClasspath +org.ow2.asm:asm-util:9.10.1=spotbugs +org.ow2.asm:asm-util:9.7.1=latestDepTestRuntimeClasspath,testRuntimeClasspath +org.ow2.asm:asm:9.10.1=buildTimeInstrumentationPlugin,compileClasspath,latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,muzzleTooling,runtimeClasspath,spotbugs,testCompileClasspath,testRuntimeClasspath +org.slf4j:jcl-over-slf4j:1.7.30=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +org.slf4j:jul-to-slf4j:1.7.30=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +org.slf4j:log4j-over-slf4j:1.7.30=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +org.slf4j:slf4j-api:1.7.30=buildTimeInstrumentationPlugin,compileClasspath,muzzleBootstrap,muzzleTooling,runtimeClasspath +org.slf4j:slf4j-api:1.7.32=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +org.slf4j:slf4j-api:2.0.17=spotbugsSlf4j +org.slf4j:slf4j-api:2.0.18=spotbugs +org.slf4j:slf4j-simple:2.0.17=spotbugsSlf4j +org.snakeyaml:snakeyaml-engine:2.9=buildTimeInstrumentationPlugin,latestDepTestRuntimeClasspath,muzzleTooling,runtimeClasspath,testRuntimeClasspath +org.spockframework:spock-bom:2.4-groovy-3.0=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +org.spockframework:spock-core:2.4-groovy-3.0=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +org.tabletest:tabletest-junit:1.2.2=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +org.tabletest:tabletest-parser:1.2.1=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath +org.xmlresolver:xmlresolver:5.3.3=spotbugs +empty=spotbugsPlugins diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java new file mode 100644 index 00000000000..8c2332ef049 --- /dev/null +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java @@ -0,0 +1,106 @@ +package datadog.trace.instrumentation.azure.functions.worker; + +import static datadog.trace.agent.tooling.bytebuddy.matcher.NameMatchers.named; +import static datadog.trace.bootstrap.instrumentation.api.AgentPropagation.extractContextAndGetSpanContext; +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activateSpan; +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.startSpan; +import static datadog.trace.bootstrap.instrumentation.api.Java8BytecodeBridge.spanFromScope; +import static datadog.trace.instrumentation.azure.functions.worker.DurableFunctionsDecorator.AZURE_FUNCTIONS_REQUEST; +import static datadog.trace.instrumentation.azure.functions.worker.DurableFunctionsDecorator.DECORATE; +import static net.bytebuddy.matcher.ElementMatchers.isMethod; +import static net.bytebuddy.matcher.ElementMatchers.isPublic; +import static net.bytebuddy.matcher.ElementMatchers.takesArgument; +import static net.bytebuddy.matcher.ElementMatchers.takesArguments; + +import com.google.auto.service.AutoService; +import com.microsoft.azure.functions.TraceContext; +import com.microsoft.azure.functions.internal.spi.middleware.MiddlewareContext; +import datadog.context.ContextScope; +import datadog.trace.agent.tooling.Instrumenter; +import datadog.trace.agent.tooling.InstrumenterModule; +import datadog.trace.bootstrap.instrumentation.api.AgentSpan; +import datadog.trace.bootstrap.instrumentation.api.AgentSpanContext; +import net.bytebuddy.asm.Advice; + +@AutoService(InstrumenterModule.class) +public final class AzureFunctionsWorkerInstrumentation extends InstrumenterModule.Tracing + implements Instrumenter.ForSingleType, Instrumenter.HasMethodAdvice { + + public AzureFunctionsWorkerInstrumentation() { + super("azure-functions"); + } + + @Override + public String instrumentedType() { + return "com.microsoft.azure.functions.worker.chain.FunctionExecutionMiddleware"; + } + + @Override + public String[] helperClassNames() { + return new String[] { + packageName + ".DurableFunctionsDecorator", + packageName + ".DurableFunctionsUtils", + packageName + ".TraceContextExtractAdapter" + }; + } + + @Override + public void methodAdvice(MethodTransformer transformer) { + transformer.applyAdvice( + isMethod() + .and(isPublic()) + .and(named("invoke")) + .and(takesArguments(2)) + .and( + takesArgument( + 0, + named( + "com.microsoft.azure.functions.internal.spi.middleware.MiddlewareContext"))) + .and( + takesArgument( + 1, + named( + "com.microsoft.azure.functions.internal.spi.middleware.MiddlewareChain"))), + AzureFunctionsWorkerInstrumentation.class.getName() + "$InvokeAdvice"); + } + + public static class InvokeAdvice { + @Advice.OnMethodEnter(suppress = Throwable.class) + public static ContextScope onEnter(@Advice.Argument(0) MiddlewareContext context) { + final String trigger = DurableFunctionsUtils.getTrigger(context); + if (trigger == null) { + return null; + } + if ("DurableOrchestration".equals(trigger) + && !DurableFunctionsUtils.shouldTraceOrchestration(context)) { + return null; + } + + final TraceContext traceContext = context.getTraceContext(); + AgentSpanContext.Extracted parent = + traceContext == null + ? null + : extractContextAndGetSpanContext(traceContext, TraceContextExtractAdapter.GETTER); + parent = DurableFunctionsUtils.reconcileSamplingPriority(parent, traceContext); + + final AgentSpan span = startSpan("azure-functions", AZURE_FUNCTIONS_REQUEST, parent); + DECORATE.afterStart(span); + DECORATE.onInvoke(span, context.getFunctionName(), trigger); + return activateSpan(span); + } + + @Advice.OnMethodExit(onThrowable = Throwable.class, suppress = Throwable.class) + public static void onExit( + @Advice.Enter ContextScope scope, @Advice.Thrown Throwable throwable) { + if (scope != null) { + final AgentSpan span = spanFromScope(scope); + if (!DurableFunctionsUtils.isReplayControlFlow(throwable)) { + DECORATE.onError(span, throwable); + } + DECORATE.beforeFinish(span); + scope.close(); + span.finish(); + } + } + } +} diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsDecorator.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsDecorator.java new file mode 100644 index 00000000000..f6ec65d193a --- /dev/null +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsDecorator.java @@ -0,0 +1,46 @@ +package datadog.trace.instrumentation.azure.functions.worker; + +import datadog.trace.api.naming.SpanNaming; +import datadog.trace.bootstrap.instrumentation.api.AgentSpan; +import datadog.trace.bootstrap.instrumentation.api.InternalSpanTypes; +import datadog.trace.bootstrap.instrumentation.api.Tags; +import datadog.trace.bootstrap.instrumentation.api.UTF8BytesString; +import datadog.trace.bootstrap.instrumentation.decorator.BaseDecorator; + +public final class DurableFunctionsDecorator extends BaseDecorator { + public static final DurableFunctionsDecorator DECORATE = new DurableFunctionsDecorator(); + public static final CharSequence AZURE_FUNCTIONS_REQUEST = + UTF8BytesString.create( + SpanNaming.instance().namingSchema().cloud().operationForFaas("azure")); + + private static final CharSequence AZURE_FUNCTIONS = UTF8BytesString.create("azure-functions"); + + private DurableFunctionsDecorator() {} + + @Override + protected String[] instrumentationNames() { + return new String[] {"azure-functions"}; + } + + @Override + protected CharSequence spanType() { + return InternalSpanTypes.SERVERLESS; + } + + @Override + protected CharSequence component() { + return AZURE_FUNCTIONS; + } + + @Override + protected void doAfterStart(AgentSpan span) { + super.doAfterStart(span); + span.setTag(Tags.SPAN_KIND, Tags.SPAN_KIND_SERVER); + } + + public void onInvoke(AgentSpan span, String functionName, String trigger) { + span.setResourceName(trigger + " " + functionName); + span.setTag("aas.function.name", functionName); + span.setTag("aas.function.trigger", trigger); + } +} diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java new file mode 100644 index 00000000000..570cb2df037 --- /dev/null +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java @@ -0,0 +1,165 @@ +package datadog.trace.instrumentation.azure.functions.worker; + +import static datadog.trace.api.sampling.PrioritySampling.SAMPLER_DROP; +import static datadog.trace.api.sampling.PrioritySampling.UNSET; + +import com.microsoft.azure.functions.TraceContext; +import com.microsoft.azure.functions.internal.spi.middleware.MiddlewareContext; +import datadog.trace.bootstrap.instrumentation.api.AgentSpanContext; +import java.util.Base64; + +public final class DurableFunctionsUtils { + private static final String ACTIVITY_ANNOTATION = "DurableActivityTrigger"; + private static final String ENTITY_ANNOTATION = "DurableEntityTrigger"; + private static final String ORCHESTRATION_ANNOTATION = "DurableOrchestrationTrigger"; + + // OrchestratorRequest protobuf field numbers. + private static final int PAST_EVENTS_FIELD = 3; + private static final int NEW_EVENTS_FIELD = 4; + + // HistoryEvent protobuf field numbers. + private static final int TASK_FAILED_FIELD = 8; + private static final int SUB_ORCHESTRATION_FAILED_FIELD = 11; + + private DurableFunctionsUtils() {} + + public static String getTrigger(MiddlewareContext context) { + if (context.getParameterName(ORCHESTRATION_ANNOTATION) != null) { + return "DurableOrchestration"; + } + if (context.getParameterName(ACTIVITY_ANNOTATION) != null) { + return "DurableActivity"; + } + if (context.getParameterName(ENTITY_ANNOTATION) != null) { + return "DurableEntity"; + } + return null; + } + + public static AgentSpanContext.Extracted reconcileSamplingPriority( + AgentSpanContext.Extracted parent, TraceContext traceContext) { + if (parent != null + && parent.getSamplingPriority() == SAMPLER_DROP + && traceContext != null + && TraceContextExtractAdapter.datadogSamplingPriority(traceContext.getTracestate()) + == UNSET) { + return parent.withSamplingPriority(UNSET); + } + return parent; + } + + public static boolean isReplayControlFlow(Throwable throwable) { + while (throwable != null) { + final String className = throwable.getClass().getName(); + if ("com.microsoft.durabletask.OrchestratorBlockedException".equals(className) + || "com.microsoft.durabletask.interruption.OrchestratorBlockedException".equals(className) + || "com.microsoft.durabletask.interruption.ContinueAsNewInterruption".equals(className)) { + return true; + } + throwable = throwable.getCause(); + } + return false; + } + + /** + * Returns whether an orchestration invocation should produce a span. + * + *

Durable orchestrators replay from the beginning whenever new history arrives. Those + * executions are implementation details rather than new logical operations, so successful replays + * are suppressed. A replay containing a newly failed activity or sub-orchestration is retained so + * the failure remains visible. If the trigger payload cannot be inspected, tracing fails open. + */ + public static boolean shouldTraceOrchestration(MiddlewareContext context) { + try { + final String parameterName = context.getParameterName(ORCHESTRATION_ANNOTATION); + final Object parameterValue = context.getParameterValue(parameterName); + if (!(parameterValue instanceof String)) { + return true; + } + + final byte[] request = Base64.getDecoder().decode((String) parameterValue); + boolean replay = false; + boolean newFailure = false; + final int[] position = {0}; + while (position[0] < request.length) { + final long tag = readVarint(request, position, request.length); + final int field = (int) (tag >>> 3); + final int wireType = (int) (tag & 7); + if (wireType == 2) { + final long length = readVarint(request, position, request.length); + if (length < 0 || length > request.length - position[0]) { + return true; + } + final int end = position[0] + (int) length; + if (field == PAST_EVENTS_FIELD) { + replay = true; + } else if (field == NEW_EVENTS_FIELD && containsFailureEvent(request, position[0], end)) { + newFailure = true; + } + position[0] = end; + } else { + skipValue(request, position, request.length, wireType); + } + } + return !replay || newFailure; + } catch (Throwable ignored) { + return true; + } + } + + private static boolean containsFailureEvent(byte[] data, int offset, int limit) { + final int[] position = {offset}; + while (position[0] < limit) { + final long tag = readVarint(data, position, limit); + final int field = (int) (tag >>> 3); + final int wireType = (int) (tag & 7); + if (wireType == 2 + && (field == TASK_FAILED_FIELD || field == SUB_ORCHESTRATION_FAILED_FIELD)) { + return true; + } + skipValue(data, position, limit, wireType); + } + return false; + } + + private static long readVarint(byte[] data, int[] position, int limit) { + long value = 0; + for (int shift = 0; shift < 64 && position[0] < limit; shift += 7) { + final int current = data[position[0]++] & 0xff; + value |= (long) (current & 0x7f) << shift; + if ((current & 0x80) == 0) { + return value; + } + } + throw new IllegalArgumentException("Invalid protobuf varint"); + } + + private static void skipValue(byte[] data, int[] position, int limit, int wireType) { + switch (wireType) { + case 0: + readVarint(data, position, limit); + return; + case 1: + if (limit - position[0] < 8) { + throw new IllegalArgumentException("Invalid fixed64 field"); + } + position[0] += 8; + return; + case 2: + final long length = readVarint(data, position, limit); + if (length < 0 || length > limit - position[0]) { + throw new IllegalArgumentException("Invalid length-delimited field"); + } + position[0] += (int) length; + return; + case 5: + if (limit - position[0] < 4) { + throw new IllegalArgumentException("Invalid fixed32 field"); + } + position[0] += 4; + return; + default: + throw new IllegalArgumentException("Unsupported protobuf wire type"); + } + } +} diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/TraceContextExtractAdapter.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/TraceContextExtractAdapter.java new file mode 100644 index 00000000000..d31e210635b --- /dev/null +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/TraceContextExtractAdapter.java @@ -0,0 +1,87 @@ +package datadog.trace.instrumentation.azure.functions.worker; + +import static datadog.trace.api.sampling.PrioritySampling.UNSET; + +import com.microsoft.azure.functions.TraceContext; +import datadog.trace.bootstrap.instrumentation.api.AgentPropagation; + +public final class TraceContextExtractAdapter + implements AgentPropagation.ContextVisitor { + public static final TraceContextExtractAdapter GETTER = new TraceContextExtractAdapter(); + + private static final String TRACEPARENT = "traceparent"; + private static final String TRACESTATE = "tracestate"; + + private TraceContextExtractAdapter() {} + + @Override + public void forEachKey(TraceContext carrier, AgentPropagation.KeyClassifier classifier) { + final String tracestate = carrier.getTracestate(); + String traceparent = carrier.getTraceparent(); + if (datadogSamplingPriority(tracestate) > 0) { + traceparent = setSampledFlag(traceparent); + } + if (traceparent != null && !classifier.accept(TRACEPARENT, traceparent)) { + return; + } + if (tracestate != null) { + classifier.accept(TRACESTATE, tracestate); + } + } + + static int datadogSamplingPriority(String tracestate) { + if (tracestate == null) { + return UNSET; + } + int memberStart = 0; + while (memberStart < tracestate.length()) { + int memberEnd = tracestate.indexOf(',', memberStart); + if (memberEnd < 0) { + memberEnd = tracestate.length(); + } + int start = memberStart; + while (start < memberEnd && tracestate.charAt(start) == ' ') { + start++; + } + if (memberEnd - start >= 3 + && tracestate.charAt(start) == 'd' + && tracestate.charAt(start + 1) == 'd' + && tracestate.charAt(start + 2) == '=') { + return parseSamplingPriority(tracestate, start + 3, memberEnd); + } + memberStart = memberEnd + 1; + } + return UNSET; + } + + private static int parseSamplingPriority(String value, int start, int end) { + while (start < end) { + int fieldEnd = value.indexOf(';', start); + if (fieldEnd < 0 || fieldEnd > end) { + fieldEnd = end; + } + if (fieldEnd - start > 2 && value.charAt(start) == 's' && value.charAt(start + 1) == ':') { + try { + return Integer.parseInt(value.substring(start + 2, fieldEnd)); + } catch (NumberFormatException ignored) { + return UNSET; + } + } + start = fieldEnd + 1; + } + return UNSET; + } + + private static String setSampledFlag(String traceparent) { + if (traceparent == null || traceparent.length() < 55) { + return traceparent; + } + final int high = Character.digit(traceparent.charAt(53), 16); + final int low = Character.digit(traceparent.charAt(54), 16); + if (high < 0 || low < 0 || (low & 1) != 0) { + return traceparent; + } + final char sampledLow = Character.forDigit(low | 1, 16); + return traceparent.substring(0, 54) + sampledLow + traceparent.substring(55); + } +} diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/groovy/AzureFunctionsWorkerTest.groovy b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/groovy/AzureFunctionsWorkerTest.groovy new file mode 100644 index 00000000000..37253f6a3b0 --- /dev/null +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/groovy/AzureFunctionsWorkerTest.groovy @@ -0,0 +1,357 @@ +import static datadog.trace.api.sampling.PrioritySampling.SAMPLER_KEEP +import static datadog.trace.api.sampling.PrioritySampling.UNSET +import static org.mockito.ArgumentMatchers.anyString +import static org.mockito.Mockito.doAnswer +import static org.mockito.Mockito.mock +import static org.mockito.Mockito.when + +import com.microsoft.azure.functions.TraceContext +import com.microsoft.azure.functions.internal.spi.middleware.MiddlewareChain +import com.microsoft.azure.functions.internal.spi.middleware.MiddlewareContext +import com.microsoft.azure.functions.worker.chain.FunctionExecutionMiddleware +import com.microsoft.durabletask.interruption.ContinueAsNewInterruption +import com.microsoft.durabletask.interruption.OrchestratorBlockedException +import datadog.trace.agent.test.naming.VersionedNamingTestBase +import datadog.trace.api.DDSpanTypes +import datadog.trace.bootstrap.instrumentation.api.AgentTracer +import datadog.trace.bootstrap.instrumentation.api.Tags +import spock.lang.Unroll + +abstract class AzureFunctionsWorkerTest extends VersionedNamingTestBase { + + @Override + String service() { + null + } + + @Unroll + def "creates a span for #trigger"() { + setup: + def context = contextFor(annotation, "MyFunction") + def chain = mock(MiddlewareChain) + + when: + new FunctionExecutionMiddleware().invoke(context, chain) + + then: + assertTraces(1) { + trace(1) { + span { + parent() + operationName operation() + resourceName "$trigger MyFunction" + spanType DDSpanTypes.SERVERLESS + errored false + tags { + defaultTags() + "$Tags.COMPONENT" "azure-functions" + "$Tags.SPAN_KIND" "$Tags.SPAN_KIND_SERVER" + "aas.function.name" "MyFunction" + "aas.function.trigger" trigger + } + } + } + } + + where: + annotation | trigger + "DurableOrchestrationTrigger" | "DurableOrchestration" + "DurableActivityTrigger" | "DurableActivity" + "DurableEntityTrigger" | "DurableEntity" + } + + def "does not create a span for a non-Durable function"() { + setup: + def context = contextFor(null, "HttpFunction") + + when: + new FunctionExecutionMiddleware().invoke(context, mock(MiddlewareChain)) + + then: + assertTraces(0) {} + } + + def "does not create a span for a successful orchestration replay"() { + setup: + def context = contextFor("DurableOrchestrationTrigger", "Orchestrator") + // OrchestratorRequest { pastEvents: {}, newEvents: HistoryEvent { taskCompleted: {} } } + when(context.getParameterValue("input")).thenReturn( + Base64.encoder.encodeToString([0x1a, 0x00, 0x22, 0x02, 0x3a, 0x00] as byte[])) + + when: + new FunctionExecutionMiddleware().invoke(context, mock(MiddlewareChain)) + + then: + assertTraces(0) {} + } + + def "does not treat a failure field with the wrong protobuf wire type as a failure"() { + setup: + def context = contextFor("DurableOrchestrationTrigger", "Orchestrator") + // OrchestratorRequest { pastEvents: {}, newEvents: HistoryEvent { field 8: varint 0 } } + when(context.getParameterValue("input")).thenReturn( + Base64.encoder.encodeToString([0x1a, 0x00, 0x22, 0x02, 0x40, 0x00] as byte[])) + + when: + new FunctionExecutionMiddleware().invoke(context, mock(MiddlewareChain)) + + then: + assertTraces(0) {} + } + + def "creates a span for an initial orchestration execution"() { + setup: + def context = contextFor("DurableOrchestrationTrigger", "Orchestrator") + // OrchestratorRequest { newEvents: HistoryEvent { executionStarted: {} } } + when(context.getParameterValue("input")).thenReturn( + Base64.encoder.encodeToString([0x22, 0x02, 0x1a, 0x00] as byte[])) + + when: + new FunctionExecutionMiddleware().invoke(context, mock(MiddlewareChain)) + + then: + assertTraces(1) { + trace(1) { + span { + operationName operation() + resourceName "DurableOrchestration Orchestrator" + spanType DDSpanTypes.SERVERLESS + errored false + tags { + defaultTags() + "$Tags.COMPONENT" "azure-functions" + "$Tags.SPAN_KIND" "$Tags.SPAN_KIND_SERVER" + "aas.function.name" "Orchestrator" + "aas.function.trigger" "DurableOrchestration" + } + } + } + } + } + + @Unroll + def "creates a span for an orchestration replay with #failure"() { + setup: + def context = contextFor("DurableOrchestrationTrigger", "Orchestrator") + // OrchestratorRequest { pastEvents: {}, newEvents: HistoryEvent { failure: {} } } + when(context.getParameterValue("input")).thenReturn( + Base64.encoder.encodeToString([0x1a, 0x00, 0x22, 0x02, failureTag, 0x00] as byte[])) + + when: + new FunctionExecutionMiddleware().invoke(context, mock(MiddlewareChain)) + + then: + assertTraces(1) { + trace(1) { + span { + operationName operation() + resourceName "DurableOrchestration Orchestrator" + spanType DDSpanTypes.SERVERLESS + errored false + tags { + defaultTags() + "$Tags.COMPONENT" "azure-functions" + "$Tags.SPAN_KIND" "$Tags.SPAN_KIND_SERVER" + "aas.function.name" "Orchestrator" + "aas.function.trigger" "DurableOrchestration" + } + } + } + } + + where: + failure | failureTag + "a failed activity" | 0x42 + "a failed sub-orchestration" | 0x5a + } + + @Unroll + def "fails open for an unreadable orchestration payload: #description"() { + setup: + def context = contextFor("DurableOrchestrationTrigger", "Orchestrator") + when(context.getParameterValue("input")).thenReturn(payload) + + when: + new FunctionExecutionMiddleware().invoke(context, mock(MiddlewareChain)) + + then: + assertTraces(1) { + trace(1) { + span { + operationName operation() + resourceName "DurableOrchestration Orchestrator" + spanType DDSpanTypes.SERVERLESS + errored false + tags { + defaultTags() + "$Tags.COMPONENT" "azure-functions" + "$Tags.SPAN_KIND" "$Tags.SPAN_KIND_SERVER" + "aas.function.name" "Orchestrator" + "aas.function.trigger" "DurableOrchestration" + } + } + } + } + + where: + description | payload + "non-base64 input" | "not base64" + "truncated protobuf" | Base64.encoder.encodeToString([0x1a, 0x02, 0x00] as byte[]) + "non-string input" | new Object() + } + + def "continues the Azure trace and ignores its implicit sampling rejection"() { + setup: + def context = contextFor("DurableActivityTrigger", "Activity") + def traceContext = mock(TraceContext) + when(traceContext.getTraceparent()).thenReturn( + "00-0000000000000000000000000000002a-000000000000002b-00") + when(traceContext.getTracestate()).thenReturn(null) + when(context.getTraceContext()).thenReturn(traceContext) + def priority = Integer.MIN_VALUE + def chain = mock(MiddlewareChain) + doAnswer { + priority = AgentTracer.activeSpan().spanContext().samplingPriority + null + }.when(chain).doNext(context) + + when: + new FunctionExecutionMiddleware().invoke(context, chain) + + then: + TEST_WRITER.waitForTraces(1) + def span = TEST_WRITER[0][0] + span.traceId.toString() == "42" + span.parentId == 43 + priority == UNSET + span.samplingPriority() == SAMPLER_KEEP + } + + def "honors a Datadog keep decision when Azure clears the W3C sampled flag"() { + setup: + def context = contextFor("DurableOrchestrationTrigger", "Orchestrator") + def traceContext = mock(TraceContext) + when(traceContext.getTraceparent()).thenReturn( + "00-0000000000000000000000000000002a-000000000000002b-00") + when(traceContext.getTracestate()).thenReturn("dd=s:2;o:rum") + when(context.getTraceContext()).thenReturn(traceContext) + def priority = Integer.MIN_VALUE + def chain = mock(MiddlewareChain) + doAnswer { + priority = AgentTracer.activeSpan().spanContext().samplingPriority + null + }.when(chain).doNext(context) + + when: + new FunctionExecutionMiddleware().invoke(context, chain) + + then: + TEST_WRITER.waitForTraces(1) + def span = TEST_WRITER[0][0] + span.traceId.toString() == "42" + span.parentId == 43 + priority > 0 + } + + @Unroll + def "does not mark #controlFlow.class.simpleName as an error"() { + setup: + def context = contextFor("DurableOrchestrationTrigger", "Orchestrator") + def chain = mock(MiddlewareChain) + doAnswer { throw new RuntimeException(controlFlow) }.when(chain).doNext(context) + + when: + new FunctionExecutionMiddleware().invoke(context, chain) + + then: + thrown(RuntimeException) + assertTraces(1) { + trace(1) { + span { + operationName operation() + resourceName "DurableOrchestration Orchestrator" + spanType DDSpanTypes.SERVERLESS + errored false + tags { + defaultTags() + "$Tags.COMPONENT" "azure-functions" + "$Tags.SPAN_KIND" "$Tags.SPAN_KIND_SERVER" + "aas.function.name" "Orchestrator" + "aas.function.trigger" "DurableOrchestration" + } + } + } + } + + where: + controlFlow << [ + new com.microsoft.durabletask.OrchestratorBlockedException(), + new OrchestratorBlockedException(), + new ContinueAsNewInterruption() + ] + } + + def "marks application failures as errors"() { + setup: + def context = contextFor("DurableActivityTrigger", "Activity") + def chain = mock(MiddlewareChain) + doAnswer { throw new IllegalStateException("failure") }.when(chain).doNext(context) + + when: + new FunctionExecutionMiddleware().invoke(context, chain) + + then: + thrown(IllegalStateException) + assertTraces(1) { + trace(1) { + span { + operationName operation() + resourceName "DurableActivity Activity" + spanType DDSpanTypes.SERVERLESS + errored true + tags { + errorTags IllegalStateException, "failure" + defaultTags() + "$Tags.COMPONENT" "azure-functions" + "$Tags.SPAN_KIND" "$Tags.SPAN_KIND_SERVER" + "aas.function.name" "Activity" + "aas.function.trigger" "DurableActivity" + } + } + } + } + } + + private static MiddlewareContext contextFor(String annotation, String functionName) { + def context = mock(MiddlewareContext) + when(context.getFunctionName()).thenReturn(functionName) + when(context.getParameterName(anyString())).thenAnswer { + it.arguments[0] == annotation ? "input" : null + } + context + } +} + +class AzureFunctionsWorkerV0ForkedTest extends AzureFunctionsWorkerTest { + @Override + int version() { + 0 + } + + @Override + String operation() { + "dd-tracer-serverless-span" + } +} + +class AzureFunctionsWorkerV1Test extends AzureFunctionsWorkerTest { + @Override + int version() { + 1 + } + + @Override + String operation() { + "azure.functions.invoke" + } +} diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/com/microsoft/azure/functions/worker/chain/FunctionExecutionMiddleware.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/com/microsoft/azure/functions/worker/chain/FunctionExecutionMiddleware.java new file mode 100644 index 00000000000..1cda6f23f27 --- /dev/null +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/com/microsoft/azure/functions/worker/chain/FunctionExecutionMiddleware.java @@ -0,0 +1,10 @@ +package com.microsoft.azure.functions.worker.chain; + +import com.microsoft.azure.functions.internal.spi.middleware.MiddlewareChain; +import com.microsoft.azure.functions.internal.spi.middleware.MiddlewareContext; + +public class FunctionExecutionMiddleware { + public void invoke(MiddlewareContext context, MiddlewareChain chain) throws Exception { + chain.doNext(context); + } +} diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/com/microsoft/durabletask/OrchestratorBlockedException.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/com/microsoft/durabletask/OrchestratorBlockedException.java new file mode 100644 index 00000000000..d8839fdf084 --- /dev/null +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/com/microsoft/durabletask/OrchestratorBlockedException.java @@ -0,0 +1,3 @@ +package com.microsoft.durabletask; + +public class OrchestratorBlockedException extends RuntimeException {} diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/com/microsoft/durabletask/interruption/ContinueAsNewInterruption.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/com/microsoft/durabletask/interruption/ContinueAsNewInterruption.java new file mode 100644 index 00000000000..e540e7519e0 --- /dev/null +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/com/microsoft/durabletask/interruption/ContinueAsNewInterruption.java @@ -0,0 +1,3 @@ +package com.microsoft.durabletask.interruption; + +public class ContinueAsNewInterruption extends RuntimeException {} diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/com/microsoft/durabletask/interruption/OrchestratorBlockedException.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/com/microsoft/durabletask/interruption/OrchestratorBlockedException.java new file mode 100644 index 00000000000..681aa96c4bd --- /dev/null +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/com/microsoft/durabletask/interruption/OrchestratorBlockedException.java @@ -0,0 +1,3 @@ +package com.microsoft.durabletask.interruption; + +public class OrchestratorBlockedException extends RuntimeException {} diff --git a/dd-trace-core/src/main/java/datadog/trace/core/propagation/ExtractedContext.java b/dd-trace-core/src/main/java/datadog/trace/core/propagation/ExtractedContext.java index af503e6a6ed..7c48ac9e452 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/propagation/ExtractedContext.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/propagation/ExtractedContext.java @@ -1,5 +1,7 @@ package datadog.trace.core.propagation; +import static datadog.trace.api.sampling.SamplingMechanism.EXTERNAL_OVERRIDE; + import datadog.trace.api.DDTraceId; import datadog.trace.api.TagMap; import datadog.trace.api.TraceConfig; @@ -117,6 +119,15 @@ public PropagationTags getPropagationTags() { return propagationTags; } + @Override + public ExtractedContext withSamplingPriority(final int samplingPriority) { + setSamplingPriority(samplingPriority); + if (propagationTags != null) { + propagationTags.updateTraceSamplingPriority(samplingPriority, EXTERNAL_OVERRIDE); + } + return this; + } + @Override public String toString() { StringBuilder builder = new StringBuilder("ExtractedContext{"); diff --git a/dd-trace-core/src/test/java/datadog/trace/core/propagation/ExtractedContextTest.java b/dd-trace-core/src/test/java/datadog/trace/core/propagation/ExtractedContextTest.java new file mode 100644 index 00000000000..130d5adef33 --- /dev/null +++ b/dd-trace-core/src/test/java/datadog/trace/core/propagation/ExtractedContextTest.java @@ -0,0 +1,40 @@ +package datadog.trace.core.propagation; + +import static datadog.trace.api.TracePropagationStyle.TRACECONTEXT; +import static datadog.trace.api.sampling.PrioritySampling.SAMPLER_DROP; +import static datadog.trace.api.sampling.PrioritySampling.UNSET; +import static datadog.trace.api.sampling.SamplingMechanism.EXTERNAL_OVERRIDE; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertSame; + +import datadog.trace.api.DDTraceId; +import org.junit.jupiter.api.Test; + +class ExtractedContextTest { + @Test + void replacesSamplingPriorityWithoutChangingTraceIdentity() { + PropagationTags propagationTags = PropagationTags.factory().empty(); + propagationTags.updateTraceSamplingPriority(SAMPLER_DROP, EXTERNAL_OVERRIDE); + ExtractedContext original = + new ExtractedContext( + DDTraceId.from(42), 43, SAMPLER_DROP, null, propagationTags, TRACECONTEXT); + + ExtractedContext updated = original.withSamplingPriority(UNSET); + + assertSame(original, updated); + assertEquals(original.getTraceId(), updated.getTraceId()); + assertEquals(original.getSpanId(), updated.getSpanId()); + assertEquals(UNSET, updated.getSamplingPriority()); + assertEquals(UNSET, updated.getPropagationTags().getSamplingPriority()); + } + + @Test + void replacesSamplingPriorityWithoutPropagationTags() { + ExtractedContext context = + new ExtractedContext(DDTraceId.from(42), 43, SAMPLER_DROP, null, null, TRACECONTEXT); + + context.withSamplingPriority(UNSET); + + assertEquals(UNSET, context.getSamplingPriority()); + } +} diff --git a/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/AgentSpanContext.java b/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/AgentSpanContext.java index 6e950925305..fcde079c1ea 100644 --- a/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/AgentSpanContext.java +++ b/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/AgentSpanContext.java @@ -72,6 +72,14 @@ default CharSequence getIntegrationName() { boolean isRemote(); interface Extracted extends AgentSpanContext { + /** + * Returns this extracted context with the supplied sampling priority when the implementation + * supports replacing propagation decisions. + */ + default Extracted withSamplingPriority(int samplingPriority) { + return this; + } + /** * Gets the span links related to the other terminated context. * diff --git a/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/TagContext.java b/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/TagContext.java index 078ffeb1625..f7098e94aa2 100644 --- a/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/TagContext.java +++ b/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/TagContext.java @@ -33,7 +33,7 @@ public class TagContext implements AgentSpanContext.Extracted { private final HttpHeaders httpHeaders; private final Map baggage; private Baggage w3cBaggage; - private final int samplingPriority; + private int samplingPriority; private final TraceConfig traceConfig; private final TracePropagationStyle propagationStyle; private final DDTraceId traceId; @@ -181,6 +181,10 @@ public final int getSamplingPriority() { return samplingPriority; } + protected final void setSamplingPriority(int samplingPriority) { + this.samplingPriority = samplingPriority; + } + public final Map getBaggage() { return baggage; } diff --git a/settings.gradle.kts b/settings.gradle.kts index e557332879d..182dfc511fc 100644 --- a/settings.gradle.kts +++ b/settings.gradle.kts @@ -322,6 +322,7 @@ include( ":dd-java-agent:instrumentation:axis2-1.3", ":dd-java-agent:instrumentation:axway-api-7.5", ":dd-java-agent:instrumentation:azure-functions-1.2.2", + ":dd-java-agent:instrumentation:azure-functions-worker-2.7", ":dd-java-agent:instrumentation:beanshell-2.0", ":dd-java-agent:instrumentation:caffeine-1.0", ":dd-java-agent:instrumentation:cdi-1.2", From f6cb9cf240f586a01fed7291b2909a78976ea3ab Mon Sep 17 00:00:00 2001 From: jcstorms1 Date: Wed, 23 Sep 2026 13:09:09 -0400 Subject: [PATCH 02/10] Add Azure Functions worker code owner --- .github/CODEOWNERS | 1 + 1 file changed, 1 insertion(+) diff --git a/.github/CODEOWNERS b/.github/CODEOWNERS index d61754b6290..c5c661a393f 100644 --- a/.github/CODEOWNERS +++ b/.github/CODEOWNERS @@ -172,6 +172,7 @@ /dd-java-agent/instrumentation/aws-java/aws-java-lambda-handler-1.2/ @DataDog/apm-serverless /dd-java-agent/instrumentation/azure-functions/ @DataDog/apm-serverless /dd-java-agent/instrumentation/azure-functions-1.2.2/ @DataDog/apm-serverless +/dd-java-agent/instrumentation/azure-functions-worker-2.7/ @DataDog/apm-serverless # @DataDog/apm-lang-platform-java /.editorconfig @DataDog/apm-lang-platform-java From 376c24d8510fef4f062deedbabd056132182f1aa Mon Sep 17 00:00:00 2001 From: jcstorms1 Date: Wed, 23 Sep 2026 13:39:55 -0400 Subject: [PATCH 03/10] Convert Azure Functions worker tests to JUnit --- .../groovy/AzureFunctionsWorkerTest.groovy | 357 ------------------ .../test/java/AzureFunctionsWorkerTest.java | 290 ++++++++++++++ .../AzureFunctionsWorkerV0ForkedTest.java | 12 + .../test/java/AzureFunctionsWorkerV1Test.java | 12 + 4 files changed, 314 insertions(+), 357 deletions(-) delete mode 100644 dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/groovy/AzureFunctionsWorkerTest.groovy create mode 100644 dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java create mode 100644 dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerV0ForkedTest.java create mode 100644 dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerV1Test.java diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/groovy/AzureFunctionsWorkerTest.groovy b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/groovy/AzureFunctionsWorkerTest.groovy deleted file mode 100644 index 37253f6a3b0..00000000000 --- a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/groovy/AzureFunctionsWorkerTest.groovy +++ /dev/null @@ -1,357 +0,0 @@ -import static datadog.trace.api.sampling.PrioritySampling.SAMPLER_KEEP -import static datadog.trace.api.sampling.PrioritySampling.UNSET -import static org.mockito.ArgumentMatchers.anyString -import static org.mockito.Mockito.doAnswer -import static org.mockito.Mockito.mock -import static org.mockito.Mockito.when - -import com.microsoft.azure.functions.TraceContext -import com.microsoft.azure.functions.internal.spi.middleware.MiddlewareChain -import com.microsoft.azure.functions.internal.spi.middleware.MiddlewareContext -import com.microsoft.azure.functions.worker.chain.FunctionExecutionMiddleware -import com.microsoft.durabletask.interruption.ContinueAsNewInterruption -import com.microsoft.durabletask.interruption.OrchestratorBlockedException -import datadog.trace.agent.test.naming.VersionedNamingTestBase -import datadog.trace.api.DDSpanTypes -import datadog.trace.bootstrap.instrumentation.api.AgentTracer -import datadog.trace.bootstrap.instrumentation.api.Tags -import spock.lang.Unroll - -abstract class AzureFunctionsWorkerTest extends VersionedNamingTestBase { - - @Override - String service() { - null - } - - @Unroll - def "creates a span for #trigger"() { - setup: - def context = contextFor(annotation, "MyFunction") - def chain = mock(MiddlewareChain) - - when: - new FunctionExecutionMiddleware().invoke(context, chain) - - then: - assertTraces(1) { - trace(1) { - span { - parent() - operationName operation() - resourceName "$trigger MyFunction" - spanType DDSpanTypes.SERVERLESS - errored false - tags { - defaultTags() - "$Tags.COMPONENT" "azure-functions" - "$Tags.SPAN_KIND" "$Tags.SPAN_KIND_SERVER" - "aas.function.name" "MyFunction" - "aas.function.trigger" trigger - } - } - } - } - - where: - annotation | trigger - "DurableOrchestrationTrigger" | "DurableOrchestration" - "DurableActivityTrigger" | "DurableActivity" - "DurableEntityTrigger" | "DurableEntity" - } - - def "does not create a span for a non-Durable function"() { - setup: - def context = contextFor(null, "HttpFunction") - - when: - new FunctionExecutionMiddleware().invoke(context, mock(MiddlewareChain)) - - then: - assertTraces(0) {} - } - - def "does not create a span for a successful orchestration replay"() { - setup: - def context = contextFor("DurableOrchestrationTrigger", "Orchestrator") - // OrchestratorRequest { pastEvents: {}, newEvents: HistoryEvent { taskCompleted: {} } } - when(context.getParameterValue("input")).thenReturn( - Base64.encoder.encodeToString([0x1a, 0x00, 0x22, 0x02, 0x3a, 0x00] as byte[])) - - when: - new FunctionExecutionMiddleware().invoke(context, mock(MiddlewareChain)) - - then: - assertTraces(0) {} - } - - def "does not treat a failure field with the wrong protobuf wire type as a failure"() { - setup: - def context = contextFor("DurableOrchestrationTrigger", "Orchestrator") - // OrchestratorRequest { pastEvents: {}, newEvents: HistoryEvent { field 8: varint 0 } } - when(context.getParameterValue("input")).thenReturn( - Base64.encoder.encodeToString([0x1a, 0x00, 0x22, 0x02, 0x40, 0x00] as byte[])) - - when: - new FunctionExecutionMiddleware().invoke(context, mock(MiddlewareChain)) - - then: - assertTraces(0) {} - } - - def "creates a span for an initial orchestration execution"() { - setup: - def context = contextFor("DurableOrchestrationTrigger", "Orchestrator") - // OrchestratorRequest { newEvents: HistoryEvent { executionStarted: {} } } - when(context.getParameterValue("input")).thenReturn( - Base64.encoder.encodeToString([0x22, 0x02, 0x1a, 0x00] as byte[])) - - when: - new FunctionExecutionMiddleware().invoke(context, mock(MiddlewareChain)) - - then: - assertTraces(1) { - trace(1) { - span { - operationName operation() - resourceName "DurableOrchestration Orchestrator" - spanType DDSpanTypes.SERVERLESS - errored false - tags { - defaultTags() - "$Tags.COMPONENT" "azure-functions" - "$Tags.SPAN_KIND" "$Tags.SPAN_KIND_SERVER" - "aas.function.name" "Orchestrator" - "aas.function.trigger" "DurableOrchestration" - } - } - } - } - } - - @Unroll - def "creates a span for an orchestration replay with #failure"() { - setup: - def context = contextFor("DurableOrchestrationTrigger", "Orchestrator") - // OrchestratorRequest { pastEvents: {}, newEvents: HistoryEvent { failure: {} } } - when(context.getParameterValue("input")).thenReturn( - Base64.encoder.encodeToString([0x1a, 0x00, 0x22, 0x02, failureTag, 0x00] as byte[])) - - when: - new FunctionExecutionMiddleware().invoke(context, mock(MiddlewareChain)) - - then: - assertTraces(1) { - trace(1) { - span { - operationName operation() - resourceName "DurableOrchestration Orchestrator" - spanType DDSpanTypes.SERVERLESS - errored false - tags { - defaultTags() - "$Tags.COMPONENT" "azure-functions" - "$Tags.SPAN_KIND" "$Tags.SPAN_KIND_SERVER" - "aas.function.name" "Orchestrator" - "aas.function.trigger" "DurableOrchestration" - } - } - } - } - - where: - failure | failureTag - "a failed activity" | 0x42 - "a failed sub-orchestration" | 0x5a - } - - @Unroll - def "fails open for an unreadable orchestration payload: #description"() { - setup: - def context = contextFor("DurableOrchestrationTrigger", "Orchestrator") - when(context.getParameterValue("input")).thenReturn(payload) - - when: - new FunctionExecutionMiddleware().invoke(context, mock(MiddlewareChain)) - - then: - assertTraces(1) { - trace(1) { - span { - operationName operation() - resourceName "DurableOrchestration Orchestrator" - spanType DDSpanTypes.SERVERLESS - errored false - tags { - defaultTags() - "$Tags.COMPONENT" "azure-functions" - "$Tags.SPAN_KIND" "$Tags.SPAN_KIND_SERVER" - "aas.function.name" "Orchestrator" - "aas.function.trigger" "DurableOrchestration" - } - } - } - } - - where: - description | payload - "non-base64 input" | "not base64" - "truncated protobuf" | Base64.encoder.encodeToString([0x1a, 0x02, 0x00] as byte[]) - "non-string input" | new Object() - } - - def "continues the Azure trace and ignores its implicit sampling rejection"() { - setup: - def context = contextFor("DurableActivityTrigger", "Activity") - def traceContext = mock(TraceContext) - when(traceContext.getTraceparent()).thenReturn( - "00-0000000000000000000000000000002a-000000000000002b-00") - when(traceContext.getTracestate()).thenReturn(null) - when(context.getTraceContext()).thenReturn(traceContext) - def priority = Integer.MIN_VALUE - def chain = mock(MiddlewareChain) - doAnswer { - priority = AgentTracer.activeSpan().spanContext().samplingPriority - null - }.when(chain).doNext(context) - - when: - new FunctionExecutionMiddleware().invoke(context, chain) - - then: - TEST_WRITER.waitForTraces(1) - def span = TEST_WRITER[0][0] - span.traceId.toString() == "42" - span.parentId == 43 - priority == UNSET - span.samplingPriority() == SAMPLER_KEEP - } - - def "honors a Datadog keep decision when Azure clears the W3C sampled flag"() { - setup: - def context = contextFor("DurableOrchestrationTrigger", "Orchestrator") - def traceContext = mock(TraceContext) - when(traceContext.getTraceparent()).thenReturn( - "00-0000000000000000000000000000002a-000000000000002b-00") - when(traceContext.getTracestate()).thenReturn("dd=s:2;o:rum") - when(context.getTraceContext()).thenReturn(traceContext) - def priority = Integer.MIN_VALUE - def chain = mock(MiddlewareChain) - doAnswer { - priority = AgentTracer.activeSpan().spanContext().samplingPriority - null - }.when(chain).doNext(context) - - when: - new FunctionExecutionMiddleware().invoke(context, chain) - - then: - TEST_WRITER.waitForTraces(1) - def span = TEST_WRITER[0][0] - span.traceId.toString() == "42" - span.parentId == 43 - priority > 0 - } - - @Unroll - def "does not mark #controlFlow.class.simpleName as an error"() { - setup: - def context = contextFor("DurableOrchestrationTrigger", "Orchestrator") - def chain = mock(MiddlewareChain) - doAnswer { throw new RuntimeException(controlFlow) }.when(chain).doNext(context) - - when: - new FunctionExecutionMiddleware().invoke(context, chain) - - then: - thrown(RuntimeException) - assertTraces(1) { - trace(1) { - span { - operationName operation() - resourceName "DurableOrchestration Orchestrator" - spanType DDSpanTypes.SERVERLESS - errored false - tags { - defaultTags() - "$Tags.COMPONENT" "azure-functions" - "$Tags.SPAN_KIND" "$Tags.SPAN_KIND_SERVER" - "aas.function.name" "Orchestrator" - "aas.function.trigger" "DurableOrchestration" - } - } - } - } - - where: - controlFlow << [ - new com.microsoft.durabletask.OrchestratorBlockedException(), - new OrchestratorBlockedException(), - new ContinueAsNewInterruption() - ] - } - - def "marks application failures as errors"() { - setup: - def context = contextFor("DurableActivityTrigger", "Activity") - def chain = mock(MiddlewareChain) - doAnswer { throw new IllegalStateException("failure") }.when(chain).doNext(context) - - when: - new FunctionExecutionMiddleware().invoke(context, chain) - - then: - thrown(IllegalStateException) - assertTraces(1) { - trace(1) { - span { - operationName operation() - resourceName "DurableActivity Activity" - spanType DDSpanTypes.SERVERLESS - errored true - tags { - errorTags IllegalStateException, "failure" - defaultTags() - "$Tags.COMPONENT" "azure-functions" - "$Tags.SPAN_KIND" "$Tags.SPAN_KIND_SERVER" - "aas.function.name" "Activity" - "aas.function.trigger" "DurableActivity" - } - } - } - } - } - - private static MiddlewareContext contextFor(String annotation, String functionName) { - def context = mock(MiddlewareContext) - when(context.getFunctionName()).thenReturn(functionName) - when(context.getParameterName(anyString())).thenAnswer { - it.arguments[0] == annotation ? "input" : null - } - context - } -} - -class AzureFunctionsWorkerV0ForkedTest extends AzureFunctionsWorkerTest { - @Override - int version() { - 0 - } - - @Override - String operation() { - "dd-tracer-serverless-span" - } -} - -class AzureFunctionsWorkerV1Test extends AzureFunctionsWorkerTest { - @Override - int version() { - 1 - } - - @Override - String operation() { - "azure.functions.invoke" - } -} diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java new file mode 100644 index 00000000000..c3b79300cab --- /dev/null +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java @@ -0,0 +1,290 @@ +import static datadog.trace.agent.test.assertions.SpanMatcher.span; +import static datadog.trace.agent.test.assertions.TagsMatcher.defaultTags; +import static datadog.trace.agent.test.assertions.TagsMatcher.error; +import static datadog.trace.agent.test.assertions.TagsMatcher.tag; +import static datadog.trace.agent.test.assertions.TraceMatcher.trace; +import static datadog.trace.api.sampling.PrioritySampling.SAMPLER_KEEP; +import static datadog.trace.api.sampling.PrioritySampling.UNSET; +import static datadog.trace.test.junit.utils.assertions.Matchers.is; +import static datadog.trace.test.junit.utils.assertions.Matchers.matches; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.junit.jupiter.params.provider.Arguments.arguments; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +import com.microsoft.azure.functions.TraceContext; +import com.microsoft.azure.functions.internal.spi.middleware.MiddlewareChain; +import com.microsoft.azure.functions.internal.spi.middleware.MiddlewareContext; +import com.microsoft.azure.functions.worker.chain.FunctionExecutionMiddleware; +import com.microsoft.durabletask.interruption.ContinueAsNewInterruption; +import com.microsoft.durabletask.interruption.OrchestratorBlockedException; +import datadog.trace.agent.test.AbstractInstrumentationTest; +import datadog.trace.agent.test.assertions.SpanMatcher; +import datadog.trace.api.DDSpanTypes; +import datadog.trace.bootstrap.instrumentation.api.AgentTracer; +import datadog.trace.bootstrap.instrumentation.api.Tags; +import datadog.trace.core.DDSpan; +import java.util.Base64; +import java.util.Objects; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.regex.Pattern; +import java.util.stream.Stream; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; +import org.tabletest.junit.TableTest; + +abstract class AzureFunctionsWorkerTest extends AbstractInstrumentationTest { + + abstract String operation(); + + // spotless:off + @TableTest({ + "scenario | annotation | trigger", + "orchestration | DurableOrchestrationTrigger | DurableOrchestration", + "activity | DurableActivityTrigger | DurableActivity", + "entity | DurableEntityTrigger | DurableEntity" + }) + // spotless:on + void createsSpanForDurableTrigger(String annotation, String trigger) throws Exception { + MiddlewareContext context = contextFor(annotation, "MyFunction"); + + new FunctionExecutionMiddleware().invoke(context, mock(MiddlewareChain.class)); + + assertTraces(trace(durableSpan("MyFunction", trigger))); + } + + @Test + void doesNotCreateSpanForNonDurableFunction() throws Exception { + MiddlewareContext context = contextFor(null, "HttpFunction"); + + new FunctionExecutionMiddleware().invoke(context, mock(MiddlewareChain.class)); + + assertTraces(); + } + + @Test + void doesNotCreateSpanForSuccessfulOrchestrationReplay() throws Exception { + MiddlewareContext context = contextFor("DurableOrchestrationTrigger", "Orchestrator"); + // OrchestratorRequest { pastEvents: {}, newEvents: HistoryEvent { taskCompleted: {} } } + when(context.getParameterValue("input")) + .thenReturn( + Base64.getEncoder().encodeToString(new byte[] {0x1a, 0x00, 0x22, 0x02, 0x3a, 0x00})); + + new FunctionExecutionMiddleware().invoke(context, mock(MiddlewareChain.class)); + + assertTraces(); + } + + @Test + void doesNotTreatFailureFieldWithWrongProtobufWireTypeAsFailure() throws Exception { + MiddlewareContext context = contextFor("DurableOrchestrationTrigger", "Orchestrator"); + // OrchestratorRequest { pastEvents: {}, newEvents: HistoryEvent { field 8: varint 0 } } + when(context.getParameterValue("input")) + .thenReturn( + Base64.getEncoder().encodeToString(new byte[] {0x1a, 0x00, 0x22, 0x02, 0x40, 0x00})); + + new FunctionExecutionMiddleware().invoke(context, mock(MiddlewareChain.class)); + + assertTraces(); + } + + @Test + void createsSpanForInitialOrchestrationExecution() throws Exception { + MiddlewareContext context = contextFor("DurableOrchestrationTrigger", "Orchestrator"); + // OrchestratorRequest { newEvents: HistoryEvent { executionStarted: {} } } + when(context.getParameterValue("input")) + .thenReturn(Base64.getEncoder().encodeToString(new byte[] {0x22, 0x02, 0x1a, 0x00})); + + new FunctionExecutionMiddleware().invoke(context, mock(MiddlewareChain.class)); + + assertTraces(trace(durableSpan("Orchestrator", "DurableOrchestration"))); + } + + // spotless:off + @TableTest({ + "scenario | failureTag", + "failed activity | 66", + "failed sub-orchestration | 90" + }) + // spotless:on + void createsSpanForOrchestrationReplayWithFailure(int failureTag) throws Exception { + MiddlewareContext context = contextFor("DurableOrchestrationTrigger", "Orchestrator"); + // OrchestratorRequest { pastEvents: {}, newEvents: HistoryEvent { failure: {} } } + when(context.getParameterValue("input")) + .thenReturn( + Base64.getEncoder() + .encodeToString(new byte[] {0x1a, 0x00, 0x22, 0x02, (byte) failureTag, 0x00})); + + new FunctionExecutionMiddleware().invoke(context, mock(MiddlewareChain.class)); + + assertTraces(trace(durableSpan("Orchestrator", "DurableOrchestration"))); + } + + @ParameterizedTest(name = "{0}") + // spotless:off + @TableTest({ + "scenario | description | payload", + "non-base64 input | non-base64 input | 'not base64'", + "truncated protobuf | truncated protobuf | GgIA" + }) + // spotless:on + @MethodSource("failsOpenForUnreadableOrchestrationPayloadArguments") + void failsOpenForUnreadableOrchestrationPayload(String description, Object payload) + throws Exception { + MiddlewareContext context = contextFor("DurableOrchestrationTrigger", "Orchestrator"); + when(context.getParameterValue("input")).thenReturn(payload); + + new FunctionExecutionMiddleware().invoke(context, mock(MiddlewareChain.class)); + + assertTraces(trace(durableSpan("Orchestrator", "DurableOrchestration"))); + } + + static Stream failsOpenForUnreadableOrchestrationPayloadArguments() { + return Stream.of(arguments("non-string input", new Object())); + } + + @Test + void continuesAzureTraceAndIgnoresImplicitSamplingRejection() throws Exception { + MiddlewareContext context = contextFor("DurableActivityTrigger", "Activity"); + TraceContext traceContext = mock(TraceContext.class); + when(traceContext.getTraceparent()) + .thenReturn("00-0000000000000000000000000000002a-000000000000002b-00"); + when(traceContext.getTracestate()).thenReturn(null); + when(context.getTraceContext()).thenReturn(traceContext); + AtomicInteger priority = new AtomicInteger(Integer.MIN_VALUE); + MiddlewareChain chain = mock(MiddlewareChain.class); + doAnswer( + invocation -> { + priority.set(AgentTracer.activeSpan().spanContext().getSamplingPriority()); + return null; + }) + .when(chain) + .doNext(context); + + new FunctionExecutionMiddleware().invoke(context, chain); + + writer.waitForTraces(1); + DDSpan span = writer.firstTrace().get(0); + assertEquals("42", span.getTraceId().toString()); + assertEquals(43L, span.getParentId()); + assertEquals(UNSET, priority.get()); + assertEquals(SAMPLER_KEEP, span.samplingPriority()); + } + + @Test + void honorsDatadogKeepDecisionWhenAzureClearsW3cSampledFlag() throws Exception { + MiddlewareContext context = contextFor("DurableOrchestrationTrigger", "Orchestrator"); + TraceContext traceContext = mock(TraceContext.class); + when(traceContext.getTraceparent()) + .thenReturn("00-0000000000000000000000000000002a-000000000000002b-00"); + when(traceContext.getTracestate()).thenReturn("dd=s:2;o:rum"); + when(context.getTraceContext()).thenReturn(traceContext); + AtomicInteger priority = new AtomicInteger(Integer.MIN_VALUE); + MiddlewareChain chain = mock(MiddlewareChain.class); + doAnswer( + invocation -> { + priority.set(AgentTracer.activeSpan().spanContext().getSamplingPriority()); + return null; + }) + .when(chain) + .doNext(context); + + new FunctionExecutionMiddleware().invoke(context, chain); + + writer.waitForTraces(1); + DDSpan span = writer.firstTrace().get(0); + assertEquals("42", span.getTraceId().toString()); + assertEquals(43L, span.getParentId()); + assertTrue(priority.get() > 0); + } + + @ParameterizedTest(name = "{0}") + @MethodSource("doesNotMarkReplayControlFlowAsErrorArguments") + void doesNotMarkReplayControlFlowAsError(String description, Throwable controlFlow) + throws Exception { + MiddlewareContext context = contextFor("DurableOrchestrationTrigger", "Orchestrator"); + MiddlewareChain chain = mock(MiddlewareChain.class); + doAnswer( + invocation -> { + throw new RuntimeException(controlFlow); + }) + .when(chain) + .doNext(context); + + assertThrows( + RuntimeException.class, () -> new FunctionExecutionMiddleware().invoke(context, chain)); + + assertTraces(trace(durableSpan("Orchestrator", "DurableOrchestration"))); + } + + static Stream doesNotMarkReplayControlFlowAsErrorArguments() { + return Stream.of( + arguments( + "legacy OrchestratorBlockedException", + new com.microsoft.durabletask.OrchestratorBlockedException()), + arguments("OrchestratorBlockedException", new OrchestratorBlockedException()), + arguments("ContinueAsNewInterruption", new ContinueAsNewInterruption())); + } + + @Test + void marksApplicationFailuresAsErrors() throws Exception { + MiddlewareContext context = contextFor("DurableActivityTrigger", "Activity"); + MiddlewareChain chain = mock(MiddlewareChain.class); + doAnswer( + invocation -> { + throw new IllegalStateException("failure"); + }) + .when(chain) + .doNext(context); + + assertThrows( + IllegalStateException.class, + () -> new FunctionExecutionMiddleware().invoke(context, chain)); + + assertTraces( + trace( + span() + .root() + .operationName(Pattern.compile(Pattern.quote(operation()))) + .resourceName("DurableActivity Activity") + .type(DDSpanTypes.SERVERLESS) + .error(true) + .tags( + defaultTags(), + error(IllegalStateException.class, "failure"), + tag(Tags.COMPONENT, matches(Pattern.quote("azure-functions"))), + tag(Tags.SPAN_KIND, is(Tags.SPAN_KIND_SERVER)), + tag("aas.function.name", is("Activity")), + tag("aas.function.trigger", is("DurableActivity"))))); + } + + private SpanMatcher durableSpan(String functionName, String trigger) { + return span() + .root() + .operationName(Pattern.compile(Pattern.quote(operation()))) + .resourceName(trigger + " " + functionName) + .type(DDSpanTypes.SERVERLESS) + .error(false) + .tags( + defaultTags(), + tag(Tags.COMPONENT, matches(Pattern.quote("azure-functions"))), + tag(Tags.SPAN_KIND, is(Tags.SPAN_KIND_SERVER)), + tag("aas.function.name", is(functionName)), + tag("aas.function.trigger", is(trigger))); + } + + private static MiddlewareContext contextFor(String annotation, String functionName) { + MiddlewareContext context = mock(MiddlewareContext.class); + when(context.getFunctionName()).thenReturn(functionName); + when(context.getParameterName(anyString())) + .thenAnswer( + invocation -> Objects.equals(invocation.getArgument(0), annotation) ? "input" : null); + return context; + } +} diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerV0ForkedTest.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerV0ForkedTest.java new file mode 100644 index 00000000000..e3d36f11ea3 --- /dev/null +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerV0ForkedTest.java @@ -0,0 +1,12 @@ +import static datadog.trace.api.config.TracerConfig.TRACE_SPAN_ATTRIBUTE_SCHEMA; + +import datadog.trace.test.junit.utils.config.WithConfig; + +@WithConfig(key = TRACE_SPAN_ATTRIBUTE_SCHEMA, value = "v0") +class AzureFunctionsWorkerV0ForkedTest extends AzureFunctionsWorkerTest { + + @Override + String operation() { + return "dd-tracer-serverless-span"; + } +} diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerV1Test.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerV1Test.java new file mode 100644 index 00000000000..ce27c2c15ca --- /dev/null +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerV1Test.java @@ -0,0 +1,12 @@ +import static datadog.trace.api.config.TracerConfig.TRACE_SPAN_ATTRIBUTE_SCHEMA; + +import datadog.trace.test.junit.utils.config.WithConfig; + +@WithConfig(key = TRACE_SPAN_ATTRIBUTE_SCHEMA, value = "v1") +class AzureFunctionsWorkerV1Test extends AzureFunctionsWorkerTest { + + @Override + String operation() { + return "azure.functions.invoke"; + } +} From 12dec943396a3b6a28afb95a9b8d5e6e7899168b Mon Sep 17 00:00:00 2001 From: jcstorms1 Date: Wed, 23 Sep 2026 15:21:32 -0400 Subject: [PATCH 04/10] Fix Durable Functions replay edge cases --- .../AzureFunctionsWorkerInstrumentation.java | 38 ++--- .../worker/Base64StringInputStream.java | 34 ++++ .../worker/DurableFunctionsUtils.java | 145 ++++++++++++++---- .../worker/TraceContextExtractAdapter.java | 14 +- .../test/java/AzureFunctionsWorkerTest.java | 58 ++++++- .../api/AgentSpanContextTest.java | 15 ++ 6 files changed, 247 insertions(+), 57 deletions(-) create mode 100644 dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/Base64StringInputStream.java create mode 100644 internal-api/src/test/java/datadog/trace/bootstrap/instrumentation/api/AgentSpanContextTest.java diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java index 8c2332ef049..b37e9e335cb 100644 --- a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java @@ -1,11 +1,7 @@ package datadog.trace.instrumentation.azure.functions.worker; import static datadog.trace.agent.tooling.bytebuddy.matcher.NameMatchers.named; -import static datadog.trace.bootstrap.instrumentation.api.AgentPropagation.extractContextAndGetSpanContext; -import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activateSpan; -import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.startSpan; import static datadog.trace.bootstrap.instrumentation.api.Java8BytecodeBridge.spanFromScope; -import static datadog.trace.instrumentation.azure.functions.worker.DurableFunctionsDecorator.AZURE_FUNCTIONS_REQUEST; import static datadog.trace.instrumentation.azure.functions.worker.DurableFunctionsDecorator.DECORATE; import static net.bytebuddy.matcher.ElementMatchers.isMethod; import static net.bytebuddy.matcher.ElementMatchers.isPublic; @@ -13,13 +9,11 @@ import static net.bytebuddy.matcher.ElementMatchers.takesArguments; import com.google.auto.service.AutoService; -import com.microsoft.azure.functions.TraceContext; import com.microsoft.azure.functions.internal.spi.middleware.MiddlewareContext; import datadog.context.ContextScope; import datadog.trace.agent.tooling.Instrumenter; import datadog.trace.agent.tooling.InstrumenterModule; import datadog.trace.bootstrap.instrumentation.api.AgentSpan; -import datadog.trace.bootstrap.instrumentation.api.AgentSpanContext; import net.bytebuddy.asm.Advice; @AutoService(InstrumenterModule.class) @@ -40,6 +34,7 @@ public String[] helperClassNames() { return new String[] { packageName + ".DurableFunctionsDecorator", packageName + ".DurableFunctionsUtils", + packageName + ".Base64StringInputStream", packageName + ".TraceContextExtractAdapter" }; } @@ -76,29 +71,30 @@ public static ContextScope onEnter(@Advice.Argument(0) MiddlewareContext context return null; } - final TraceContext traceContext = context.getTraceContext(); - AgentSpanContext.Extracted parent = - traceContext == null - ? null - : extractContextAndGetSpanContext(traceContext, TraceContextExtractAdapter.GETTER); - parent = DurableFunctionsUtils.reconcileSamplingPriority(parent, traceContext); - - final AgentSpan span = startSpan("azure-functions", AZURE_FUNCTIONS_REQUEST, parent); - DECORATE.afterStart(span); - DECORATE.onInvoke(span, context.getFunctionName(), trigger); - return activateSpan(span); + return DurableFunctionsUtils.startSpanScope(context, trigger); } @Advice.OnMethodExit(onThrowable = Throwable.class, suppress = Throwable.class) public static void onExit( - @Advice.Enter ContextScope scope, @Advice.Thrown Throwable throwable) { - if (scope != null) { - final AgentSpan span = spanFromScope(scope); + @Advice.Argument(0) MiddlewareContext context, + @Advice.Enter ContextScope scope, + @Advice.Thrown Throwable throwable) { + ContextScope activeScope = scope; + if (activeScope == null + && throwable != null + && !DurableFunctionsUtils.isReplayControlFlow(throwable)) { + final String trigger = DurableFunctionsUtils.getTrigger(context); + if ("DurableOrchestration".equals(trigger)) { + activeScope = DurableFunctionsUtils.startSpanScope(context, trigger); + } + } + if (activeScope != null) { + final AgentSpan span = spanFromScope(activeScope); if (!DurableFunctionsUtils.isReplayControlFlow(throwable)) { DECORATE.onError(span, throwable); } DECORATE.beforeFinish(span); - scope.close(); + activeScope.close(); span.finish(); } } diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/Base64StringInputStream.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/Base64StringInputStream.java new file mode 100644 index 00000000000..9ee4c414e75 --- /dev/null +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/Base64StringInputStream.java @@ -0,0 +1,34 @@ +package datadog.trace.instrumentation.azure.functions.worker; + +import java.io.InputStream; + +/** Exposes an ASCII Base64 string as a stream without copying it into a byte array. */ +public final class Base64StringInputStream extends InputStream { + private final String value; + private int position; + + public Base64StringInputStream(String value) { + this.value = value; + } + + @Override + public int read() { + return position < value.length() ? value.charAt(position++) : -1; + } + + @Override + public int read(byte[] buffer, int offset, int length) { + if (length == 0) { + return 0; + } + if (position >= value.length()) { + return -1; + } + final int count = Math.min(length, value.length() - position); + for (int index = 0; index < count; index++) { + buffer[offset + index] = (byte) value.charAt(position + index); + } + position += count; + return count; + } +} diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java index 570cb2df037..ff0cbcc3c7f 100644 --- a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java @@ -2,10 +2,19 @@ import static datadog.trace.api.sampling.PrioritySampling.SAMPLER_DROP; import static datadog.trace.api.sampling.PrioritySampling.UNSET; +import static datadog.trace.bootstrap.instrumentation.api.AgentPropagation.extractContextAndGetSpanContext; +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activateSpan; +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.startSpan; +import static datadog.trace.instrumentation.azure.functions.worker.DurableFunctionsDecorator.AZURE_FUNCTIONS_REQUEST; +import static datadog.trace.instrumentation.azure.functions.worker.DurableFunctionsDecorator.DECORATE; import com.microsoft.azure.functions.TraceContext; import com.microsoft.azure.functions.internal.spi.middleware.MiddlewareContext; +import datadog.context.ContextScope; +import datadog.trace.bootstrap.instrumentation.api.AgentSpan; import datadog.trace.bootstrap.instrumentation.api.AgentSpanContext; +import java.io.IOException; +import java.io.InputStream; import java.util.Base64; public final class DurableFunctionsUtils { @@ -48,6 +57,20 @@ public static AgentSpanContext.Extracted reconcileSamplingPriority( return parent; } + public static ContextScope startSpanScope(MiddlewareContext context, String trigger) { + final TraceContext traceContext = context.getTraceContext(); + AgentSpanContext.Extracted parent = + traceContext == null + ? null + : extractContextAndGetSpanContext(traceContext, TraceContextExtractAdapter.GETTER); + parent = reconcileSamplingPriority(parent, traceContext); + + final AgentSpan span = startSpan("azure-functions", AZURE_FUNCTIONS_REQUEST, parent); + DECORATE.afterStart(span); + DECORATE.onInvoke(span, context.getFunctionName(), trigger); + return activateSpan(span); + } + public static boolean isReplayControlFlow(Throwable throwable) { while (throwable != null) { final String className = throwable.getClass().getName(); @@ -77,28 +100,42 @@ public static boolean shouldTraceOrchestration(MiddlewareContext context) { return true; } - final byte[] request = Base64.getDecoder().decode((String) parameterValue); + final InputStream request = + Base64.getDecoder().wrap(new Base64StringInputStream((String) parameterValue)); + final byte[] skipBuffer = new byte[512]; boolean replay = false; boolean newFailure = false; - final int[] position = {0}; - while (position[0] < request.length) { - final long tag = readVarint(request, position, request.length); + final long[] position = {0}; + while (true) { + final int first = request.read(); + if (first < 0) { + break; + } + position[0]++; + final long tag = readVarint(request, position, Long.MAX_VALUE, first); final int field = (int) (tag >>> 3); final int wireType = (int) (tag & 7); if (wireType == 2) { - final long length = readVarint(request, position, request.length); - if (length < 0 || length > request.length - position[0]) { - return true; - } - final int end = position[0] + (int) length; + final long length = readVarint(request, position, Long.MAX_VALUE); + final long end = checkedEnd(position[0], length); if (field == PAST_EVENTS_FIELD) { replay = true; - } else if (field == NEW_EVENTS_FIELD && containsFailureEvent(request, position[0], end)) { - newFailure = true; + skipFully(request, position, length, skipBuffer); + } else if (field == NEW_EVENTS_FIELD) { + if (containsFailureEvent(request, position, end, skipBuffer)) { + newFailure = true; + } + } else { + skipFully(request, position, length, skipBuffer); + } + if (position[0] != end) { + throw new IllegalArgumentException("Invalid length-delimited field"); } - position[0] = end; } else { - skipValue(request, position, request.length, wireType); + skipValue(request, position, Long.MAX_VALUE, wireType, skipBuffer); + } + if (replay && newFailure) { + return true; } } return !replay || newFailure; @@ -107,59 +144,103 @@ public static boolean shouldTraceOrchestration(MiddlewareContext context) { } } - private static boolean containsFailureEvent(byte[] data, int offset, int limit) { - final int[] position = {offset}; + private static boolean containsFailureEvent( + InputStream data, long[] position, long limit, byte[] skipBuffer) throws IOException { while (position[0] < limit) { final long tag = readVarint(data, position, limit); final int field = (int) (tag >>> 3); final int wireType = (int) (tag & 7); if (wireType == 2 && (field == TASK_FAILED_FIELD || field == SUB_ORCHESTRATION_FAILED_FIELD)) { + skipFully(data, position, limit - position[0], skipBuffer); return true; } - skipValue(data, position, limit, wireType); + skipValue(data, position, limit, wireType, skipBuffer); } return false; } - private static long readVarint(byte[] data, int[] position, int limit) { + private static long readVarint(InputStream data, long[] position, long limit) throws IOException { + if (position[0] >= limit) { + throw new IllegalArgumentException("Invalid protobuf varint"); + } + final int first = data.read(); + if (first < 0) { + throw new IllegalArgumentException("Invalid protobuf varint"); + } + position[0]++; + return readVarint(data, position, limit, first); + } + + private static long readVarint(InputStream data, long[] position, long limit, int first) + throws IOException { long value = 0; - for (int shift = 0; shift < 64 && position[0] < limit; shift += 7) { - final int current = data[position[0]++] & 0xff; + int current = first; + for (int shift = 0; shift < 64; shift += 7) { value |= (long) (current & 0x7f) << shift; if ((current & 0x80) == 0) { return value; } + if (position[0] >= limit) { + break; + } + current = data.read(); + if (current < 0) { + break; + } + position[0]++; } throw new IllegalArgumentException("Invalid protobuf varint"); } - private static void skipValue(byte[] data, int[] position, int limit, int wireType) { + private static void skipValue( + InputStream data, long[] position, long limit, int wireType, byte[] skipBuffer) + throws IOException { switch (wireType) { case 0: readVarint(data, position, limit); return; case 1: - if (limit - position[0] < 8) { - throw new IllegalArgumentException("Invalid fixed64 field"); - } - position[0] += 8; + skipWithinLimit(data, position, limit, 8, skipBuffer, "Invalid fixed64 field"); return; case 2: final long length = readVarint(data, position, limit); - if (length < 0 || length > limit - position[0]) { - throw new IllegalArgumentException("Invalid length-delimited field"); - } - position[0] += (int) length; + skipWithinLimit( + data, position, limit, length, skipBuffer, "Invalid length-delimited field"); return; case 5: - if (limit - position[0] < 4) { - throw new IllegalArgumentException("Invalid fixed32 field"); - } - position[0] += 4; + skipWithinLimit(data, position, limit, 4, skipBuffer, "Invalid fixed32 field"); return; default: throw new IllegalArgumentException("Unsupported protobuf wire type"); } } + + private static void skipWithinLimit( + InputStream data, long[] position, long limit, long length, byte[] skipBuffer, String error) + throws IOException { + if (length < 0 || length > limit - position[0]) { + throw new IllegalArgumentException(error); + } + skipFully(data, position, length, skipBuffer); + } + + private static void skipFully(InputStream data, long[] position, long length, byte[] skipBuffer) + throws IOException { + while (length > 0) { + final int read = data.read(skipBuffer, 0, (int) Math.min(length, skipBuffer.length)); + if (read < 0) { + throw new IllegalArgumentException("Truncated protobuf field"); + } + position[0] += read; + length -= read; + } + } + + private static long checkedEnd(long position, long length) { + if (length < 0 || length > Long.MAX_VALUE - position) { + throw new IllegalArgumentException("Invalid length-delimited field"); + } + return position + length; + } } diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/TraceContextExtractAdapter.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/TraceContextExtractAdapter.java index d31e210635b..3ea6c52d457 100644 --- a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/TraceContextExtractAdapter.java +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/TraceContextExtractAdapter.java @@ -40,20 +40,28 @@ static int datadogSamplingPriority(String tracestate) { memberEnd = tracestate.length(); } int start = memberStart; - while (start < memberEnd && tracestate.charAt(start) == ' ') { + while (start < memberEnd && isOptionalWhitespace(tracestate.charAt(start))) { start++; } - if (memberEnd - start >= 3 + int end = memberEnd; + while (end > start && isOptionalWhitespace(tracestate.charAt(end - 1))) { + end--; + } + if (end - start >= 3 && tracestate.charAt(start) == 'd' && tracestate.charAt(start + 1) == 'd' && tracestate.charAt(start + 2) == '=') { - return parseSamplingPriority(tracestate, start + 3, memberEnd); + return parseSamplingPriority(tracestate, start + 3, end); } memberStart = memberEnd + 1; } return UNSET; } + private static boolean isOptionalWhitespace(char value) { + return value == ' ' || value == '\t'; + } + private static int parseSamplingPriority(String value, int start, int end) { while (start < end) { int fieldEnd = value.indexOf(';', start); diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java index c3b79300cab..6352a6d0b99 100644 --- a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java @@ -94,6 +94,26 @@ void doesNotTreatFailureFieldWithWrongProtobufWireTypeAsFailure() throws Excepti assertTraces(); } + @Test + void suppressesReplayWithLargeHistory() throws Exception { + MiddlewareContext context = contextFor("DurableOrchestrationTrigger", "Orchestrator"); + byte[] request = new byte[4103]; + // OrchestratorRequest { pastEvents: <4096 bytes>, newEvents: taskCompleted {} } + request[0] = 0x1a; + request[1] = (byte) 0x80; + request[2] = 0x20; + request[4099] = 0x22; + request[4100] = 0x02; + request[4101] = 0x3a; + request[4102] = 0x00; + when(context.getParameterValue("input")) + .thenReturn(Base64.getEncoder().encodeToString(request)); + + new FunctionExecutionMiddleware().invoke(context, mock(MiddlewareChain.class)); + + assertTraces(); + } + @Test void createsSpanForInitialOrchestrationExecution() throws Exception { MiddlewareContext context = contextFor("DurableOrchestrationTrigger", "Orchestrator"); @@ -126,6 +146,42 @@ void createsSpanForOrchestrationReplayWithFailure(int failureTag) throws Excepti assertTraces(trace(durableSpan("Orchestrator", "DurableOrchestration"))); } + @Test + void createsErrorSpanWhenSuppressedOrchestrationReplayFails() throws Exception { + MiddlewareContext context = contextFor("DurableOrchestrationTrigger", "Orchestrator"); + // OrchestratorRequest { pastEvents: {}, newEvents: HistoryEvent { taskCompleted: {} } } + when(context.getParameterValue("input")) + .thenReturn( + Base64.getEncoder().encodeToString(new byte[] {0x1a, 0x00, 0x22, 0x02, 0x3a, 0x00})); + MiddlewareChain chain = mock(MiddlewareChain.class); + doAnswer( + invocation -> { + throw new IllegalStateException("replay failure"); + }) + .when(chain) + .doNext(context); + + assertThrows( + IllegalStateException.class, + () -> new FunctionExecutionMiddleware().invoke(context, chain)); + + assertTraces( + trace( + span() + .root() + .operationName(Pattern.compile(Pattern.quote(operation()))) + .resourceName("DurableOrchestration Orchestrator") + .type(DDSpanTypes.SERVERLESS) + .error(true) + .tags( + defaultTags(), + error(IllegalStateException.class, "replay failure"), + tag(Tags.COMPONENT, matches(Pattern.quote("azure-functions"))), + tag(Tags.SPAN_KIND, is(Tags.SPAN_KIND_SERVER)), + tag("aas.function.name", is("Orchestrator")), + tag("aas.function.trigger", is("DurableOrchestration"))))); + } + @ParameterizedTest(name = "{0}") // spotless:off @TableTest({ @@ -183,7 +239,7 @@ void honorsDatadogKeepDecisionWhenAzureClearsW3cSampledFlag() throws Exception { TraceContext traceContext = mock(TraceContext.class); when(traceContext.getTraceparent()) .thenReturn("00-0000000000000000000000000000002a-000000000000002b-00"); - when(traceContext.getTracestate()).thenReturn("dd=s:2;o:rum"); + when(traceContext.getTracestate()).thenReturn("vendor=value,\tdd=s:2;o:rum\t"); when(context.getTraceContext()).thenReturn(traceContext); AtomicInteger priority = new AtomicInteger(Integer.MIN_VALUE); MiddlewareChain chain = mock(MiddlewareChain.class); diff --git a/internal-api/src/test/java/datadog/trace/bootstrap/instrumentation/api/AgentSpanContextTest.java b/internal-api/src/test/java/datadog/trace/bootstrap/instrumentation/api/AgentSpanContextTest.java new file mode 100644 index 00000000000..aea6c034aae --- /dev/null +++ b/internal-api/src/test/java/datadog/trace/bootstrap/instrumentation/api/AgentSpanContextTest.java @@ -0,0 +1,15 @@ +package datadog.trace.bootstrap.instrumentation.api; + +import static datadog.trace.api.sampling.PrioritySampling.UNSET; +import static org.junit.jupiter.api.Assertions.assertSame; + +import org.junit.jupiter.api.Test; + +class AgentSpanContextTest { + @Test + void extractedContextReturnsItselfWhenSamplingPriorityCannotBeReplaced() { + AgentSpanContext.Extracted context = NoopSpanContext.INSTANCE; + + assertSame(context, context.withSamplingPriority(UNSET)); + } +} From 8cd74a37b776afd96bf3f3dc510d1ae97c84c9ce Mon Sep 17 00:00:00 2001 From: jcstorms1 Date: Wed, 23 Sep 2026 16:36:31 -0400 Subject: [PATCH 05/10] Fix Azure Functions replay error handling --- .../AzureFunctionsWorkerInstrumentation.java | 9 +++- .../worker/DurableFunctionsUtils.java | 18 +++++++- .../test/java/AzureFunctionsWorkerTest.java | 41 +++++++++++++++++++ 3 files changed, 64 insertions(+), 4 deletions(-) diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java index b37e9e335cb..cd0051a7930 100644 --- a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java @@ -3,6 +3,7 @@ import static datadog.trace.agent.tooling.bytebuddy.matcher.NameMatchers.named; import static datadog.trace.bootstrap.instrumentation.api.Java8BytecodeBridge.spanFromScope; import static datadog.trace.instrumentation.azure.functions.worker.DurableFunctionsDecorator.DECORATE; +import static java.util.concurrent.TimeUnit.MILLISECONDS; import static net.bytebuddy.matcher.ElementMatchers.isMethod; import static net.bytebuddy.matcher.ElementMatchers.isPublic; import static net.bytebuddy.matcher.ElementMatchers.takesArgument; @@ -61,13 +62,16 @@ public void methodAdvice(MethodTransformer transformer) { public static class InvokeAdvice { @Advice.OnMethodEnter(suppress = Throwable.class) - public static ContextScope onEnter(@Advice.Argument(0) MiddlewareContext context) { + public static ContextScope onEnter( + @Advice.Argument(0) MiddlewareContext context, + @Advice.Local("startTimeMicros") long startTimeMicros) { final String trigger = DurableFunctionsUtils.getTrigger(context); if (trigger == null) { return null; } if ("DurableOrchestration".equals(trigger) && !DurableFunctionsUtils.shouldTraceOrchestration(context)) { + startTimeMicros = MILLISECONDS.toMicros(System.currentTimeMillis()); return null; } @@ -78,6 +82,7 @@ public static ContextScope onEnter(@Advice.Argument(0) MiddlewareContext context public static void onExit( @Advice.Argument(0) MiddlewareContext context, @Advice.Enter ContextScope scope, + @Advice.Local("startTimeMicros") long startTimeMicros, @Advice.Thrown Throwable throwable) { ContextScope activeScope = scope; if (activeScope == null @@ -85,7 +90,7 @@ public static void onExit( && !DurableFunctionsUtils.isReplayControlFlow(throwable)) { final String trigger = DurableFunctionsUtils.getTrigger(context); if ("DurableOrchestration".equals(trigger)) { - activeScope = DurableFunctionsUtils.startSpanScope(context, trigger); + activeScope = DurableFunctionsUtils.startSpanScope(context, trigger, startTimeMicros); } } if (activeScope != null) { diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java index ff0cbcc3c7f..6b48c3051e0 100644 --- a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java @@ -16,6 +16,7 @@ import java.io.IOException; import java.io.InputStream; import java.util.Base64; +import java.util.IdentityHashMap; public final class DurableFunctionsUtils { private static final String ACTIVITY_ANNOTATION = "DurableActivityTrigger"; @@ -58,6 +59,11 @@ public static AgentSpanContext.Extracted reconcileSamplingPriority( } public static ContextScope startSpanScope(MiddlewareContext context, String trigger) { + return startSpanScope(context, trigger, 0); + } + + public static ContextScope startSpanScope( + MiddlewareContext context, String trigger, long startTimeMicros) { final TraceContext traceContext = context.getTraceContext(); AgentSpanContext.Extracted parent = traceContext == null @@ -65,14 +71,22 @@ public static ContextScope startSpanScope(MiddlewareContext context, String trig : extractContextAndGetSpanContext(traceContext, TraceContextExtractAdapter.GETTER); parent = reconcileSamplingPriority(parent, traceContext); - final AgentSpan span = startSpan("azure-functions", AZURE_FUNCTIONS_REQUEST, parent); + final AgentSpan span = + startTimeMicros > 0 + ? startSpan("azure-functions", AZURE_FUNCTIONS_REQUEST, parent, startTimeMicros) + : startSpan("azure-functions", AZURE_FUNCTIONS_REQUEST, parent); DECORATE.afterStart(span); DECORATE.onInvoke(span, context.getFunctionName(), trigger); return activateSpan(span); } public static boolean isReplayControlFlow(Throwable throwable) { - while (throwable != null) { + if (throwable == null) { + return false; + } + + final IdentityHashMap seen = new IdentityHashMap<>(); + while (throwable != null && seen.put(throwable, Boolean.TRUE) == null) { final String className = throwable.getClass().getName(); if ("com.microsoft.durabletask.OrchestratorBlockedException".equals(className) || "com.microsoft.durabletask.interruption.OrchestratorBlockedException".equals(className) diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java index 6352a6d0b99..7c5fd9ed982 100644 --- a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java @@ -7,7 +7,9 @@ import static datadog.trace.api.sampling.PrioritySampling.UNSET; import static datadog.trace.test.junit.utils.assertions.Matchers.is; import static datadog.trace.test.junit.utils.assertions.Matchers.matches; +import static java.util.concurrent.TimeUnit.MILLISECONDS; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.junit.jupiter.params.provider.Arguments.arguments; @@ -28,9 +30,11 @@ import datadog.trace.bootstrap.instrumentation.api.AgentTracer; import datadog.trace.bootstrap.instrumentation.api.Tags; import datadog.trace.core.DDSpan; +import datadog.trace.instrumentation.azure.functions.worker.DurableFunctionsUtils; import java.util.Base64; import java.util.Objects; import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicLong; import java.util.regex.Pattern; import java.util.stream.Stream; import org.junit.jupiter.api.Test; @@ -154,8 +158,11 @@ void createsErrorSpanWhenSuppressedOrchestrationReplayFails() throws Exception { .thenReturn( Base64.getEncoder().encodeToString(new byte[] {0x1a, 0x00, 0x22, 0x02, 0x3a, 0x00})); MiddlewareChain chain = mock(MiddlewareChain.class); + AtomicLong invocationStartMillis = new AtomicLong(); doAnswer( invocation -> { + invocationStartMillis.set(System.currentTimeMillis()); + Thread.sleep(25); throw new IllegalStateException("replay failure"); }) .when(chain) @@ -165,6 +172,11 @@ void createsErrorSpanWhenSuppressedOrchestrationReplayFails() throws Exception { IllegalStateException.class, () -> new FunctionExecutionMiddleware().invoke(context, chain)); + writer.waitForTraces(1); + DDSpan errorSpan = writer.firstTrace().get(0); + assertTrue(errorSpan.getStartTime() <= MILLISECONDS.toNanos(invocationStartMillis.get())); + assertTrue(errorSpan.getDurationNano() >= MILLISECONDS.toNanos(20)); + assertTraces( trace( span() @@ -182,6 +194,35 @@ void createsErrorSpanWhenSuppressedOrchestrationReplayFails() throws Exception { tag("aas.function.trigger", is("DurableOrchestration"))))); } + @Test + void stopsTraversingRepeatedExceptionCause() { + AtomicInteger causeReads = new AtomicInteger(); + Throwable[] cycle = new Throwable[2]; + cycle[0] = + new RuntimeException("first") { + @Override + public Throwable getCause() { + if (causeReads.incrementAndGet() > 2) { + throw new AssertionError("cause traversal did not terminate"); + } + return cycle[1]; + } + }; + cycle[1] = + new RuntimeException("second") { + @Override + public Throwable getCause() { + if (causeReads.incrementAndGet() > 2) { + throw new AssertionError("cause traversal did not terminate"); + } + return cycle[0]; + } + }; + + assertFalse(DurableFunctionsUtils.isReplayControlFlow(cycle[0])); + assertEquals(2, causeReads.get()); + } + @ParameterizedTest(name = "{0}") // spotless:off @TableTest({ From e159a3d17365a64fe336c626e46dfe6d780f4213 Mon Sep 17 00:00:00 2001 From: jcstorms1 Date: Wed, 23 Sep 2026 17:29:21 -0400 Subject: [PATCH 06/10] Address Azure Functions review feedback --- .../azure-functions-worker-2.7/build.gradle | 26 ----------- .../build.gradle.kts | 38 ++++++++++++++++ ...tream.java => AsciiStringInputStream.java} | 6 +-- .../AzureFunctionsWorkerInstrumentation.java | 2 +- .../worker/DurableFunctionsUtils.java | 2 +- .../test/java/AzureFunctionsWorkerTest.java | 45 +++++++------------ 6 files changed, 60 insertions(+), 59 deletions(-) delete mode 100644 dd-java-agent/instrumentation/azure-functions-worker-2.7/build.gradle create mode 100644 dd-java-agent/instrumentation/azure-functions-worker-2.7/build.gradle.kts rename dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/{Base64StringInputStream.java => AsciiStringInputStream.java} (77%) diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/build.gradle b/dd-java-agent/instrumentation/azure-functions-worker-2.7/build.gradle deleted file mode 100644 index 930115252af..00000000000 --- a/dd-java-agent/instrumentation/azure-functions-worker-2.7/build.gradle +++ /dev/null @@ -1,26 +0,0 @@ -plugins { - id 'dd-trace-java.module.instrumentation' -} - -muzzle { - pass { - group = 'com.microsoft.azure.functions' - module = 'azure-functions-java-spi' - versions = '[1.0.0,)' - extraDependency 'com.microsoft.azure.functions:azure-functions-java-core-library:1.2.0' - } -} - -addTestSuiteForDir('latestDepTest', 'test') - -dependencies { - compileOnly group: 'com.microsoft.azure.functions', name: 'azure-functions-java-core-library', version: '1.2.0' - compileOnly group: 'com.microsoft.azure.functions', name: 'azure-functions-java-spi', version: '1.0.0' - - testImplementation group: 'com.microsoft.azure.functions', name: 'azure-functions-java-core-library', version: '1.2.0' - testImplementation group: 'com.microsoft.azure.functions', name: 'azure-functions-java-spi', version: '1.0.0' - testImplementation libs.bundles.mockito - - latestDepTestImplementation group: 'com.microsoft.azure.functions', name: 'azure-functions-java-core-library', version: '+' - latestDepTestImplementation group: 'com.microsoft.azure.functions', name: 'azure-functions-java-spi', version: '+' -} diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/build.gradle.kts b/dd-java-agent/instrumentation/azure-functions-worker-2.7/build.gradle.kts new file mode 100644 index 00000000000..944def40f19 --- /dev/null +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/build.gradle.kts @@ -0,0 +1,38 @@ +import groovy.lang.Closure + +plugins { + id("dd-trace-java.module.instrumentation") +} + +muzzle { + pass { + group = "com.microsoft.azure.functions" + module = "azure-functions-java-spi" + versions = "[1.0.0,)" + extraDependency("com.microsoft.azure.functions:azure-functions-java-core-library:1.2.0") + } +} + +fun addTestSuiteForDir(name: String, directory: String) { + (project.extra["addTestSuiteForDir"] as Closure<*>).call(name, directory) +} + +addTestSuiteForDir("latestDepTest", "test") + +dependencies { + compileOnly("com.microsoft.azure.functions:azure-functions-java-core-library:1.2.0") + compileOnly("com.microsoft.azure.functions:azure-functions-java-spi:1.0.0") + + testImplementation("com.microsoft.azure.functions:azure-functions-java-core-library:1.2.0") + testImplementation("com.microsoft.azure.functions:azure-functions-java-spi:1.0.0") + testImplementation(libs.bundles.mockito) + + add( + "latestDepTestImplementation", + "com.microsoft.azure.functions:azure-functions-java-core-library:+" + ) + add( + "latestDepTestImplementation", + "com.microsoft.azure.functions:azure-functions-java-spi:+" + ) +} diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/Base64StringInputStream.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AsciiStringInputStream.java similarity index 77% rename from dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/Base64StringInputStream.java rename to dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AsciiStringInputStream.java index 9ee4c414e75..e26d6e9803f 100644 --- a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/Base64StringInputStream.java +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AsciiStringInputStream.java @@ -2,12 +2,12 @@ import java.io.InputStream; -/** Exposes an ASCII Base64 string as a stream without copying it into a byte array. */ -public final class Base64StringInputStream extends InputStream { +/** Exposes an ASCII string as a stream without copying it into a byte array. */ +public final class AsciiStringInputStream extends InputStream { private final String value; private int position; - public Base64StringInputStream(String value) { + public AsciiStringInputStream(String value) { this.value = value; } diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java index cd0051a7930..ed827072bb7 100644 --- a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java @@ -35,7 +35,7 @@ public String[] helperClassNames() { return new String[] { packageName + ".DurableFunctionsDecorator", packageName + ".DurableFunctionsUtils", - packageName + ".Base64StringInputStream", + packageName + ".AsciiStringInputStream", packageName + ".TraceContextExtractAdapter" }; } diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java index 6b48c3051e0..3269eede800 100644 --- a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java @@ -115,7 +115,7 @@ public static boolean shouldTraceOrchestration(MiddlewareContext context) { } final InputStream request = - Base64.getDecoder().wrap(new Base64StringInputStream((String) parameterValue)); + Base64.getDecoder().wrap(new AsciiStringInputStream((String) parameterValue)); final byte[] skipBuffer = new byte[512]; boolean replay = false; boolean newFailure = false; diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java index 7c5fd9ed982..a7b7f59f89c 100644 --- a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java @@ -22,8 +22,6 @@ import com.microsoft.azure.functions.internal.spi.middleware.MiddlewareChain; import com.microsoft.azure.functions.internal.spi.middleware.MiddlewareContext; import com.microsoft.azure.functions.worker.chain.FunctionExecutionMiddleware; -import com.microsoft.durabletask.interruption.ContinueAsNewInterruption; -import com.microsoft.durabletask.interruption.OrchestratorBlockedException; import datadog.trace.agent.test.AbstractInstrumentationTest; import datadog.trace.agent.test.assertions.SpanMatcher; import datadog.trace.api.DDSpanTypes; @@ -47,14 +45,12 @@ abstract class AzureFunctionsWorkerTest extends AbstractInstrumentationTest { abstract String operation(); - // spotless:off @TableTest({ - "scenario | annotation | trigger", + "scenario | annotation | trigger ", "orchestration | DurableOrchestrationTrigger | DurableOrchestration", - "activity | DurableActivityTrigger | DurableActivity", - "entity | DurableEntityTrigger | DurableEntity" + "activity | DurableActivityTrigger | DurableActivity ", + "entity | DurableEntityTrigger | DurableEntity " }) - // spotless:on void createsSpanForDurableTrigger(String annotation, String trigger) throws Exception { MiddlewareContext context = contextFor(annotation, "MyFunction"); @@ -102,10 +98,11 @@ void doesNotTreatFailureFieldWithWrongProtobufWireTypeAsFailure() throws Excepti void suppressesReplayWithLargeHistory() throws Exception { MiddlewareContext context = contextFor("DurableOrchestrationTrigger", "Orchestrator"); byte[] request = new byte[4103]; - // OrchestratorRequest { pastEvents: <4096 bytes>, newEvents: taskCompleted {} } + // Field 3 (pastEvents), length-delimited with a two-byte varint length of 4096. request[0] = 0x1a; request[1] = (byte) 0x80; request[2] = 0x20; + // Field 4 (newEvents), containing HistoryEvent field 7 (taskCompleted) with an empty payload. request[4099] = 0x22; request[4100] = 0x02; request[4101] = 0x3a; @@ -130,13 +127,11 @@ void createsSpanForInitialOrchestrationExecution() throws Exception { assertTraces(trace(durableSpan("Orchestrator", "DurableOrchestration"))); } - // spotless:off @TableTest({ "scenario | failureTag", - "failed activity | 66", - "failed sub-orchestration | 90" + "failed activity | 66 ", + "failed sub-orchestration | 90 " }) - // spotless:on void createsSpanForOrchestrationReplayWithFailure(int failureTag) throws Exception { MiddlewareContext context = contextFor("DurableOrchestrationTrigger", "Orchestrator"); // OrchestratorRequest { pastEvents: {}, newEvents: HistoryEvent { failure: {} } } @@ -224,13 +219,11 @@ public Throwable getCause() { } @ParameterizedTest(name = "{0}") - // spotless:off @TableTest({ - "scenario | description | payload", + "scenario | description | payload ", "non-base64 input | non-base64 input | 'not base64'", - "truncated protobuf | truncated protobuf | GgIA" + "truncated protobuf | truncated protobuf | GgIA " }) - // spotless:on @MethodSource("failsOpenForUnreadableOrchestrationPayloadArguments") void failsOpenForUnreadableOrchestrationPayload(String description, Object payload) throws Exception { @@ -301,10 +294,15 @@ void honorsDatadogKeepDecisionWhenAzureClearsW3cSampledFlag() throws Exception { assertTrue(priority.get() > 0); } - @ParameterizedTest(name = "{0}") - @MethodSource("doesNotMarkReplayControlFlowAsErrorArguments") - void doesNotMarkReplayControlFlowAsError(String description, Throwable controlFlow) + @TableTest({ + "scenario | controlFlowType ", + "legacy OrchestratorBlockedException | com.microsoft.durabletask.OrchestratorBlockedException ", + "OrchestratorBlockedException | com.microsoft.durabletask.interruption.OrchestratorBlockedException", + "ContinueAsNewInterruption | com.microsoft.durabletask.interruption.ContinueAsNewInterruption " + }) + void doesNotMarkReplayControlFlowAsError(Class controlFlowType) throws Exception { + Throwable controlFlow = controlFlowType.getDeclaredConstructor().newInstance(); MiddlewareContext context = contextFor("DurableOrchestrationTrigger", "Orchestrator"); MiddlewareChain chain = mock(MiddlewareChain.class); doAnswer( @@ -320,15 +318,6 @@ void doesNotMarkReplayControlFlowAsError(String description, Throwable controlFl assertTraces(trace(durableSpan("Orchestrator", "DurableOrchestration"))); } - static Stream doesNotMarkReplayControlFlowAsErrorArguments() { - return Stream.of( - arguments( - "legacy OrchestratorBlockedException", - new com.microsoft.durabletask.OrchestratorBlockedException()), - arguments("OrchestratorBlockedException", new OrchestratorBlockedException()), - arguments("ContinueAsNewInterruption", new ContinueAsNewInterruption())); - } - @Test void marksApplicationFailuresAsErrors() throws Exception { MiddlewareContext context = contextFor("DurableActivityTrigger", "Activity"); From ab2c863fe914bcc54561a51f218559d592aa428d Mon Sep 17 00:00:00 2001 From: jcstorms1 Date: Wed, 23 Sep 2026 17:34:54 -0400 Subject: [PATCH 07/10] Sort Azure Functions helper classes --- .../functions/worker/AzureFunctionsWorkerInstrumentation.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java index ed827072bb7..f676499197e 100644 --- a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java @@ -33,9 +33,9 @@ public String instrumentedType() { @Override public String[] helperClassNames() { return new String[] { + packageName + ".AsciiStringInputStream", packageName + ".DurableFunctionsDecorator", packageName + ".DurableFunctionsUtils", - packageName + ".AsciiStringInputStream", packageName + ".TraceContextExtractAdapter" }; } From 868f14c316528b7423c46b4fb32f8b60f7f6badb Mon Sep 17 00:00:00 2001 From: jcstorms1 Date: Thu, 24 Sep 2026 11:24:52 -0400 Subject: [PATCH 08/10] Scope Azure sampling workaround to durable functions --- .../worker/DurableFunctionsUtils.java | 17 +------- .../worker/TraceContextExtractAdapter.java | 6 ++- .../test/java/AzureFunctionsWorkerTest.java | 34 ++++++++++++++-- .../core/propagation/ExtractedContext.java | 11 ----- .../propagation/ExtractedContextTest.java | 40 ------------------- .../instrumentation/api/AgentSpanContext.java | 8 ---- .../instrumentation/api/TagContext.java | 6 +-- .../api/AgentSpanContextTest.java | 15 ------- 8 files changed, 38 insertions(+), 99 deletions(-) delete mode 100644 dd-trace-core/src/test/java/datadog/trace/core/propagation/ExtractedContextTest.java delete mode 100644 internal-api/src/test/java/datadog/trace/bootstrap/instrumentation/api/AgentSpanContextTest.java diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java index 3269eede800..665307c77a0 100644 --- a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java @@ -1,7 +1,5 @@ package datadog.trace.instrumentation.azure.functions.worker; -import static datadog.trace.api.sampling.PrioritySampling.SAMPLER_DROP; -import static datadog.trace.api.sampling.PrioritySampling.UNSET; import static datadog.trace.bootstrap.instrumentation.api.AgentPropagation.extractContextAndGetSpanContext; import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activateSpan; import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.startSpan; @@ -46,18 +44,6 @@ public static String getTrigger(MiddlewareContext context) { return null; } - public static AgentSpanContext.Extracted reconcileSamplingPriority( - AgentSpanContext.Extracted parent, TraceContext traceContext) { - if (parent != null - && parent.getSamplingPriority() == SAMPLER_DROP - && traceContext != null - && TraceContextExtractAdapter.datadogSamplingPriority(traceContext.getTracestate()) - == UNSET) { - return parent.withSamplingPriority(UNSET); - } - return parent; - } - public static ContextScope startSpanScope(MiddlewareContext context, String trigger) { return startSpanScope(context, trigger, 0); } @@ -65,11 +51,10 @@ public static ContextScope startSpanScope(MiddlewareContext context, String trig public static ContextScope startSpanScope( MiddlewareContext context, String trigger, long startTimeMicros) { final TraceContext traceContext = context.getTraceContext(); - AgentSpanContext.Extracted parent = + final AgentSpanContext.Extracted parent = traceContext == null ? null : extractContextAndGetSpanContext(traceContext, TraceContextExtractAdapter.GETTER); - parent = reconcileSamplingPriority(parent, traceContext); final AgentSpan span = startTimeMicros > 0 diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/TraceContextExtractAdapter.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/TraceContextExtractAdapter.java index 3ea6c52d457..d210f9b66c7 100644 --- a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/TraceContextExtractAdapter.java +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/TraceContextExtractAdapter.java @@ -18,7 +18,11 @@ private TraceContextExtractAdapter() {} public void forEachKey(TraceContext carrier, AgentPropagation.KeyClassifier classifier) { final String tracestate = carrier.getTracestate(); String traceparent = carrier.getTraceparent(); - if (datadogSamplingPriority(tracestate) > 0) { + final int datadogSamplingPriority = datadogSamplingPriority(tracestate); + if (datadogSamplingPriority == UNSET || datadogSamplingPriority > 0) { + // The Azure host can clear the sampled flag when host telemetry is disabled even though the + // trace identifiers are intended to correlate durable worker invocations. Restore the flag + // only in this durable-trigger extraction path, while preserving explicit Datadog drops. traceparent = setSampledFlag(traceparent); } if (traceparent != null && !classifier.accept(TRACEPARENT, traceparent)) { diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java index a7b7f59f89c..de80b4b18a1 100644 --- a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java @@ -3,8 +3,8 @@ import static datadog.trace.agent.test.assertions.TagsMatcher.error; import static datadog.trace.agent.test.assertions.TagsMatcher.tag; import static datadog.trace.agent.test.assertions.TraceMatcher.trace; +import static datadog.trace.api.sampling.PrioritySampling.SAMPLER_DROP; import static datadog.trace.api.sampling.PrioritySampling.SAMPLER_KEEP; -import static datadog.trace.api.sampling.PrioritySampling.UNSET; import static datadog.trace.test.junit.utils.assertions.Matchers.is; import static datadog.trace.test.junit.utils.assertions.Matchers.matches; import static java.util.concurrent.TimeUnit.MILLISECONDS; @@ -240,7 +240,7 @@ static Stream failsOpenForUnreadableOrchestrationPayloadArguments() { } @Test - void continuesAzureTraceAndIgnoresImplicitSamplingRejection() throws Exception { + void continuesAzureTraceWhenHostClearsW3cSampledFlag() throws Exception { MiddlewareContext context = contextFor("DurableActivityTrigger", "Activity"); TraceContext traceContext = mock(TraceContext.class); when(traceContext.getTraceparent()) @@ -263,10 +263,38 @@ void continuesAzureTraceAndIgnoresImplicitSamplingRejection() throws Exception { DDSpan span = writer.firstTrace().get(0); assertEquals("42", span.getTraceId().toString()); assertEquals(43L, span.getParentId()); - assertEquals(UNSET, priority.get()); + assertEquals(SAMPLER_KEEP, priority.get()); assertEquals(SAMPLER_KEEP, span.samplingPriority()); } + @Test + void honorsDatadogDropDecisionWhenAzureClearsW3cSampledFlag() throws Exception { + MiddlewareContext context = contextFor("DurableActivityTrigger", "Activity"); + TraceContext traceContext = mock(TraceContext.class); + when(traceContext.getTraceparent()) + .thenReturn("00-0000000000000000000000000000002a-000000000000002b-00"); + when(traceContext.getTracestate()).thenReturn("dd=s:0"); + when(context.getTraceContext()).thenReturn(traceContext); + AtomicInteger priority = new AtomicInteger(Integer.MIN_VALUE); + MiddlewareChain chain = mock(MiddlewareChain.class); + doAnswer( + invocation -> { + priority.set(AgentTracer.activeSpan().spanContext().getSamplingPriority()); + return null; + }) + .when(chain) + .doNext(context); + + new FunctionExecutionMiddleware().invoke(context, chain); + + assertEquals(SAMPLER_DROP, priority.get()); + writer.waitForTraces(1); + DDSpan span = writer.firstTrace().get(0); + assertEquals("42", span.getTraceId().toString()); + assertEquals(43L, span.getParentId()); + assertEquals(SAMPLER_DROP, span.samplingPriority()); + } + @Test void honorsDatadogKeepDecisionWhenAzureClearsW3cSampledFlag() throws Exception { MiddlewareContext context = contextFor("DurableOrchestrationTrigger", "Orchestrator"); diff --git a/dd-trace-core/src/main/java/datadog/trace/core/propagation/ExtractedContext.java b/dd-trace-core/src/main/java/datadog/trace/core/propagation/ExtractedContext.java index 7c48ac9e452..af503e6a6ed 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/propagation/ExtractedContext.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/propagation/ExtractedContext.java @@ -1,7 +1,5 @@ package datadog.trace.core.propagation; -import static datadog.trace.api.sampling.SamplingMechanism.EXTERNAL_OVERRIDE; - import datadog.trace.api.DDTraceId; import datadog.trace.api.TagMap; import datadog.trace.api.TraceConfig; @@ -119,15 +117,6 @@ public PropagationTags getPropagationTags() { return propagationTags; } - @Override - public ExtractedContext withSamplingPriority(final int samplingPriority) { - setSamplingPriority(samplingPriority); - if (propagationTags != null) { - propagationTags.updateTraceSamplingPriority(samplingPriority, EXTERNAL_OVERRIDE); - } - return this; - } - @Override public String toString() { StringBuilder builder = new StringBuilder("ExtractedContext{"); diff --git a/dd-trace-core/src/test/java/datadog/trace/core/propagation/ExtractedContextTest.java b/dd-trace-core/src/test/java/datadog/trace/core/propagation/ExtractedContextTest.java deleted file mode 100644 index 130d5adef33..00000000000 --- a/dd-trace-core/src/test/java/datadog/trace/core/propagation/ExtractedContextTest.java +++ /dev/null @@ -1,40 +0,0 @@ -package datadog.trace.core.propagation; - -import static datadog.trace.api.TracePropagationStyle.TRACECONTEXT; -import static datadog.trace.api.sampling.PrioritySampling.SAMPLER_DROP; -import static datadog.trace.api.sampling.PrioritySampling.UNSET; -import static datadog.trace.api.sampling.SamplingMechanism.EXTERNAL_OVERRIDE; -import static org.junit.jupiter.api.Assertions.assertEquals; -import static org.junit.jupiter.api.Assertions.assertSame; - -import datadog.trace.api.DDTraceId; -import org.junit.jupiter.api.Test; - -class ExtractedContextTest { - @Test - void replacesSamplingPriorityWithoutChangingTraceIdentity() { - PropagationTags propagationTags = PropagationTags.factory().empty(); - propagationTags.updateTraceSamplingPriority(SAMPLER_DROP, EXTERNAL_OVERRIDE); - ExtractedContext original = - new ExtractedContext( - DDTraceId.from(42), 43, SAMPLER_DROP, null, propagationTags, TRACECONTEXT); - - ExtractedContext updated = original.withSamplingPriority(UNSET); - - assertSame(original, updated); - assertEquals(original.getTraceId(), updated.getTraceId()); - assertEquals(original.getSpanId(), updated.getSpanId()); - assertEquals(UNSET, updated.getSamplingPriority()); - assertEquals(UNSET, updated.getPropagationTags().getSamplingPriority()); - } - - @Test - void replacesSamplingPriorityWithoutPropagationTags() { - ExtractedContext context = - new ExtractedContext(DDTraceId.from(42), 43, SAMPLER_DROP, null, null, TRACECONTEXT); - - context.withSamplingPriority(UNSET); - - assertEquals(UNSET, context.getSamplingPriority()); - } -} diff --git a/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/AgentSpanContext.java b/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/AgentSpanContext.java index fcde079c1ea..6e950925305 100644 --- a/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/AgentSpanContext.java +++ b/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/AgentSpanContext.java @@ -72,14 +72,6 @@ default CharSequence getIntegrationName() { boolean isRemote(); interface Extracted extends AgentSpanContext { - /** - * Returns this extracted context with the supplied sampling priority when the implementation - * supports replacing propagation decisions. - */ - default Extracted withSamplingPriority(int samplingPriority) { - return this; - } - /** * Gets the span links related to the other terminated context. * diff --git a/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/TagContext.java b/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/TagContext.java index f7098e94aa2..078ffeb1625 100644 --- a/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/TagContext.java +++ b/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/TagContext.java @@ -33,7 +33,7 @@ public class TagContext implements AgentSpanContext.Extracted { private final HttpHeaders httpHeaders; private final Map baggage; private Baggage w3cBaggage; - private int samplingPriority; + private final int samplingPriority; private final TraceConfig traceConfig; private final TracePropagationStyle propagationStyle; private final DDTraceId traceId; @@ -181,10 +181,6 @@ public final int getSamplingPriority() { return samplingPriority; } - protected final void setSamplingPriority(int samplingPriority) { - this.samplingPriority = samplingPriority; - } - public final Map getBaggage() { return baggage; } diff --git a/internal-api/src/test/java/datadog/trace/bootstrap/instrumentation/api/AgentSpanContextTest.java b/internal-api/src/test/java/datadog/trace/bootstrap/instrumentation/api/AgentSpanContextTest.java deleted file mode 100644 index aea6c034aae..00000000000 --- a/internal-api/src/test/java/datadog/trace/bootstrap/instrumentation/api/AgentSpanContextTest.java +++ /dev/null @@ -1,15 +0,0 @@ -package datadog.trace.bootstrap.instrumentation.api; - -import static datadog.trace.api.sampling.PrioritySampling.UNSET; -import static org.junit.jupiter.api.Assertions.assertSame; - -import org.junit.jupiter.api.Test; - -class AgentSpanContextTest { - @Test - void extractedContextReturnsItselfWhenSamplingPriorityCannotBeReplaced() { - AgentSpanContext.Extracted context = NoopSpanContext.INSTANCE; - - assertSame(context, context.withSamplingPriority(UNSET)); - } -} From 675b1ca459e31fe5cd5c6ac5de3a69d0344d4fba Mon Sep 17 00:00:00 2001 From: jcstorms1 Date: Thu, 24 Sep 2026 14:35:40 -0400 Subject: [PATCH 09/10] Align Azure worker test fixture with SPI --- .../functions/worker/chain/FunctionExecutionMiddleware.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/com/microsoft/azure/functions/worker/chain/FunctionExecutionMiddleware.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/com/microsoft/azure/functions/worker/chain/FunctionExecutionMiddleware.java index 1cda6f23f27..5b8b9822d87 100644 --- a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/com/microsoft/azure/functions/worker/chain/FunctionExecutionMiddleware.java +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/com/microsoft/azure/functions/worker/chain/FunctionExecutionMiddleware.java @@ -1,9 +1,11 @@ package com.microsoft.azure.functions.worker.chain; +import com.microsoft.azure.functions.internal.spi.middleware.Middleware; import com.microsoft.azure.functions.internal.spi.middleware.MiddlewareChain; import com.microsoft.azure.functions.internal.spi.middleware.MiddlewareContext; -public class FunctionExecutionMiddleware { +public class FunctionExecutionMiddleware implements Middleware { + @Override public void invoke(MiddlewareContext context, MiddlewareChain chain) throws Exception { chain.doNext(context); } From 96637fd5557e476f76a254bd43d81fb141d6b708 Mon Sep 17 00:00:00 2001 From: jcstorms1 Date: Thu, 24 Sep 2026 15:38:34 -0400 Subject: [PATCH 10/10] Harden durable function replay tracing --- .../AzureFunctionsWorkerInstrumentation.java | 12 ++--- .../worker/DurableFunctionsUtils.java | 11 ++--- .../test/java/AzureFunctionsWorkerTest.java | 44 +++++++++++++++++-- 3 files changed, 51 insertions(+), 16 deletions(-) diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java index f676499197e..efe806deff2 100644 --- a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/AzureFunctionsWorkerInstrumentation.java @@ -64,15 +64,17 @@ public static class InvokeAdvice { @Advice.OnMethodEnter(suppress = Throwable.class) public static ContextScope onEnter( @Advice.Argument(0) MiddlewareContext context, + @Advice.Local("trigger") String trigger, @Advice.Local("startTimeMicros") long startTimeMicros) { - final String trigger = DurableFunctionsUtils.getTrigger(context); + trigger = DurableFunctionsUtils.getTrigger(context); if (trigger == null) { return null; } - if ("DurableOrchestration".equals(trigger) - && !DurableFunctionsUtils.shouldTraceOrchestration(context)) { + if ("DurableOrchestration".equals(trigger)) { startTimeMicros = MILLISECONDS.toMicros(System.currentTimeMillis()); - return null; + if (!DurableFunctionsUtils.shouldTraceOrchestration(context)) { + return null; + } } return DurableFunctionsUtils.startSpanScope(context, trigger); @@ -82,13 +84,13 @@ public static ContextScope onEnter( public static void onExit( @Advice.Argument(0) MiddlewareContext context, @Advice.Enter ContextScope scope, + @Advice.Local("trigger") String trigger, @Advice.Local("startTimeMicros") long startTimeMicros, @Advice.Thrown Throwable throwable) { ContextScope activeScope = scope; if (activeScope == null && throwable != null && !DurableFunctionsUtils.isReplayControlFlow(throwable)) { - final String trigger = DurableFunctionsUtils.getTrigger(context); if ("DurableOrchestration".equals(trigger)) { activeScope = DurableFunctionsUtils.startSpanScope(context, trigger, startTimeMicros); } diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java index 665307c77a0..aab77d76563 100644 --- a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/main/java/datadog/trace/instrumentation/azure/functions/worker/DurableFunctionsUtils.java @@ -127,9 +127,6 @@ public static boolean shouldTraceOrchestration(MiddlewareContext context) { } else { skipFully(request, position, length, skipBuffer); } - if (position[0] != end) { - throw new IllegalArgumentException("Invalid length-delimited field"); - } } else { skipValue(request, position, Long.MAX_VALUE, wireType, skipBuffer); } @@ -138,25 +135,25 @@ public static boolean shouldTraceOrchestration(MiddlewareContext context) { } } return !replay || newFailure; - } catch (Throwable ignored) { + } catch (Exception ignored) { return true; } } private static boolean containsFailureEvent( InputStream data, long[] position, long limit, byte[] skipBuffer) throws IOException { + boolean failure = false; while (position[0] < limit) { final long tag = readVarint(data, position, limit); final int field = (int) (tag >>> 3); final int wireType = (int) (tag & 7); if (wireType == 2 && (field == TASK_FAILED_FIELD || field == SUB_ORCHESTRATION_FAILED_FIELD)) { - skipFully(data, position, limit - position[0], skipBuffer); - return true; + failure = true; } skipValue(data, position, limit, wireType, skipBuffer); } - return false; + return failure; } private static long readVarint(InputStream data, long[] position, long limit) throws IOException { diff --git a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java index de80b4b18a1..e5269da97a8 100644 --- a/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java +++ b/dd-java-agent/instrumentation/azure-functions-worker-2.7/src/test/java/AzureFunctionsWorkerTest.java @@ -146,12 +146,38 @@ void createsSpanForOrchestrationReplayWithFailure(int failureTag) throws Excepti } @Test - void createsErrorSpanWhenSuppressedOrchestrationReplayFails() throws Exception { + void createsSpanWhenFailurePrecedesAnotherHistoryField() throws Exception { MiddlewareContext context = contextFor("DurableOrchestrationTrigger", "Orchestrator"); - // OrchestratorRequest { pastEvents: {}, newEvents: HistoryEvent { taskCompleted: {} } } + // OrchestratorRequest { + // pastEvents: {}, + // newEvents: HistoryEvent { taskFailed: {}, field 20: 5 } + // } when(context.getParameterValue("input")) .thenReturn( - Base64.getEncoder().encodeToString(new byte[] {0x1a, 0x00, 0x22, 0x02, 0x3a, 0x00})); + Base64.getEncoder() + .encodeToString( + new byte[] {0x1a, 0x00, 0x22, 0x05, 0x42, 0x00, (byte) 0xa0, 0x01, 0x05})); + + new FunctionExecutionMiddleware().invoke(context, mock(MiddlewareChain.class)); + + assertTraces(trace(durableSpan("Orchestrator", "DurableOrchestration"))); + } + + @Test + void createsErrorSpanWhenSuppressedOrchestrationReplayFails() throws Exception { + MiddlewareContext context = contextFor("DurableOrchestrationTrigger", "Orchestrator"); + // OrchestratorRequest { pastEvents: {}, newEvents: HistoryEvent { taskCompleted: {} } } + String payload = + Base64.getEncoder().encodeToString(new byte[] {0x1a, 0x00, 0x22, 0x02, 0x3a, 0x00}); + AtomicLong payloadReadStartMillis = new AtomicLong(); + doAnswer( + invocation -> { + payloadReadStartMillis.set(System.currentTimeMillis()); + Thread.sleep(25); + return payload; + }) + .when(context) + .getParameterValue("input"); MiddlewareChain chain = mock(MiddlewareChain.class); AtomicLong invocationStartMillis = new AtomicLong(); doAnswer( @@ -169,8 +195,9 @@ void createsErrorSpanWhenSuppressedOrchestrationReplayFails() throws Exception { writer.waitForTraces(1); DDSpan errorSpan = writer.firstTrace().get(0); + assertTrue(errorSpan.getStartTime() <= MILLISECONDS.toNanos(payloadReadStartMillis.get())); assertTrue(errorSpan.getStartTime() <= MILLISECONDS.toNanos(invocationStartMillis.get())); - assertTrue(errorSpan.getDurationNano() >= MILLISECONDS.toNanos(20)); + assertTrue(errorSpan.getDurationNano() >= MILLISECONDS.toNanos(45)); assertTraces( trace( @@ -239,6 +266,15 @@ static Stream failsOpenForUnreadableOrchestrationPayloadArguments() { return Stream.of(arguments("non-string input", new Object())); } + @Test + void doesNotSwallowErrorsWhileInspectingOrchestrationPayload() { + MiddlewareContext context = contextFor("DurableOrchestrationTrigger", "Orchestrator"); + when(context.getParameterValue("input")).thenThrow(new AssertionError("failure")); + + assertThrows( + AssertionError.class, () -> DurableFunctionsUtils.shouldTraceOrchestration(context)); + } + @Test void continuesAzureTraceWhenHostClearsW3cSampledFlag() throws Exception { MiddlewareContext context = contextFor("DurableActivityTrigger", "Activity");