Skip to content

Commit 073ad18

Browse files
authored
Tag consumer spans with the pathway hash on every data streams checkpoint (#11808)
Tag consumer spans with the pathway hash on every data streams checkpoint The pathway.hash span tag was only set on the produce/inject path (DataStreamsPropagator.inject). Consumers that checkpoint without injecting (e.g. RabbitMQ) had no pathway.hash, unlike the JS and Python tracers which tag the span on every checkpoint. Set it centrally in DefaultDataStreamsMonitoring.setCheckpoint so all consume-side integrations get it. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> Remove explanatory comment on pathway hash tagging Per review on #11808; the rationale lives in the PR description. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> Address review on #11808 - cache pathwayContext.getHash() in a local instead of calling it twice - use a negative hash in the test so it actually exercises Long.toUnsignedString Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> Expect pathway.hash on the inferred-proxy server span when DSM is enabled Now that DefaultDataStreamsMonitoring.setCheckpoint tags the span on every checkpoint, the inferred-proxy (API Gateway) HTTP server span gets a pathway.hash like any other DSM-enabled server span. The shared HttpServerTest already asserts this conditionally; SpringBootBasedTest's hand-rolled inferred-proxy assertions were missing it. Mirror the same guard. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> Merge branch 'master' into eric.firth/dsm-consumer-pathway-hash Co-authored-by: eric.firth <eric.firth@datadoghq.com>
1 parent 41fffe8 commit 073ad18

3 files changed

Lines changed: 76 additions & 0 deletions

File tree

dd-java-agent/instrumentation/spring/spring-webmvc/spring-webmvc-3.1/src/test/groovy/test/boot/SpringBootBasedTest.groovy

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -570,6 +570,9 @@ class SpringBootBasedTest extends HttpServerTest<ConfigurableApplicationContext>
570570
"$Tags.HTTP_ROUTE" "/success"
571571
"servlet.context" "/boot-context"
572572
"servlet.path" "/success"
573+
if ({ isDataStreamsEnabled() }) {
574+
"$DDTags.PATHWAY_HASH" { String }
575+
}
573576
defaultTags()
574577
}
575578
}

dd-trace-core/src/main/java/datadog/trace/core/datastreams/DefaultDataStreamsMonitoring.java

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
package datadog.trace.core.datastreams;
22

33
import static datadog.communication.ddagent.DDAgentFeaturesDiscovery.V01_DATASTREAMS_ENDPOINT;
4+
import static datadog.trace.api.DDTags.PATHWAY_HASH;
45
import static datadog.trace.api.datastreams.DataStreamsContext.fromTags;
56
import static datadog.trace.api.datastreams.DataStreamsTags.Direction.INBOUND;
67
import static datadog.trace.api.datastreams.DataStreamsTags.Direction.OUTBOUND;
@@ -309,6 +310,10 @@ public void setCheckpoint(AgentSpan span, DataStreamsContext context) {
309310
PathwayContext pathwayContext = span.spanContext().getPathwayContext();
310311
if (pathwayContext != null) {
311312
pathwayContext.setCheckpoint(context, this::add);
313+
long pathwayHash = pathwayContext.getHash();
314+
if (pathwayHash != 0) {
315+
span.setTag(PATHWAY_HASH, Long.toUnsignedString(pathwayHash));
316+
}
312317
}
313318
}
314319

dd-trace-core/src/test/java/datadog/trace/core/datastreams/DefaultDataStreamsMonitoringTest.java

Lines changed: 68 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,22 +10,29 @@
1010
import static org.junit.jupiter.api.Assertions.assertNotNull;
1111
import static org.junit.jupiter.api.Assertions.assertTrue;
1212
import static org.mockito.ArgumentMatchers.any;
13+
import static org.mockito.ArgumentMatchers.eq;
1314
import static org.mockito.Mockito.RETURNS_SMART_NULLS;
1415
import static org.mockito.Mockito.mock;
16+
import static org.mockito.Mockito.never;
1517
import static org.mockito.Mockito.verify;
1618
import static org.mockito.Mockito.when;
1719

