From 572b1f6e7e3bca4a844cce3a434884632fc09c0e Mon Sep 17 00:00:00 2001 From: Vikentiy Fesunov Date: Tue, 11 Aug 2026 12:37:47 +0200 Subject: [PATCH] Prevent payload corruption on buffer overflow If caller tries to serialize a metric larger than the payload size limit, we throw BufferOverflow exception, and the user is requested to call `resetMetric`. Previously, it failed to correctly rewind everything, in particular interners and delta encoders, to the state they had prior to the errant metric, and subsequent entries using dictionary elements of the failed metric would produce dangling references. This patch adds commit/revert methods to Interner and DeltaEncoder. When the record in progress is added to the final payload, we make the state of inerner and encoder permanent. On failure, revert resets to the last committed state. --- .../http/serializer/DeltaEncoder.java | 10 +++ .../dogstatsd/http/serializer/Interner.java | 17 +++++ .../http/serializer/PayloadBuilder.java | 28 ++++++++ .../http/serializer/PayloadBuilderTest.java | 69 +++++++++++++++++++ .../dogstatsd/http/serializer/TestUtil.java | 13 +++- 5 files changed, 135 insertions(+), 2 deletions(-) diff --git a/dogstatsd-http/core/src/main/java/com/datadoghq/dogstatsd/http/serializer/DeltaEncoder.java b/dogstatsd-http/core/src/main/java/com/datadoghq/dogstatsd/http/serializer/DeltaEncoder.java index 1668b401..f4af894e 100644 --- a/dogstatsd-http/core/src/main/java/com/datadoghq/dogstatsd/http/serializer/DeltaEncoder.java +++ b/dogstatsd-http/core/src/main/java/com/datadoghq/dogstatsd/http/serializer/DeltaEncoder.java @@ -9,6 +9,7 @@ class DeltaEncoder { private long prev = 0; + private long undo = 0; long encode(long v) { long r = v - prev; @@ -16,7 +17,16 @@ long encode(long v) { return r; } + void commit() { + undo = prev; + } + + void revert() { + prev = undo; + } + void clear() { prev = 0; + undo = 0; } } diff --git a/dogstatsd-http/core/src/main/java/com/datadoghq/dogstatsd/http/serializer/Interner.java b/dogstatsd-http/core/src/main/java/com/datadoghq/dogstatsd/http/serializer/Interner.java index bef16438..eb8edd20 100644 --- a/dogstatsd-http/core/src/main/java/com/datadoghq/dogstatsd/http/serializer/Interner.java +++ b/dogstatsd-http/core/src/main/java/com/datadoghq/dogstatsd/http/serializer/Interner.java @@ -7,6 +7,7 @@ package com.datadoghq.dogstatsd.http.serializer; +import java.util.ArrayList; import java.util.HashMap; class Interner { @@ -15,6 +16,8 @@ interface Encoder { } private HashMap inner = new HashMap<>(); + private ArrayList undo = new ArrayList<>(); + private long lastId = 0; private final Encoder encoder; private final T empty; @@ -38,12 +41,26 @@ long intern(T val) { encoder.encode(val); lastId++; + undo.add(val); inner.put(val, lastId); return lastId; } + void commit() { + undo.clear(); + } + + void revert() { + for (T it : undo) { + inner.remove(it); + } + lastId -= undo.size(); + undo.clear(); + } + void clear() { inner.clear(); + undo.clear(); lastId = 0; } } diff --git a/dogstatsd-http/core/src/main/java/com/datadoghq/dogstatsd/http/serializer/PayloadBuilder.java b/dogstatsd-http/core/src/main/java/com/datadoghq/dogstatsd/http/serializer/PayloadBuilder.java index a21983f7..889377c5 100644 --- a/dogstatsd-http/core/src/main/java/com/datadoghq/dogstatsd/http/serializer/PayloadBuilder.java +++ b/dogstatsd-http/core/src/main/java/com/datadoghq/dogstatsd/http/serializer/PayloadBuilder.java @@ -230,7 +230,21 @@ void endMetric() { m.clearDependentFields(); m.encodeDependentFields(); } + payload.put(record); + + nameStrInterner.commit(); + tagStrInterner.commit(); + tagsInterner.commit(); + resourceStrInterner.commit(); + resourcesInterner.commit(); + originInfoInterner.commit(); + + nameRefsDelta.commit(); + tagsetRefsDelta.commit(); + resourcesRefsDelta.commit(); + originInfoRefsDelta.commit(); + timestampsDelta.commit(); } finally { resetMetric(); } @@ -247,6 +261,20 @@ public void resetMetric() { timestamps.clear(); values.clear(); counts.clear(); + + nameStrInterner.revert(); + tagStrInterner.revert(); + tagsInterner.revert(); + resourceStrInterner.revert(); + resourcesInterner.revert(); + originInfoInterner.revert(); + + nameRefsDelta.revert(); + tagsetRefsDelta.revert(); + resourcesRefsDelta.revert(); + originInfoRefsDelta.revert(); + timestampsDelta.revert(); + metricInProgress = null; } diff --git a/dogstatsd-http/core/src/test/java/com/datadoghq/dogstatsd/http/serializer/PayloadBuilderTest.java b/dogstatsd-http/core/src/test/java/com/datadoghq/dogstatsd/http/serializer/PayloadBuilderTest.java index b07deb42..6f868fe6 100644 --- a/dogstatsd-http/core/src/test/java/com/datadoghq/dogstatsd/http/serializer/PayloadBuilderTest.java +++ b/dogstatsd-http/core/src/test/java/com/datadoghq/dogstatsd/http/serializer/PayloadBuilderTest.java @@ -8,12 +8,15 @@ package com.datadoghq.dogstatsd.http.serializer; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertThrows; import static org.junit.Assert.assertTrue; import com.datadoghq.dogstatsd.Sketch; +import java.nio.BufferOverflowException; import java.util.ArrayList; import java.util.Arrays; import org.junit.Test; +import org.junit.function.ThrowingRunnable; public class PayloadBuilderTest { @Test @@ -370,4 +373,70 @@ public void handle(byte[] p) { assertTrue(p.length <= b.maxPayloadSize); } } + + @Test + // Make sure we leave builder in consistent state on error. + public void rollback() { + final ArrayList payloads1 = new ArrayList<>(); + PayloadBuilder pb1 = + new PayloadBuilder( + new PayloadConsumer() { + @Override + public void handle(byte[] p) { + payloads1.add(p); + } + }); + + final ArrayList payloads2 = new ArrayList<>(); + PayloadBuilder pb2 = + new PayloadBuilder( + new PayloadConsumer() { + @Override + public void handle(byte[] p) { + payloads2.add(p); + } + }); + + pb1.count("m.before") + .setTags(Arrays.asList(new String[] {"env:prod"})) + .addPoint(100, 1) + .close(); + pb2.count("m.before") + .setTags(Arrays.asList(new String[] {"env:prod"})) + .addPoint(100, 1) + .close(); + + // 40k * 8 bytes = 320 KB > 256 KB payload limit. + final ScalarMetric huge = + pb2.gauge("m.huge").setTags(Arrays.asList(new String[] {"env:prod", "kind:huge"})); + for (int i = 0; i < 40000; i++) { + huge.addPoint(1000, 3.14); + } + assertThrows( + BufferOverflowException.class, + new ThrowingRunnable() { + @Override + public void run() { + huge.close(); + } + }); + pb2.resetMetric(); + + pb1.gauge("m.huge") + .setTags(Arrays.asList(new String[] {"env:prod", "kind:huge"})) + .addPoint(200, 2) + .close(); + pb2.gauge("m.huge") + .setTags(Arrays.asList(new String[] {"env:prod", "kind:huge"})) + .addPoint(200, 2) + .close(); + + pb1.close(); + pb2.close(); + + assertEquals(payloads1.size(), payloads2.size()); + for (int i = 0; i < payloads1.size(); i++) { + TestUtil.assertPayload(payloads2.get(i), payloads1.get(i)); + } + } } diff --git a/dogstatsd-http/core/src/test/java/com/datadoghq/dogstatsd/http/serializer/TestUtil.java b/dogstatsd-http/core/src/test/java/com/datadoghq/dogstatsd/http/serializer/TestUtil.java index 0de6793f..f4db0530 100644 --- a/dogstatsd-http/core/src/test/java/com/datadoghq/dogstatsd/http/serializer/TestUtil.java +++ b/dogstatsd-http/core/src/test/java/com/datadoghq/dogstatsd/http/serializer/TestUtil.java @@ -103,9 +103,14 @@ static String protodump(byte[] p) { } static String protodump(int[] p) { + return protodump(toBytes(p)); + } + + // int[] lets tests write unsigned byte values without conversion. + static byte[] toBytes(int[] p) { byte[] pb = new byte[p.length]; for (int i = 0; i < p.length; i++) pb[i] = (byte) p[i]; - return protodump(pb); + return pb; } static void formatTwoCols(Formatter out, String hl, String hr, String dl, String dr) { @@ -129,9 +134,13 @@ static void formatTwoCols(Formatter out, String hl, String hr, String dl, String // expected is int[] to be able to write unsigned byte values wtihout conversion. static void assertPayload(byte[] got, int[] expected) { + assertPayload(got, toBytes(expected)); + } + + static void assertPayload(byte[] got, byte[] expected) { boolean same = got.length == expected.length; for (int i = 0; same && i < got.length; i++) { - same &= got[i] == (byte) expected[i]; + same &= got[i] == expected[i]; } if (!same) { Formatter out = new Formatter();