Skip to content
Draft
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
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
import static datadog.trace.instrumentation.kafka_clients.KafkaDecorator.KAFKA_PRODUCED_KEY;

import datadog.trace.api.Config;
import datadog.trace.api.Functions;
import datadog.trace.bootstrap.instrumentation.api.AgentPropagation;
import datadog.trace.bootstrap.instrumentation.api.AgentPropagation.ContextVisitor;
import java.nio.ByteBuffer;
Expand All @@ -20,14 +21,21 @@ public class TextMapExtractAdapter implements ContextVisitor<Headers> {
private static final Logger log = LoggerFactory.getLogger(TextMapExtractAdapter.class);

public static final TextMapExtractAdapter GETTER =
new TextMapExtractAdapter(Config.get().isKafkaClientBase64DecodingEnabled());
new TextMapExtractAdapter(
Config.get().isKafkaClientBase64DecodingEnabled(),
Config.get().isKafkaClientBase64DecodingGuardEnabled());

private final Function<byte[], String> headerValueTransformer;
private final Base64.Decoder decoder;

public TextMapExtractAdapter(boolean decodeBase64Headers) {
this(decodeBase64Headers, true);
}

public TextMapExtractAdapter(boolean decodeBase64Headers, boolean guardEnabled) {
if (decodeBase64Headers) {
this.headerValueTransformer = BASE64_DECODE;
this.headerValueTransformer =
guardEnabled ? new Functions.GuardedBase64Decode()::tryApply : BASE64_DECODE;
Comment on lines +37 to +38

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Guard time-in-queue Base64 decoding too

When Base64 decoding and time-in-queue extraction are enabled, a malformed x_datadog_kafka_produced header still reaches extractTimeInQueueStart, which directly calls decoder.decode(header.value()) (lines 65-74) on every consumed record. That path bypasses the new latch entirely, so it continues to allocate and throw an IllegalArgumentException repeatedly—the failure mode this guard is intended to eliminate. The same direct decode remains in the Kafka 3.8 adapter.

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I believe this comment is wrong, x_datadog_kafka_produced can be a valid valid value, e.g. an timestamp (8 bytes), but this cannot be decoded as Base64, so it's rather an an encoding mismatch, not malformed data.

this.decoder = Base64.getDecoder();
} else {
this.headerValueTransformer = UTF8_BYTES_TO_STRING;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
import static datadog.trace.api.telemetry.LogCollector.EXCLUDE_TELEMETRY;

import datadog.trace.api.Config;
import datadog.trace.api.Functions;
import datadog.trace.bootstrap.instrumentation.api.AgentPropagation;
import datadog.trace.bootstrap.instrumentation.api.AgentPropagation.ContextVisitor;
import java.nio.ByteBuffer;
Expand All @@ -19,14 +20,21 @@ public class TextMapExtractAdapter implements ContextVisitor<Headers> {
private static final Logger log = LoggerFactory.getLogger(TextMapExtractAdapter.class);

public static final TextMapExtractAdapter GETTER =
new TextMapExtractAdapter(Config.get().isKafkaClientBase64DecodingEnabled());
new TextMapExtractAdapter(
Config.get().isKafkaClientBase64DecodingEnabled(),
Config.get().isKafkaClientBase64DecodingGuardEnabled());

private final Function<byte[], String> headerValueTransformer;
private final Base64.Decoder decoder;

public TextMapExtractAdapter(boolean decodeBase64Headers) {
this(decodeBase64Headers, true);
}

public TextMapExtractAdapter(boolean decodeBase64Headers, boolean guardEnabled) {
if (decodeBase64Headers) {
this.headerValueTransformer = BASE64_DECODE;
this.headerValueTransformer =
guardEnabled ? new Functions.GuardedBase64Decode()::tryApply : BASE64_DECODE;
this.decoder = Base64.getDecoder();
} else {
this.headerValueTransformer = UTF8_BYTES_TO_STRING;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -113,6 +113,8 @@ public final class TraceInstrumentationConfig {
"kafka.client.propagation.disabled.topics";
public static final String KAFKA_CLIENT_BASE64_DECODING_ENABLED =
"kafka.client.base64.decoding.enabled";
public static final String KAFKA_CLIENT_BASE64_DECODING_GUARD_ENABLED =
"kafka.client.base64.decoding.guard.enabled";

public static final String JMS_PROPAGATION_DISABLED_TOPICS = "jms.propagation.disabled.topics";
public static final String JMS_PROPAGATION_DISABLED_QUEUES = "jms.propagation.disabled.queues";
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,137 @@
package datadog.trace.api;

import static java.nio.charset.StandardCharsets.UTF_8;

import datadog.trace.api.Functions.GuardedBase64Decode;
import java.util.Base64;
import java.util.Random;
import org.openjdk.jmh.annotations.Benchmark;
import org.openjdk.jmh.annotations.Fork;
import org.openjdk.jmh.annotations.Measurement;
import org.openjdk.jmh.annotations.Param;
import org.openjdk.jmh.annotations.Scope;
import org.openjdk.jmh.annotations.Setup;
import org.openjdk.jmh.annotations.State;
import org.openjdk.jmh.annotations.Threads;
import org.openjdk.jmh.annotations.Warmup;

/**
* {@link GuardedBase64Decode#decodeOrNull}, the exception-free decoder, against {@link
* Base64#getDecoder()}. These are the two costs that set the guard's {@code closeAfter}: the
* exception-free decoder's overhead on valid input (the cost of staying engaged), and the JDK
* decoder's failure on invalid input (the cost of disengaging too early).
*
* <ul>
* <li>{@code size}: {@code header} is a trace-header-sized value, {@code long} is about 1.4 KB of
* Base64, where the JDK's block decoding, and any intrinsic for it, has room to pay off.
* <li>{@code depth}: the stack the JDK decoder's exception fills in; a consumer thread's stack is
* deeper than a benchmark thread's.
* </ul>
*
* <p>Run with {@code ./gradlew :internal-api:jmh -Pjmh.includes=Base64DecodeBenchmark
* -Pjmh.profilers=gc}, and with {@code -PtestJvm=17} for a newer JDK.
*
* <p>Results, one run each: Zulu 8.0.382 and Zulu 17.0.7 (HotSpot), MacBook M1, single thread, 2
* forks of 5 one-second iterations, on a laptop with normal background activity. x86 is not
* measured. ns/op is derived from ops/s; B/op is from {@code -prof gc}.
*
* <pre>
* ns/op (B/op) JDK 8 JDK 17
* header long header long
* jdkValid, depth 0 87.7 (160) 3631 53.4 (104) 1598
* exceptionFreeValid 86.3 (160) 3904 80.3 (104) 3072
* jdkInvalid, depth 0 889 (928) 926 921 (1024) 956
* jdkInvalid, depth 50 2080 (1904) 2189 2498 (2384) 2512
* exceptionFreeInvalid 6.4 (0) 6.4 6.2 (0) 6.3
* </pre>
*
* On JDK 8 the exception-free decoder is about as fast as the JDK's on header-sized input, and
* about 7% slower on long input. On JDK 17 the JDK's decoder, which decodes in blocks and has an
* intrinsic on this platform, is about 27 ns faster on header-sized input and about twice as fast
* on long input; that difference is why the guard switches between the two rather than always using
* the exception-free one. Invalid input costs the exception-free decoder about 6 ns and no
* allocation, at any length, where the JDK's costs about 0.9 us, or 2.1 to 2.5 us at depth 50.
Comment on lines +52 to +53

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nitpick: Scope the ~6 ns / 0 B result to corruption in the first unit, since late corruption and bad padding can allocate after a clean first unit.

Suggested change
* the exception-free one. Invalid input costs the exception-free decoder about 6 ns and no
* allocation, at any length, where the JDK's costs about 0.9 us, or 2.1 to 2.5 us at depth 50.
* the exception-free one. The invalid-input rows above corrupt the fourth byte, so rejection
* happens before output allocation: about 6 ns and 0 B/op for both measured sizes. Late corruption
* and malformed padding require separate measurements; they can scan and allocate output first. The
* JDK's early-rejection cost was about 0.9 us, or 2.1 to 2.5 us at depth 50.

*
* <p>The exception-free decoder allocates its output only once the first unit is clean. Against a
* copy that allocated up front, in the same run, that cost about 7 ns (JDK 17) to 12 ns (JDK 8) on
* valid header-sized input at depth 0, and nothing measurable at depth 50 or on long input, while
* saving 2 to 4 ns and 40 B on invalid header-sized input and about 55 ns and 1 KB on invalid long
* input. It runs only while the guard is engaged. The {@code jdk*} rows are from an earlier run of
* the same day, the {@code exceptionFree*} rows from the later one.
*/
Comment on lines +18 to +61

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

praise: Looks like a good javadoc

@Fork(2)
@Warmup(iterations = 3, time = 1)
@Measurement(iterations = 5, time = 1)
@Threads(1)
@State(Scope.Benchmark)
public class Base64DecodeBenchmark {

@Param({"header", "long"})
String size;

@Param({"0", "50"})
int depth;

byte[] valid;
byte[] invalid;

@Setup
public void setup() {
byte[] data;
if ("header".equals(size)) {
data = "1234567890123456789".getBytes(UTF_8);
} else {
data = new byte[1024];
new Random(12672).nextBytes(data);
}
valid = Base64.getEncoder().encode(data);
// a plain-text value where Base64 was expected, bad from its fourth byte
invalid = valid.clone();
invalid[3] = '-';
String expected = new String(data, UTF_8);
if (!expected.equals(GuardedBase64Decode.decodeOrNull(valid))
|| !expected.equals(jdk(valid))
|| GuardedBase64Decode.decodeOrNull(invalid) != null
|| jdk(invalid) != null) {
throw new IllegalStateException("the two decoders must agree on the benchmark inputs");
}
}
Comment on lines +75 to +98

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

suggestion: I propose to introduce cases with a late corruption and an invalid-padding. The current fourth-byte corruption currently exercises only rejection before output allocation.

Suggested change
byte[] valid;
byte[] invalid;
@Setup
public void setup() {
byte[] data;
if ("header".equals(size)) {
data = "1234567890123456789".getBytes(UTF_8);
} else {
data = new byte[1024];
new Random(12672).nextBytes(data);
}
valid = Base64.getEncoder().encode(data);
// a plain-text value where Base64 was expected, bad from its fourth byte
invalid = valid.clone();
invalid[3] = '-';
String expected = new String(data, UTF_8);
if (!expected.equals(GuardedBase64Decode.decodeOrNull(valid))
|| !expected.equals(jdk(valid))
|| GuardedBase64Decode.decodeOrNull(invalid) != null
|| jdk(invalid) != null) {
throw new IllegalStateException("the two decoders must agree on the benchmark inputs");
}
}
@Param({"early", "late", "padding"})
String invalidCase;
byte[] valid;
byte[] invalid;
@Setup
public void setup() {
byte[] data;
if ("header".equals(size)) {
data = "1234567890123456789".getBytes(UTF_8);
} else {
data = new byte[1024];
new Random(12672).nextBytes(data);
}
valid = Base64.getEncoder().encode(data);
invalid = valid.clone();
switch (invalidCase) {
case "early":
invalid[3] = '-';
break;
case "late":
invalid[invalid.length - 3] = '-';
break;
case "padding":
invalid[invalid.length - 2] = '=';
invalid[invalid.length - 1] = 'A';
break;
default:
throw new IllegalArgumentException("unknown invalid case: " + invalidCase);
}
String expected = new String(data, UTF_8);
if (!expected.equals(GuardedBase64Decode.decodeOrNull(valid))
|| !expected.equals(jdk(valid))
|| GuardedBase64Decode.decodeOrNull(invalid) != null
|| jdk(invalid) != null) {
throw new IllegalStateException("the two decoders must agree on the benchmark inputs");
}
}


static String jdk(byte[] src) {
try {
return new String(Base64.getDecoder().decode(src), UTF_8);
} catch (IllegalArgumentException e) {
return null;
}
}

@Benchmark
public Object jdkValid() {
return jdkAt(depth, valid);
}

@Benchmark
public Object exceptionFreeValid() {
return exceptionFreeAt(depth, valid);
}

@Benchmark
public Object jdkInvalid() {
return jdkAt(depth, invalid);
}

@Benchmark
public Object exceptionFreeInvalid() {
return exceptionFreeAt(depth, invalid);
}

private static Object jdkAt(int remaining, byte[] src) {
return remaining > 0 ? jdkAt(remaining - 1, src) : jdk(src);
}

private static Object exceptionFreeAt(int remaining, byte[] src) {
return remaining > 0
? exceptionFreeAt(remaining - 1, src)
: GuardedBase64Decode.decodeOrNull(src);
}
}
Loading
Loading