diff --git a/src/lib/langfuse-transport.ts b/src/lib/langfuse-transport.ts index afc9f5b..ac85b68 100644 --- a/src/lib/langfuse-transport.ts +++ b/src/lib/langfuse-transport.ts @@ -1,8 +1,50 @@ import { NodeTracerProvider } from "@opentelemetry/sdk-trace-node"; import { LangfuseSpanProcessor } from "@langfuse/otel"; +/** Trace name used for HTTP transport in Langfuse. */ +export const TRACE_NAME_HTTP = "ask262-http"; + +/** Trace name used for stdio transport in Langfuse. */ +export const TRACE_NAME_STDIO = "ask262-stdio"; + +/** Langfuse attribute key for setting the trace name on a span. */ +export const LANGFUSE_TRACE_NAME_ATTR = "langfuse.trace.name"; + +/** Langfuse attribute key for observation input. */ +export const LANGFUSE_OBSERVATION_INPUT_ATTR = "langfuse.observation.input"; + +/** Langfuse attribute key for observation output. */ +export const LANGFUSE_OBSERVATION_OUTPUT_ATTR = "langfuse.observation.output"; + +/** Langfuse attribute key for trace-level input. */ +export const LANGFUSE_TRACE_INPUT_ATTR = "langfuse.trace.input"; + +/** Langfuse attribute key for trace-level output. */ +export const LANGFUSE_TRACE_OUTPUT_ATTR = "langfuse.trace.output"; + let provider: NodeTracerProvider | null = null; +/** + * Extract MCP tool call information from a JSON-RPC request body. + * Used to populate trace-level input in Langfuse. + */ +export function extractMcpToolInfo( + body: unknown, +): { method?: string; tool?: string; input?: unknown } { + if (typeof body !== "object" || body === null) return {}; + const b = body as Record; + const rpcMethod = b.method as string | undefined; + if (rpcMethod === "tools/call") { + const params = b.params as Record | undefined; + return { + method: rpcMethod, + tool: params?.name as string | undefined, + input: params?.arguments, + }; + } + return { method: rpcMethod }; +} + /** * Initialize Langfuse OTel span processor. * Called once at server startup before any spans are created. diff --git a/src/mcp-server-http.ts b/src/mcp-server-http.ts index 6b525e2..9ec9189 100644 --- a/src/mcp-server-http.ts +++ b/src/mcp-server-http.ts @@ -36,6 +36,14 @@ import { import { DEFAULT_PORT, STORAGE_DIR as STORAGE_DIR_REL } from "./constants.js"; import { createEmbeddings } from "./lib/embeddings-factory.js"; import { LogOperation, logger } from "./lib/logger.js"; +import { trace } from "@opentelemetry/api"; +import { + LANGFUSE_OBSERVATION_INPUT_ATTR, + LANGFUSE_OBSERVATION_OUTPUT_ATTR, + LANGFUSE_TRACE_NAME_ATTR, + LANGFUSE_TRACE_OUTPUT_ATTR, + TRACE_NAME_HTTP, +} from "./lib/langfuse-transport.js"; import { getSessionMetadata, setupTracing, withSpan } from "./lib/tracing.js"; // Resolve storage path relative to this script's directory @@ -93,12 +101,24 @@ async function createMcpServer() { return await withSpan( "ask262_search_spec_sections", { - "langfuse.observation.input": JSON.stringify({ query }), + [LANGFUSE_OBSERVATION_INPUT_ATTR]: JSON.stringify({ query }), tool: searchSpecToolName, query, }, async () => { const result = await searchSpecTool({ query }); + trace + .getActiveSpan() + ?.setAttribute( + LANGFUSE_TRACE_OUTPUT_ATTR, + JSON.stringify(result), + ); + trace + .getActiveSpan() + ?.setAttribute( + LANGFUSE_OBSERVATION_OUTPUT_ATTR, + JSON.stringify(result), + ); return { content: [{ type: "text", text: JSON.stringify(result, null, 2) }], structuredContent: result, @@ -126,7 +146,7 @@ async function createMcpServer() { return await withSpan( "ask262_get_section_content", { - "langfuse.observation.input": JSON.stringify({ + [LANGFUSE_OBSERVATION_INPUT_ATTR]: JSON.stringify({ sectionIds, recursive, }), @@ -135,6 +155,18 @@ async function createMcpServer() { }, async () => { const result = await getSectionContentTool({ sectionIds, recursive }); + trace + .getActiveSpan() + ?.setAttribute( + LANGFUSE_TRACE_OUTPUT_ATTR, + JSON.stringify(result), + ); + trace + .getActiveSpan() + ?.setAttribute( + LANGFUSE_OBSERVATION_OUTPUT_ATTR, + JSON.stringify(result), + ); return { content: [{ type: "text", text: JSON.stringify(result, null, 2) }], structuredContent: result, @@ -162,14 +194,24 @@ async function createMcpServer() { return await withSpan( "ask262_evaluate_in_engine262", { - "langfuse.observation.input": JSON.stringify({ - code: code.slice(0, 200), - }), + [LANGFUSE_OBSERVATION_INPUT_ATTR]: JSON.stringify({ code }), tool: evaluateToolName, code_length: code.length, }, async () => { const result = await evaluateTool({ code }); + trace + .getActiveSpan() + ?.setAttribute( + LANGFUSE_TRACE_OUTPUT_ATTR, + JSON.stringify(result), + ); + trace + .getActiveSpan() + ?.setAttribute( + LANGFUSE_OBSERVATION_OUTPUT_ATTR, + JSON.stringify(result), + ); const isError = result.error !== undefined; const text = isError ? result.error @@ -259,8 +301,17 @@ export async function main() { // Handle request within trace context (passing trace ID from header if available) const sessionMetadata = getSessionMetadata("http"); return await withSpan( - LogOperation.HANDLING_MCP_HTTP_REQUEST, - { method: c.req.method, client_ip: clientIp }, + "mcp_http_request", + { + [LANGFUSE_TRACE_NAME_ATTR]: TRACE_NAME_HTTP, + [LANGFUSE_OBSERVATION_INPUT_ATTR]: JSON.stringify({ + method: c.req.method, + endpoint: "/mcp", + client_ip: clientIp, + }), + method: c.req.method, + client_ip: clientIp, + }, async () => { const op = log.start(LogOperation.HANDLING_MCP_HTTP_REQUEST, { method: c.req.method, @@ -284,7 +335,6 @@ export async function main() { op.end({ status: "success" }); - // Return the Web Standard Response directly return response; } catch (err) { const error = err instanceof Error ? err : new Error(String(err)); diff --git a/src/mcp-server-stdio.ts b/src/mcp-server-stdio.ts index 94ce499..1f25649 100644 --- a/src/mcp-server-stdio.ts +++ b/src/mcp-server-stdio.ts @@ -37,6 +37,15 @@ import { import { STORAGE_DIR as STORAGE_DIR_REL } from "./constants.js"; import { createEmbeddings } from "./lib/embeddings-factory.js"; import { LogOperation, logger } from "./lib/logger.js"; +import { trace } from "@opentelemetry/api"; +import { + LANGFUSE_OBSERVATION_INPUT_ATTR, + LANGFUSE_OBSERVATION_OUTPUT_ATTR, + LANGFUSE_TRACE_INPUT_ATTR, + LANGFUSE_TRACE_NAME_ATTR, + LANGFUSE_TRACE_OUTPUT_ATTR, + TRACE_NAME_STDIO, +} from "./lib/langfuse-transport.js"; import { createProcessScopedTrace, getSessionMetadata, @@ -149,12 +158,29 @@ export async function main() { sessionTraceId, "ask262_search_spec_sections", { - "langfuse.observation.input": JSON.stringify({ query }), + LANGFUSE_TRACE_NAME_ATTR: TRACE_NAME_STDIO, + LANGFUSE_TRACE_INPUT_ATTR: JSON.stringify({ + tool: searchSpecToolName, + input: { query }, + }), + LANGFUSE_OBSERVATION_INPUT_ATTR: JSON.stringify({ query }), tool: searchSpecToolName, query, }, async () => { const result = await searchSpecTool({ query }); + trace + .getActiveSpan() + ?.setAttribute( + LANGFUSE_OBSERVATION_OUTPUT_ATTR, + JSON.stringify(result), + ); + trace + .getActiveSpan() + ?.setAttribute( + LANGFUSE_TRACE_OUTPUT_ATTR, + JSON.stringify(result), + ); return { content: [{ type: "text", text: JSON.stringify(result, null, 2) }], structuredContent: result, @@ -187,7 +213,12 @@ export async function main() { sessionTraceId, "ask262_get_section_content", { - "langfuse.observation.input": JSON.stringify({ + LANGFUSE_TRACE_NAME_ATTR: TRACE_NAME_STDIO, + LANGFUSE_TRACE_INPUT_ATTR: JSON.stringify({ + tool: sectionContentToolName, + input: { sectionIds, recursive }, + }), + LANGFUSE_OBSERVATION_INPUT_ATTR: JSON.stringify({ sectionIds, recursive, }), @@ -196,6 +227,18 @@ export async function main() { }, async () => { const result = await getSectionContentTool({ sectionIds, recursive }); + trace + .getActiveSpan() + ?.setAttribute( + LANGFUSE_OBSERVATION_OUTPUT_ATTR, + JSON.stringify(result), + ); + trace + .getActiveSpan() + ?.setAttribute( + LANGFUSE_TRACE_OUTPUT_ATTR, + JSON.stringify(result), + ); return { content: [{ type: "text", text: JSON.stringify(result, null, 2) }], structuredContent: result, @@ -225,14 +268,29 @@ export async function main() { sessionTraceId, "ask262_evaluate_in_engine262", { - "langfuse.observation.input": JSON.stringify({ - code: code.slice(0, 200), + LANGFUSE_TRACE_NAME_ATTR: TRACE_NAME_STDIO, + LANGFUSE_TRACE_INPUT_ATTR: JSON.stringify({ + tool: evaluateToolName, + input: { code }, }), + LANGFUSE_OBSERVATION_INPUT_ATTR: JSON.stringify({ code }), tool: evaluateToolName, code_length: code.length, }, async () => { const result = await evaluateTool({ code }); + trace + .getActiveSpan() + ?.setAttribute( + LANGFUSE_OBSERVATION_OUTPUT_ATTR, + JSON.stringify(result), + ); + trace + .getActiveSpan() + ?.setAttribute( + LANGFUSE_TRACE_OUTPUT_ATTR, + JSON.stringify(result), + ); const isError = result.error !== undefined; const text = isError ? result.error : JSON.stringify(result, null, 2); return {