1820
import datadog.communication.ddagent.DDAgentFeaturesDiscovery;
1921
import datadog.metrics.api.Histograms;
2022
import datadog.metrics.impl.DDSketchHistograms;
2123
import datadog.trace.api.Config;
24+
import datadog.trace.api.DDTags;
2225
import datadog.trace.api.TraceConfig;
26+
import datadog.trace.api.datastreams.DataStreamsContext;
2327
import datadog.trace.api.datastreams.DataStreamsTags;
2428
import datadog.trace.api.datastreams.KafkaConfigReport;
29+
import datadog.trace.api.datastreams.PathwayContext;
2530
import datadog.trace.api.datastreams.SchemaRegistryUsage;
2631
import datadog.trace.api.datastreams.StatsPoint;
2732
import datadog.trace.api.experimental.DataStreamsContextCarrier;
2833
import datadog.trace.api.time.ControllableTimeSource;
34+
import datadog.trace.bootstrap.instrumentation.api.AgentSpan;
35+
import datadog.trace.bootstrap.instrumentation.api.AgentSpanContext;
2936
import datadog.trace.common.metrics.EventListener;
3037
import datadog.trace.common.metrics.Sink;
3138
import datadog.trace.core.DDCoreJavaSpecification;
@@ -1559,6 +1566,67 @@ void close() {
15591566
}
15601567
}
15611568

1569+
private DefaultDataStreamsMonitoring newDataStreamsMonitoring() {
1570+
return new DefaultDataStreamsMonitoring(
1571+
mock(Sink.class),
1572+
stubFeatures(true),
1573+
new ControllableTimeSource(),
1574+
() -> stubTraceConfig(true),
1575+
mock(DatastreamsPayloadWriter.class),
1576+
DEFAULT_BUCKET_DURATION_NANOS);
1577+
}
1578+
1579+
@Test
1580+
void setCheckpointTagsTheSpanWithThePathwayHash() {
1581+
DefaultDataStreamsMonitoring dataStreams = newDataStreamsMonitoring();
1582+
1583+
// Negative so the signed and unsigned string representations differ, proving the tag uses
1584+
// Long.toUnsignedString (pathway hashes are unsigned 64-bit values).
1585+
long hash = -1234567890123456789L;
1586+
PathwayContext pathwayContext = mock(PathwayContext.class);
1587+
when(pathwayContext.getHash()).thenReturn(hash);
1588+
AgentSpanContext spanContext = mock(AgentSpanContext.class);
1589+
when(spanContext.getPathwayContext()).thenReturn(pathwayContext);
1590+
AgentSpan span = mock(AgentSpan.class);
1591+
when(span.spanContext()).thenReturn(spanContext);
1592+
DataStreamsContext context = mock(DataStreamsContext.class);
1593+
1594+
dataStreams.setCheckpoint(span, context);
1595+
1596+
verify(pathwayContext).setCheckpoint(eq(context), any());
1597+
verify(span).setTag(DDTags.PATHWAY_HASH, Long.toUnsignedString(hash));
1598+
}
1599+
1600+
@Test
1601+
void setCheckpointDoesNotTagTheSpanWhenThePathwayHashIsZero() {
1602+
DefaultDataStreamsMonitoring dataStreams = newDataStreamsMonitoring();
1603+
1604+
PathwayContext pathwayContext = mock(PathwayContext.class);
1605+
when(pathwayContext.getHash()).thenReturn(0L);
1606+
AgentSpanContext spanContext = mock(AgentSpanContext.class);
1607+
when(spanContext.getPathwayContext()).thenReturn(pathwayContext);
1608+
AgentSpan span = mock(AgentSpan.class);
1609+
when(span.spanContext()).thenReturn(spanContext);
1610+
1611+
dataStreams.setCheckpoint(span, mock(DataStreamsContext.class));
1612+
1613+
verify(span, never()).setTag(eq(DDTags.PATHWAY_HASH), any(String.class));
1614+
}
1615+
1616+
@Test
1617+
void setCheckpointDoesNotTagTheSpanWhenThereIsNoPathwayContext() {
1618+
DefaultDataStreamsMonitoring dataStreams = newDataStreamsMonitoring();
1619+
1620+
AgentSpanContext spanContext = mock(AgentSpanContext.class);
1621+
when(spanContext.getPathwayContext()).thenReturn(null);
1622+
AgentSpan span = mock(AgentSpan.class);
1623+
when(span.spanContext()).thenReturn(spanContext);
1624+
1625+
dataStreams.setCheckpoint(span, mock(DataStreamsContext.class));
1626+
1627+
verify(span, never()).setTag(eq(DDTags.PATHWAY_HASH), any(String.class));
1628+
}
1629+
15621630
static class CustomContextCarrier implements DataStreamsContextCarrier {
15631631
private final Map<String, Object> data = new HashMap<>();
15641632

0 commit comments

Comments
 (0)