Skip to content

Avoid repeated exceptions from malformed Kafka Base64 headers with AdaptiveLatch ("quick" fix) - #12672

Draft
dougqh wants to merge 13 commits into
masterfrom
experiment/base64-guard-benchmark
Draft

dougqh wants to merge 13 commits into
masterfrom
experiment/base64-guard-benchmark

Conversation

@dougqh

@dougqh dougqh commented Sep 28, 2026 •

Copy link
Copy Markdown
Contributor

What Does This Do

Stops Kafka header extraction from paying for a thrown, stack-trace-filled IllegalArgumentException on every malformed header value when kafka.client.base64.decoding.enabled is 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:

    • apply is the optimistic path. It's fastest on good input but throws X on bad input.
    • applySafely is the cautious path. It must not throw X, and it reports bad input by returning reject(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 closeAfter consecutive inputs pass it cleanly. closeAfter is passed to the constructor, which rejects values below 1. A pre-check in front of apply, an exception-free re-implementation, and a repair step are all valid cautious paths. It's the input-dependent sibling of Latch (Stop repeated NoSuchFieldError in Jackson 2.16 IAST interner lookup (quick fix) #12670) and ClassLatch (Skip repeated AbstractMethodError from JDBC getClientInfo #12702), with the same apply / fallback / tryApply naming.

  • Functions.GuardedBase64Decode switches between Base64.getDecoder() and decodeOrNull, a new exception-free decoder. decodeOrNull accepts exactly what the JDK's basic decoder accepts, and returns null where the JDK would throw. closeAfter is 64.

  • Kafka wiring (kafka-clients-0.11 and kafka-clients-3.8 TextMapExtractAdapter): uses the guard when kafka.client.base64.decoding.guard.enabled is on (default true), and falls back to the plain BASE64_DECODE when it's off.

public static final class GuardedBase64Decode
    extends AdaptiveLatch<byte[], String, IllegalArgumentException> {
  public GuardedBase64Decode() {
    super(IllegalArgumentException.class, CLOSE_AFTER); // 64
  }

  @Override
  protected String apply(byte[] bytes) {
    return new String(Base64.getDecoder().decode(bytes), UTF_8);
  }

  @Override
  protected String applySafely(byte[] bytes) {
    String decoded = decodeOrNull(bytes);
    return decoded != null ? decoded : reject(bytes);
  }
}

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_DECODE just catches it and returns null.

A Latch or ClassLatch can'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.safeParse among 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.

  • Staying engaged costs the cautious path's overhead on each good input. Disengaging costs one optimistic failure when the next bad input arrives.
  • Disengaging after about failureCost / cautiousOverhead good inputs is never worse than 2× the best possible schedule, whatever the input pattern.
  • For Base64 on JDK 17, a failure costs about 913 ns more than a rejection at depth 0, and about 2,468 ns at depth 50. The cautious overhead is about 27 ns, so the break-even is 34 to 91. A consumer's stack is deeper than a benchmark's, so it's set to 64.
  • The rule is documented on AdaptiveLatch. Every site passes its own value to the constructor; there's no default.

decodeOrNull correctness.

  • 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:

    • every input up to 8 characters over a reduced alphabet covering padding, low and high bits, and an illegal byte
    • 20,000 random encodings, padded and unpadded, plus a mutation, a truncation and random noise for each
    • a @TableTest of named edge cases

    The 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. decodeOrNull allocates 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):

Base64DecodeBenchmark compares ours with the JDK decoder directly:

ns/op (B/op) JDK 8 header JDK 8 long JDK 17 header JDK 17 long
JDK decoder, valid 87.7 (160) 3631 53.4 (104) 1598
ours, valid 86.3 (160) 3904 80.3 (104) 3072
JDK decoder, invalid, depth 0 889 (928) 926 921 (1024) 956
JDK decoder, invalid, depth 50 2080 (1904) 2189 2498 (2384) 2512
ours, invalid 6.4 (0) 6.4 6.2 (0) 6.3

The JDK-decoder rows come from an earlier run the same day than the "ours" rows. The allocation comparison above was a single run.

