Skip to content

Commit 1331243

Browse files
asimurkacursoragent
andcommitted
LCORE-3755: emit eval turn attrs on llm.inference vs root
Split turn-summary OTEL fields on top of the unified span hierarchy: JSON tool/RAG payloads, model/provider, and tokens go on llm.inference; input/output/session/compacted stay on the root span. Co-authored-by: Cursor <cursoragent@cursor.com>
1 parent fa4bb69 commit 1331243

20 files changed

Lines changed: 443 additions & 210 deletions

‎src/app/endpoints/a2a.py‎

Lines changed: 18 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -158,11 +158,11 @@ def _record_model_span(span: trace.Span, model_id: str) -> None:
158158
span: The active OpenTelemetry span.
159159
model_id: Full model identifier in "provider/model" format.
160160
"""
161-
provider_id, _ = extract_provider_and_model_from_model_id(model_id)
161+
provider_id, bare_model_id = extract_provider_and_model_from_model_id(model_id)
162162
set_span_attributes(
163163
span,
164164
{
165-
SpanAttributes.LLM_MODEL_ID: model_id,
165+
SpanAttributes.LLM_MODEL_ID: bare_model_id,
166166
SpanAttributes.LLM_PROVIDER_ID: provider_id,
167167
},
168168
)
@@ -172,13 +172,17 @@ def _record_execution_span(
172172
span: trace.Span,
173173
tool_call_names: list[str],
174174
run_result: Optional[AgentRunResult[str]],
175+
compacted: bool,
176+
inference_time: float,
175177
) -> None:
176178
"""Record tool-call metrics, token usage, and output on an a2a.execute span.
177179
178180
Parameters:
179181
span: The active OpenTelemetry span.
180182
tool_call_names: Tool names collected during streaming.
181183
run_result: Completed agent run result, or None.
184+
compacted: Whether the turn used compacted conversation context.
185+
inference_time: Request processing duration in seconds.
182186
"""
183187
if tool_call_names:
184188
set_span_attributes(
@@ -205,6 +209,9 @@ def _record_execution_span(
205209
if output_text:
206210
span.set_attribute(SpanAttributes.OUTPUT, output_text)
207211

212+
span.set_attribute(SpanAttributes.COMPACTED, compacted)
213+
span.set_attribute(SpanAttributes.INFERENCE_TIME, inference_time)
214+
208215

209216
async def _persist_compacted_a2a_turn(
210217
client: Any,
@@ -385,7 +392,7 @@ async def execute(
385392
"Failed to publish failure event: %s", enqueue_error, exc_info=True
386393
)
387394

388-
async def _process_task_streaming( # pylint: disable=too-many-locals
395+
async def _process_task_streaming( # pylint: disable=too-many-locals,too-many-statements
389396
self,
390397
context: RequestContext,
391398
task_updater: TaskUpdater,
@@ -405,6 +412,7 @@ async def _process_task_streaming( # pylint: disable=too-many-locals
405412

406413
with tracer.start_as_current_span("a2a.execute") as span:
407414
span.set_attribute(SpanAttributes.SESSION_ID, context_id)
415+
started_at = datetime.now(UTC)
408416

409417
# Extract user input using SDK utility
410418
user_input = context.get_user_input()
@@ -578,7 +586,13 @@ async def _process_task_streaming( # pylint: disable=too-many-locals
578586
client, responses_params, compaction, agent, task_id
579587
)
580588

581-
_record_execution_span(span, self._tool_call_names, self._run_result)
589+
_record_execution_span(
590+
span,
591+
self._tool_call_names,
592+
self._run_result,
593+
compacted=compaction.compacted,
594+
inference_time=(datetime.now(UTC) - started_at).total_seconds(),
595+
)
582596

583597
# Publish the final task result event
584598
if aggregator.task_state == TaskState.working:

‎src/app/endpoints/info.py‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,7 @@
2020
)
2121
from models.api.responses.successful import InfoResponse
2222
from models.config import Action
23-
from utils.otel_tracing import set_span_attributes
23+
from utils.otel_tracing import SpanAttributes, set_span_attributes
2424
from utils.types import Responses
2525
from version import __version__
2626

@@ -85,8 +85,8 @@ async def info_endpoint_handler(
8585
set_span_attributes(
8686
span,
8787
{
88-
"service.name": configuration.configuration.name,
89-
"service.version": __version__,
88+
SpanAttributes.SERVICE_NAME: configuration.configuration.name,
89+
SpanAttributes.SERVICE_VERSION: __version__,
9090
},
9191
)
9292
return InfoResponse(

‎src/app/endpoints/query.py‎

Lines changed: 22 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -44,6 +44,7 @@
4444
SpanEvents,
4545
add_span_event,
4646
anonymize_value,
47+
root_span_turn_attributes,
4748
set_span_attributes,
4849
)
4950
from utils.query import (
@@ -153,7 +154,8 @@ async def _handle_query_with_tracing(
153154
"""
154155
check_configuration_loaded(configuration)
155156

