diff --git a/dd-java-agent/instrumentation/spring/spring-webflux/spring-webflux-5.0/src/test/groovy/dd/trace/instrumentation/springwebflux/client/SpringWebfluxHttpClientBase.groovy b/dd-java-agent/instrumentation/spring/spring-webflux/spring-webflux-5.0/src/test/groovy/dd/trace/instrumentation/springwebflux/client/SpringWebfluxHttpClientBase.groovy index 3b8fca1f2a9..066f37474e6 100644 --- a/dd-java-agent/instrumentation/spring/spring-webflux/spring-webflux-5.0/src/test/groovy/dd/trace/instrumentation/springwebflux/client/SpringWebfluxHttpClientBase.groovy +++ b/dd-java-agent/instrumentation/spring/spring-webflux/spring-webflux-5.0/src/test/groovy/dd/trace/instrumentation/springwebflux/client/SpringWebfluxHttpClientBase.groovy @@ -9,6 +9,7 @@ import datadog.trace.bootstrap.instrumentation.api.Tags import datadog.trace.bootstrap.instrumentation.api.URIUtils import datadog.trace.instrumentation.netty41.client.NettyHttpClientDecorator import datadog.trace.instrumentation.springwebflux.client.SpringWebfluxHttpClientDecorator +import io.netty.channel.Channel import org.springframework.http.HttpMethod import org.springframework.web.reactive.function.client.ClientRequest import org.springframework.web.reactive.function.client.ClientResponse @@ -44,7 +45,28 @@ abstract class SpringWebfluxHttpClientBase extends HttpClientTest implements Tes check() - response.statusCode().value() + consumeResponse(response) + } + + protected static int consumeResponse(ClientResponse response) { + int status = response.statusCode().value() + Channel channel = responseChannel(response) + response.bodyToMono(Void).block() + // Body completion can occur inside channelRead, before channelReadComplete is fired. + channel.eventLoop().submit({} as Runnable).sync() + return status + } + + private static Channel responseChannel(ClientResponse response) { + def clientHttpResponse = field(response, "response") + def nettyResponse = field(clientHttpResponse, "response") + return nettyResponse.channel() + } + + private static Object field(Object target, String name) { + def field = target.class.getDeclaredField(name) + field.accessible = true + return field.get(target) } @Override diff --git a/dd-java-agent/instrumentation/spring/spring-webflux/spring-webflux-5.0/src/test/groovy/dd/trace/instrumentation/springwebflux/client/SpringWebfluxHttpClientDoAfterTerminateTest.groovy b/dd-java-agent/instrumentation/spring/spring-webflux/spring-webflux-5.0/src/test/groovy/dd/trace/instrumentation/springwebflux/client/SpringWebfluxHttpClientDoAfterTerminateTest.groovy index c03f9a18082..8355d53d18d 100644 --- a/dd-java-agent/instrumentation/spring/spring-webflux/spring-webflux-5.0/src/test/groovy/dd/trace/instrumentation/springwebflux/client/SpringWebfluxHttpClientDoAfterTerminateTest.groovy +++ b/dd-java-agent/instrumentation/spring/spring-webflux/spring-webflux-5.0/src/test/groovy/dd/trace/instrumentation/springwebflux/client/SpringWebfluxHttpClientDoAfterTerminateTest.groovy @@ -29,7 +29,7 @@ class SpringWebfluxHttpClientDoAfterTerminateTest extends SpringWebfluxHttpClien check() - response.statusCode().value() + consumeResponse(response) } @Override diff --git a/dd-java-agent/instrumentation/spring/spring-webflux/spring-webflux-5.0/src/test/groovy/dd/trace/instrumentation/springwebflux/client/SpringWebfluxHttpClientDoOnSuccessOrErrorTest.groovy b/dd-java-agent/instrumentation/spring/spring-webflux/spring-webflux-5.0/src/test/groovy/dd/trace/instrumentation/springwebflux/client/SpringWebfluxHttpClientDoOnSuccessOrErrorTest.groovy index eb99848e9d5..366e87d6df1 100644 --- a/dd-java-agent/instrumentation/spring/spring-webflux/spring-webflux-5.0/src/test/groovy/dd/trace/instrumentation/springwebflux/client/SpringWebfluxHttpClientDoOnSuccessOrErrorTest.groovy +++ b/dd-java-agent/instrumentation/spring/spring-webflux/spring-webflux-5.0/src/test/groovy/dd/trace/instrumentation/springwebflux/client/SpringWebfluxHttpClientDoOnSuccessOrErrorTest.groovy @@ -29,7 +29,7 @@ abstract class SpringWebfluxHttpClientDoOnSuccessOrErrorTest extends SpringWebfl check() - response.statusCode().value() + consumeResponse(response) } @Override diff --git a/dd-java-agent/instrumentation/spring/spring-webflux/spring-webflux-5.0/src/test/groovy/dd/trace/instrumentation/springwebflux/client/SpringWebfluxHttpClientDoOnSuccessTest.groovy b/dd-java-agent/instrumentation/spring/spring-webflux/spring-webflux-5.0/src/test/groovy/dd/trace/instrumentation/springwebflux/client/SpringWebfluxHttpClientDoOnSuccessTest.groovy index 82efca874c9..d6fbb96bc75 100644 --- a/dd-java-agent/instrumentation/spring/spring-webflux/spring-webflux-5.0/src/test/groovy/dd/trace/instrumentation/springwebflux/client/SpringWebfluxHttpClientDoOnSuccessTest.groovy +++ b/dd-java-agent/instrumentation/spring/spring-webflux/spring-webflux-5.0/src/test/groovy/dd/trace/instrumentation/springwebflux/client/SpringWebfluxHttpClientDoOnSuccessTest.groovy @@ -29,7 +29,7 @@ class SpringWebfluxHttpClientDoOnSuccessTest extends SpringWebfluxHttpClientBase check() - response.statusCode().value() + consumeResponse(response) } @Override diff --git a/dd-java-agent/instrumentation/spring/spring-webflux/spring-webflux-5.0/src/test/groovy/dd/trace/instrumentation/springwebflux/client/SpringWebfluxHttpClientDoOnTerminateTest.groovy b/dd-java-agent/instrumentation/spring/spring-webflux/spring-webflux-5.0/src/test/groovy/dd/trace/instrumentation/springwebflux/client/SpringWebfluxHttpClientDoOnTerminateTest.groovy index d7a0620dca9..b6b05def34d 100644 --- a/dd-java-agent/instrumentation/spring/spring-webflux/spring-webflux-5.0/src/test/groovy/dd/trace/instrumentation/springwebflux/client/SpringWebfluxHttpClientDoOnTerminateTest.groovy +++ b/dd-java-agent/instrumentation/spring/spring-webflux/spring-webflux-5.0/src/test/groovy/dd/trace/instrumentation/springwebflux/client/SpringWebfluxHttpClientDoOnTerminateTest.groovy @@ -29,7 +29,7 @@ class SpringWebfluxHttpClientDoOnTerminateTest extends SpringWebfluxHttpClientBa check() - response.statusCode().value() + consumeResponse(response) } @Override