AdaptiveLatchBenchmark measures the latch itself, on Integer.parseInt with a digit-scan pre-check, at depth 0 and 50, on JDK 17:

  • While engaged, bad input costs about 2.6 ns instead of about 908 ns (31 ns against 2.4 µs at depth 50), and allocates nothing instead of 880 B.
  • While disengaged, the latch costs about 3–4.5 ns over no latch on good input, which is more than one field read should. That hasn't been investigated.
  • At one bad input in 100, 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.

FunctionsBase64Benchmark keeps its earlier results from the benchmark-local sketch. They weren't re-measured, and the Javadoc says so.

Tests.

  • AdaptiveLatchTest:
    • each path is taken only in its state
    • an optimistic failure engages the latch and retries cautiously, including a repair case
    • the latch disengages after closeAfter clean calls, and a rejection restarts the count
    • a closeAfter below 1 is rejected
    • other exceptions propagate
    • reject yields the fallback, while a real null counts as clean
  • GuardedBase64DecodeTest: the differential tests above, plus engaging and disengaging behavior.
  • TextMapExtractAdapterTest passes for both Kafka modules, and checkConfigurations passes. All JMH sources compile.

Contributor Checklist

🤖 Generated with Claude Code

@dougqh dougqh added comp: core Tracer core tag: no release notes Changes to exclude from release notes tag: experimental Experimental changes tag: ai generated Largely based on code generated by an AI or LLM labels Sep 28, 2026
@dd-octo-sts

dd-octo-sts Bot commented Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

🟢 Java Benchmark SLOs — All performance SLOs passed

Suite Status
Startup 🟢 pass

SLO thresholds are defined here based on automatically generated metrics. A warning is raised when results are within 5% of the threshold.

PR vs. master results
Scenario Candidate master Δ (95% CI of mean)
startup:insecure-bank:iast:Agent 14.08 s 14.07 s [-0.7%; +0.9%] (no difference)
startup:insecure-bank:tracing:Agent 12.97 s 13.07 s [-1.5%; -0.1%] (maybe better)
startup:petclinic:appsec:Agent 17.21 s 17.15 s [-0.5%; +1.2%] (no difference)
startup:petclinic:iast:Agent 17.00 s 17.17 s [-1.9%; -0.1%] (maybe better)
startup:petclinic:profiling:Agent 16.83 s 16.74 s [-0.6%; +1.6%] (no difference)
startup:petclinic:sca:Agent 17.17 s 16.59 s [-1.2%; +8.1%] (no difference)
startup:petclinic:tracing:Agent 16.30 s 16.33 s [-1.5%; +1.0%] (no difference)

Commit: b9c3c997 · CI Pipeline · Benchmarking Platform UI


Load and DaCapo benchmarks can be triggered manually in the GitLab pipeline. Results will appear in the Benchmarking Platform UI after completion.

@dougqh dougqh changed the title Exploratory: JMH benchmark for a generic guarded-decode/breaker pattern Exploratory: JMH benchmark for a DynamicLatch guarded-decode pattern Sep 30, 2026
@datadog-prod-us1-5

datadog-prod-us1-5 Bot commented Sep 30, 2026 •

Copy link
Copy Markdown
Contributor

🎯 Code Coverage (details)
• Patch Coverage: 92.86%
• Overall Coverage: 59.32% (-0.01%)

This comment will be updated automatically if new data arrives.
🔗 Commit SHA: b9c3c99 | Docs | Give us feedback!

@dougqh dougqh changed the title Exploratory: JMH benchmark for a DynamicLatch guarded-decode pattern Exploratory: JMH benchmark for an AdaptiveLatch guarded-decode pattern Oct 2, 2026
dougqh and others added 8 commits October 5, 2026 13:43
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>
@dougqh
dougqh force-pushed the experiment/base64-guard-benchmark branch from 59c56f5 to 140a424 Compare October 5, 2026 17:55
@dougqh dougqh added inst: kafka Kafka instrumentation and removed tag: no release notes Changes to exclude from release notes tag: experimental Experimental changes labels Oct 5, 2026
@dougqh dougqh changed the title Exploratory: JMH benchmark for an AdaptiveLatch guarded-decode pattern Avoid repeated exceptions from malformed Kafka Base64 headers with AdaptiveLatch Oct 5, 2026
@dougqh dougqh added the type: bug fix Bug fix label Oct 5, 2026
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>
@pr-commenter

