Identifying otel http calls (#5918)
Co-authored-by: jack-berg <34418638+jack-berg@users.noreply.github.com>
This commit is contained in:
parent
f9be6821a5
commit
f99e4961cb
|
@ -0,0 +1,40 @@
|
||||||
|
/*
|
||||||
|
* Copyright The OpenTelemetry Authors
|
||||||
|
* SPDX-License-Identifier: Apache-2.0
|
||||||
|
*/
|
||||||
|
|
||||||
|
package io.opentelemetry.exporter.internal;
|
||||||
|
|
||||||
|
import io.opentelemetry.context.Context;
|
||||||
|
import io.opentelemetry.context.ContextKey;
|
||||||
|
import java.util.Objects;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* This class is internal and is hence not for public use. Its APIs are unstable and can change at
|
||||||
|
* any time.
|
||||||
|
*/
|
||||||
|
public final class InstrumentationUtil {
|
||||||
|
private static final ContextKey<Boolean> SUPPRESS_INSTRUMENTATION_KEY =
|
||||||
|
ContextKey.named("suppress_internal_exporter_instrumentation");
|
||||||
|
|
||||||
|
private InstrumentationUtil() {}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Adds a Context boolean key that will allow to identify HTTP calls coming from OTel exporters.
|
||||||
|
* The key later be checked by an automatic instrumentation to avoid tracing OTel exporter's
|
||||||
|
* calls.
|
||||||
|
*/
|
||||||
|
public static void suppressInstrumentation(Runnable runnable) {
|
||||||
|
Context.current().with(SUPPRESS_INSTRUMENTATION_KEY, true).wrap(runnable).run();
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Checks if an automatic instrumentation should be suppressed with the provided Context.
|
||||||
|
*
|
||||||
|
* @return TRUE to suppress the automatic instrumentation, FALSE to continue with the
|
||||||
|
* instrumentation.
|
||||||
|
*/
|
||||||
|
public static boolean shouldSuppressInstrumentation(Context context) {
|
||||||
|
return Objects.equals(context.get(SUPPRESS_INSTRUMENTATION_KEY), true);
|
||||||
|
}
|
||||||
|
}
|
|
@ -0,0 +1,27 @@
|
||||||
|
/*
|
||||||
|
* Copyright The OpenTelemetry Authors
|
||||||
|
* SPDX-License-Identifier: Apache-2.0
|
||||||
|
*/
|
||||||
|
|
||||||
|
package io.opentelemetry.exporter.internal;
|
||||||
|
|
||||||
|
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||||
|
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||||
|
|
||||||
|
import io.opentelemetry.context.Context;
|
||||||
|
import org.junit.jupiter.api.Test;
|
||||||
|
|
||||||
|
class InstrumentationUtilTest {
|
||||||
|
@Test
|
||||||
|
void verifySuppressInstrumentation() {
|
||||||
|
// Should be false by default.
|
||||||
|
assertFalse(InstrumentationUtil.shouldSuppressInstrumentation(Context.current()));
|
||||||
|
|
||||||
|
// Should be true inside the Runnable passed to InstrumentationUtil.suppressInstrumentation.
|
||||||
|
InstrumentationUtil.suppressInstrumentation(
|
||||||
|
() -> assertTrue(InstrumentationUtil.shouldSuppressInstrumentation(Context.current())));
|
||||||
|
|
||||||
|
// Should be false after the runnable finishes.
|
||||||
|
assertFalse(InstrumentationUtil.shouldSuppressInstrumentation(Context.current()));
|
||||||
|
}
|
||||||
|
}
|
|
@ -23,6 +23,7 @@
|
||||||
|
|
||||||
package io.opentelemetry.exporter.sender.okhttp.internal;
|
package io.opentelemetry.exporter.sender.okhttp.internal;
|
||||||
|
|
||||||
|
import io.opentelemetry.exporter.internal.InstrumentationUtil;
|
||||||
import io.opentelemetry.exporter.internal.RetryUtil;
|
import io.opentelemetry.exporter.internal.RetryUtil;
|
||||||
import io.opentelemetry.exporter.internal.grpc.GrpcExporterUtil;
|
import io.opentelemetry.exporter.internal.grpc.GrpcExporterUtil;
|
||||||
import io.opentelemetry.exporter.internal.grpc.GrpcResponse;
|
import io.opentelemetry.exporter.internal.grpc.GrpcResponse;
|
||||||
|
@ -112,51 +113,53 @@ public final class OkHttpGrpcSender<T extends Marshaler> implements GrpcSender<T
|
||||||
RequestBody requestBody = new GrpcRequestBody(request, compressionEnabled);
|
RequestBody requestBody = new GrpcRequestBody(request, compressionEnabled);
|
||||||
requestBuilder.post(requestBody);
|
requestBuilder.post(requestBody);
|
||||||
|
|
||||||
client
|
InstrumentationUtil.suppressInstrumentation(
|
||||||
.newCall(requestBuilder.build())
|
() ->
|
||||||
.enqueue(
|
client
|
||||||
new Callback() {
|
.newCall(requestBuilder.build())
|
||||||
@Override
|
.enqueue(
|
||||||
public void onFailure(Call call, IOException e) {
|
new Callback() {
|
||||||
String description = e.getMessage();
|
@Override
|
||||||
if (description == null) {
|
public void onFailure(Call call, IOException e) {
|
||||||
description = "";
|
String description = e.getMessage();
|
||||||
}
|
if (description == null) {
|
||||||
onError.accept(GrpcResponse.create(2 /* UNKNOWN */, description), e);
|
description = "";
|
||||||
}
|
}
|
||||||
|
onError.accept(GrpcResponse.create(2 /* UNKNOWN */, description), e);
|
||||||
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public void onResponse(Call call, Response response) {
|
public void onResponse(Call call, Response response) {
|
||||||
// Response body is empty but must be consumed to access trailers.
|
// Response body is empty but must be consumed to access trailers.
|
||||||
try {
|
try {
|
||||||
response.body().bytes();
|
response.body().bytes();
|
||||||
} catch (IOException e) {
|
} catch (IOException e) {
|
||||||
onError.accept(
|
onError.accept(
|
||||||
GrpcResponse.create(
|
GrpcResponse.create(
|
||||||
GrpcExporterUtil.GRPC_STATUS_UNKNOWN,
|
GrpcExporterUtil.GRPC_STATUS_UNKNOWN,
|
||||||
"Could not consume server response."),
|
"Could not consume server response."),
|
||||||
e);
|
e);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
String status = grpcStatus(response);
|
String status = grpcStatus(response);
|
||||||
if ("0".equals(status)) {
|
if ("0".equals(status)) {
|
||||||
onSuccess.run();
|
onSuccess.run();
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
String errorMessage = grpcMessage(response);
|
String errorMessage = grpcMessage(response);
|
||||||
int statusCode;
|
int statusCode;
|
||||||
try {
|
try {
|
||||||
statusCode = Integer.parseInt(status);
|
statusCode = Integer.parseInt(status);
|
||||||
} catch (NumberFormatException ex) {
|
} catch (NumberFormatException ex) {
|
||||||
statusCode = GrpcExporterUtil.GRPC_STATUS_UNKNOWN;
|
statusCode = GrpcExporterUtil.GRPC_STATUS_UNKNOWN;
|
||||||
}
|
}
|
||||||
onError.accept(
|
onError.accept(
|
||||||
GrpcResponse.create(statusCode, errorMessage),
|
GrpcResponse.create(statusCode, errorMessage),
|
||||||
new IllegalStateException(errorMessage));
|
new IllegalStateException(errorMessage));
|
||||||
}
|
}
|
||||||
});
|
}));
|
||||||
}
|
}
|
||||||
|
|
||||||
@Nullable
|
@Nullable
|
||||||
|
|
|
@ -5,6 +5,7 @@
|
||||||
|
|
||||||
package io.opentelemetry.exporter.sender.okhttp.internal;
|
package io.opentelemetry.exporter.sender.okhttp.internal;
|
||||||
|
|
||||||
|
import io.opentelemetry.exporter.internal.InstrumentationUtil;
|
||||||
import io.opentelemetry.exporter.internal.RetryUtil;
|
import io.opentelemetry.exporter.internal.RetryUtil;
|
||||||
import io.opentelemetry.exporter.internal.auth.Authenticator;
|
import io.opentelemetry.exporter.internal.auth.Authenticator;
|
||||||
import io.opentelemetry.exporter.internal.http.HttpSender;
|
import io.opentelemetry.exporter.internal.http.HttpSender;
|
||||||
|
@ -101,38 +102,40 @@ public final class OkHttpHttpSender implements HttpSender {
|
||||||
requestBuilder.post(body);
|
requestBuilder.post(body);
|
||||||
}
|
}
|
||||||
|
|
||||||
client
|
InstrumentationUtil.suppressInstrumentation(
|
||||||
.newCall(requestBuilder.build())
|
() ->
|
||||||
.enqueue(
|
client
|
||||||
new Callback() {
|
.newCall(requestBuilder.build())
|
||||||
@Override
|
.enqueue(
|
||||||
public void onFailure(Call call, IOException e) {
|
new Callback() {
|
||||||
onError.accept(e);
|
@Override
|
||||||
}
|
public void onFailure(Call call, IOException e) {
|
||||||
|
onError.accept(e);
|
||||||
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public void onResponse(Call call, okhttp3.Response response) {
|
public void onResponse(Call call, okhttp3.Response response) {
|
||||||
try (ResponseBody body = response.body()) {
|
try (ResponseBody body = response.body()) {
|
||||||
onResponse.accept(
|
onResponse.accept(
|
||||||
new Response() {
|
new Response() {
|
||||||
@Override
|
@Override
|
||||||
public int statusCode() {
|
public int statusCode() {
|
||||||
return response.code();
|
return response.code();
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public String statusMessage() {
|
public String statusMessage() {
|
||||||
return response.message();
|
return response.message();
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public byte[] responseBody() throws IOException {
|
public byte[] responseBody() throws IOException {
|
||||||
return body.bytes();
|
return body.bytes();
|
||||||
|
}
|
||||||
|
});
|
||||||
}
|
}
|
||||||
});
|
}
|
||||||
}
|
}));
|
||||||
}
|
|
||||||
});
|
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
|
|
|
@ -18,6 +18,13 @@ import okhttp3.Dispatcher;
|
||||||
* at any time.
|
* at any time.
|
||||||
*/
|
*/
|
||||||
public final class OkHttpUtil {
|
public final class OkHttpUtil {
|
||||||
|
@SuppressWarnings("NonFinalStaticField")
|
||||||
|
private static boolean propagateContextForTestingInDispatcher = false;
|
||||||
|
|
||||||
|
public static void setPropagateContextForTestingInDispatcher(
|
||||||
|
boolean propagateContextForTestingInDispatcher) {
|
||||||
|
OkHttpUtil.propagateContextForTestingInDispatcher = propagateContextForTestingInDispatcher;
|
||||||
|
}
|
||||||
|
|
||||||
/** Returns a {@link Dispatcher} using daemon threads, otherwise matching the OkHttp default. */
|
/** Returns a {@link Dispatcher} using daemon threads, otherwise matching the OkHttp default. */
|
||||||
public static Dispatcher newDispatcher() {
|
public static Dispatcher newDispatcher() {
|
||||||
|
@ -28,7 +35,7 @@ public final class OkHttpUtil {
|
||||||
60,
|
60,
|
||||||
TimeUnit.SECONDS,
|
TimeUnit.SECONDS,
|
||||||
new SynchronousQueue<>(),
|
new SynchronousQueue<>(),
|
||||||
new DaemonThreadFactory("okhttp-dispatch")));
|
new DaemonThreadFactory("okhttp-dispatch", propagateContextForTestingInDispatcher)));
|
||||||
}
|
}
|
||||||
|
|
||||||
private OkHttpUtil() {}
|
private OkHttpUtil() {}
|
||||||
|
|
|
@ -0,0 +1,58 @@
|
||||||
|
/*
|
||||||
|
* Copyright The OpenTelemetry Authors
|
||||||
|
* SPDX-License-Identifier: Apache-2.0
|
||||||
|
*/
|
||||||
|
|
||||||
|
package io.opentelemetry.exporter.sender.okhttp.internal;
|
||||||
|
|
||||||
|
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||||
|
|
||||||
|
import io.opentelemetry.context.Context;
|
||||||
|
import io.opentelemetry.exporter.internal.InstrumentationUtil;
|
||||||
|
import java.util.concurrent.CountDownLatch;
|
||||||
|
import java.util.concurrent.atomic.AtomicBoolean;
|
||||||
|
import org.junit.jupiter.api.AfterEach;
|
||||||
|
import org.junit.jupiter.api.Assertions;
|
||||||
|
import org.junit.jupiter.api.BeforeEach;
|
||||||
|
import org.junit.jupiter.api.Test;
|
||||||
|
|
||||||
|
abstract class AbstractOkHttpSuppressionTest<T> {
|
||||||
|
|
||||||
|
@BeforeEach
|
||||||
|
void setUp() {
|
||||||
|
OkHttpUtil.setPropagateContextForTestingInDispatcher(true);
|
||||||
|
}
|
||||||
|
|
||||||
|
@AfterEach
|
||||||
|
void tearDown() {
|
||||||
|
OkHttpUtil.setPropagateContextForTestingInDispatcher(false);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void testSuppressInstrumentation() throws InterruptedException {
|
||||||
|
CountDownLatch latch = new CountDownLatch(1);
|
||||||
|
AtomicBoolean suppressInstrumentation = new AtomicBoolean(false);
|
||||||
|
|
||||||
|
Runnable onSuccess = Assertions::fail;
|
||||||
|
Runnable onFailure =
|
||||||
|
() -> {
|
||||||
|
suppressInstrumentation.set(
|
||||||
|
InstrumentationUtil.shouldSuppressInstrumentation(Context.current()));
|
||||||
|
latch.countDown();
|
||||||
|
};
|
||||||
|
|
||||||
|
send(getSender(), onSuccess, onFailure);
|
||||||
|
|
||||||
|
latch.await();
|
||||||
|
|
||||||
|
assertTrue(suppressInstrumentation.get());
|
||||||
|
}
|
||||||
|
|
||||||
|
abstract void send(T sender, Runnable onSuccess, Runnable onFailure);
|
||||||
|
|
||||||
|
private T getSender() {
|
||||||
|
return createSender("https://none");
|
||||||
|
}
|
||||||
|
|
||||||
|
abstract T createSender(String endpoint);
|
||||||
|
}
|
|
@ -0,0 +1,36 @@
|
||||||
|
/*
|
||||||
|
* Copyright The OpenTelemetry Authors
|
||||||
|
* SPDX-License-Identifier: Apache-2.0
|
||||||
|
*/
|
||||||
|
|
||||||
|
package io.opentelemetry.exporter.sender.okhttp.internal;
|
||||||
|
|
||||||
|
import io.opentelemetry.exporter.internal.marshal.MarshalerWithSize;
|
||||||
|
import io.opentelemetry.exporter.internal.marshal.Serializer;
|
||||||
|
import java.util.Collections;
|
||||||
|
|
||||||
|
class OkHttpGrpcSuppressionTest
|
||||||
|
extends AbstractOkHttpSuppressionTest<
|
||||||
|
OkHttpGrpcSender<OkHttpGrpcSuppressionTest.DummyMarshaler>> {
|
||||||
|
|
||||||
|
@Override
|
||||||
|
void send(OkHttpGrpcSender<DummyMarshaler> sender, Runnable onSuccess, Runnable onFailure) {
|
||||||
|
sender.send(new DummyMarshaler(), onSuccess, (grpcResponse, throwable) -> onFailure.run());
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
OkHttpGrpcSender<DummyMarshaler> createSender(String endpoint) {
|
||||||
|
return new OkHttpGrpcSender<>(
|
||||||
|
"https://localhost", false, 10L, Collections.emptyMap(), null, null, null);
|
||||||
|
}
|
||||||
|
|
||||||
|
protected static class DummyMarshaler extends MarshalerWithSize {
|
||||||
|
|
||||||
|
protected DummyMarshaler() {
|
||||||
|
super(0);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
protected void writeTo(Serializer output) {}
|
||||||
|
}
|
||||||
|
}
|
|
@ -0,0 +1,39 @@
|
||||||
|
/*
|
||||||
|
* Copyright The OpenTelemetry Authors
|
||||||
|
* SPDX-License-Identifier: Apache-2.0
|
||||||
|
*/
|
||||||
|
|
||||||
|
package io.opentelemetry.exporter.sender.okhttp.internal;
|
||||||
|
|
||||||
|
import java.io.IOException;
|
||||||
|
import java.io.OutputStream;
|
||||||
|
import java.nio.charset.StandardCharsets;
|
||||||
|
import java.util.Collections;
|
||||||
|
import java.util.function.Consumer;
|
||||||
|
|
||||||
|
class OkHttpHttpSuppressionTest extends AbstractOkHttpSuppressionTest<OkHttpHttpSender> {
|
||||||
|
|
||||||
|
@Override
|
||||||
|
void send(OkHttpHttpSender sender, Runnable onSuccess, Runnable onFailure) {
|
||||||
|
byte[] content = "A".getBytes(StandardCharsets.UTF_8);
|
||||||
|
Consumer<OutputStream> outputStreamConsumer =
|
||||||
|
outputStream -> {
|
||||||
|
try {
|
||||||
|
outputStream.write(content);
|
||||||
|
} catch (IOException e) {
|
||||||
|
throw new RuntimeException(e);
|
||||||
|
}
|
||||||
|
};
|
||||||
|
sender.send(
|
||||||
|
outputStreamConsumer,
|
||||||
|
content.length,
|
||||||
|
(response) -> onSuccess.run(),
|
||||||
|
(error) -> onFailure.run());
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
OkHttpHttpSender createSender(String endpoint) {
|
||||||
|
return new OkHttpHttpSender(
|
||||||
|
endpoint, false, "text/plain", 10L, Collections::emptyMap, null, null, null, null);
|
||||||
|
}
|
||||||
|
}
|
|
@ -5,6 +5,7 @@
|
||||||
|
|
||||||
package io.opentelemetry.sdk.internal;
|
package io.opentelemetry.sdk.internal;
|
||||||
|
|
||||||
|
import io.opentelemetry.context.Context;
|
||||||
import java.util.concurrent.Executors;
|
import java.util.concurrent.Executors;
|
||||||
import java.util.concurrent.ThreadFactory;
|
import java.util.concurrent.ThreadFactory;
|
||||||
import java.util.concurrent.atomic.AtomicInteger;
|
import java.util.concurrent.atomic.AtomicInteger;
|
||||||
|
@ -20,14 +21,30 @@ public final class DaemonThreadFactory implements ThreadFactory {
|
||||||
private final String namePrefix;
|
private final String namePrefix;
|
||||||
private final AtomicInteger counter = new AtomicInteger();
|
private final AtomicInteger counter = new AtomicInteger();
|
||||||
private final ThreadFactory delegate = Executors.defaultThreadFactory();
|
private final ThreadFactory delegate = Executors.defaultThreadFactory();
|
||||||
|
private final boolean propagateContextForTesting;
|
||||||
|
|
||||||
public DaemonThreadFactory(String namePrefix) {
|
public DaemonThreadFactory(String namePrefix) {
|
||||||
|
this(namePrefix, /* propagateContextForTesting= */ false);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* {@link DaemonThreadFactory}'s constructor.
|
||||||
|
*
|
||||||
|
* @param namePrefix Used when setting the new thread's name.
|
||||||
|
* @param propagateContextForTesting For tests only. When enabled, the current thread's {@link
|
||||||
|
* Context} will be passed over to the new threads, this is useful for validating scenarios
|
||||||
|
* where context propagation is available through bytecode instrumentation.
|
||||||
|
*/
|
||||||
|
public DaemonThreadFactory(String namePrefix, boolean propagateContextForTesting) {
|
||||||
this.namePrefix = namePrefix;
|
this.namePrefix = namePrefix;
|
||||||
|
this.propagateContextForTesting = propagateContextForTesting;
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public Thread newThread(Runnable runnable) {
|
public Thread newThread(Runnable runnable) {
|
||||||
Thread t = delegate.newThread(runnable);
|
Thread t =
|
||||||
|
delegate.newThread(
|
||||||
|
propagateContextForTesting ? Context.current().wrap(runnable) : runnable);
|
||||||
try {
|
try {
|
||||||
t.setDaemon(true);
|
t.setDaemon(true);
|
||||||
t.setName(namePrefix + "-" + counter.incrementAndGet());
|
t.setName(namePrefix + "-" + counter.incrementAndGet());
|
||||||
|
|
Loading…
Reference in New Issue