156-
started_at = datetime.datetime.now(datetime.UTC).strftime("%Y-%m-%dT%H:%M:%SZ")
157+
started_at_dt = datetime.datetime.now(datetime.UTC)
158+
started_at = started_at_dt.strftime("%Y-%m-%dT%H:%M:%SZ")
157159
user_id, _, _skip_userid_check, token = auth
158160

159161
# Set initial span attributes
@@ -275,20 +277,21 @@ async def _handle_query_with_tracing(
275277
shield_ids=query_request.shield_ids,
276278
no_tools=bool(query_request.no_tools),
277279
image_attachments=image_attachments,
280+
extra_rag_chunks=(
281+
inline_rag_context.rag_chunks
282+
if moderation_result.decision == "passed"
283+
else None
284+
),
278285
)
279286

280287
if moderation_result.decision == "passed":
281-
# Combine inline RAG results (BYOK + Solr) with tool-based RAG results for the transcript
282-
rag_chunks = inline_rag_context.rag_chunks
283-
tool_rag_chunks = turn_summary.rag_chunks
284-
logger.info("RAG as a tool retrieved %d chunks", len(tool_rag_chunks))
285-
turn_summary.rag_chunks = rag_chunks + tool_rag_chunks
286-
287-
# Add tool-based RAG documents and chunks
288-
rag_documents = inline_rag_context.referenced_documents
289-
tool_rag_documents = turn_summary.referenced_documents
288+
# Combine inline RAG documents with tool-based RAG documents for the transcript
289+
logger.info(
290+
"RAG as a tool retrieved %d chunks",
291+
len(turn_summary.rag_chunks) - len(inline_rag_context.rag_chunks),
292+
)
290293
turn_summary.referenced_documents = deduplicate_referenced_documents(
291-
rag_documents + tool_rag_documents
294+
inline_rag_context.referenced_documents + turn_summary.referenced_documents
292295
)
293296

294297
# Get topic summary for new conversation
@@ -314,7 +317,8 @@ async def _handle_query_with_tracing(
314317
quota_limiters=configuration.quota_limiters, user_id=user_id
315318
)
316319

317-
completed_at = datetime.datetime.now(datetime.UTC).strftime("%Y-%m-%dT%H:%M:%SZ")
320+
completed_at_dt = datetime.datetime.now(datetime.UTC)
321+
completed_at = completed_at_dt.strftime("%Y-%m-%dT%H:%M:%SZ")
318322
conversation_id = normalize_conversation_id(responses_params.conversation)
319323