pr-commenter Bot commented Oct 5, 2026 •

Copy link
Copy Markdown

Kafka / producer-benchmark

Parameters

Baseline Candidate
baseline_or_candidate baseline candidate
git_branch master experiment/base64-guard-benchmark
git_commit_date 1791236323 1791297201
git_commit_sha 1380b85 b9c3c99
See matching parameters
Baseline Candidate
ci_job_date 1791298555 1791298555
ci_job_id 2113260119 2113260119
ci_pipeline_id 142746646 142746646
cpu_model Intel(R) Xeon(R) Platinum 8259CL CPU @ 2.50GHz Intel(R) Xeon(R) Platinum 8259CL CPU @ 2.50GHz
jdkVersion 11.0.31 11.0.31
jmhVersion 1.36 1.36
jvm /usr/lib/jvm/java-11-openjdk-amd64/bin/java /usr/lib/jvm/java-11-openjdk-amd64/bin/java
jvmArgs -Dhttp.proxyHost=127.0.0.1 -Dhttp.proxyPort=15002 -Dhttps.proxyHost=127.0.0.1 -Dhttps.proxyPort=15002 -Dhttp.nonProxyHosts=localhost *.localhost
kernel_version Linux runner-zfyrx7zua-project-304-concurrent-0-cp129fwy 6.8.0-1031-aws #33~22.04.1-Ubuntu SMP Thu Jun 26 14:22:30 UTC 2025 x86_64 x86_64 x86_64 GNU/Linux Linux runner-zfyrx7zua-project-304-concurrent-0-cp129fwy 6.8.0-1031-aws #33~22.04.1-Ubuntu SMP Thu Jun 26 14:22:30 UTC 2025 x86_64 x86_64 x86_64 GNU/Linux
vmName OpenJDK 64-Bit Server VM OpenJDK 64-Bit Server VM
vmVersion 11.0.31+11-post-1ubuntu1-22.04.2-Ubuntu 11.0.31+11-post-1ubuntu1-22.04.2-Ubuntu

Summary

Found 0 performance improvements and 0 performance regressions! Performance is the same for 3 metrics, 0 unstable metrics.

See unchanged results
scenario Δ mean throughput
scenario:not-instrumented/KafkaProduceBenchmark.benchProduce same
scenario:only-tracing-dsm-disabled-benchmarks/KafkaProduceBenchmark.benchProduce same
scenario:only-tracing-dsm-enabled-benchmarks/KafkaProduceBenchmark.benchProduce same

@pr-commenter

pr-commenter Bot commented Oct 5, 2026 •

Copy link
Copy Markdown

Kafka / consumer-benchmark

Parameters

Baseline Candidate
baseline_or_candidate baseline candidate
git_branch master experiment/base64-guard-benchmark
git_commit_date 1791236323 1791297201
git_commit_sha 1380b85 b9c3c99
See matching parameters
Baseline Candidate
ci_job_date 1791298628 1791298628
ci_job_id 2113260123 2113260123
ci_pipeline_id 142746646 142746646
cpu_model Intel(R) Xeon(R) Platinum 8259CL CPU @ 2.50GHz Intel(R) Xeon(R) Platinum 8259CL CPU @ 2.50GHz
jdkVersion 11.0.31 11.0.31
jmhVersion 1.36 1.36
jvm /usr/lib/jvm/java-11-openjdk-amd64/bin/java /usr/lib/jvm/java-11-openjdk-amd64/bin/java
jvmArgs -Dhttp.proxyHost=127.0.0.1 -Dhttp.proxyPort=15002 -Dhttps.proxyHost=127.0.0.1 -Dhttps.proxyPort=15002 -Dhttp.nonProxyHosts=localhost *.localhost
kernel_version Linux runner-zfyrx7zua-project-304-concurrent-0-815h19wg 6.8.0-1031-aws #33~22.04.1-Ubuntu SMP Thu Jun 26 14:22:30 UTC 2025 x86_64 x86_64 x86_64 GNU/Linux Linux runner-zfyrx7zua-project-304-concurrent-0-815h19wg 6.8.0-1031-aws #33~22.04.1-Ubuntu SMP Thu Jun 26 14:22:30 UTC 2025 x86_64 x86_64 x86_64 GNU/Linux
vmName OpenJDK 64-Bit Server VM OpenJDK 64-Bit Server VM
vmVersion 11.0.31+11-post-1ubuntu1-22.04.2-Ubuntu 11.0.31+11-post-1ubuntu1-22.04.2-Ubuntu

