Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions dd-trace-core/build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -120,6 +120,8 @@ dependencies {
testImplementation group: 'commons-codec', name: 'commons-codec', version: '1.3'
testImplementation group: 'com.amazonaws', name: 'aws-lambda-java-events', version:'3.11.0'
testImplementation group: 'com.google.protobuf', name: 'protobuf-java', version: '3.14.0'
testImplementation libs.jnr.unixsocket
testImplementation libs.okhttp3.mockwebserver
testImplementation libs.testcontainers
testImplementation project(':utils:test-junit-utils')
testImplementation project(':utils:test-junit-converter-utils')
Expand Down
25 changes: 13 additions & 12 deletions dd-trace-core/gradle.lockfile
Original file line number Diff line number Diff line change
Expand Up @@ -21,14 +21,14 @@ com.github.docker-java:docker-java-api:3.4.2=jmhRuntimeClasspath,testCompileClas
com.github.docker-java:docker-java-transport-zerodep:3.4.2=jmhRuntimeClasspath,testCompileClasspath,testRuntimeClasspath,traceAgentTestCompileClasspath,traceAgentTestRuntimeClasspath
com.github.docker-java:docker-java-transport:3.4.2=jmhRuntimeClasspath,testCompileClasspath,testRuntimeClasspath,traceAgentTestCompileClasspath,traceAgentTestRuntimeClasspath
com.github.javaparser:javaparser-core:3.25.6=codenarc
com.github.jnr:jffi:1.3.15=jmhRuntimeClasspath,runtimeClasspath,testRuntimeClasspath,traceAgentTestRuntimeClasspath
com.github.jnr:jnr-a64asm:1.0.0=jmhRuntimeClasspath,runtimeClasspath,testRuntimeClasspath,traceAgentTestRuntimeClasspath
com.github.jnr:jnr-constants:0.10.4=jmhRuntimeClasspath,runtimeClasspath,testRuntimeClasspath,traceAgentTestRuntimeClasspath
com.github.jnr:jnr-enxio:0.32.20=jmhRuntimeClasspath,runtimeClasspath,testRuntimeClasspath,traceAgentTestRuntimeClasspath
com.github.jnr:jnr-ffi:2.2.19=jmhRuntimeClasspath,runtimeClasspath,testRuntimeClasspath,traceAgentTestRuntimeClasspath
com.github.jnr:jnr-posix:3.1.22=jmhRuntimeClasspath,runtimeClasspath,testRuntimeClasspath,traceAgentTestRuntimeClasspath
com.github.jnr:jnr-unixsocket:0.38.25=jmhRuntimeClasspath,runtimeClasspath,testRuntimeClasspath,traceAgentTestRuntimeClasspath
com.github.jnr:jnr-x86asm:1.0.2=jmhRuntimeClasspath,runtimeClasspath,testRuntimeClasspath,traceAgentTestRuntimeClasspath
com.github.jnr:jffi:1.3.15=jmhRuntimeClasspath,runtimeClasspath,testCompileClasspath,testRuntimeClasspath,traceAgentTestCompileClasspath,traceAgentTestRuntimeClasspath
com.github.jnr:jnr-a64asm:1.0.0=jmhRuntimeClasspath,runtimeClasspath,testCompileClasspath,testRuntimeClasspath,traceAgentTestCompileClasspath,traceAgentTestRuntimeClasspath
com.github.jnr:jnr-constants:0.10.4=jmhRuntimeClasspath,runtimeClasspath,testCompileClasspath,testRuntimeClasspath,traceAgentTestCompileClasspath,traceAgentTestRuntimeClasspath
com.github.jnr:jnr-enxio:0.32.20=jmhRuntimeClasspath,runtimeClasspath,testCompileClasspath,testRuntimeClasspath,traceAgentTestCompileClasspath,traceAgentTestRuntimeClasspath
com.github.jnr:jnr-ffi:2.2.19=jmhRuntimeClasspath,runtimeClasspath,testCompileClasspath,testRuntimeClasspath,traceAgentTestCompileClasspath,traceAgentTestRuntimeClasspath
com.github.jnr:jnr-posix:3.1.22=jmhRuntimeClasspath,runtimeClasspath,testCompileClasspath,testRuntimeClasspath,traceAgentTestCompileClasspath,traceAgentTestRuntimeClasspath
com.github.jnr:jnr-unixsocket:0.38.25=jmhRuntimeClasspath,runtimeClasspath,testCompileClasspath,testRuntimeClasspath,traceAgentTestCompileClasspath,traceAgentTestRuntimeClasspath
com.github.jnr:jnr-x86asm:1.0.2=jmhRuntimeClasspath,runtimeClasspath,testCompileClasspath,testRuntimeClasspath,traceAgentTestCompileClasspath,traceAgentTestRuntimeClasspath
com.github.spotbugs:spotbugs-annotations:4.10.4=compileClasspath,jmhCompileClasspath,spotbugs
com.github.spotbugs:spotbugs:4.10.4=spotbugs
com.github.stephenc.jcip:jcip-annotations:1.0-1=spotbugs
Expand All @@ -51,6 +51,7 @@ com.google.protobuf:protobuf-java:3.14.0=jmhRuntimeClasspath,testCompileClasspat
com.google.re2j:re2j:1.8=compileClasspath,jmhCompileClasspath,jmhRuntimeClasspath,runtimeClasspath,testCompileClasspath,testRuntimeClasspath,traceAgentTestCompileClasspath,traceAgentTestRuntimeClasspath
com.squareup.moshi:moshi:1.11.0=compileClasspath,jmhCompileClasspath,jmhRuntimeClasspath,runtimeClasspath,testCompileClasspath,testRuntimeClasspath,traceAgentTestCompileClasspath,traceAgentTestRuntimeClasspath
com.squareup.okhttp3:logging-interceptor:3.12.12=jmhRuntimeClasspath,testCompileClasspath,testRuntimeClasspath,traceAgentTestCompileClasspath,traceAgentTestRuntimeClasspath
com.squareup.okhttp3:mockwebserver:3.12.12=jmhRuntimeClasspath,testCompileClasspath,testRuntimeClasspath,traceAgentTestCompileClasspath,traceAgentTestRuntimeClasspath
com.squareup.okhttp3:okhttp:3.12.12=jmhRuntimeClasspath,testCompileClasspath,testRuntimeClasspath,traceAgentTestCompileClasspath,traceAgentTestRuntimeClasspath
com.squareup.okio:okio:1.17.5=compileClasspath,jmhCompileClasspath,jmhRuntimeClasspath,runtimeClasspath,testCompileClasspath,testRuntimeClasspath,traceAgentTestCompileClasspath,traceAgentTestRuntimeClasspath
com.thoughtworks.qdox:qdox:1.12.1=codenarc
Expand Down Expand Up @@ -127,13 +128,13 @@ org.openjdk.jmh:jmh-generator-reflection:1.37=jmh,jmhCompileClasspath,jmhRuntime
org.openjdk.jol:jol-core:0.17=jmhRuntimeClasspath,testCompileClasspath,testRuntimeClasspath,traceAgentTestCompileClasspath,traceAgentTestRuntimeClasspath
org.opentest4j:opentest4j:1.3.0=jmhRuntimeClasspath,testCompileClasspath,testRuntimeClasspath,traceAgentTestCompileClasspath,traceAgentTestRuntimeClasspath
org.ow2.asm:asm-analysis:9.10.1=spotbugs
org.ow2.asm:asm-analysis:9.7.1=jmhRuntimeClasspath,runtimeClasspath,testRuntimeClasspath,traceAgentTestRuntimeClasspath
org.ow2.asm:asm-analysis:9.7.1=jmhRuntimeClasspath,runtimeClasspath,testCompileClasspath,testRuntimeClasspath,traceAgentTestCompileClasspath,traceAgentTestRuntimeClasspath
org.ow2.asm:asm-commons:9.10.1=jacocoAnt,spotbugs
org.ow2.asm:asm-commons:9.7.1=jmhRuntimeClasspath,runtimeClasspath,testRuntimeClasspath,traceAgentTestRuntimeClasspath
org.ow2.asm:asm-commons:9.7.1=jmhRuntimeClasspath,runtimeClasspath,testCompileClasspath,testRuntimeClasspath,traceAgentTestCompileClasspath,traceAgentTestRuntimeClasspath
org.ow2.asm:asm-tree:9.10.1=jacocoAnt,spotbugs
org.ow2.asm:asm-tree:9.7.1=jmhRuntimeClasspath,runtimeClasspath,testRuntimeClasspath,traceAgentTestRuntimeClasspath
org.ow2.asm:asm-tree:9.7.1=jmhRuntimeClasspath,runtimeClasspath,testCompileClasspath,testRuntimeClasspath,traceAgentTestCompileClasspath,traceAgentTestRuntimeClasspath
org.ow2.asm:asm-util:9.10.1=spotbugs
org.ow2.asm:asm-util:9.7.1=jmhRuntimeClasspath,runtimeClasspath,testRuntimeClasspath,traceAgentTestRuntimeClasspath
org.ow2.asm:asm-util:9.7.1=jmhRuntimeClasspath,runtimeClasspath,testCompileClasspath,testRuntimeClasspath,traceAgentTestCompileClasspath,traceAgentTestRuntimeClasspath
org.ow2.asm:asm:9.0=jmh,jmhCompileClasspath
org.ow2.asm:asm:9.10.1=jacocoAnt,jmhRuntimeClasspath,spotbugs,testCompileClasspath,testRuntimeClasspath,traceAgentTestCompileClasspath,traceAgentTestRuntimeClasspath
org.ow2.asm:asm:9.7.1=runtimeClasspath
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,109 @@
package datadog.common.socket;

import java.io.File;
import java.io.IOException;
import java.net.InetAddress;
import java.net.InetSocketAddress;
import java.net.ServerSocket;
import java.net.Socket;
import java.net.SocketAddress;
import java.net.SocketException;
import java.nio.channels.ClosedChannelException;
import javax.net.ServerSocketFactory;
import jnr.unixsocket.UnixServerSocketChannel;
import jnr.unixsocket.UnixSocketAddress;
import jnr.unixsocket.UnixSocketChannel;

/**
* Adapts a JNR Unix-domain server channel to APIs such as MockWebServer that require a {@link
* ServerSocket}. Adapted from OkHttp's <a
* href="https://github.com/square/okhttp/blob/master/samples/unixdomainsockets/src/main/java/okhttp3/unixdomainsockets/UnixDomainServerSocketFactory.java">Unix-domain
* socket sample</a>.
*/
public final class UnixDomainServerSocketFactory extends ServerSocketFactory {
private final File path;

public UnixDomainServerSocketFactory(File path) {
this.path = path;
}

@Override
public ServerSocket createServerSocket() throws IOException {
return new UnixDomainServerSocket();
}

@Override
public ServerSocket createServerSocket(int port) throws IOException {
return createServerSocket();
}

@Override
public ServerSocket createServerSocket(int port, int backlog) throws IOException {
return createServerSocket();
}

@Override
public ServerSocket createServerSocket(int port, int backlog, InetAddress inetAddress)
throws IOException {
return createServerSocket();
}

private final class UnixDomainServerSocket extends ServerSocket {
private UnixServerSocketChannel serverSocketChannel;
private InetSocketAddress endpoint;

private UnixDomainServerSocket() throws IOException {}

@Override
public void bind(SocketAddress endpoint, int backlog) throws IOException {
this.endpoint = (InetSocketAddress) endpoint;
UnixServerSocketChannel channel = UnixServerSocketChannel.open();
boolean bound = false;
try {
channel.configureBlocking(true);
channel.socket().bind(new UnixSocketAddress(path));
serverSocketChannel = channel;
bound = true;
} finally {
if (!bound) {
channel.close();
}
}
}

@Override
public void setReuseAddress(boolean on) {
// MockWebServer configures this TCP option before binding. It has no UDS equivalent.
}

@Override
public int getLocalPort() {
return 1; // MockWebServer requires a port even though a UDS has none.
}

@Override
public SocketAddress getLocalSocketAddress() {
return endpoint;
}

@Override
public Socket accept() throws IOException {
try {
UnixSocketChannel channel = serverSocketChannel.accept();
return new TunnelingUnixSocket(path, channel, endpoint);
} catch (ClosedChannelException e) {
SocketException socketException = new SocketException("Socket is closed");
socketException.initCause(e);
throw socketException;
}
}

@Override
public void close() throws IOException {
super.close();
if (serverSocketChannel != null) {
serverSocketChannel.close();
}
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -2,9 +2,17 @@

import static datadog.trace.api.ProtocolVersion.V0_5;
import static datadog.trace.api.config.GeneralConfig.EXPERIMENTAL_PROPAGATE_PROCESS_TAGS_ENABLED;
import static datadog.trace.api.config.GeneralConfig.JDK_SOCKET_ENABLED;
import static datadog.trace.common.writer.ddagent.Prioritization.ENSURE_TRACE;
import static okhttp3.mockwebserver.SocketPolicy.DISCONNECT_AT_END;
import static okhttp3.mockwebserver.SocketPolicy.NO_RESPONSE;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertSame;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.junit.jupiter.api.condition.JRE.JAVA_16;
import static org.junit.jupiter.api.condition.OS.LINUX;
import static org.junit.jupiter.api.condition.OS.MAC;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyInt;
import static org.mockito.ArgumentMatchers.anyLong;
Expand All @@ -18,6 +26,7 @@
import static org.mockito.Mockito.verifyNoMoreInteractions;
import static org.mockito.Mockito.when;

import datadog.common.socket.UnixDomainServerSocketFactory;
import datadog.communication.ddagent.DDAgentFeaturesDiscovery;
import datadog.communication.http.OkHttpUtils;
import datadog.communication.serialization.FlushingBuffer;
Expand All @@ -39,18 +48,26 @@
import datadog.trace.test.junit.utils.config.WithConfig;
import datadog.trace.test.util.Flaky;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.Collections;
import java.util.List;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.Phaser;
import java.util.concurrent.Semaphore;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference;
import okhttp3.HttpUrl;
import okhttp3.mockwebserver.MockResponse;
import okhttp3.mockwebserver.MockWebServer;
import okhttp3.mockwebserver.RecordedRequest;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.Timeout;
import org.junit.jupiter.api.condition.EnabledForJreRange;
import org.junit.jupiter.api.condition.EnabledOnOs;
import org.mockito.Mockito;
import org.tabletest.junit.TableTest;

Expand Down Expand Up @@ -339,12 +356,11 @@ void monitorHappyPath(String agentVersion) {
List<DDSpan> minimalTrace = createMinimalTrace();

// DQH -- need to set-up a dummy agent for the final send callback to work
JavaTestHttpServer agent =
try (JavaTestHttpServer agent =
JavaTestHttpServer.httpServer(
server ->
server.handlers(
h -> h.put(agentVersion, api -> api.getResponse().status(200).send())));
try {
h -> h.put(agentVersion, api -> api.getResponse().status(200).send())))) {
HttpUrl agentUrl = HttpUrl.get(agent.getAddress());
okhttp3.OkHttpClient client = OkHttpUtils.buildHttpClient(agentUrl, 1000);
DDAgentFeaturesDiscovery discovery =
Expand Down Expand Up @@ -383,8 +399,6 @@ void monitorHappyPath(String agentVersion) {
writer.close();

verify(healthMetrics, times(1)).onShutdown(true);
} finally {
agent.close();
}
}

Expand All @@ -399,25 +413,11 @@ void monitorAgentReturnsError(String agentVersion) {
List<DDSpan> minimalTrace = createMinimalTrace();

// DQH -- need to set-up a dummy agent for the final send callback to work
final boolean[] first = {true};
JavaTestHttpServer agent =
try (JavaTestHttpServer agent =
JavaTestHttpServer.httpServer(
server ->
server.handlers(
h ->
h.put(
agentVersion,
api -> {
// DQH - DDApi sniffs for end point existence, so respond with 200 the
// first time
if (first[0]) {
api.getResponse().status(200).send();
first[0] = false;
} else {
api.getResponse().status(500).send();
}
})));
try {
h -> h.put(agentVersion, api -> api.getResponse().status(500).send())))) {
HttpUrl agentUrl = HttpUrl.get(agent.getAddress());
okhttp3.OkHttpClient client = OkHttpUtils.buildHttpClient(agentUrl, 1000);
DDAgentFeaturesDiscovery discovery =
Expand Down Expand Up @@ -456,8 +456,86 @@ void monitorAgentReturnsError(String agentVersion) {
writer.close();

verify(healthMetrics, times(1)).onShutdown(true);
}
}

@Test
@WithConfig(key = JDK_SOCKET_ENABLED, value = "true")
@EnabledForJreRange(min = JAVA_16)
@EnabledOnOs({LINUX, MAC})
void unixSocketTimeoutKeepsWorkerAliveAndReconnects() throws Exception {
assertTrue(Config.get().isJdkSocketEnabled());

Path socketPath = Files.createTempFile("dd-trace-agent-", ".sock");
Files.delete(socketPath);

HealthMetrics healthMetrics = mock(HealthMetrics.class);
AtomicReference<Thread> failedSendThread = new AtomicReference<>();
AtomicReference<Thread> successfulSendThread = new AtomicReference<>();
CountDownLatch failedSend = new CountDownLatch(1);
CountDownLatch successfulSend = new CountDownLatch(1);
doAnswer(
invocation -> {
failedSendThread.set(Thread.currentThread());
failedSend.countDown();
return null;
})
.when(healthMetrics)
.onFailedSend(anyInt(), anyInt(), any());
doAnswer(
invocation -> {
successfulSendThread.set(Thread.currentThread());
successfulSend.countDown();
return null;
})
.when(healthMetrics)
.onSend(anyInt(), anyInt(), any());

DDAgentFeaturesDiscovery discovery = mock(DDAgentFeaturesDiscovery.class);
when(discovery.getTraceEndpoint()).thenReturn("v0.4/traces");

try (MockWebServer server = new MockWebServer();
DDAgentWriter writer =
DDAgentWriter.builder()
.featureDiscovery(discovery)
.unixDomainSocket(socketPath.toString())
.timeoutMillis(500)
.monitoring(monitoring)
.healthMetrics(healthMetrics)
.flushIntervalMilliseconds(10)
.flushTimeout(5, TimeUnit.SECONDS)
.build()) {
server.setServerSocketFactory(new UnixDomainServerSocketFactory(socketPath.toFile()));
// Read the first request fully, then withhold the response to trigger a header-read timeout.
server.enqueue(new MockResponse().setSocketPolicy(NO_RESPONSE));
server.enqueue(new MockResponse().setResponseCode(200).setSocketPolicy(DISCONNECT_AT_END));
server.start();

writer.start();

writer.write(createMinimalTrace());
assertTrue(
failedSend.await(5, TimeUnit.SECONDS), "The periodic flush did not report failure");

RecordedRequest failedRequest = server.takeRequest(5, TimeUnit.SECONDS);
assertNotNull(failedRequest);
assertEquals(0, failedRequest.getSequenceNumber());
assertEquals(1, server.getRequestCount());
verify(healthMetrics, times(1)).onFailedSend(anyInt(), anyInt(), any());

writer.write(createMinimalTrace());
assertTrue(
successfulSend.await(5, TimeUnit.SECONDS), "The worker did not send the next payload");

RecordedRequest successfulRequest = server.takeRequest(5, TimeUnit.SECONDS);
assertNotNull(successfulRequest);
// Sequence numbers are per connection; zero again proves that this used a fresh socket.
assertEquals(0, successfulRequest.getSequenceNumber());
assertEquals(2, server.getRequestCount());
verify(healthMetrics, times(1)).onSend(anyInt(), anyInt(), any());
assertSame(failedSendThread.get(), successfulSendThread.get());
} finally {
agent.close();
Files.deleteIfExists(socketPath);
}
}

Expand Down
12 changes: 11 additions & 1 deletion utils/socket-utils/build.gradle.kts
Original file line number Diff line number Diff line change
Expand Up @@ -4,21 +4,31 @@ plugins {
`java-library`
idea
id("dd-trace-java.module.internal-library")
id("dd-trace-java.jmh-conventions")
}

extensions.getByName("tracerJava").withGroovyBuilder {
invokeMethod("addSourceSetFor", arrayOf(JavaVersion.VERSION_17, mapOf("compileOnly" to true)))
}

dependencies {
add("main_java17CompileOnly", project(":components:annotations"))
implementation(project(":components:environment"))
implementation(project(":utils:logging-utils"))
implementation(libs.slf4j)
implementation(libs.jnr.unixsocket)
testImplementation(files(sourceSets["main_java17"].output))
jmhImplementation(files(sourceSets["main_java17"].output))
}

listOf("compileMain_java17Java", "compileTestJava").forEach {
jmh {
jmhVersion = libs.versions.jmh.get()
includeTests = false
resultFormat = "JSON"
failOnError = true
}

listOf("compileMain_java17Java", "compileTestJava", "compileJmhJava").forEach {
tasks.named<JavaCompile>(it) {
// The Java 17 implementation can lift this offset, but compileTestJava must first be split if
// the remaining socket tests still need to run on Java 8.
Expand Down
Loading
Loading