Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,129 @@
package datadog.trace.llmobs;

import datadog.context.Context;
import datadog.context.ContextScope;
import datadog.context.propagation.Propagators;
import datadog.trace.api.llmobs.LLMObs;
import datadog.trace.api.llmobs.LLMObsContext;
import datadog.trace.api.llmobs.LLMObsSpan;
import datadog.trace.api.llmobs.LLMObsTags;
import datadog.trace.bootstrap.instrumentation.api.AgentScope;
import datadog.trace.bootstrap.instrumentation.api.AgentSpan;
import datadog.trace.bootstrap.instrumentation.api.AgentSpanContext;
import datadog.trace.bootstrap.instrumentation.api.AgentTracer;
import datadog.trace.llmobs.domain.DDLLMObsSpan;
import java.io.Closeable;
import java.util.LinkedHashMap;
import java.util.Map;
import java.util.Objects;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

/**
* Explicit, manual distributed tracing propagation for LLM Observability, for boundaries that
* automatic instrumentation doesn't cover — e.g. an SQS worker reading its own message attributes.
*
* <p>The standard APM trace context (trace id, parent id, sampling, {@code x-datadog-tags}, ...) is
* injected/extracted via the normal {@link Propagators}. LLMObs-specific values (ml_app,
* session_id, agent attribution) ride as additional {@code _dd.p.llmobs_*} tags appended to the
* same {@code x-datadog-tags} carrier entry — the wire container dd-trace-py/js/go already use for
* these tags — so a mixed-language pipeline can still join a trace across this hop.
*/
public class DDLLMObsPropagator implements LLMObs.LLMObsPropagator {
private static final Logger LOGGER = LoggerFactory.getLogger(DDLLMObsPropagator.class);

// Package-private (not exposed by dd-trace-core) but a long-stable wire header name; also
// listed in datadog.trace.util.PropagationUtils#KNOWN_PROPAGATION_HEADERS.
private static final String DATADOG_TAGS_HEADER = "x-datadog-tags";

@Override
public Map<String, String> injectDistributedHeaders(
LLMObsSpan span, Map<String, String> headers) {
Objects.requireNonNull(span, "span");
Objects.requireNonNull(headers, "headers");
if (!(span instanceof DDLLMObsSpan)) {
LOGGER.debug(
"injectDistributedHeaders requires a span started by the LLM Observability SDK, got {}; ignoring",
span.getClass());
return headers;
}
DDLLMObsSpan llmObsSpan = (DDLLMObsSpan) span;
AgentSpan agentSpan = llmObsSpan.getAgentSpan();

Propagators.defaultPropagator().inject(agentSpan, headers, Map::put);

StringBuilder llmObsTags = new StringBuilder();
appendTag(llmObsTags, LLMObsTags.PROPAGATED_ML_APP, llmObsSpan.getMlApp());
appendTag(llmObsTags, LLMObsTags.PROPAGATED_SESSION_ID, llmObsSpan.getSessionId());
appendTag(llmObsTags, LLMObsTags.PROPAGATED_PAGENT_SPAN_ID, llmObsSpan.getParentAgentSpanId());
appendTag(llmObsTags, LLMObsTags.PROPAGATED_PAGENT_NAME, llmObsSpan.getParentAgentName());

if (llmObsTags.length() > 0) {
String existing = headers.get(DATADOG_TAGS_HEADER);
headers.put(
DATADOG_TAGS_HEADER,
existing == null || existing.isEmpty()
? llmObsTags.toString()
: existing + "," + llmObsTags);
}
return headers;
}

@Override
public Closeable activateDistributedHeaders(Map<String, String> headers) {
Objects.requireNonNull(headers, "headers");

Context extracted =
Propagators.defaultPropagator()
.extract(Context.root(), headers, (carrier, visitor) -> carrier.forEach(visitor));
AgentSpan extractedSpan = AgentSpan.fromContext(extracted);
if (extractedSpan == null) {
LOGGER.debug(
"no distributed trace context found in headers; activateDistributedHeaders is a no-op");
return () -> {};
}

Map<String, String> llmObsTags = parseLlmObsTags(headers.get(DATADOG_TAGS_HEADER));
String sessionId = llmObsTags.get(LLMObsTags.PROPAGATED_SESSION_ID);
String pagentSpanId = llmObsTags.get(LLMObsTags.PROPAGATED_PAGENT_SPAN_ID);
String pagentName = llmObsTags.get(LLMObsTags.PROPAGATED_PAGENT_NAME);

AgentScope apmScope = AgentTracer.get().activateSpan(extractedSpan);
AgentSpanContext extractedContext = extractedSpan.spanContext();
ContextScope llmObsScope =
LLMObsContext.attach(extractedContext, sessionId, null, pagentSpanId, pagentName);

return () -> {
llmObsScope.close();
apmScope.close();
};
}

private static void appendTag(StringBuilder sb, String key, String value) {
if (value == null || value.isEmpty()) {
return;
}
if (sb.length() > 0) {
sb.append(',');
}
sb.append(key).append('=').append(value);
}

private static Map<String, String> parseLlmObsTags(String xDatadogTags) {
Map<String, String> tags = new LinkedHashMap<>();
if (xDatadogTags == null || xDatadogTags.isEmpty()) {
return tags;
}
for (String pair : xDatadogTags.split(",")) {
int eq = pair.indexOf('=');
if (eq <= 0) {
continue;
}
String key = pair.substring(0, eq);
if (key.startsWith("_dd.p.llmobs_")) {
tags.put(key, pair.substring(eq + 1));
}
}
return tags;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,8 @@ public static void start(Instrumentation inst, SharedCommunicationObjects sco) {
LLMObsInternal.setEvalProcessor(new LLMObsCustomEvalProcessor(mlApp, sco, config));

LLMObsInternal.setFeedbackProcessor(new LLMObsCustomFeedbackProcessor(mlApp, sco, config));

LLMObsInternal.setPropagator(new DDLLMObsPropagator());
}

private static class LLMObsCustomFeedbackProcessor implements LLMObs.LLMObsFeedbackProcessor {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,9 @@ public class DDLLMObsSpan implements LLMObsSpan {
private final String spanKind;
private final String mlApp;
private final boolean hasSessionId;
private final String sessionId;
private final String parentAgentSpanId;
private final String parentAgentName;
private final ContextScope scope;
// Non-null only for agent-kind spans started without an ambient APM root. Activating the
// agent's APM span keeps children in the same APM trace so the trace-ID gate passes and
Expand Down Expand Up @@ -160,6 +163,7 @@ public DDLLMObsSpan(
}

this.hasSessionId = sessionId != null && !sessionId.isEmpty();
this.sessionId = this.hasSessionId ? sessionId : null;
if (this.hasSessionId) {
span.setTag(LLMOBS_TAG_PREFIX + LLMObsTags.SESSION_ID, sessionId);
}
Expand Down Expand Up @@ -197,6 +201,8 @@ public DDLLMObsSpan(
span.setTag(PAGENT_NAME_TAG_INTERNAL, resolvedParentAgentName);
}
}
this.parentAgentSpanId = resolvedParentAgentSpanId;
this.parentAgentName = resolvedParentAgentName;

// Propagate the effective sessionId and agent attribution to descendant LLMObs spans.
scope =
Expand Down Expand Up @@ -681,4 +687,35 @@ public DDTraceId getTraceId() {
public long getSpanId() {
return span.getSpanId();
}

/** Internal accessor for the underlying APM span, used by {@code DDLLMObsPropagator}. */
public AgentSpan getAgentSpan() {
return span;
}

/** Internal accessor for this span's effective ml_app, used by {@code DDLLMObsPropagator}. */
public String getMlApp() {
return mlApp;
}

/**
* Internal accessor for this span's effective session_id (including one inherited from an
* enclosing LLMObs span), used by {@code DDLLMObsPropagator}. May be {@code null}.
*/
public String getSessionId() {
return sessionId;
}

/**
* Internal accessor for this span's effective agent attribution, used by {@code
* DDLLMObsPropagator}. May be {@code null}.
*/
public String getParentAgentSpanId() {
return parentAgentSpanId;
}

/** See {@link #getParentAgentSpanId()}. May be {@code null}. */
public String getParentAgentName() {
return parentAgentName;
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,124 @@
package datadog.trace.llmobs;

import static org.junit.jupiter.api.Assertions.assertEquals;
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 datadog.trace.agent.tooling.TracerInstaller;
import datadog.trace.api.WellKnownTags;
import datadog.trace.api.llmobs.LLMObsContext;
import datadog.trace.bootstrap.instrumentation.api.AgentScope;
import datadog.trace.bootstrap.instrumentation.api.AgentSpan;
import datadog.trace.bootstrap.instrumentation.api.AgentTracer;
import datadog.trace.bootstrap.instrumentation.api.Tags;
import datadog.trace.core.CoreTracer;
import datadog.trace.llmobs.domain.DDLLMObsSpan;
import java.io.Closeable;
import java.util.HashMap;
import java.util.Map;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;

/**
* Round-trips {@link DDLLMObsPropagator} through a plain {@code Map<String, String>} carrier — the
* shape a customer's own SQS message-attribute map would take.
*/
class DDLLMObsPropagatorTest {

private static CoreTracer tracer;
private final DDLLMObsPropagator propagator = new DDLLMObsPropagator();

@BeforeAll
static void installTracer() {
tracer = CoreTracer.builder().build();
TracerInstaller.forceInstallGlobalTracer(tracer);
}

@AfterAll
static void closeTracer() {
TracerInstaller.forceInstallGlobalTracer(null);
tracer.close();
}

private static DDLLMObsSpan newAgentSpan(String name, String mlApp, String sessionId) {
return newSpan(Tags.LLMOBS_AGENT_SPAN_KIND, name, mlApp, sessionId);
}

private static DDLLMObsSpan newToolSpan(String name, String mlApp, String sessionId) {
return newSpan(Tags.LLMOBS_TOOL_SPAN_KIND, name, mlApp, sessionId);
}

private static DDLLMObsSpan newSpan(String kind, String name, String mlApp, String sessionId) {
WellKnownTags tags =
new WellKnownTags("runtime-id", "hostname", "test", "service", "version", "java");
return new DDLLMObsSpan(kind, name, mlApp, sessionId, "service", tags);
}

private static AgentScope startRootApmScope() {
AgentSpan root = AgentTracer.get().buildSpan("apm", "sqs.produce").start();
return AgentTracer.activateSpan(root);
}

@Test
void injectRequiresNonNullSpanAndHeaders() {
assertThrows(
NullPointerException.class,
() -> propagator.injectDistributedHeaders(null, new HashMap<>()));
}

@Test
void activateRequiresNonNullHeaders() {
assertThrows(NullPointerException.class, () -> propagator.activateDistributedHeaders(null));
}

@Test
void activateWithoutTraceContextIsNoOp() throws Exception {
try (Closeable scope = propagator.activateDistributedHeaders(new HashMap<>())) {
assertNull(LLMObsContext.current());
}
}

@Test
void injectThenActivateJoinsSameTraceAndPropagatesLlmObsContext() throws Exception {
Map<String, String> headers = new HashMap<>();

try (AgentScope apmScope = startRootApmScope()) {
DDLLMObsSpan producerAgent = newAgentSpan("producer-agent", "my-ml-app", "session-123");
long producerTraceId = producerAgent.getTraceId().toLong();
try {
propagator.injectDistributedHeaders(producerAgent, headers);
} finally {
producerAgent.finish();
}

// Simulate the consumer side: a different message handled with no ambient context.
try (Closeable consumerScope = propagator.activateDistributedHeaders(headers)) {
DDLLMObsSpan consumerTool = newToolSpan("consumer-tool", "my-ml-app", null);
try {
assertEquals(producerTraceId, consumerTool.getTraceId().toLong());
assertEquals("session-123", consumerTool.getSessionId());
} finally {
consumerTool.finish();
}
}
}
}

@Test
void injectAlwaysIncludesMlAppEvenWithoutSessionIdOrAttribution() {
Map<String, String> headers = new HashMap<>();
try (AgentScope apmScope = startRootApmScope()) {
DDLLMObsSpan standaloneTool = newToolSpan("standalone-tool", "my-ml-app", null);
try {
propagator.injectDistributedHeaders(standaloneTool, headers);
String xDatadogTags = headers.get("x-datadog-tags");
// ml_app is always present since it's required on every LLMObs span.
assertTrue(xDatadogTags != null && xDatadogTags.contains("_dd.p.llmobs_ml_app=my-ml-app"));
} finally {
standaloneTool.finish();
}
}
}
}
Loading
Loading