diff --git a/README.md b/README.md index 21cd531..3f33592 100644 --- a/README.md +++ b/README.md @@ -43,6 +43,20 @@ pnpm dev The server will be available at `http://localhost:3000`. +### OpenTelemetry (GenAI) tracing + +`responses.js` emits OpenTelemetry spans for response execution (`gen_ai.operation.name=chat`) and tool execution (`gen_ai.operation.name=execute_tool`). + +- Parent trace context is extracted from incoming propagation headers (`traceparent`/`tracestate`), so spans attach to your upstream instrumentation. +- Tool metadata is always traced (`gen_ai.tool.name`, `gen_ai.tool.type`, `gen_ai.tool.call.id`, `mcp.server_label`). +- Tool arguments/results are optional and controlled by: + +```bash +OTEL_GENAI_CAPTURE_TOOL_CONTENT=true +``` + +Set it to `false` (or unset) to avoid collecting tool arguments/results in span attributes. + ### Running Examples Explore the various capabilities with our example scripts located in the [./examples](./examples) folder: diff --git a/package.json b/package.json index 661aac3..8e75569 100644 --- a/package.json +++ b/package.json @@ -68,6 +68,7 @@ "author": "Hugging Face", "license": "MIT", "dependencies": { + "@opentelemetry/api": "^1.9.0", "@modelcontextprotocol/sdk": "^1.26.0", "express": "^4.22.1", "openai": "^5.8.2", diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 8dd7946..4da6cb1 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -19,6 +19,9 @@ importers: '@modelcontextprotocol/sdk': specifier: ^1.26.0 version: 1.26.0(zod@3.25.76) + '@opentelemetry/api': + specifier: ^1.9.0 + version: 1.9.0 express: specifier: ^4.22.1 version: 4.22.1 @@ -311,6 +314,10 @@ packages: '@cfworker/json-schema': optional: true + '@opentelemetry/api@1.9.0': + resolution: {integrity: sha512-3giAOQvZiH5F9bMlMiv8+GSPMeqg0dbaeo58/0SlA9sxSqZhnUtxzX9/2FzyhS9sWQf5S0GJE0AKBrFqjpeYcg==} + engines: {node: '>=8.0.0'} + '@pkgjs/parseargs@0.11.0': resolution: {integrity: sha512-+1VkjdD0QBLPodGrJUeqarH8VAIvQODIbwh9XpP5Syisf7YoQgsJKPNFoqqLQlu+VQ/tVSshMR6loPMn8U+dPg==} engines: {node: '>=14'} @@ -1844,6 +1851,8 @@ snapshots: transitivePeerDependencies: - supports-color + '@opentelemetry/api@1.9.0': {} + '@pkgjs/parseargs@0.11.0': optional: true diff --git a/src/routes/responses.ts b/src/routes/responses.ts index 3adfd59..9326730 100644 --- a/src/routes/responses.ts +++ b/src/routes/responses.ts @@ -3,6 +3,7 @@ import { type ValidatedRequest } from "../middleware/validation.js"; import type { CreateResponseParams, McpServerParams, McpApprovalRequestParams } from "../schemas.js"; import { generateUniqueId } from "../lib/generateUniqueId.js"; import { OpenAI } from "openai"; +import { context, propagation, SpanStatusCode, trace, type Attributes, type Context, type Span } from "@opentelemetry/api"; import type { Response, ResponseContentPartAddedEvent, @@ -34,6 +35,11 @@ class StreamingError extends Error { type IncompleteResponse = Omit; const SEQUENCE_NUMBER_PLACEHOLDER = -1; +const tracer = trace.getTracer("responses.js.routes.responses"); + +const OTEL_GENAI_CAPTURE_TOOL_CONTENT = + process.env.OTEL_GENAI_CAPTURE_TOOL_CONTENT === "1" || + process.env.OTEL_GENAI_CAPTURE_TOOL_CONTENT?.toLowerCase() === "true"; // All headers are forwarded by default, except these ones. const NOT_FORWARDED_HEADERS = new Set([ @@ -52,6 +58,38 @@ const NOT_FORWARDED_HEADERS = new Set([ "upgrade", ]); +const buildJsonAttribute = (value: unknown): string => { + if (typeof value === "string") { + return value; + } + try { + return JSON.stringify(value); + } catch { + return String(value); + } +}; + +const getRequestTraceContext = (req: ValidatedRequest): Context => { + const carrier: Record = {}; + for (const [key, value] of Object.entries(req.headers)) { + if (typeof value === "string") { + carrier[key] = value; + } else if (Array.isArray(value)) { + carrier[key] = value.join(","); + } + } + + return propagation.extract(context.active(), carrier); +}; + +const recordError = (span: Span, error: unknown): void => { + span.recordException(error instanceof Error ? error : new Error(String(error))); + span.setStatus({ + code: SpanStatusCode.ERROR, + message: error instanceof Error ? error.message : String(error), + }); +}; + export const postCreateResponse = async ( req: ValidatedRequest, res: ExpressResponse @@ -90,6 +128,23 @@ async function* runCreateResponseStream( req: ValidatedRequest, res: ExpressResponse ): AsyncGenerator { + const requestContext = getRequestTraceContext(req); + const requestSpan = tracer.startSpan( + "responses.create", + { + attributes: { + "gen_ai.operation.name": "chat", + "gen_ai.request.model": req.body.model, + "gen_ai.request.max_tokens": req.body.max_output_tokens ?? undefined, + "gen_ai.request.temperature": req.body.temperature ?? undefined, + "gen_ai.request.top_p": req.body.top_p ?? undefined, + "gen_ai.response.id": undefined, + }, + }, + requestContext + ); + const traceContext = trace.setSpan(requestContext, requestSpan); + let sequenceNumber = 0; // Prepare response object that will be iteratively populated const responseObject: IncompleteResponse = { @@ -117,6 +172,13 @@ async function* runCreateResponseStream( total_tokens: 0, }, }; + requestSpan.setAttribute("gen_ai.response.id", responseObject.id); + // if (req.body.instructions) { + // requestSpan.setAttribute( + // "gen_ai.system_instructions", + // buildJsonAttribute([{ type: "text", content: req.body.instructions }]) + // ); + // } // Response created event yield { @@ -132,49 +194,65 @@ async function* runCreateResponseStream( sequence_number: sequenceNumber++, }; - // Any events (LLM call, MCP call, list tools, etc.) try { - for await (const event of innerRunStream(req, res, responseObject)) { - yield { ...event, sequence_number: sequenceNumber++ }; + // Any events (LLM call, MCP call, list tools, etc.) + try { + for await (const event of innerRunStream(req, res, responseObject, traceContext)) { + yield { ...event, sequence_number: sequenceNumber++ }; + } + } catch (error) { + // Error event => stop + console.error("Error in stream:", error); + + const message = + typeof error === "object" && + error && + "message" in error && + typeof (error as { message: unknown }).message === "string" + ? (error as { message: string }).message + : "An error occurred in stream"; + + responseObject.status = "failed"; + responseObject.error = { + code: "server_error", + message, + }; + recordError(requestSpan, error); + yield { + type: "response.failed", + response: responseObject as Response, + sequence_number: sequenceNumber++, + }; + return; } - } catch (error) { - // Error event => stop - console.error("Error in stream:", error); - - const message = - typeof error === "object" && - error && - "message" in error && - typeof (error as { message: unknown }).message === "string" - ? (error as { message: string }).message - : "An error occurred in stream"; - - responseObject.status = "failed"; - responseObject.error = { - code: "server_error", - message, - }; + + // Response completed event + responseObject.status = "completed"; + if (responseObject.usage) { + requestSpan.setAttributes({ + "gen_ai.usage.input_tokens": responseObject.usage.input_tokens, + "gen_ai.usage.output_tokens": responseObject.usage.output_tokens, + }); + } + requestSpan.setAttributes({ + "gen_ai.response.model": responseObject.model, + "response.status": responseObject.status, + }); yield { - type: "response.failed", + type: "response.completed", response: responseObject as Response, sequence_number: sequenceNumber++, }; - return; + } finally { + requestSpan.end(); } - - // Response completed event - responseObject.status = "completed"; - yield { - type: "response.completed", - response: responseObject as Response, - sequence_number: sequenceNumber++, - }; } async function* innerRunStream( req: ValidatedRequest, res: ExpressResponse, - responseObject: IncompleteResponse + responseObject: IncompleteResponse, + traceContext: Context ): AsyncGenerator { // Retrieve API key from headers const apiKey = req.headers.authorization?.split(" ")[1]; @@ -196,6 +274,43 @@ async function* innerRunStream( throw new Error(`Not implemented: only 'auto' summary is supported. Got '${req.body.reasoning?.summary}'`); } + // Trace function tool calls provided by the client in input history + if (Array.isArray(req.body.input)) { + for (const item of req.body.input) { + if (item.type !== "function_call") { + continue; + } + + const matchingOutput = req.body.input.find( + inputItem => inputItem.type === "function_call_output" && inputItem.call_id === item.call_id + ) as Extract[number], { type: "function_call_output" }> | undefined; + + const functionCallSpanAttributes: Attributes = { + "gen_ai.operation.name": "execute_tool", + "gen_ai.tool.type": "function", + "gen_ai.tool.call.id": item.call_id, + "gen_ai.tool.name": item.name ?? "unknown_function", + }; + + if (OTEL_GENAI_CAPTURE_TOOL_CONTENT) { + if (item.arguments) { + functionCallSpanAttributes["gen_ai.tool.call.arguments"] = buildJsonAttribute(item.arguments); + } + if (matchingOutput?.output) { + functionCallSpanAttributes["gen_ai.tool.call.result"] = buildJsonAttribute(matchingOutput.output); + } + } + + const functionCallSpan = tracer.startSpan( + "gen_ai.execute_tool", + { attributes: functionCallSpanAttributes }, + traceContext + ); + functionCallSpan.setAttribute("tool.status", matchingOutput ? "ok" : "requested"); + functionCallSpan.end(); + } + } + // List MCP tools from server (if required) + prepare tools for the LLM let tools: ChatCompletionTool[] | undefined = []; const mcpToolsMapping: Record = {}; @@ -228,7 +343,7 @@ async function* innerRunStream( } // Otherwise, list tools from MCP server if (!mcpListTools) { - for await (const event of listMcpToolsStream(tool, responseObject)) { + for await (const event of listMcpToolsStream(tool, responseObject, traceContext)) { yield event; } mcpListTools = responseObject.output.at(-1) as ResponseOutputItem.McpListTools; @@ -433,7 +548,8 @@ async function* innerRunStream( approvalRequest, mcpToolsMapping, responseObject, - payload + payload, + traceContext )) { yield event; } @@ -451,7 +567,7 @@ async function* innerRunStream( do { previousMessageCount = currentMessageCount; - for await (const event of handleOneTurnStream(apiKey, payload, responseObject, mcpToolsMapping, defaultHeaders)) { + for await (const event of handleOneTurnStream(apiKey, payload, responseObject, mcpToolsMapping, defaultHeaders, traceContext)) { yield event; } @@ -462,8 +578,21 @@ async function* innerRunStream( async function* listMcpToolsStream( tool: McpServerParams, - responseObject: IncompleteResponse + responseObject: IncompleteResponse, + traceContext: Context ): AsyncGenerator { + const span = tracer.startSpan( + "gen_ai.execute_tool", + { + attributes: { + "gen_ai.operation.name": "execute_tool", + "gen_ai.tool.name": "mcp.list_tools", + "gen_ai.tool.type": "extension", + "mcp.server_label": tool.server_label, + }, + }, + traceContext + ); const outputObject: ResponseOutputItem.McpListTools = { id: generateUniqueId("mcpl"), type: "mcp_list_tools", @@ -501,6 +630,7 @@ async function* listMcpToolsStream( annotations: mcpTool.annotations, description: mcpTool.description, })); + span.setAttribute("mcp.tools.count", outputObject.tools.length); yield { type: "response.output_item.done", output_index: responseObject.output.length - 1, @@ -510,6 +640,7 @@ async function* listMcpToolsStream( } catch (error) { const errorMessage = `Failed to list tools from MCP server '${tool.server_label}': ${error instanceof Error ? error.message : "Unknown error"}`; console.error(errorMessage); + recordError(span, error); yield { type: "response.mcp_list_tools.failed", item_id: outputObject.id, @@ -517,6 +648,8 @@ async function* listMcpToolsStream( sequence_number: SEQUENCE_NUMBER_PLACEHOLDER, }; throw new Error(errorMessage); + } finally { + span.end(); } } @@ -528,8 +661,23 @@ async function* handleOneTurnStream( payload: ChatCompletionCreateParamsStreaming, responseObject: IncompleteResponse, mcpToolsMapping: Record, - defaultHeaders: Record + defaultHeaders: Record, + traceContext: Context ): AsyncGenerator { + const llmSpan = tracer.startSpan( + "gen_ai.chat", + { + attributes: { + "gen_ai.operation.name": "chat", + "gen_ai.request.model": payload.model, + "gen_ai.request.max_tokens": payload.max_tokens ?? undefined, + "gen_ai.request.temperature": payload.temperature ?? undefined, + "gen_ai.request.top_p": payload.top_p ?? undefined, + }, + }, + traceContext + ); + const client = new OpenAI({ baseURL: process.env.OPENAI_BASE_URL ?? "https://router.huggingface.co/v1", apiKey: apiKey, @@ -541,7 +689,8 @@ async function* handleOneTurnStream( let previousTotalTokens = responseObject.usage?.total_tokens ?? 0; let currentTextMode: "text" | "reasoning" = "text"; - for await (const chunk of stream) { + try { + for await (const chunk of stream) { if (chunk.usage) { // Overwrite usage with the latest chunk's usage responseObject.usage = { @@ -566,14 +715,14 @@ async function* handleOneTurnStream( // If start or end of reasoning, skip token and update the current text mode if (reasoningText) { if (currentTextMode === "text") { - for await (const event of closeLastOutputItem(responseObject, payload, mcpToolsMapping)) { + for await (const event of closeLastOutputItem(responseObject, payload, mcpToolsMapping, traceContext)) { yield event; } } currentTextMode = "reasoning"; } else if (delta.content) { if (currentTextMode === "reasoning") { - for await (const event of closeLastOutputItem(responseObject, payload, mcpToolsMapping)) { + for await (const event of closeLastOutputItem(responseObject, payload, mcpToolsMapping, traceContext)) { yield event; } } @@ -773,10 +922,22 @@ async function* handleOneTurnStream( } } } - } + } - for await (const event of closeLastOutputItem(responseObject, payload, mcpToolsMapping)) { - yield event; + for await (const event of closeLastOutputItem(responseObject, payload, mcpToolsMapping, traceContext)) { + yield event; + } + } catch (error) { + recordError(llmSpan, error); + throw error; + } finally { + if (responseObject.usage) { + llmSpan.setAttributes({ + "gen_ai.usage.input_tokens": responseObject.usage.input_tokens, + "gen_ai.usage.output_tokens": responseObject.usage.output_tokens, + }); + } + llmSpan.end(); } } @@ -789,7 +950,8 @@ async function* callApprovedMCPToolStream( approvalRequest: McpApprovalRequestParams | undefined, mcpToolsMapping: Record, responseObject: IncompleteResponse, - payload: ChatCompletionCreateParamsStreaming + payload: ChatCompletionCreateParamsStreaming, + traceContext: Context ): AsyncGenerator { if (!approvalRequest) { throw new Error(`MCP approval request '${approval_request_id}' not found`); @@ -819,11 +981,40 @@ async function* callApprovedMCPToolStream( sequence_number: SEQUENCE_NUMBER_PLACEHOLDER, }; + const toolSpan = tracer.startSpan( + "gen_ai.execute_tool", + { + attributes: { + "gen_ai.operation.name": "execute_tool", + "gen_ai.tool.name": approvalRequest.name, + "gen_ai.tool.type": "extension", + "gen_ai.tool.call.id": outputObject.id, + "mcp.server_label": approvalRequest.server_label, + ...(OTEL_GENAI_CAPTURE_TOOL_CONTENT + ? { + "gen_ai.tool.call.arguments": buildJsonAttribute(approvalRequest.arguments), + } + : {}), + }, + }, + traceContext + ); + const toolParams = mcpToolsMapping[approvalRequest.name]; - const toolResult = await callMcpTool(toolParams, approvalRequest.name, approvalRequest.arguments); + let toolResult; + try { + toolResult = await callMcpTool(toolParams, approvalRequest.name, approvalRequest.arguments); + } catch (error) { + recordError(toolSpan, error); + toolSpan.end(); + throw error; + } if (toolResult.error) { outputObject.error = toolResult.error; + toolSpan.setAttribute("tool.status", "error"); + toolSpan.setAttribute("tool.error", toolResult.error); + recordError(toolSpan, new Error(toolResult.error)); yield { type: "response.mcp_call.failed", item_id: outputObject.id, @@ -832,6 +1023,10 @@ async function* callApprovedMCPToolStream( }; } else { outputObject.output = toolResult.output; + toolSpan.setAttribute("tool.status", "ok"); + if (OTEL_GENAI_CAPTURE_TOOL_CONTENT) { + toolSpan.setAttribute("gen_ai.tool.call.result", buildJsonAttribute(toolResult.output)); + } yield { type: "response.mcp_call.completed", item_id: outputObject.id, @@ -870,6 +1065,8 @@ async function* callApprovedMCPToolStream( content: outputObject.output ? outputObject.output : outputObject.error ? `Error: ${outputObject.error}` : "", } ); + + toolSpan.end(); } function requiresApproval(toolName: string, mcpToolsMapping: Record): boolean { @@ -888,7 +1085,8 @@ function requiresApproval(toolName: string, mcpToolsMapping: Record + mcpToolsMapping: Record, + traceContext: Context ): AsyncGenerator { const lastOutputItem = responseObject.output.at(-1); if (lastOutputItem) { @@ -955,6 +1153,21 @@ async function* closeLastOutputItem( sequence_number: SEQUENCE_NUMBER_PLACEHOLDER, }; } else if (lastOutputItem?.type === "function_call") { + const functionCallSpanAttributes: Attributes = { + "gen_ai.operation.name": "execute_tool", + "gen_ai.tool.name": lastOutputItem.name, + "gen_ai.tool.type": "function", + "gen_ai.tool.call.id": lastOutputItem.call_id || lastOutputItem.id, + }; + if (OTEL_GENAI_CAPTURE_TOOL_CONTENT) { + functionCallSpanAttributes["gen_ai.tool.call.arguments"] = buildJsonAttribute(lastOutputItem.arguments); + } + const functionCallSpan = tracer.startSpan( + "gen_ai.execute_tool", + { attributes: functionCallSpanAttributes }, + traceContext + ); + yield { type: "response.function_call_arguments.done", item_id: lastOutputItem.id as string, @@ -964,12 +1177,14 @@ async function* closeLastOutputItem( }; lastOutputItem.status = "completed"; + functionCallSpan.setAttribute("tool.status", "requested"); yield { type: "response.output_item.done", output_index: responseObject.output.length - 1, item: lastOutputItem, sequence_number: SEQUENCE_NUMBER_PLACEHOLDER, }; + functionCallSpan.end(); } else if (lastOutputItem?.type === "mcp_call") { yield { type: "response.mcp_call_arguments.done", @@ -981,9 +1196,31 @@ async function* closeLastOutputItem( // Call MCP tool const toolParams = mcpToolsMapping[lastOutputItem.name]; - const toolResult = await callMcpTool(toolParams, lastOutputItem.name, lastOutputItem.arguments); + const toolSpanAttributes: Attributes = { + "gen_ai.operation.name": "execute_tool", + "gen_ai.tool.name": lastOutputItem.name, + "gen_ai.tool.type": "extension", + "gen_ai.tool.call.id": lastOutputItem.id, + "mcp.server_label": lastOutputItem.server_label, + }; + if (OTEL_GENAI_CAPTURE_TOOL_CONTENT) { + toolSpanAttributes["gen_ai.tool.call.arguments"] = buildJsonAttribute(lastOutputItem.arguments); + } + const toolSpan = tracer.startSpan("gen_ai.execute_tool", { attributes: toolSpanAttributes }, traceContext); + + let toolResult; + try { + toolResult = await callMcpTool(toolParams, lastOutputItem.name, lastOutputItem.arguments); + } catch (error) { + recordError(toolSpan, error); + toolSpan.end(); + throw error; + } if (toolResult.error) { lastOutputItem.error = toolResult.error; + toolSpan.setAttribute("tool.status", "error"); + toolSpan.setAttribute("tool.error", toolResult.error); + recordError(toolSpan, new Error(toolResult.error)); yield { type: "response.mcp_call.failed", item_id: lastOutputItem.id as string, @@ -992,6 +1229,10 @@ async function* closeLastOutputItem( }; } else { lastOutputItem.output = toolResult.output; + toolSpan.setAttribute("tool.status", "ok"); + if (OTEL_GENAI_CAPTURE_TOOL_CONTENT) { + toolSpan.setAttribute("gen_ai.tool.call.result", buildJsonAttribute(toolResult.output)); + } yield { type: "response.mcp_call.completed", item_id: lastOutputItem.id as string, @@ -999,6 +1240,7 @@ async function* closeLastOutputItem( sequence_number: SEQUENCE_NUMBER_PLACEHOLDER, }; } + toolSpan.end(); yield { type: "response.output_item.done",