Repository navigation
Avoid repeated exceptions from malformed Kafka Base64 headers with AdaptiveLatch ("quick" fix) - #12672
Avoid repeated exceptions from malformed Kafka Base64 headers with AdaptiveLatch ("quick" fix)#12672dougqh wants to merge 13 commits into
Conversation
🟢 Java Benchmark SLOs — All performance SLOs passed
PR vs. master results
Commit: Load and DaCapo benchmarks can be triggered manually in the GitLab pipeline. Results will appear in the Benchmarking Platform UI after completion. |
Follow-up to the Kafka header Base64-decode guard fix. Compares the hand-specialized GuardedBase64Decode/Breaker shape against a generic ParseHandler + Breaker abstraction (construction-time vs. call-time strategy composition) to see what devirtualization costs, if any, a reusable version of this pattern would carry. Motivated by APMLP-1513's children (APMLP-1760, APMLP-1769, APMLP-1772, APMLP-1783), which show enough independent exception/parsing-cost cases to justify considering a shared primitive rather than one-off copies. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Replace the ParseHandler plus GenericBreakerCtor/GenericBreakerParam pair with one abstract DynamicLatch sketch, local to the benchmark, whose subclass is the strategy: an optimistic parse, a cheap correct pre-check and a stackless failure. That removes the constructor-versus-call-time question, since there is no separate strategy object. The latch has two flavors: get lets the failure flow to the caller (throwing a stackless stand-in while engaged), and tryGetOrNull converts it to null and builds no exception at all while engaged. Latches are static final fields of a named final subclass. The hand-written Breaker stays as the specialized baseline, and a setup self-check fails fast if the sketch misbehaves. Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
Rename the optimistic hook parse to handle, so it reads the same across the Latch family, and the pre-check isDefinitelyInvalid to isKnownToFail, which is not specific to parsers. Naming only. Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
Record a five-fork run on Zulu 17 (M1) in the benchmark's Javadoc: the single-input arms and the mixed-input arms at five invalid rates. The sketch matches the hand-written breaker; the adaptive form is a tradeoff against always pre-checking at high invalid rates. Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
Brings the sketch in line with Latch and ClassLatch: tryGetOrNull now returns fallback(input), null unless overridden, for a failed or known-to-fail input. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
handle becomes apply and tryGetOrNull becomes tryApply. Also re-aligns the results table after the AdaptiveLatch rename. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Customer (~180K/day) hits repeated IllegalArgumentException from Functions.BASE64_DECODE when a Kafka producer sends non-Base64 header values, paying full stack-trace fill-in on every throw. Adds a guarded decode path (GuardedBase64Decode) that fast-fails with a stack-trace-free exception once a real failure is seen, with hysteresis to re-check after a run of valid input, wired behind a new kafka.client.base64.decoding.guard.enabled config flag (default true). Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
…n it AdaptiveLatch moves from the benchmark into datadog.trace.util with only the converting form, tryApply; the flow-through get has no caller. Functions.GuardedBase64Decode becomes a subclass of it, so a malformed header no longer builds even a stackless exception while engaged. The pre-check no longer turns away unpadded or empty input: the basic decoder accepts both, so only a byte outside the alphabet, or a final unit of one character, is known to fail. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
59c56f5 to
140a424
Compare
Measures the primitive itself, on Integer.parseInt rather than Base64: the status quo, an always-on pre-check, and the latch disengaged, engaged and in a mixed stream, at stack depth 0 and 50. Records one run on Zulu 17, and narrows AdaptiveLatch's Javadoc to what was measured: the disengaged path is not free, and the latch only helps while bad input arrives more often than once per closeAfter calls. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Kafka / producer-benchmarkParameters
See matching parameters
SummaryFound 0 performance improvements and 0 performance regressions! Performance is the same for 3 metrics, 0 unstable metrics. See unchanged results
|
Kafka / consumer-benchmarkParameters
See matching parameters
SummaryFound 0 performance improvements and 0 performance regressions! Performance is the same for 3 metrics, 0 unstable metrics. See unchanged results
|
… one AdaptiveLatch now switches between two paths instead of guarding one with a pre-check: apply, the optimistic path, which may throw, and applySafely, a cautious path that must not, which reports bad input by returning reject(input). A pre-check is one way to build the cautious path; an exception-free implementation is another. closeAfter becomes abstract, and its Javadoc gives the rent-or-buy rule for choosing it. GuardedBase64Decode's cautious path is now decodeOrNull, an exception-free decoder that accepts exactly what Base64.getDecoder() accepts, checked against it by differential tests on JDK 8, 17 and 25. It allocates its output only once the first unit is clean, so input that is bad from the start costs no allocation. closeAfter is 64, from Base64DecodeBenchmark's costs on JDK 17. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
It is fixed policy for a call site, like failureType, and is read only on the failure path, so a final field costs nothing over an inlined constant. The constructor rejects a value below 1, which would engage and immediately disengage the latch, and a closeAfter() accessor remains. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
| * final unit, and nothing may follow it; a final unit of one character is rejected. | ||
| */ | ||
| @Nullable | ||
| static String decodeOrNull(byte[] src) { |
There was a problem hiding this comment.
I'm debating whether this is the right level for this. Curious what others think?
On newer JVMs using the stock decoder is faster for the happy path because it is intrinsified and takes advantage of hardware SIMD support. The exception free version is only better for invalid values.
And the adaptive approach is only really worth it because header parsing is one of the hottest things in the customer application
There was a problem hiding this comment.
I don't have a practical answer. Maybe this can be done as a later improvement. Also, it would be interesting to know on what value this code is failing.
IIUC Kafka header values are arbitrary bytes, and forEachKey attempts Base64 decoding on every visited non-null value before classification. A valid UTF-8, other encoding or binary application header can therefore fail Base64 decoding without being malformed.
Its state is a plain counter by design: a stale or lost update costs one more cautious call or one more optimistic failure, never a wrong result. SpotBugs reports nothing for it, so no suppression is needed. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Codex Review SummaryThis comment shows the latest Codex review activity on this pull request.
ℹ️ About Codex in GitHubYour team has set up Codex to review pull requests in this repo. Reviews are triggered when you
Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings. |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 3c7c25fc04
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| this.headerValueTransformer = | ||
| guardEnabled ? new Functions.GuardedBase64Decode()::tryApply : BASE64_DECODE; |
There was a problem hiding this comment.
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 👍 / 👎.
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
| /** | ||
| * {@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. | ||
| * | ||
| * <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. | ||
| */ |
There was a problem hiding this comment.
praise: Looks like a good javadoc
| return result; | ||
| } | ||
|
|
||
| // Same hysteresis, but preserves throw-based failure semantics: a precheck-known failure |
There was a problem hiding this comment.
suggestion: Same hysteresis, ... suggests we have read that file as book, but that might not be the case, i.e. maybe repeat what this hysteresis is.
| import static datadog.trace.api.config.TraceInstrumentationConfig.JMS_PROPAGATION_DISABLED_TOPICS; | ||
| import static datadog.trace.api.config.TraceInstrumentationConfig.JMS_UNACKNOWLEDGED_MAX_AGE; | ||
| import static datadog.trace.api.config.TraceInstrumentationConfig.KAFKA_CLIENT_BASE64_DECODING_ENABLED; | ||
| import static datadog.trace.api.config.TraceInstrumentationConfig.KAFKA_CLIENT_BASE64_DECODING_GUARD_ENABLED; |
There was a problem hiding this comment.
note: FYI this config key is missing form the registry
| /** | ||
| * Base64-decodes bytes that are occasionally not Base64 at all, for example header values from a | ||
| * misconfigured producer that mixes encodings, without paying for {@link Base64}'s decoder to | ||
| * throw, and fill in a stack trace, on every one of them. Returns {@code null} for input that | ||
| * does not decode, like {@link #BASE64_DECODE}. | ||
| * | ||
| * <p>Decodes with {@link Base64#getDecoder()} until it fails, then with {@link | ||
| * #decodeOrNull(byte[])}, an exception-free decoder that accepts exactly what {@link | ||
| * Base64#getDecoder()} accepts, until {@code CLOSE_AFTER} consecutive inputs decode again (see | ||
| * {@link AdaptiveLatch}). | ||
| * | ||
| * <p>One instance can be shared across threads: its state is advisory, so a stale read costs one | ||
| * extra cautious decode or one extra real failure, never a wrong result. | ||
| */ |
There was a problem hiding this comment.
praise: neat doc, did clarify doc skill helped there?
| private static final int CLOSE_THRESHOLD = 20; | ||
|
|
||
| /** | ||
| * Thrown only when the cheap alphabet-scan precheck already knows the input can't be valid | ||
| * Base64, so we never call into {@link Base64}'s decoder (or pay for its exception's stack-trace | ||
| * capture) at all. Deliberately extends {@link IllegalArgumentException} so it's a drop-in for | ||
| * existing catch sites around the real decoder; callers should not rely on a stack trace being | ||
| * present for this specific instance — a real decode failure still throws the JDK's own | ||
| * exception, unmodified, with its own message and stack trace. | ||
| * | ||
| * <p>A new instance is built per failure, never shared, so suppression cannot accumulate on it. A | ||
| * shared instance would also have to disable suppression and forbid {@code initCause}. | ||
| */ | ||
| static final class FastFailBase64Exception extends IllegalArgumentException { | ||
| private static final String MESSAGE = "Header value is not valid Base64"; | ||
|
|
||
| FastFailBase64Exception() { | ||
| super(MESSAGE); | ||
| } | ||
|
|
||
| // IllegalArgumentException has no writableStackTrace-suppressing constructor of its own, so | ||
| // skip the stack walk here instead. | ||
| @Override | ||
| public synchronized Throwable fillInStackTrace() { | ||
| return this; | ||
| } | ||
| } | ||
|
|
||
| // Single-word countdown: 0 == closed (no precheck), >0 == guarded, counting down to close. | ||
| // Plain int on purpose: this is advisory hysteresis, not correctness-critical state, so a lost | ||
| // update or a stale read across threads just means one extra precheck or one extra exception. | ||
| static final class Breaker { | ||
| int state; | ||
|
|
||
| String decode(byte[] bytes) { | ||
| if (state > 0 && !looksLikeBase64(bytes)) { | ||
| return null; | ||
| } | ||
| String result = BASE64_DECODE.apply(bytes); | ||
| if (result == null) { | ||
| state = CLOSE_THRESHOLD; | ||
| } else if (state > 0) { | ||
| state--; | ||
| } | ||
| return result; | ||
| } | ||
|
|
||
| // Same hysteresis, but preserves throw-based failure semantics: a precheck-known failure | ||
| // throws our stack-trace-free stand-in, a real decode failure lets the JDK's own | ||
| // IllegalArgumentException (with its own message and stack trace) propagate untouched. | ||
| String decodeOrThrow(byte[] bytes) { | ||
| if (state > 0 && !looksLikeBase64(bytes)) { | ||
| throw new FastFailBase64Exception(); | ||
| } | ||
| try { | ||
| String result = new String(Base64.getDecoder().decode(bytes), StandardCharsets.UTF_8); | ||
| if (state > 0) { | ||
| state--; | ||
| } | ||
| return result; | ||
| } catch (IllegalArgumentException e) { | ||
| state = CLOSE_THRESHOLD; | ||
| throw e; | ||
| } | ||
| } | ||
| } |
There was a problem hiding this comment.
issue: This is from codex:
The converting Breaker uses 20 successful decodes across intervening failures, while production uses 64 consecutive successes. Over 1,000 alternating calls, this causes 25 JDK failures versus 1, so the mixed results conflate policy and latch overhead. Use the same countdown and cautious decoder for that baseline; the throwing arm remains a separate strategy.
| private static final int CLOSE_THRESHOLD = 20; | |
| /** | |
| * Thrown only when the cheap alphabet-scan precheck already knows the input can't be valid | |
| * Base64, so we never call into {@link Base64}'s decoder (or pay for its exception's stack-trace | |
| * capture) at all. Deliberately extends {@link IllegalArgumentException} so it's a drop-in for | |
| * existing catch sites around the real decoder; callers should not rely on a stack trace being | |
| * present for this specific instance — a real decode failure still throws the JDK's own | |
| * exception, unmodified, with its own message and stack trace. | |
| * | |
| * <p>A new instance is built per failure, never shared, so suppression cannot accumulate on it. A | |
| * shared instance would also have to disable suppression and forbid {@code initCause}. | |
| */ | |
| static final class FastFailBase64Exception extends IllegalArgumentException { | |
| private static final String MESSAGE = "Header value is not valid Base64"; | |
| FastFailBase64Exception() { | |
| super(MESSAGE); | |
| } | |
| // IllegalArgumentException has no writableStackTrace-suppressing constructor of its own, so | |
| // skip the stack walk here instead. | |
| @Override | |
| public synchronized Throwable fillInStackTrace() { | |
| return this; | |
| } | |
| } | |
| // Single-word countdown: 0 == closed (no precheck), >0 == guarded, counting down to close. | |
| // Plain int on purpose: this is advisory hysteresis, not correctness-critical state, so a lost | |
| // update or a stale read across threads just means one extra precheck or one extra exception. | |
| static final class Breaker { | |
| int state; | |
| String decode(byte[] bytes) { | |
| if (state > 0 && !looksLikeBase64(bytes)) { | |
| return null; | |
| } | |
| String result = BASE64_DECODE.apply(bytes); | |
| if (result == null) { | |
| state = CLOSE_THRESHOLD; | |
| } else if (state > 0) { | |
| state--; | |
| } | |
| return result; | |
| } | |
| // Same hysteresis, but preserves throw-based failure semantics: a precheck-known failure | |
| // throws our stack-trace-free stand-in, a real decode failure lets the JDK's own | |
| // IllegalArgumentException (with its own message and stack trace) propagate untouched. | |
| String decodeOrThrow(byte[] bytes) { | |
| if (state > 0 && !looksLikeBase64(bytes)) { | |
| throw new FastFailBase64Exception(); | |
| } | |
| try { | |
| String result = new String(Base64.getDecoder().decode(bytes), StandardCharsets.UTF_8); | |
| if (state > 0) { | |
| state--; | |
| } | |
| return result; | |
| } catch (IllegalArgumentException e) { | |
| state = CLOSE_THRESHOLD; | |
| throw e; | |
| } | |
| } | |
| } | |
| private static final int CLOSE_THRESHOLD = 64; | |
| /** | |
| * Thrown only when the cheap alphabet-scan precheck already knows the input can't be valid | |
| * Base64, so we never call into {@link Base64}'s decoder (or pay for its exception's stack-trace | |
| * capture) at all. Deliberately extends {@link IllegalArgumentException} so it's a drop-in for | |
| * existing catch sites around the real decoder; callers should not rely on a stack trace being | |
| * present for this specific instance — a real decode failure still throws the JDK's own | |
| * exception, unmodified, with its own message and stack trace. | |
| * | |
| * <p>A new instance is built per failure, never shared, so suppression cannot accumulate on it. A | |
| * shared instance would also have to disable suppression and forbid {@code initCause}. | |
| */ | |
| static final class FastFailBase64Exception extends IllegalArgumentException { | |
| private static final String MESSAGE = "Header value is not valid Base64"; | |
| FastFailBase64Exception() { | |
| super(MESSAGE); | |
| } | |
| // IllegalArgumentException has no writableStackTrace-suppressing constructor of its own, so | |
| // skip the stack walk here instead. | |
| @Override | |
| public synchronized Throwable fillInStackTrace() { | |
| return this; | |
| } | |
| } | |
| // Single-word countdown: 0 == closed (no precheck), >0 == guarded, counting down to close. | |
| // Plain int on purpose: this is advisory hysteresis, not correctness-critical state, so a lost | |
| // update or a stale read across threads just means one extra precheck or one extra exception. | |
| static final class Breaker { | |
| int state; | |
| String decode(byte[] bytes) { | |
| if (state == 0) { | |
| String result = BASE64_DECODE.apply(bytes); | |
| if (result != null) { | |
| return result; | |
| } | |
| state = CLOSE_THRESHOLD; | |
| } else { | |
| state--; | |
| } | |
| String result = GuardedBase64Decode.decodeOrNull(bytes); | |
| if (result == null) { | |
| state = CLOSE_THRESHOLD; | |
| } | |
| return result; | |
| } | |
| // Same hysteresis, but preserves throw-based failure semantics: a precheck-known failure | |
| // throws our stack-trace-free stand-in, a real decode failure lets the JDK's own | |
| // IllegalArgumentException (with its own message and stack trace) propagate untouched. | |
| String decodeOrThrow(byte[] bytes) { | |
| if (state > 0 && !looksLikeBase64(bytes)) { | |
| state = CLOSE_THRESHOLD; | |
| throw new FastFailBase64Exception(); | |
| } | |
| try { | |
| String result = new String(Base64.getDecoder().decode(bytes), StandardCharsets.UTF_8); | |
| if (state > 0) { | |
| state--; | |
| } | |
| return result; | |
| } catch (IllegalArgumentException e) { | |
| state = CLOSE_THRESHOLD; | |
| throw e; | |
| } | |
| } | |
| } |
| 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"); | ||
| } | ||
| } |
There was a problem hiding this comment.
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.
| 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"); | |
| } | |
| } |
| * The sketch matches the hand-written breaker: within 0.2 ns in the single-input arms, and within | ||
| * about 4 ns in the mixed ones (the largest gap is flow-through at one in 10,000, 39.6 ns against | ||
| * 35.3 ns). While engaged, the converting flavor costs about 2.8 ns where the status quo costs | ||
| * about 906 ns; the flow-through flavor costs about 11.5 ns, the price of building a stackless | ||
| * exception. Always pre-checking more than doubles the cost of valid input (71 ns against 30 ns), | ||
| * which is what the adaptive form avoids. | ||
| * | ||
| * <p>It is a tradeoff, not a free win. With one invalid input in 100 or rarer, the guard costs 1 to | ||
| * 4 ns over doing nothing and is about 24 to 34 ns cheaper than always pre-checking. With a high | ||
| * rate (one in 2 or one in 10), always pre-checking is faster than the adaptive guard (36.6 ns | ||
| * against about 47 to 52 ns at one in 2, and 65 ns against about 70 to 72 ns at one in 10). The | ||
| * cause was not investigated. |
There was a problem hiding this comment.
suggestion: Keep the historical table, but avoid presenting the earlier sketch's timings as results for the final 64-consecutive-success guard or the Kafka adapter.
| * The sketch matches the hand-written breaker: within 0.2 ns in the single-input arms, and within | |
| * about 4 ns in the mixed ones (the largest gap is flow-through at one in 10,000, 39.6 ns against | |
| * 35.3 ns). While engaged, the converting flavor costs about 2.8 ns where the status quo costs | |
| * about 906 ns; the flow-through flavor costs about 11.5 ns, the price of building a stackless | |
| * exception. Always pre-checking more than doubles the cost of valid input (71 ns against 30 ns), | |
| * which is what the adaptive form avoids. | |
| * | |
| * <p>It is a tradeoff, not a free win. With one invalid input in 100 or rarer, the guard costs 1 to | |
| * 4 ns over doing nothing and is about 24 to 34 ns cheaper than always pre-checking. With a high | |
| * rate (one in 2 or one in 10), always pre-checking is faster than the adaptive guard (36.6 ns | |
| * against about 47 to 52 ns at one in 2, and 65 ns against about 70 to 72 ns at one in 10). The | |
| * cause was not investigated. | |
| * These historical results describe only the benchmark-local sketch and its pre-check strategy. | |
| * They do not establish the overhead or mixed-input tradeoff of the production guard, which uses an | |
| * exception-free decoder and disengages after 64 consecutive successful decodes. Re-run the current | |
| * arms with equivalent policies before drawing conclusions about that implementation. | |
| * | |
| * <p>The current latch arms call an exact-type {@code static final} directly. They do not measure | |
| * the Kafka adapter's captured {@code Function}, header iteration, or concurrent use of shared | |
| * state. |
| * 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. |
There was a problem hiding this comment.
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.
| * 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. |
| * final unit, and nothing may follow it; a final unit of one character is rejected. | ||
| */ | ||
| @Nullable | ||
| static String decodeOrNull(byte[] src) { |
There was a problem hiding this comment.
I don't have a practical answer. Maybe this can be done as a later improvement. Also, it would be interesting to know on what value this code is failing.
IIUC Kafka header values are arbitrary bytes, and forEachKey attempts Base64 decoding on every visited non-null value before classification. A valid UTF-8, other encoding or binary application header can therefore fail Base64 decoding without being malformed.
| this.headerValueTransformer = | ||
| guardEnabled ? new Functions.GuardedBase64Decode()::tryApply : BASE64_DECODE; |
There was a problem hiding this comment.
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.
Brought over unchanged, with its tests, for use in URL-to-URI conversion. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
What Does This Do
Stops Kafka header extraction from paying for a thrown, stack-trace-filled
IllegalArgumentExceptionon every malformed header value whenkafka.client.base64.decoding.enabledis on. It does this with a new shared primitive,AdaptiveLatch.AdaptiveLatch(internal-api,datadog.trace.util) is a self-resetting switch between two ways of doing an operation whose failure depends on the input:applyis the optimistic path. It's fastest on good input but throwsXon bad input.applySafelyis the cautious path. It must not throwX, and it reports bad input by returningreject(input).While disengaged, calls take the optimistic path. A failure there engages the latch and retries the same input cautiously. Calls then stay on the cautious path until
closeAfterconsecutive inputs pass it cleanly.closeAfteris passed to the constructor, which rejects values below 1. A pre-check in front ofapply, an exception-free re-implementation, and a repair step are all valid cautious paths. It's the input-dependent sibling ofLatch(Stop repeated NoSuchFieldError in Jackson 2.16 IAST interner lookup (quick fix) #12670) andClassLatch(Skip repeated AbstractMethodError from JDBC getClientInfo #12702), with the sameapply/fallback/tryApplynaming.Functions.GuardedBase64Decodeswitches betweenBase64.getDecoder()anddecodeOrNull, a new exception-free decoder.decodeOrNullaccepts exactly what the JDK's basic decoder accepts, and returnsnullwhere the JDK would throw.closeAfteris 64.Kafka wiring (
kafka-clients-0.11andkafka-clients-3.8TextMapExtractAdapter): uses the guard whenkafka.client.base64.decoding.guard.enabledis on (defaulttrue), and falls back to the plainBASE64_DECODEwhen it's off.Motivation
A customer sees about 180K/day of this exact throw from Kafka header extraction: non-Base64 header values from a producer that mixes encodings. The JDK decoder fills in a full stack trace on every one, and the existing
BASE64_DECODEjust catches it and returnsnull.A
LatchorClassLatchcan't help here, because the failure is per value: a bad header says nothing about the next one. Enough other input-dependent exception and parsing-cost cases exist,URIUtils.safeParseamong them, to justify a shared primitive rather than another one-off breaker.This supersedes #12671, which had a hand-written breaker for the same fix, and folds in its Kafka wiring and config flag.
Additional Notes
Why switch between decoders rather than always use the exception-free one. On JDK 8 the two decoders cost about the same on header-sized input. On JDK 17 the JDK decoder, which decodes in blocks and has an intrinsic on aarch64, is about 27 ns faster on header-sized input and about 2× faster on long input. Always using ours would add that to every header on newer JDKs, and the adapter decodes every header value on each message, not just Datadog's. Switching keeps JDK speed for valid input and ours for bad input.
Choosing
closeAfter: a rent-or-buy decision.failureCost / cautiousOverheadgood inputs is never worse than 2× the best possible schedule, whatever the input pattern.AdaptiveLatch. Every site passes its own value to the constructor; there's no default.decodeOrNullcorrectness.It follows the JDK's rules: padding is optional, but if present it must complete the final unit, nothing may follow it, and a single dangling character is rejected.
Differential tests check it against
Base64.getDecoder()on JDK 8, 17 and 25:@TableTestof named edge casesThe tests also assert that the JDK decoder fails only with
IllegalArgumentException, which the latch relies on.This fixes a bug from Fix repeated exception cost from malformed Kafka header Base64 decoding #12671. Its pre-check treated empty input, and any length that isn't a multiple of 4, as invalid.
Base64.getDecoder()accepts both ("YQ"→a), so Fix repeated exception cost from malformed Kafka header Base64 decoding #12671 would have dropped valid unpadded headers while engaged.Allocation on bad input.
decodeOrNullallocates its output buffer only once the first 4-character unit is clean. Input that's bad from the start therefore costs about 6 ns and no allocation at any length, compared with 60–70 ns and 1 KB for long garbage when allocating up front. A same-run comparison showed this costs about 7 ns (JDK 17) to 12 ns (JDK 8) on valid header-sized input at depth 0, and nothing measurable otherwise. It only runs while the guard is engaged.Benchmarks (one run each on Zulu 8.0.382 and 17.0.7, M1, single thread, 2 forks):
Base64DecodeBenchmarkcompares ours with the JDK decoder directly:The JDK-decoder rows come from an earlier run the same day than the "ours" rows. The allocation comparison above was a single run.
AdaptiveLatchBenchmarkmeasures the latch itself, onInteger.parseIntwith a digit-scan pre-check, at depth 0 and 50, on JDK 17:closeAfter= 20 is far too small for that operation (its break-even is about 290). The latch saves nothing there, which matches what the cost model predicts.FunctionsBase64Benchmarkkeeps its earlier results from the benchmark-local sketch. They weren't re-measured, and the Javadoc says so.Tests.
AdaptiveLatchTest:closeAfterclean calls, and a rejection restarts the countcloseAfterbelow 1 is rejectedrejectyields the fallback, while a realnullcounts as cleanGuardedBase64DecodeTest: the differential tests above, plus engaging and disengaging behavior.TextMapExtractAdapterTestpasses for both Kafka modules, andcheckConfigurationspasses. All JMH sources compile.Contributor Checklist
type:and (comp:orinst:) labels in addition to any other useful labelsclose,fix, or any linking keywords when referencing an issueUse
solvesinstead, and assign the PR milestone to the issue🤖 Generated with Claude Code