320324
logger.info("Storing query results")
@@ -335,13 +339,14 @@ async def _handle_query_with_tracing(
335339

336340
logger.info("Building final response")
337341

338-
# Set final span attributes
342+
# Set final root-span attributes (llm.* attrs live on the llm.inference span)
339343
set_span_attributes(
340344
root_span,
341-
{
342-
SpanAttributes.SESSION_ID: conversation_id,
343-
SpanAttributes.OUTPUT: turn_summary.llm_response,
344-
},
345+
root_span_turn_attributes(
346+
turn_summary,
347+
session_id=conversation_id,
348+
compacted=compaction.compacted,
349+
),
345350
)
346351

347352
# Emit LLM response completed event

‎src/app/endpoints/responses.py‎

Lines changed: 44 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -79,7 +79,9 @@
7979
SpanEvents,
8080
add_span_event,
8181
anonymize_value,
82+
llm_inference_span_attributes,
8283
record_exception,
84+
root_span_turn_attributes,
8385
set_span_attributes,
8486
)
8587
from utils.prompts import get_system_prompt
@@ -158,18 +160,18 @@ def _count_request_attachments(response_input: ResponseInput) -> int:
158160
def _finalize_responses_root_span(
159161
root_span: trace.Span,
160162
turn_summary: TurnSummary,
163+
compacted: bool = False,
161164
) -> None:
162165
"""Set final root-span attributes and completion events for /responses.
163166
164167
Args:
165168
root_span: OpenTelemetry root span for the request.
166-
turn_summary: Completed turn summary with output text.
169+
turn_summary: Completed turn summary with the LLM response text.
170+
compacted: Whether the turn used compacted conversation context.
167171
"""
168172
set_span_attributes(
169173
root_span,
170-
{
171-
SpanAttributes.OUTPUT: turn_summary.llm_response,
172-
},
174+
root_span_turn_attributes(turn_summary, compacted=compacted),
173175
)
174176
add_span_event(root_span, SpanEvents.LLM_RESPONSE_COMPLETED)
175177

