diff --git a/components/json/src/main/java/datadog/json/JsonWriter.java b/components/json/src/main/java/datadog/json/JsonWriter.java index 0fcde05a629..0420f7b7f7b 100644 --- a/components/json/src/main/java/datadog/json/JsonWriter.java +++ b/components/json/src/main/java/datadog/json/JsonWriter.java @@ -19,6 +19,7 @@ public final class JsonWriter implements Flushable, AutoCloseable { private final JsonStructure structure; private boolean requireComma; + private int bytesWritten; /** Creates a writer with structure check. */ public JsonWriter() { @@ -245,6 +246,11 @@ public void flush() { } } + /** Approximate number of bytes written so far. */ + public int sizeInBytes() { + return this.bytesWritten; + } + @Override public void close() { try { @@ -268,6 +274,7 @@ private void endsValue() { private void write(char ch) { try { this.writer.write(ch); + this.bytesWritten++; } catch (IOException ignored) { } } @@ -275,6 +282,7 @@ private void write(char ch) { private void writeStringLiteral(String str) { try { this.writer.write('"'); + int count = 1; for (int i = 0; i < str.length(); ++i) { char c = str.charAt(i); @@ -286,6 +294,7 @@ private void writeStringLiteral(String str) { this.writer.write(HEX_DIGITS[(c >>> 8) & 0xF]); this.writer.write(HEX_DIGITS[(c >>> 4) & 0xF]); this.writer.write(HEX_DIGITS[c & 0xF]); + count += 6; } else { switch (c) { case '"': // Quotation mark @@ -293,26 +302,32 @@ private void writeStringLiteral(String str) { case '/': // Solidus this.writer.write('\\'); this.writer.write(c); + count += 2; break; case '\b': // Backspace this.writer.write('\\'); this.writer.write('b'); + count += 2; break; case '\f': // Form feed this.writer.write('\\'); this.writer.write('f'); + count += 2; break; case '\n': // Line feed this.writer.write('\\'); this.writer.write('n'); + count += 2; break; case '\r': // Carriage return this.writer.write('\\'); this.writer.write('r'); + count += 2; break; case '\t': // Horizontal tab this.writer.write('\\'); this.writer.write('t'); + count += 2; break; default: if (c < 0x20) { @@ -322,8 +337,10 @@ private void writeStringLiteral(String str) { this.writer.write('0'); this.writer.write(HEX_DIGITS[(c >>> 4) & 0xF]); this.writer.write(HEX_DIGITS[c & 0xF]); + count += 6; } else { this.writer.write(c); + count += 1; } break; } @@ -331,6 +348,9 @@ private void writeStringLiteral(String str) { } this.writer.write('"'); + count += 1; + + this.bytesWritten += count; } catch (IOException ignored) { } } @@ -338,6 +358,7 @@ private void writeStringLiteral(String str) { private void writeStringRaw(String str) { try { this.writer.write(str); + this.bytesWritten += str.length(); // exact if ASCII, estimate otherwise } catch (IOException ignored) { } } diff --git a/components/json/src/test/java/datadog/json/JsonWriterTest.java b/components/json/src/test/java/datadog/json/JsonWriterTest.java index d6239a053cc..68624cfa936 100644 --- a/components/json/src/test/java/datadog/json/JsonWriterTest.java +++ b/components/json/src/test/java/datadog/json/JsonWriterTest.java @@ -165,6 +165,54 @@ void testCompleteObject() { } } + @Test + void testSizeInBytes() { + try (JsonWriter writer = new JsonWriter()) { + assertEquals(0, writer.sizeInBytes(), "Check empty writer size"); + + writer.beginArray(); + assertSizeInBytes(writer, "Check size after plain ASCII string"); + writer.value("bar"); + + assertSizeInBytes(writer, "Check size after escaped characters"); + writer.value("\"\\/"); + + assertSizeInBytes(writer, "Check size after named escapes"); + writer.value("\b\f\n\r\t"); + + assertSizeInBytes(writer, "Check size after control character escape"); + writer.value("\u0001"); + + assertSizeInBytes(writer, "Check size after non-ASCII character escape"); + writer.value("café"); + + assertSizeInBytes(writer, "Check size after int value"); + writer.value(3); + + assertSizeInBytes(writer, "Check size after long value"); + writer.value(3456789123L); + + assertSizeInBytes(writer, "Check size after float value"); + writer.value(3.142f); + + assertSizeInBytes(writer, "Check size after double value"); + writer.value(PI); + + assertSizeInBytes(writer, "Check size after boolean value"); + writer.value(true); + + assertSizeInBytes(writer, "Check size after null value"); + writer.nullValue(); + + writer.endArray(); + assertSizeInBytes(writer, "Check final size matches written bytes"); + } + } + + private static void assertSizeInBytes(JsonWriter writer, String message) { + assertEquals(writer.toByteArray().length, writer.sizeInBytes(), message); + } + @Test void testCompleteArray() { try (JsonWriter writer = new JsonWriter()) { diff --git a/dd-trace-core/src/main/java/datadog/trace/common/writer/OtlpPayloadDispatcher.java b/dd-trace-core/src/main/java/datadog/trace/common/writer/OtlpPayloadDispatcher.java index b2817b46dbc..50eed88b627 100644 --- a/dd-trace-core/src/main/java/datadog/trace/common/writer/OtlpPayloadDispatcher.java +++ b/dd-trace-core/src/main/java/datadog/trace/common/writer/OtlpPayloadDispatcher.java @@ -7,8 +7,14 @@ import java.util.Collection; import java.util.Collections; import java.util.List; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; final class OtlpPayloadDispatcher implements PayloadDispatcher { + private static final Logger log = LoggerFactory.getLogger(OtlpPayloadDispatcher.class); + + private static final int FLUSH_THRESHOLD_BYTES = 5 << 20; // 5 MiB + private final OtlpTraceCollector collector; private final OtlpSender sender; @@ -20,13 +26,21 @@ final class OtlpPayloadDispatcher implements PayloadDispatcher { @Override public void addTrace(List> trace) { collector.addTrace(trace); + // flush proactively to keep payload size bounded + if (collector.sizeInBytes() >= FLUSH_THRESHOLD_BYTES) { + flush(); + } } @Override public void flush() { - OtlpPayload payload = collector.collectTraces(); - if (payload != OtlpPayload.EMPTY) { - sender.send(payload); + try { + OtlpPayload payload = collector.collectTraces(); + if (payload != OtlpPayload.EMPTY) { + sender.send(payload); + } + } catch (RuntimeException e) { // don't catch severe Errors + log.debug("Failed to send OTLP payload", e); } } diff --git a/dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpProtoBuffer.java b/dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpProtoBuffer.java index eab87bf6096..81196bfaa23 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpProtoBuffer.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpProtoBuffer.java @@ -18,12 +18,23 @@ * @see GrowableBuffer */ public final class OtlpProtoBuffer { + // hard limit to avoid unbounded buffering; matches OTLP spec's recommended default + public static final int MAX_CAPACITY_BYTES = 64 << 20; // 64 MiB + private final int initialCapacity; private ByteBuffer buffer; private int remaining; public OtlpProtoBuffer(int requiredCapacity) { this.initialCapacity = nextPowerOfTwo(requiredCapacity); + if (this.initialCapacity > MAX_CAPACITY_BYTES) { + throw new IllegalArgumentException( + "OTLP payload initial capacity of " + + this.initialCapacity + + " bytes exceeds maximum buffer size of " + + MAX_CAPACITY_BYTES + + " bytes"); + } this.buffer = ByteBuffer.allocate(initialCapacity); this.remaining = initialCapacity; } @@ -96,6 +107,11 @@ public ByteBuffer flip() { return buffer; } + /** Returns the number of bytes currently recorded in the buffer. */ + public int sizeInBytes() { + return buffer.capacity() - remaining; + } + /** * Returns an {@link OtlpPayload} containing the protobuf encoded content. * @@ -123,10 +139,21 @@ private void checkCapacity(int required) { ByteBuffer oldBuffer = flip(); int oldSize = oldBuffer.remaining(); // round up to next multiple of initialCapacity that can accommodate required - int newSize = (oldSize + required + initialCapacity - 1) & -initialCapacity; - ByteBuffer newBuffer = ByteBuffer.allocate(newSize); + // (uses long arithmetic so overflow can be detected before allocating) + long newSize = ((long) oldSize + required + initialCapacity - 1) & -initialCapacity; + if (newSize > MAX_CAPACITY_BYTES) { + throw new IllegalStateException( + "OTLP payload exceeds maximum buffer size of " + + MAX_CAPACITY_BYTES + + " bytes: " + + oldSize + + " bytes buffered, " + + required + + " more requested"); + } + ByteBuffer newBuffer = ByteBuffer.allocate((int) newSize); // copy over old content so it stays at the far end - remaining = newSize - oldSize; + remaining = (int) newSize - oldSize; newBuffer.position(remaining); newBuffer.put(oldBuffer); buffer = newBuffer; diff --git a/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceCollector.java b/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceCollector.java index dda40ff74b2..127e27af467 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceCollector.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceCollector.java @@ -15,6 +15,9 @@ public abstract class OtlpTraceCollector { /** Collects all spans added since the last collection. */ public abstract OtlpPayload collectTraces(); + /** Returns the number of bytes buffered since the last collection. */ + public abstract int sizeInBytes(); + protected final boolean shouldExport(CoreSpan span) { return span.samplingPriority() > 0 // trace-level sampling priority || span.getTag(SPAN_SAMPLING_MECHANISM_TAG) != null; // span-level sampling priority diff --git a/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceJsonCollector.java b/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceJsonCollector.java index c5da9bcddaf..94fe3ed5706 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceJsonCollector.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceJsonCollector.java @@ -2,6 +2,7 @@ import static datadog.trace.core.otlp.common.OtlpCommonJson.writeScopeAndSchema; import static datadog.trace.core.otlp.common.OtlpPayload.JSON_CONTENT_TYPE; +import static datadog.trace.core.otlp.common.OtlpProtoBuffer.MAX_CAPACITY_BYTES; import static datadog.trace.core.otlp.common.OtlpResourceJson.TRACE_RESOURCE_FRAGMENT; import static datadog.trace.core.otlp.trace.OtlpTraceJson.writeSpan; @@ -54,11 +55,22 @@ public void addTrace(List> spans) { payloadStarted = true; } - for (CoreSpan span : spans) { - visitSpan(span); + try { + for (CoreSpan span : spans) { + visitSpan(span); + } + } catch (Throwable e) { + // reset the buffer for subsequent traces + stop(); + throw e; } } + @Override + public int sizeInBytes() { + return writer == null ? 0 : writer.sizeInBytes(); + } + /** * Marshals the traces collected so far into a JSON payload. * @@ -178,5 +190,10 @@ private void completeSpan() { // reset temporary elements for next span currentSpan = null; currentSpanLinks = Collections.emptyList(); + + if (writer.sizeInBytes() > MAX_CAPACITY_BYTES) { + throw new IllegalStateException( + "OTLP payload exceeds maximum buffer size of " + MAX_CAPACITY_BYTES + " bytes"); + } } } diff --git a/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceProtoCollector.java b/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceProtoCollector.java index fa1efa87fba..10f89056210 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceProtoCollector.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceProtoCollector.java @@ -56,12 +56,23 @@ public void addTrace(List> spans) { payloadStarted = true; } - // OtlpProtoBuffer collects spans in reverse - for (int i = spans.size() - 1; i >= 0; i--) { - visitSpan(spans.get(i)); + try { + // OtlpProtoBuffer collects spans in reverse + for (int i = spans.size() - 1; i >= 0; i--) { + visitSpan(spans.get(i)); + } + } catch (Throwable e) { + // reset the buffer for subsequent traces + stop(); + throw e; } } + @Override + public int sizeInBytes() { + return protobuf.sizeInBytes(); + } + /** * Marshals the traces collected so far into a chunked payload. * diff --git a/dd-trace-core/src/test/java/datadog/trace/common/writer/OtlpPayloadDispatcherTest.java b/dd-trace-core/src/test/java/datadog/trace/common/writer/OtlpPayloadDispatcherTest.java index 6e55ebeadaa..0e0b44a08a3 100644 --- a/dd-trace-core/src/test/java/datadog/trace/common/writer/OtlpPayloadDispatcherTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/common/writer/OtlpPayloadDispatcherTest.java @@ -5,6 +5,7 @@ import static java.util.Collections.singletonList; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.any; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.verifyNoInteractions; @@ -95,6 +96,68 @@ void emptyTraceForwardsNothing() { verifyNoInteractions(sender); } + @Test + void largeTraceTriggersProactiveFlushWithoutExplicitFlush() { + OtlpPayloadDispatcher dispatcher = new OtlpPayloadDispatcher(sender, collector); + collector.fakeSizeInBytes = Integer.MAX_VALUE; + + dispatcher.addTrace(singletonList(sampledSpan())); + + // no explicit dispatcher.flush() call + ArgumentCaptor captor = ArgumentCaptor.forClass(OtlpPayload.class); + verify(sender).send(captor.capture()); + assertEquals(1 /*spans*/, captor.getValue().getContentLength()); + } + + @Test + void belowFlushThresholdDoesNotTriggerProactiveFlush() { + OtlpPayloadDispatcher dispatcher = new OtlpPayloadDispatcher(sender, collector); + collector.fakeSizeInBytes = (5 << 20) - 1; // one byte under FLUSH_THRESHOLD_BYTES + + dispatcher.addTrace(singletonList(sampledSpan())); + + verifyNoInteractions(sender); + } + + @Test + void atFlushThresholdTriggersProactiveFlush() { + OtlpPayloadDispatcher dispatcher = new OtlpPayloadDispatcher(sender, collector); + collector.fakeSizeInBytes = 5 << 20; // exactly FLUSH_THRESHOLD_BYTES + + dispatcher.addTrace(singletonList(sampledSpan())); + + verify(sender).send(any(OtlpPayload.class)); + } + + @Test + void flushSwallowsCollectorFailureInsteadOfPropagating() { + // e.g. a held-back span pushes the payload over the buffer's hard cap only once + // collectTraces() finalizes it during a scheduled flush, not during addTrace() + OtlpPayloadDispatcher dispatcher = new OtlpPayloadDispatcher(sender, collector); + dispatcher.addTrace(singletonList(sampledSpan())); + collector.throwOnCollect = true; + + dispatcher.flush(); + + verifyNoInteractions(sender); + } + + @Test + void dispatcherRemainsUsableAfterCollectorFailure() { + OtlpPayloadDispatcher dispatcher = new OtlpPayloadDispatcher(sender, collector); + dispatcher.addTrace(singletonList(sampledSpan())); + collector.throwOnCollect = true; + dispatcher.flush(); + + collector.throwOnCollect = false; + dispatcher.addTrace(singletonList(sampledSpan())); + dispatcher.flush(); + + ArgumentCaptor captor = ArgumentCaptor.forClass(OtlpPayload.class); + verify(sender).send(captor.capture()); + assertEquals(1 /*spans*/, captor.getValue().getContentLength()); + } + @Test void getApisIsEmpty() { OtlpPayloadDispatcher dispatcher = new OtlpPayloadDispatcher(sender, collector); @@ -143,6 +206,12 @@ private static CoreSpan unsetSpan() { private static class TestCollector extends OtlpTraceCollector { final List> spansToExport = new ArrayList<>(); + // lets tests drive the proactive-flush threshold independently of spansToExport + int fakeSizeInBytes; + + // lets tests simulate a collectTraces() failure, e.g. a buffer overflow on the held-back span + boolean throwOnCollect; + @Override public void addTrace(List> spans) { for (CoreSpan span : spans) { @@ -154,6 +223,11 @@ public void addTrace(List> spans) { @Override public OtlpPayload collectTraces() { + if (throwOnCollect) { + // mirrors the real collectors, which always reset state via a finally block + spansToExport.clear(); + throw new IllegalStateException("simulated buffer overflow"); + } if (spansToExport.isEmpty()) { return OtlpPayload.EMPTY; } @@ -165,5 +239,10 @@ public OtlpPayload collectTraces() { spansToExport.clear(); } } + + @Override + public int sizeInBytes() { + return fakeSizeInBytes; + } } } diff --git a/dd-trace-core/src/test/java/datadog/trace/core/otlp/common/OtlpProtoBufferTest.java b/dd-trace-core/src/test/java/datadog/trace/core/otlp/common/OtlpProtoBufferTest.java index 5cf7a820016..da75ef39f91 100644 --- a/dd-trace-core/src/test/java/datadog/trace/core/otlp/common/OtlpProtoBufferTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/core/otlp/common/OtlpProtoBufferTest.java @@ -60,6 +60,14 @@ void initialPayloadIsEmpty() { assertEquals("application/x-protobuf", payload.getContentType()); } + @Test + void constructorRejectsCapacityExceedingMaxCapacity() { + // rounds up to a power of two above MAX_CAPACITY_BYTES; must reject before allocating + assertThrows( + IllegalArgumentException.class, + () -> new OtlpProtoBuffer(OtlpProtoBuffer.MAX_CAPACITY_BYTES + 1)); + } + // ─── recordMessage(GrowableBuffer, int) ────────────────────────────────── @Test diff --git a/dd-trace-core/src/test/java/datadog/trace/core/otlp/trace/OtlpTraceJsonCollectorTest.java b/dd-trace-core/src/test/java/datadog/trace/core/otlp/trace/OtlpTraceJsonCollectorTest.java index 2f7db58c6fd..17ec0c0d17f 100644 --- a/dd-trace-core/src/test/java/datadog/trace/core/otlp/trace/OtlpTraceJsonCollectorTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/core/otlp/trace/OtlpTraceJsonCollectorTest.java @@ -10,7 +10,10 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; import datadog.json.JsonMapper; import datadog.trace.api.DDTraceId; @@ -206,6 +209,30 @@ void multipleSpansInATraceAreAllWritten() throws IOException { assertEquals(hexSpanId(((DDSpan) parent).getSpanId()), parsedChild.get("parentSpanId")); } + @Test + void poisonedSpanResetsCollectorForNextTrace() throws IOException { + // mid-trace exception (e.g. from a malformed span) must not leave partial state behind + DDSpan realSpan = startAndFinish("op.first", "GET /first", null); + + CoreSpan poison = mock(CoreSpan.class); + when(poison.samplingPriority()).thenReturn(1); + when(poison.getTraceId()).thenThrow(new RuntimeException("boom")); + + List> poisonedTrace = new ArrayList<>(); + poisonedTrace.add((CoreSpan) realSpan); + poisonedTrace.add(poison); + + OtlpTraceJsonCollector collector = new OtlpTraceJsonCollector(); + assertThrows(RuntimeException.class, () -> collector.addTrace(poisonedTrace)); + + // a normal trace collected afterwards must not see any leftover state from the poisoned one + DDSpan normalSpan = startAndFinish("op.normal", "GET /normal", null); + collector.addTrace(asList((CoreSpan) normalSpan)); + Map parsedSpan = onlySpan(collector.collectTraces()); + + assertEquals("GET /normal", parsedSpan.get("name")); + } + @Test void spanLinkOmitsTraceStateWhenEmpty() throws IOException { AgentSpan linked = TRACER.startSpan("test", "op.linked"); diff --git a/dd-trace-core/src/test/java/datadog/trace/core/otlp/trace/OtlpTraceProtoTest.java b/dd-trace-core/src/test/java/datadog/trace/core/otlp/trace/OtlpTraceProtoTest.java index 0acca35beb2..aa9d7c7022b 100644 --- a/dd-trace-core/src/test/java/datadog/trace/core/otlp/trace/OtlpTraceProtoTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/core/otlp/trace/OtlpTraceProtoTest.java @@ -16,7 +16,10 @@ import static org.junit.jupiter.api.Assertions.assertNotEquals; import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; import com.google.protobuf.CodedInputStream; import com.google.protobuf.WireFormat; @@ -29,6 +32,7 @@ import datadog.trace.bootstrap.instrumentation.api.SpanAttributes; import datadog.trace.bootstrap.instrumentation.api.SpanLink; import datadog.trace.common.writer.LoggingWriter; +import datadog.trace.core.CoreSpan; import datadog.trace.core.CoreTracer; import datadog.trace.core.DDSpan; import datadog.trace.core.otlp.common.OtlpPayload; @@ -628,6 +632,30 @@ void testCollectMultipleTraces() throws IOException { "payload must contain spans with all three distinct trace IDs"); } + @Test + void poisonedSpanResetsCollectorForNextTrace() { + // mid-trace exception (e.g. from a malformed span) must not leave partial state behind + DDSpan realSpan = buildSpans(asList(span("first.span", "op.first", "web"))).get(0); + + CoreSpan poison = mock(CoreSpan.class); + when(poison.samplingPriority()).thenReturn(1); + when(poison.getTraceId()).thenThrow(new RuntimeException("boom")); + + List> poisonedTrace = new ArrayList<>(); + poisonedTrace.add(poison); + poisonedTrace.add(realSpan); + + OtlpTraceProtoCollector collector = new OtlpTraceProtoCollector(); + assertThrows(RuntimeException.class, () -> collector.addTrace(poisonedTrace)); + + // a normal trace collected afterwards must not see any leftover state from the poisoned one + List normalTrace = buildSpans(asList(span("normal.op", "op.normal", "web"))); + collector.addTrace(normalTrace); + OtlpPayload payload = collector.collectTraces(); + + assertTrue(payload.getContentLength() > 0, "normal trace after reset must still export"); + } + @Test void testSpanOrderInTracePreserved() throws IOException { // Verifies that spans appear in the payload with the same order as the original trace