Summary

Found 0 performance improvements and 0 performance regressions! Performance is the same for 3 metrics, 0 unstable metrics.

See unchanged results
scenario Δ mean throughput
scenario:not-instrumented/KafkaConsumerBenchmark.benchConsume unsure
[+251.557op/s; +10622.730op/s] or [+0.085%; +3.583%]
scenario:only-tracing-dsm-disabled-benchmarks/KafkaConsumerBenchmark.benchConsume same
scenario:only-tracing-dsm-enabled-benchmarks/KafkaConsumerBenchmark.benchConsume same

dougqh and others added 2 commits October 5, 2026 15:25
… 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) {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

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

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 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>
@dougqh
dougqh marked this pull request as ready for review October 6, 2026 14:33
@dougqh
dougqh requested review from a team as code owners October 6, 2026 14:33
@dougqh
dougqh requested review from ValentinZakharov and bric3 and removed request for a team October 6, 2026 14:33
@chatgpt-codex-connector

chatgpt-codex-connector Bot commented Oct 6, 2026 •

Copy link
Copy Markdown

Codex Review Summary

This comment shows the latest Codex review activity on this pull request.

Review Status Commit Review trigger
📝 Code Review ✅ Completed 2026-10-06T14:42:25.589171Z 3c7c25f Draft marked ready
ℹ️ About Codex in GitHub

Your team has set up Codex to 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" or "@codex security review".

Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings.

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

💡 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".

Comment on lines +37 to +38
this.headerValueTransformer =
guardEnabled ? new Functions.GuardedBase64Decode()::tryApply : BASE64_DECODE;

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.

@datadog-prod-us1-5 datadog-prod-us1-5 Bot left a comment

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.

Bits Code Review: PASS

More details

Focused validation supports the guarded decoder’s latch behavior and Base64 handling. Kafka 0.11 tests remain unverified because its Confluent dependency fetch returned a sandbox 502.

Was this helpful? React 👍 or 👎

Open Bits AI session

🤖 Bits Code Review · Commit 3c7c25f · @DataDog review to ask questions

@dougqh dougqh changed the title Avoid repeated exceptions from malformed Kafka Base64 headers with AdaptiveLatch Avoid repeated exceptions from malformed Kafka Base64 headers with AdaptiveLatch ("quick" fix) Oct 7, 2026
Comment on lines +18 to +61
/**
* {@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.
*/

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

return result;
}

// Same hysteresis, but preserves throw-based failure semantics: a precheck-known failure

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: 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;

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.

note: FYI this config key is missing form the registry

Comment on lines +195 to +208
/**
* 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.
*/

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: neat doc, did clarify doc skill helped there?

Comment on lines +118 to +183
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;
}
}
}

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.

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.

Suggested change
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;
}
}
}

Comment on lines +75 to +98
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");
}
}

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");
}
}

Comment on lines +67 to +78
* 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.

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: 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.

Suggested change
* 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.

Comment on lines +52 to +53
* 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.

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.

* final unit, and nothing may follow it; a final unit of one character is rejected.
*/
@Nullable
static String decodeOrNull(byte[] src) {

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 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.

Comment on lines +37 to +38
this.headerValueTransformer =
guardEnabled ? new Functions.GuardedBase64Decode()::tryApply : BASE64_DECODE;

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.

@dougqh
dougqh marked this pull request as draft October 7, 2026 19:53
dougqh added a commit that referenced this pull request Oct 7, 2026
Brought over unchanged, with its tests, for use in URL-to-URI conversion.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

comp: core Tracer core inst: kafka Kafka instrumentation tag: ai generated Largely based on code generated by an AI or LLM type: bug fix Bug fix

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants