Fix kafka clusterId gathering possible to crash#5837
Conversation
This improves the kafka client performance as well as fixing a bug in case the clusterId gathering would have errored. In that case, the library would have ended with an unhandled rejection.
Overall package sizeSelf size: 9.6 MB Dependency sizes| name | version | self size | total size | |------|---------|-----------|------------| | @datadog/libdatadog | 0.6.0 | 30.47 MB | 30.47 MB | | @datadog/native-appsec | 8.5.2 | 19.33 MB | 19.34 MB | | @datadog/pprof | 5.8.0 | 12.55 MB | 12.92 MB | | @datadog/native-iast-taint-tracking | 4.0.0 | 11.72 MB | 11.73 MB | | @opentelemetry/core | 1.30.1 | 908.66 kB | 7.16 MB | | protobufjs | 7.5.3 | 2.95 MB | 5.6 MB | | @datadog/wasm-js-rewriter | 4.0.1 | 2.85 MB | 3.58 MB | | @datadog/native-metrics | 3.1.1 | 1.02 MB | 1.43 MB | | @opentelemetry/api | 1.8.0 | 1.21 MB | 1.21 MB | | import-in-the-middle | 1.14.0 | 120.58 kB | 841.68 kB | | source-map | 0.7.4 | 226 kB | 226 kB | | opentracing | 0.14.7 | 194.81 kB | 194.81 kB | | lru-cache | 7.18.3 | 133.92 kB | 133.92 kB | | pprof-format | 2.1.0 | 111.69 kB | 111.69 kB | | @datadog/sketches-js | 2.1.1 | 109.9 kB | 109.9 kB | | lodash.sortby | 4.7.0 | 75.76 kB | 75.76 kB | | ignore | 5.3.2 | 53.63 kB | 53.63 kB | | istanbul-lib-coverage | 3.2.2 | 34.37 kB | 34.37 kB | | rfdc | 1.4.1 | 27.15 kB | 27.15 kB | | @isaacs/ttlcache | 1.4.1 | 25.2 kB | 25.2 kB | | dc-polyfill | 0.1.9 | 25.11 kB | 25.11 kB | | tlhunter-sorted-set | 0.1.0 | 24.94 kB | 24.94 kB | | shell-quote | 1.8.2 | 23.54 kB | 23.54 kB | | limiter | 1.1.5 | 23.17 kB | 23.17 kB | | retry | 0.13.1 | 18.85 kB | 18.85 kB | | semifies | 1.0.0 | 15.84 kB | 15.84 kB | | jest-docblock | 29.7.0 | 8.99 kB | 12.76 kB | | crypto-randomuuid | 1.0.0 | 11.18 kB | 11.18 kB | | ttl-set | 1.0.0 | 4.61 kB | 9.69 kB | | mutexify | 1.4.0 | 5.71 kB | 8.74 kB | | path-to-regexp | 0.1.12 | 6.6 kB | 6.6 kB | | koalas | 1.0.2 | 6.47 kB | 6.47 kB | | module-details-from-path | 1.0.4 | 3.96 kB | 3.96 kB |🤖 This report was automatically generated by heaviest-objects-in-the-universe |
Codecov ReportAll modified and coverable lines are covered by tests ✅
Additional details and impacted files@@ Coverage Diff @@
## master #5837 +/- ##
==========================================
- Coverage 79.19% 78.96% -0.24%
==========================================
Files 522 516 -6
Lines 24060 23867 -193
==========================================
- Hits 19054 18846 -208
- Misses 5006 5021 +15 ☔ View full report in Codecov by Sentry. 🚀 New features to boost your workflow:
|
BenchmarksBenchmark execution time: 2025-06-05 22:53:25 Comparing candidate commit dfe2511 in PR branch Found 0 performance improvements and 0 performance regressions! Performance is the same for 1272 metrics, 51 unstable metrics. |
Datadog ReportBranch report: ✅ 0 Failed, 1257 Passed, 0 Skipped, 20m 50.67s Total Time |
This stops injecting headers into users input objects. Instead copy these. In addition add a couple comments how to improve the current logic to not add additional requests to the broker. Instead, just instrument the library to receive the necessary information.
| const dataStreamsContext = this.tracer.setCheckpoint(edgeTags, span, payloadSize) | ||
| if (!disableHeaderInjection) { | ||
| // Message headers are not supported for kafka broker versions <0.11 | ||
| this.tracer.inject(span, 'text_map', message.headers) |
There was a problem hiding this comment.
I am not sure if this move is actually correct. That would have an impact on other changes in here as well.
|
@BridgeAR Is this PR still relevant? |
|
@rochdev yes, thanks for the ping. I will update it soon. |
…er input Three coordinated changes to the kafka producer wrap: 1. The cluster id discovery used a separate admin connection whose `describeCluster()` rejection wasn't awaited, surfacing as `process.unhandledRejection`. Read `cluster.brokerPool.metadata.clusterId` instead, priming it on first send via `cluster.refreshMetadataIfNecessary()`; `sharedPromiseTo` collapses our call with kafkajs's internal one, so total latency is unchanged. 2. Header injection mutated caller-owned `message.headers` in place. The boundary now hands kafkajs and the plugin a shallow clone instead. 3. The Produce request didn't carry headers until v3 (Kafka 0.11), so injecting on older brokers either wastes work or trips `UNKNOWN_SERVER_ERROR` on a per-leader version mismatch. Read `cluster.brokerPool.versions[0].maxVersion` and skip injection when the broker negotiated <v3; the existing UNKNOWN handler stays as a safety net for mixed-version clusters. Refs: #5837 Refs: #5270
…JS-compat path The KafkaJS-compat producer's `send` wrapper handed the caller's `payload.messages` array straight to the underlying client, so both trace-header injection and the client's own `headers: null` post-publish fixup mutated those objects in place. Shallow-clone the array on every send so neither path touches caller-owned objects; the native producer path already builds its own wrapper objects and is unaffected. The disabled-header tracking moves from a module-level WeakSet to a per-producer closure so the state's lifetime matches the producer it describes. Refs: #5837
…er input Three coordinated changes to the kafka producer wrap: 1. The cluster id discovery used a separate admin connection whose `describeCluster()` rejection wasn't awaited, surfacing as `process.unhandledRejection`. Read `cluster.brokerPool.metadata.clusterId` instead, priming it on first send via `cluster.refreshMetadataIfNecessary()`; `sharedPromiseTo` collapses our call with kafkajs's internal one, so total latency is unchanged. 2. Header injection mutated caller-owned `message.headers` in place. The boundary now hands kafkajs and the plugin a shallow clone instead. 3. The Produce request didn't carry headers until v3 (Kafka 0.11), so injecting on older brokers either wastes work or trips `UNKNOWN_SERVER_ERROR` on a per-leader version mismatch. Read `cluster.brokerPool.versions[0].maxVersion` and skip injection when the broker negotiated <v3; the existing UNKNOWN handler stays as a safety net for mixed-version clusters. Refs: #5837 Refs: #5270
…JS-compat path The KafkaJS-compat producer's `send` wrapper handed the caller's `payload.messages` array straight to the underlying client, so both trace-header injection and the client's own `headers: null` post-publish fixup mutated those objects in place. Shallow-clone the array on every send so neither path touches caller-owned objects; the native producer path already builds its own wrapper objects and is unaffected. The disabled-header tracking moves from a module-level WeakSet to a per-producer closure so the state's lifetime matches the producer it describes. Refs: #5837
|
Superseded |
…er input Three coordinated changes to the kafka producer wrap: 1. The cluster id discovery used a separate admin connection whose `describeCluster()` rejection wasn't awaited, surfacing as `process.unhandledRejection`. Read `cluster.brokerPool.metadata.clusterId` instead, priming it on first send via `cluster.refreshMetadataIfNecessary()`; `sharedPromiseTo` collapses our call with kafkajs's internal one, so total latency is unchanged. 2. Header injection mutated caller-owned `message.headers` in place. The boundary now hands kafkajs and the plugin a shallow clone instead. 3. The Produce request didn't carry headers until v3 (Kafka 0.11), so injecting on older brokers either wastes work or trips `UNKNOWN_SERVER_ERROR` on a per-leader version mismatch. Read `cluster.brokerPool.versions[0].maxVersion` and skip injection when the broker negotiated <v3; the existing UNKNOWN handler stays as a safety net for mixed-version clusters. Refs: #5837 Refs: #5270
…JS-compat path The KafkaJS-compat producer's `send` wrapper handed the caller's `payload.messages` array straight to the underlying client, so both trace-header injection and the client's own `headers: null` post-publish fixup mutated those objects in place. Shallow-clone the array on every send so neither path touches caller-owned objects; the native producer path already builds its own wrapper objects and is unaffected. The disabled-header tracking moves from a module-level WeakSet to a per-producer closure so the state's lifetime matches the producer it describes. Refs: #5837
…er input Three coordinated changes to the kafka producer wrap: 1. The cluster id discovery used a separate admin connection whose `describeCluster()` rejection wasn't awaited, surfacing as `process.unhandledRejection`. Read `cluster.brokerPool.metadata.clusterId` instead, priming it on first send via `cluster.refreshMetadataIfNecessary()`; `sharedPromiseTo` collapses our call with kafkajs's internal one, so total latency is unchanged. 2. Header injection mutated caller-owned `message.headers` in place. The boundary now hands kafkajs and the plugin a shallow clone instead. 3. The Produce request didn't carry headers until v3 (Kafka 0.11), so injecting on older brokers either wastes work or trips `UNKNOWN_SERVER_ERROR` on a per-leader version mismatch. Read `cluster.brokerPool.versions[0].maxVersion` and skip injection when the broker negotiated <v3; the existing UNKNOWN handler stays as a safety net for mixed-version clusters. Refs: #5837 Refs: #5270
…JS-compat path The KafkaJS-compat producer's `send` wrapper handed the caller's `payload.messages` array straight to the underlying client, so both trace-header injection and the client's own `headers: null` post-publish fixup mutated those objects in place. Shallow-clone the array on every send so neither path touches caller-owned objects; the native producer path already builds its own wrapper objects and is unaffected. The disabled-header tracking moves from a module-level WeakSet to a per-producer closure so the state's lifetime matches the producer it describes. Refs: #5837
This improves the kafka client performance as well as fixing a bug in case the clusterId gathering would have errored. In that case, the library would have ended with an unhandled rejection.