@@ -206,33 +208,29 @@ def _start_llm_inference_span(
206208
def _complete_llm_inference_span(
207209
span: trace.Span,
208210
turn_summary: TurnSummary,
211+
model: str,
212+
inference_time: float,
209213
) -> None:
210-
"""Record usage/tool attrs and completion event, then end an inference span.
214+
"""Record turn-summary attrs and completion event, then end an inference span.
211215
212216
Args:
213217
span: The ``llm.inference`` span to finalize.
214-
turn_summary: Completed turn summary providing token usage and tool calls.
218+
turn_summary: Completed turn summary with tools, RAG, and tokens.
219+
model: Composite model identifier in ``provider/model`` format.
220+
inference_time: Inference duration in seconds.
215221
"""
222+
provider_id, bare_model_id = extract_provider_and_model_from_model_id(model)
216223
set_span_attributes(
217224
span,
218-
{
219-
SpanAttributes.LLM_USAGE_INPUT_TOKENS: (
220-
turn_summary.token_usage.input_tokens
221-
),
222-
SpanAttributes.LLM_USAGE_OUTPUT_TOKENS: (
223-
turn_summary.token_usage.output_tokens
224-
),
225-
},
225+
llm_inference_span_attributes(
226+
turn_summary,
227+
bare_model_id,
228+
provider_id,
229+
inference_time,
230+
),
226231
)
227-
tool_names = [tc.name for tc in turn_summary.tool_calls]
228-
if tool_names:
229-
set_span_attributes(
230-
span,
231-
{
232-
SpanAttributes.TOOL_CALLS_COUNT: len(tool_names),
233-
SpanAttributes.TOOL_CALLS_NAMES: tool_names,
234-
},
235-
)
232+
if turn_summary.tool_calls:
233+
tool_names = [tc.name for tc in turn_summary.tool_calls]
236234
add_span_event(
237235
span,
238236
SpanEvents.TOOL_EXECUTION_COMPLETED,
@@ -570,7 +568,7 @@ async def handle_responses_with_tracing( # pylint: disable=too-many-locals
570568
)
571569
attachments_count = _count_request_attachments(original_request.input)
572570

573-
span_attributes: dict[str, Any] = {
571+
span_attributes: dict[SpanAttributes, Any] = {
574572
SpanAttributes.USER_ID: anonymize_value(user_id),
575573
SpanAttributes.INPUT: input_text,
576574
SpanAttributes.REQUEST_ATTACHMENTS_COUNT: attachments_count,
@@ -1231,7 +1229,8 @@ async def response_generator(
12311229
)
12321230
raise
12331231

1234-
# Populate tools before closing llm.inference so tool attrs land on that span.
1232+
# Extract response metadata from final response object before closing
1233+
# the inference span so tool/RAG attrs can be recorded on it.
12351234
if latest_response_object:
12361235
_populate_turn_summary(
12371236
latest_response_object,
@@ -1243,6 +1242,8 @@ async def response_generator(
12431242
_complete_llm_inference_span(
12441243
inference_span,
12451244
turn_summary,
1245+
api_params.model,
1246+
time.monotonic() - inference_start_time,
12461247
)
12471248

12481249
# Explicitly append the turn to conversation if context passed by previous response
@@ -1302,7 +1303,11 @@ async def generate_response(
13021303
completed_at,
13031304
turn_summary.llm_response,
13041305
)
1305-
_finalize_responses_root_span(root_span, turn_summary)
1306+
_finalize_responses_root_span(
1307+
root_span,
1308+
turn_summary,
1309+
context.compacted_original_input is not None,
1310+
)
13061311
# Persist conversation state before clients can close the stream.
13071312
yield "data: [DONE]\n\n"
13081313
finally:
@@ -1325,9 +1330,10 @@ async def handle_non_streaming_response(
13251330
"""
13261331
root_span = context.root_span
13271332
user_id = context.auth[0]
1333+
inference_span: Optional[trace.Span] = None
1334+
inference_start_time: Optional[float] = None
13281335

13291336
# Fork: Get response object (blocked vs normal)
1330-
inference_span: Optional[trace.Span] = None
13311337
if context.moderation_result.decision == "blocked":
13321338
output_text = context.moderation_result.message
13331339
api_response = OpenAIResponseObject.model_construct(
@@ -1367,6 +1373,8 @@ async def handle_non_streaming_response(
13671373
token_usage = extract_token_usage(
13681374
api_response.usage, api_params.model, context.endpoint_path
13691375
)
1376+
# Keep inference span open until turn_summary is built below so
1377+
# tool/RAG attributes can be recorded on llm.inference.
13701378
logger.info("Consuming tokens")
13711379
consume_query_tokens(
13721380
user_id=user_id,
@@ -1411,11 +1419,13 @@ async def handle_non_streaming_response(
14111419
)
14121420
turn_summary.rag_chunks.extend(context.inline_rag_context.rag_chunks)
14131421

1414-
# Close llm.inference after tools are known so usage + tool attrs share one span.
1415-
if inference_span is not None:
1422+
# Close llm.inference after tools/RAG are known so usage + eval attrs share one span.
1423+
if inference_span is not None and inference_start_time is not None:
14161424
_complete_llm_inference_span(
14171425
inference_span,
14181426
turn_summary,
1427+
api_params.model,
1428+
time.monotonic() - inference_start_time,
14191429
)
14201430

14211431
# Get available quotas
@@ -1446,7 +1456,11 @@ async def handle_non_streaming_response(
14461456
completed_at,
14471457
output_text,
14481458
)
1449-
_finalize_responses_root_span(root_span, turn_summary)
1459+
_finalize_responses_root_span(
1460+
root_span,
1461+
turn_summary,
1462+
context.compacted_original_input is not None,
1463+
)
14501464
configured_mcp_labels = {s.name for s in configuration.mcp_servers}
14511465
response_dict = (
14521466
api_response.model_dump(exclude_none=True)

0 commit comments

Comments
 (0)