diff --git a/apps/ai/src/mcp/tools/__tests__/agent-tools.test.ts b/apps/ai/src/mcp/tools/__tests__/agent-tools.test.ts index d9f1efe3f1..44c5985c61 100644 --- a/apps/ai/src/mcp/tools/__tests__/agent-tools.test.ts +++ b/apps/ai/src/mcp/tools/__tests__/agent-tools.test.ts @@ -127,10 +127,10 @@ const variantRows = [ const breakdownPairRows = [{ model: "claude-sonnet-4", service: "api", calls: 9 }] -// Exactly what the read returns for a payload it cut: `AI_TOOL_ERROR_PAYLOAD_MAX` -// characters, with the row reporting the true size the span carried. -const LONG_ARGUMENTS = `{"query":"${"x".repeat(AI_TOOL_ERROR_PAYLOAD_MAX - 12)}"}` +// Arguments longer than the read keeps: it cuts them to `AI_TOOL_ERROR_PAYLOAD_MAX` +// characters and reports the true size the span carried. const LONG_ARGUMENTS_BYTES = 12_000 +const LONG_ARGUMENTS = `{"query":"${"x".repeat(LONG_ARGUMENTS_BYTES - 12)}"}` /** The payload rows below are joined to these by trace and span id. */ const occurrence = (ids: [string, string], sessionId: string, timestamp: string, durationNs: number) => ({ @@ -153,20 +153,14 @@ const occurrenceRows = [ ] /** `statusCode` is Title case on the wire, like every other Maple span status. */ -const payload = (ids: [string, string], args: string, argumentsBytes: number) => ({ +const payload = (ids: [string, string], args: string, result = "TimeoutError: upstream timed out") => ({ traceId: ids[0].repeat(32), spanId: ids[1].repeat(16), statusCode: "Error", - arguments: args, - argumentsBytes, - result: "TimeoutError: upstream timed out", - resultBytes: 32, + spanAttributes: { "gen_ai.tool.call.arguments": args, "gen_ai.tool.call.result": result }, }) -const payloadRows = [ - payload(["a", "b"], LONG_ARGUMENTS, LONG_ARGUMENTS_BYTES), - payload(["c", "d"], `{"query":"short"}`, 17), -] +const payloadRows = [payload(["a", "b"], LONG_ARGUMENTS), payload(["c", "d"], `{"query":"short"}`)] /** The payload read, answering with rows of the test's own shape. */ const payloadsAre = (rows: ReadonlyArray): FixtureRule[] => [ @@ -478,18 +472,14 @@ describe("get_agent_tool_error rendering", () => { }) it("tells a payload the span left empty from one it never carried", async () => { - const rows = payloadRows.map((row) => ({ ...row, arguments: "", argumentsBytes: 0 })) + const rows = [payload(["a", "b"], ""), payload(["c", "d"], "")] const rendered = await renderWith(payloadsAre(rows), ERROR_DETAIL, GROUP) expect(rendered).toContain("(empty)") expect(rendered).not.toContain("(not available: the span was not retained)") }) it("fences a sample payload that carries a code fence of its own", async () => { - const rows = payloadRows.map((row) => ({ - ...row, - result: FENCED_RESULT, - resultBytes: FENCED_RESULT.length, - })) + const rows = [payload(["a", "b"], "{}", FENCED_RESULT), payload(["c", "d"], "{}", FENCED_RESULT)] const rendered = await renderWith(payloadsAre(rows), ERROR_DETAIL, GROUP) // One backtick longer than the longest run inside the payload, which is // left intact. diff --git a/apps/api/src/routes/internal/ai-sessions.http.test.ts b/apps/api/src/routes/internal/ai-sessions.http.test.ts index 09bd337525..6bb37c93e7 100644 --- a/apps/api/src/routes/internal/ai-sessions.http.test.ts +++ b/apps/api/src/routes/internal/ai-sessions.http.test.ts @@ -1668,10 +1668,11 @@ describe("POST /internal/ai-sessions/tools/error-samples", () => { traceId: occurrence.traceId, spanId: occurrence.spanId, statusCode: "Ok", - arguments: "{}", - argumentsBytes: 2, - result: occurrence.message, - resultBytes: 58, + spanAttributes: { + "maple_ai.vendor.id": "maple", + "gen_ai.tool.call.arguments": "{}", + "gen_ai.tool.call.result": occurrence.message, + }, }, ], ) @@ -1708,7 +1709,7 @@ describe("POST /internal/ai-sessions/tools/error-samples", () => { arguments: "{}", argumentsBytes: 2, result: occurrence.message, - resultBytes: 58, + resultBytes: occurrence.message.length, }, ], }) diff --git a/packages/backend/src/services/ai-sessions/ai-session-reads.ts b/packages/backend/src/services/ai-sessions/ai-session-reads.ts index 752230dae2..1d6acf9313 100644 --- a/packages/backend/src/services/ai-sessions/ai-session-reads.ts +++ b/packages/backend/src/services/ai-sessions/ai-session-reads.ts @@ -749,7 +749,11 @@ export const readAiToolErrorSamples = Effect.fn("aiSessions.toolErrorSamples")(f ), { profile: "list", context: "aiToolsErrorPayloads" }, ) - const payloadBySpan = new Map(payloads.map((row) => [`${row.traceId}:${row.spanId}`, row] as const)) + const payloadBySpan = new Map( + payloads.map( + (row) => [`${row.traceId}:${row.spanId}`, Integrations.aiToolErrorPayload(row)] as const, + ), + ) return new AiToolErrorSamplesResponse({ ...(nextCursor !== undefined && { nextCursor }), occurrences: occurrences.map((row) => { diff --git a/packages/backend/src/services/warehouse/ai-tools.clickhouse.e2e.test.ts b/packages/backend/src/services/warehouse/ai-tools.clickhouse.e2e.test.ts index 88f5a215e9..babe0b48dd 100644 --- a/packages/backend/src/services/warehouse/ai-tools.clickhouse.e2e.test.ts +++ b/packages/backend/src/services/warehouse/ai-tools.clickhouse.e2e.test.ts @@ -96,6 +96,9 @@ const TRACE_IDS_AT_1 = missingKey('["evidence"][1]["traceIds"]') const LOG_PATTERNS_AT_0 = missingKey('["evidence"][0]["logPatterns"]') /** A failure that says why in its status message alone. */ const REFUSED = "sandbox refused the command" +/** `flaky_tool`'s parameter schema. */ +const FLAKY_SCHEMA = + '{"properties": {"retries": {"type": "integer"}}, "required": ["retries"], "type": "object"}' interface SeedSpan { readonly traceId: string @@ -280,10 +283,17 @@ const SEED_SPANS: ReadonlyArray = [ durationNs: 8_000_000, status: "Error", statusMessage: "upstream returned 503", + // An OpenInference tool span whose GenAI dual-write copied the parameter + // schema into the arguments slot: the payload read shows `input.value`, + // as the session page does. attrs: agentSpan({ + [MAPLE_AI_VENDOR_ID_ATTR]: "openai_agents_sdk", + "openinference.span.kind": "TOOL", "gen_ai.operation.name": "execute_tool", "gen_ai.tool.name": "flaky_tool", - "gen_ai.tool.call.arguments": '{"retries":1}', + "tool.parameters": FLAKY_SCHEMA, + "gen_ai.tool.call.arguments": FLAKY_SCHEMA, + "input.value": '{"retries":1}', "gen_ai.tool.call.result": '{"error":"503"}', ...aiGatewayStamps({ toolCall: true, @@ -825,13 +835,15 @@ describe.skipIf(!clickhouseE2eEnabled)("agent tools reads", () => { ) assert.isFalse(payloads.sql.includes(flakyWindow.startTime)) assert.deepStrictEqual( - Effect.runSync(payloads.decodeRows(await runJson(payloads.sql))).map((row) => ({ - spanId: row.spanId, - statusCode: row.statusCode, - arguments: row.arguments, - argumentsBytes: row.argumentsBytes, - resultBytes: row.resultBytes, - })), + Effect.runSync(payloads.decodeRows(await runJson(payloads.sql))) + .map(Integrations.aiToolErrorPayload) + .map((row) => ({ + spanId: row.spanId, + statusCode: row.statusCode, + arguments: row.arguments, + argumentsBytes: row.argumentsBytes, + resultBytes: row.resultBytes, + })), [ { spanId: "tools-flaky-1", @@ -870,7 +882,9 @@ describe.skipIf(!clickhouseE2eEnabled)("agent tools reads", () => { { orgId: ORG_ID, ...slice }, { rowSchema: Integrations.aiToolErrorPayloadsRowSchema }, ) - const rows = Effect.runSync(payloads.decodeRows(await runJson(payloads.sql))) + const rows = Effect.runSync(payloads.decodeRows(await runJson(payloads.sql))).map( + Integrations.aiToolErrorPayload, + ) assert.deepStrictEqual( [...rows] .sort((a, b) => a.spanId.localeCompare(b.spanId)) diff --git a/packages/query-engine-integrations/src/__sql_baseline__/integrations.sql b/packages/query-engine-integrations/src/__sql_baseline__/integrations.sql index 586a92035c..eeee5326b7 100644 --- a/packages/query-engine-integrations/src/__sql_baseline__/integrations.sql +++ b/packages/query-engine-integrations/src/__sql_baseline__/integrations.sql @@ -1170,10 +1170,8 @@ SELECT TraceId AS traceId, SpanId AS spanId, StatusCode AS statusCode, - leftUTF8(coalesce(nullIf(SpanAttributes['gen_ai.tool.call.arguments'], ''), nullIf(SpanAttributes['ai.toolCall.args'], ''), ''), 4000) AS arguments, - length(coalesce(nullIf(SpanAttributes['gen_ai.tool.call.arguments'], ''), nullIf(SpanAttributes['ai.toolCall.args'], ''), '')) AS argumentsBytes, - leftUTF8(coalesce(nullIf(SpanAttributes['gen_ai.tool.call.result'], ''), nullIf(SpanAttributes['ai.toolCall.result'], ''), ''), 4000) AS result, - length(coalesce(nullIf(SpanAttributes['gen_ai.tool.call.result'], ''), nullIf(SpanAttributes['ai.toolCall.result'], ''), '')) AS resultBytes + mapApply((k, v) -> (k, leftUTF8(v, 16384)), mapFilter((k, v) -> (((k IN ('maple_ai.session.id', 'maple_ai.vendor.id', 'maple_ai.vendor.version', 'maple_ai.agent.name', 'gen_ai.operation.name', 'gen_ai.provider.name', 'gen_ai.system', 'gen_ai.request.model', 'gen_ai.request.max_tokens', 'gen_ai.request.choice.count', 'gen_ai.request.temperature', 'gen_ai.request.top_p', 'gen_ai.request.top_k', 'gen_ai.request.stop_sequences', 'gen_ai.request.frequency_penalty', 'gen_ai.request.presence_penalty', 'gen_ai.request.encoding_formats', 'gen_ai.request.seed', 'gen_ai.openai.request.seed', 'gen_ai.request.stream', 'gen_ai.request.reasoning.level', 'gen_ai.request.previous_response.id', 'gen_ai.request.stream_cursor', 'gen_ai.response.id', 'gen_ai.response.model', 'gen_ai.response.finish_reasons', 'gen_ai.response.finish_reason', 'gen_ai.response.status', 'gen_ai.response.time_to_first_chunk', 'gen_ai.output.type', 'gen_ai.usage.input_tokens', 'gen_ai.usage.prompt_tokens', 'gen_ai.usage.cache_read.input_tokens', 'gen_ai.usage.input_tokens.cached', 'gen_ai.usage.cache_creation.input_tokens', 'gen_ai.usage.cache_write.input_tokens', 'gen_ai.usage.output_tokens', 'gen_ai.usage.completion_tokens', 'gen_ai.usage.reasoning.output_tokens', 'gen_ai.usage.output_tokens.reasoning', 'gen_ai.usage.cost', 'gen_ai.usage.total_cost', 'maple_ai.llm_call', 'maple_ai.tool_call', 'maple_ai.error', 'maple_ai.usage.input_tokens', 'maple_ai.usage.cache_read_tokens', 'maple_ai.usage.cache_write_tokens', 'maple_ai.usage.output_tokens', 'maple_ai.usage.reasoning_tokens', 'maple_ai.usage.cost', 'gen_ai.conversation.id', 'gen_ai.conversation.compacted', 'gen_ai.agent.id', 'gen_ai.agent.name', 'gen_ai.agent.description', 'gen_ai.agent.version', 'gen_ai.tool.name', 'gen_ai.tool.call.id', 'gen_ai.tool.description', 'gen_ai.tool.type', 'gen_ai.tool.call.arguments', 'gen_ai.tool.call.result', 'gen_ai.tool.definitions', 'gen_ai.system_instructions', 'gen_ai.input.messages', 'gen_ai.prompt', 'gen_ai.output.messages', 'gen_ai.completion', 'gen_ai.data_source.id', 'gen_ai.retrieval.query.text', 'gen_ai.retrieval.top_k', 'gen_ai.retrieval.documents', 'gen_ai.memory.store.id', 'gen_ai.memory.record.id', 'gen_ai.memory.record.count', 'gen_ai.memory.query.text', 'gen_ai.memory.records', 'gen_ai.embeddings.dimension.count', 'gen_ai.evaluation.name', 'gen_ai.evaluation.score.value', 'gen_ai.evaluation.score.label', 'gen_ai.evaluation.explanation', 'gen_ai.prompt.name', 'gen_ai.prompt.version', 'gen_ai.workflow.name', 'span.metadata.attempt_index', 'span.metadata.status_code', 'trace.metadata.openrouter.provider_name', 'error.type', 'server.address', 'server.port', 'ai.model.provider', 'ai.model.id', 'ai.response.id', 'ai.response.model', 'ai.response.finishReason', 'gen_ai.client.operation.time_to_first_chunk', 'ai.usage.inputTokens', 'ai.usage.promptTokens', 'ai.usage.cachedInputTokens', 'ai.usage.inputTokenDetails.cacheReadTokens', 'ai.usage.inputTokenDetails.cacheWriteTokens', 'ai.usage.outputTokens', 'ai.usage.completionTokens', 'ai.usage.reasoningTokens', 'ai.usage.outputTokenDetails.reasoningTokens', 'ai.telemetry.functionId', 'ai.toolCall.name', 'ai.toolCall.id', 'ai.toolCall.args', 'ai.toolCall.result', 'ai.prompt.tools', 'ai.prompt.messages', 'ai.prompt', 'llm.provider', 'llm.system', 'llm.model_name', 'llm.finish_reason', 'llm.token_count.prompt', 'llm.token_count.prompt_details.cache_read', 'llm.token_count.completion', 'llm.token_count.completion_details.reasoning', 'llm.cost.total', 'tool.name', 'tool.description', 'llm.tools', 'openinference.span.kind', 'tool.parameters', 'input.value', 'output.value', 'eve.turn.id', 'maple_ai.turn.id') OR k LIKE 'gen_ai.prompt.variable.%') OR k LIKE 'llm.input_messages.%') OR k LIKE 'llm.output_messages.%'), SpanAttributes)) AS spanAttributes, + mapApply((k, v) -> (k, length(v)), mapFilter((k, v) -> lengthUTF8(v) > 16384, mapFilter((k, v) -> (((k IN ('maple_ai.session.id', 'maple_ai.vendor.id', 'maple_ai.vendor.version', 'maple_ai.agent.name', 'gen_ai.operation.name', 'gen_ai.provider.name', 'gen_ai.system', 'gen_ai.request.model', 'gen_ai.request.max_tokens', 'gen_ai.request.choice.count', 'gen_ai.request.temperature', 'gen_ai.request.top_p', 'gen_ai.request.top_k', 'gen_ai.request.stop_sequences', 'gen_ai.request.frequency_penalty', 'gen_ai.request.presence_penalty', 'gen_ai.request.encoding_formats', 'gen_ai.request.seed', 'gen_ai.openai.request.seed', 'gen_ai.request.stream', 'gen_ai.request.reasoning.level', 'gen_ai.request.previous_response.id', 'gen_ai.request.stream_cursor', 'gen_ai.response.id', 'gen_ai.response.model', 'gen_ai.response.finish_reasons', 'gen_ai.response.finish_reason', 'gen_ai.response.status', 'gen_ai.response.time_to_first_chunk', 'gen_ai.output.type', 'gen_ai.usage.input_tokens', 'gen_ai.usage.prompt_tokens', 'gen_ai.usage.cache_read.input_tokens', 'gen_ai.usage.input_tokens.cached', 'gen_ai.usage.cache_creation.input_tokens', 'gen_ai.usage.cache_write.input_tokens', 'gen_ai.usage.output_tokens', 'gen_ai.usage.completion_tokens', 'gen_ai.usage.reasoning.output_tokens', 'gen_ai.usage.output_tokens.reasoning', 'gen_ai.usage.cost', 'gen_ai.usage.total_cost', 'maple_ai.llm_call', 'maple_ai.tool_call', 'maple_ai.error', 'maple_ai.usage.input_tokens', 'maple_ai.usage.cache_read_tokens', 'maple_ai.usage.cache_write_tokens', 'maple_ai.usage.output_tokens', 'maple_ai.usage.reasoning_tokens', 'maple_ai.usage.cost', 'gen_ai.conversation.id', 'gen_ai.conversation.compacted', 'gen_ai.agent.id', 'gen_ai.agent.name', 'gen_ai.agent.description', 'gen_ai.agent.version', 'gen_ai.tool.name', 'gen_ai.tool.call.id', 'gen_ai.tool.description', 'gen_ai.tool.type', 'gen_ai.tool.call.arguments', 'gen_ai.tool.call.result', 'gen_ai.tool.definitions', 'gen_ai.system_instructions', 'gen_ai.input.messages', 'gen_ai.prompt', 'gen_ai.output.messages', 'gen_ai.completion', 'gen_ai.data_source.id', 'gen_ai.retrieval.query.text', 'gen_ai.retrieval.top_k', 'gen_ai.retrieval.documents', 'gen_ai.memory.store.id', 'gen_ai.memory.record.id', 'gen_ai.memory.record.count', 'gen_ai.memory.query.text', 'gen_ai.memory.records', 'gen_ai.embeddings.dimension.count', 'gen_ai.evaluation.name', 'gen_ai.evaluation.score.value', 'gen_ai.evaluation.score.label', 'gen_ai.evaluation.explanation', 'gen_ai.prompt.name', 'gen_ai.prompt.version', 'gen_ai.workflow.name', 'span.metadata.attempt_index', 'span.metadata.status_code', 'trace.metadata.openrouter.provider_name', 'error.type', 'server.address', 'server.port', 'ai.model.provider', 'ai.model.id', 'ai.response.id', 'ai.response.model', 'ai.response.finishReason', 'gen_ai.client.operation.time_to_first_chunk', 'ai.usage.inputTokens', 'ai.usage.promptTokens', 'ai.usage.cachedInputTokens', 'ai.usage.inputTokenDetails.cacheReadTokens', 'ai.usage.inputTokenDetails.cacheWriteTokens', 'ai.usage.outputTokens', 'ai.usage.completionTokens', 'ai.usage.reasoningTokens', 'ai.usage.outputTokenDetails.reasoningTokens', 'ai.telemetry.functionId', 'ai.toolCall.name', 'ai.toolCall.id', 'ai.toolCall.args', 'ai.toolCall.result', 'ai.prompt.tools', 'ai.prompt.messages', 'ai.prompt', 'llm.provider', 'llm.system', 'llm.model_name', 'llm.finish_reason', 'llm.token_count.prompt', 'llm.token_count.prompt_details.cache_read', 'llm.token_count.completion', 'llm.token_count.completion_details.reasoning', 'llm.cost.total', 'tool.name', 'tool.description', 'llm.tools', 'openinference.span.kind', 'tool.parameters', 'input.value', 'output.value', 'eve.turn.id', 'maple_ai.turn.id') OR k LIKE 'gen_ai.prompt.variable.%') OR k LIKE 'llm.input_messages.%') OR k LIKE 'llm.output_messages.%'), SpanAttributes))) AS cutAttributeBytes FROM trace_detail_spans WHERE OrgId = 'org_sql_catalog' AND Timestamp >= '2026-01-02 11:15:00.000000000' diff --git a/packages/query-engine-integrations/src/ai/ai-integrations.ts b/packages/query-engine-integrations/src/ai/ai-integrations.ts index 0ee97cc0ee..9e676901f0 100644 --- a/packages/query-engine-integrations/src/ai/ai-integrations.ts +++ b/packages/query-engine-integrations/src/ai/ai-integrations.ts @@ -32,7 +32,6 @@ import { AI_VENDOR_INTEGRATIONS } from "./ai-vendors" import { isRecord, unwrapMessages, unwrapOutputMessages, unwrapToolMessage } from "./ai-messages" export interface AiRefineContext { - readonly row: AiSessionSpansOutput /** The span's own attributes — the map the source key lists read. */ readonly attributes: Record /** One attribute decoded the way the mapper decodes `field`; `undefined` @@ -348,16 +347,13 @@ export const AI_NON_SIGNAL_FIELDS: ReadonlySet = new Set([...AI_CO const hasAiSignal = (values: MutableAiGenAiValues): boolean => Object.keys(values).some((field) => !AI_NON_SIGNAL_FIELDS.has(field as AiGenAiField)) -export const mapAiSpan = (row: AiSessionSpansOutput): AiAgentSpan => { - // Span attributes only, envelope and source keys alike. The gateway strips - // `maple_ai.*` from span attributes before stamping its own verdict, so a - // span-level value is authoritative — and it does not touch resource - // attributes, where one forged `gen_ai.*` or `maple_ai.*` key would mark - // every span in the service as an AI span. - const attributes = row.spanAttributes - const vendorId = readAttribute(attributes, MAPLE_AI_VENDOR_ID_ATTR) +/** Every catalog field of one span, through the integration its vendor stamp + * selects: the source keys in order, then the refine hooks. */ +const decodeGenAi = ( + attributes: Record, + vendorId: string | undefined, +): MutableAiGenAiValues => { const integration = resolveAiIntegration(vendorId) - // SAFETY: the catalog correlates each field with its value type, but a loop // over the field union cannot carry that correlation. `decodeAttribute` is // driven by the same catalog entry as the field it is written under, so the @@ -377,11 +373,38 @@ export const mapAiSpan = (row: AiSessionSpansOutput): AiAgentSpan => { break } } - integration.refine?.(genAi, { row, attributes, read }) + integration.refine?.(genAi, { attributes, read }) // The agent the ingest gateway named, which the list and its facets show: it // reads names no dialect key carries (OpenAI Agents' graph node). const stampedAgent = readAttribute(attributes, MAPLE_AI_STAMP_ATTRS.agentName) if (stampedAgent !== undefined) genAi.agentName = stampedAgent + return genAi +} + +/** + * What a tool call was called with and what came back, decoded exactly as the + * session page decodes the span — so a view that reads one tool span on its + * own shows the payload the transcript shows: an OpenInference span's real + * `input.value` rather than the parameter schema its GenAI dual-write copied + * into `gen_ai.tool.call.arguments`, a LangChain `ToolMessage` unwrapped to + * its content. `undefined` where the span captured none. + */ +export const aiToolCallPayload = ( + attributes: Record, +): { readonly arguments: unknown; readonly result: unknown } => { + const genAi = decodeGenAi(attributes, readAttribute(attributes, MAPLE_AI_VENDOR_ID_ATTR)) + return { arguments: genAi.toolCallArguments, result: genAi.toolCallResult } +} + +export const mapAiSpan = (row: AiSessionSpansOutput): AiAgentSpan => { + // Span attributes only, envelope and source keys alike. The gateway strips + // `maple_ai.*` from span attributes before stamping its own verdict, so a + // span-level value is authoritative — and it does not touch resource + // attributes, where one forged `gen_ai.*` or `maple_ai.*` key would mark + // every span in the service as an AI span. + const attributes = row.spanAttributes + const vendorId = readAttribute(attributes, MAPLE_AI_VENDOR_ID_ATTR) + const genAi = decodeGenAi(attributes, vendorId) const promptVariables = collectPromptVariables(attributes) const sessionId = readAttribute(attributes, MAPLE_AI_SESSION_ID_ATTR) diff --git a/packages/query-engine-integrations/src/ai/ai-sessions.ts b/packages/query-engine-integrations/src/ai/ai-sessions.ts index 35a7298482..4d84244d37 100644 --- a/packages/query-engine-integrations/src/ai/ai-sessions.ts +++ b/packages/query-engine-integrations/src/ai/ai-sessions.ts @@ -1358,6 +1358,23 @@ export const aiSessionSpansRowSchema: CompiledQueryRowSchema>, +): CH.Expr> => + mapFilterKeys(attributes, (key) => + aiSpanAttributePrefixes.reduce( + (matched, prefix) => matched.or(key.like(`${prefix}%`)), + key.in_(...aiSpanAttributeKeys), + ), + ) + /** Shared by both span reads, so a session keyed by id and one keyed by trace * cannot drift apart in shape — {@link aiSessionSpansRowSchema} decodes both. */ const spanProjection = ($: ColumnAccessor) => ({ @@ -1371,17 +1388,7 @@ const spanProjection = ($: ColumnAccessor) => ( statusCode: $.StatusCode, statusMessage: $.StatusMessage, timestamp: CH.toString_($.Timestamp), - // The map cut down to what `mapAiSpan` reads. Measured on production's - // largest sessions, the whole map is dominated by keys the mapper never - // touches (`db.query.text` alone was half of one session's bytes), and - // `ResourceAttributes` — which the mapper deliberately ignores, see - // `mapAiSpan` — was another 60% on top. Neither is read any more. - spanAttributes: mapFilterKeys($.SpanAttributes, (key) => - aiSpanAttributePrefixes.reduce( - (matched, prefix) => matched.or(key.like(`${prefix}%`)), - key.in_(...aiSpanAttributeKeys), - ), - ), + spanAttributes: aiSpanAttributes($.SpanAttributes), }) /** diff --git a/packages/query-engine-integrations/src/ai/ai-tools.test.ts b/packages/query-engine-integrations/src/ai/ai-tools.test.ts index a7a3eb3fd4..f7f69cb2f5 100644 --- a/packages/query-engine-integrations/src/ai/ai-tools.test.ts +++ b/packages/query-engine-integrations/src/ai/ai-tools.test.ts @@ -8,6 +8,7 @@ import { aiToolErrorBreakdownRowSchema, aiToolErrorOccurrencesQuery, aiToolErrorOccurrencesRowSchema, + aiToolErrorPayload, aiToolErrorPayloadSlice, aiToolErrorPayloadsQuery, aiToolErrorPayloadsRowSchema, @@ -23,6 +24,8 @@ import { aiToolsTotalsQuery, AI_TOOLS_BREAKDOWN_LIMIT, AI_TOOLS_SERIES_MAX_KEYS, + AI_TOOL_ERROR_ATTRIBUTE_MAX, + AI_TOOL_ERROR_PAYLOAD_MAX, AI_TOOL_OCCURRENCES_LIMIT, type AiToolErrorCallKey, } from "./ai-tools" @@ -655,25 +658,161 @@ describe("aiToolErrorPayloadsQuery", () => { expect(compiled.sql).not.toContain("2026-08-19 23:59:59") }) - it("truncates payloads by codepoint and reports their size in bytes", () => { - // `left` counts BYTES and would cut a multi-byte codepoint in half. - expect(compiled.sql).not.toContain("left(") - expect(compiled.sql).toContain("leftUTF8(") - expect(compiled.sql).toContain("AS argumentsBytes") - expect(compiled.sql).toContain("AS resultBytes") + it("reads the attributes the session page decodes, not a payload of its own", () => { + // Which attribute holds the arguments is the integrations' call, so the + // read hands back the same projection the session span read does. + expect(compiled.sql).toContain("mapFilter((k, v) ->") + expect(compiled.sql).toContain("'input.value'") + expect(compiled.sql).toContain("'tool.parameters'") + expect(compiled.sql).toContain("AS spanAttributes") + expect(compiled.sql).not.toContain("coalesce(") + }) + + it("cuts every value in SQL and reports the true size of the ones it cut", () => { + expect(compiled.sql).toContain( + `mapApply((k, v) -> (k, leftUTF8(v, ${AI_TOOL_ERROR_ATTRIBUTE_MAX})), mapFilter(`, + ) + expect(compiled.sql).toContain( + `mapApply((k, v) -> (k, length(v)), mapFilter((k, v) -> lengthUTF8(v) > ${AI_TOOL_ERROR_ATTRIBUTE_MAX}, mapFilter(`, + ) + expect(compiled.sql).toContain("AS cutAttributeBytes") expect( decodeRows(compiled, [ { traceId: "t1", spanId: "s1", statusCode: "Error", - arguments: "{}", - argumentsBytes: "2", - result: "", - resultBytes: 0, + spanAttributes: { "input.value": "{}" }, + // UInt64 arrives quoted under FORMAT JSON. + cutAttributeBytes: { "output.value": "40000" }, }, - ])[0], - ).toMatchObject({ argumentsBytes: 2, resultBytes: 0 }) + ])[0]?.cutAttributeBytes, + ).toEqual({ "output.value": 40_000 }) + }) +}) + +describe("aiToolErrorPayload", () => { + const payload = ( + spanAttributes: Record, + cutAttributeBytes: Record = {}, + ) => + aiToolErrorPayload({ + traceId: "t1", + spanId: "s1", + statusCode: "Error", + spanAttributes, + cutAttributeBytes, + }) + + // `update_seat` as the OpenAI Agents SDK's Python OpenInference instrumentor + // emitted it in production: the GenAI dual-write put the parameter schema in + // `gen_ai.tool.call.arguments` (session `verify-oa-py-nofix1-1`) where the + // baseline run (`verify-oa-py-base-1`) had the real arguments. + const UPDATE_SEAT_SCHEMA = + '{"properties": {"confirmation_number": {"description": "The confirmation number for the flight.", "title": "Confirmation Number", "type": "string"}, "new_seat": {"description": "The new seat to update to.", "title": "New Seat", "type": "string"}}, "required": ["confirmation_number", "new_seat"], "title": "update_seat_args", "type": "object", "additionalProperties": false}' + const UPDATE_SEAT_ARGS = '{"confirmation_number":"ABC123","new_seat":"14C"}' + const UPDATE_SEAT_ERROR = + "An error occurred while running the tool. Please try again. Error: Seat 14C is temporarily locked, retry once" + const updateSeat = (dualWrittenArguments: string) => ({ + "maple_ai.vendor.id": "openai_agents_sdk", + "openinference.span.kind": "TOOL", + "tool.name": "update_seat", + "gen_ai.tool.name": "update_seat", + "tool.parameters": UPDATE_SEAT_SCHEMA, + "gen_ai.tool.call.arguments": dualWrittenArguments, + "input.value": UPDATE_SEAT_ARGS, + "output.value": UPDATE_SEAT_ERROR, + "gen_ai.tool.call.result": UPDATE_SEAT_ERROR, + }) + + it("shows an OpenInference tool's arguments where the dual-write copied its schema", () => { + expect(payload(updateSeat(UPDATE_SEAT_SCHEMA))).toMatchObject({ + arguments: UPDATE_SEAT_ARGS, + argumentsBytes: UPDATE_SEAT_ARGS.length, + result: UPDATE_SEAT_ERROR, + }) + // The run whose dual-write carried the arguments reads the same. + expect(payload(updateSeat(UPDATE_SEAT_ARGS)).arguments).toBe(UPDATE_SEAT_ARGS) + }) + + // LlamaIndex `get_weather`: with the GenAI dual-write on (`verify-li-sem1-1`) + // the arguments slot holds the schema, with it off (`verify-li-sem0-1`) it is + // empty, and `input.value` has the call either way. + const WEATHER_SCHEMA = + '{"properties": {"city": {"title": "City", "type": "string"}, "unit": {"title": "Unit", "type": "string"}}, "required": ["city", "unit"], "type": "object"}' + const WEATHER_OUTPUT = + '{"blocks":[{"text":"Weather in Berlin: 18 degrees celsius, light rain."}],"tool_name":"get_weather","raw_input":{"args":[],"kwargs":{"city":"Berlin","unit":"celsius"}},"raw_output":"Weather in Berlin: 18 degrees celsius, light rain.","is_error":false}' + const getWeather = (dualWrite: Record) => ({ + "maple_ai.vendor.id": "llamaindex", + "openinference.span.kind": "TOOL", + "tool.name": "get_weather", + "tool.parameters": WEATHER_SCHEMA, + "input.value": '{"kwargs": {"city": "Berlin", "unit": "celsius"}}', + "output.value": WEATHER_OUTPUT, + ...dualWrite, + }) + + it("reads a LlamaIndex tool's input whether or not the dual-write ran", () => { + const withDualWrite = payload( + getWeather({ + "gen_ai.tool.call.arguments": WEATHER_SCHEMA, + "gen_ai.tool.call.result": WEATHER_OUTPUT, + }), + ) + const without = payload(getWeather({})) + // Re-serialised from the decoded value, as the session page renders it. + const args = '{"kwargs":{"city":"Berlin","unit":"celsius"}}' + expect(withDualWrite.arguments).toBe(args) + expect(without.arguments).toBe(args) + expect(without.result).toBe(WEATHER_OUTPUT) + }) + + it("unwraps a LangChain ToolMessage result to what the tool returned", () => { + expect( + payload({ + "gen_ai.tool.name": "get_weather", + "gen_ai.tool.call.result": + '{"type": "tool", "data": {"content": "{\\"city\\": \\"Berlin\\"}", "type": "tool", "name": "get_weather", "tool_call_id": "call_1", "status": "success"}}', + }).result, + ).toBe('{"city": "Berlin"}') + }) + + it("truncates by codepoint and reports the full size in bytes", () => { + const long = `"${"é".repeat(AI_TOOL_ERROR_PAYLOAD_MAX + 10)}"` + const cut = payload({ "gen_ai.tool.call.arguments": "{}", "gen_ai.tool.call.result": long }) + expect(cut).toMatchObject({ arguments: "{}", argumentsBytes: 2 }) + // A JSON string decodes to its text; every `é` is two bytes. + expect(Array.from(cut.result)).toHaveLength(AI_TOOL_ERROR_PAYLOAD_MAX) + expect(cut.resultBytes).toBe((AI_TOOL_ERROR_PAYLOAD_MAX + 10) * 2) + expect(payload({})).toMatchObject({ arguments: "", argumentsBytes: 0, result: "", resultBytes: 0 }) + }) + + it("shows a value the read cut as its raw text, at its true size", () => { + // As the read returns a 40 KB JSON result: cut mid-document, so it no + // longer parses. + const whole = JSON.stringify({ rows: "x".repeat(40_000) }) + const cut = whole.slice(0, AI_TOOL_ERROR_ATTRIBUTE_MAX) + const shown = payload( + { "gen_ai.tool.call.arguments": "{}", "gen_ai.tool.call.result": cut }, + { "gen_ai.tool.call.result": whole.length }, + ) + expect(shown).toMatchObject({ arguments: "{}", argumentsBytes: 2, resultBytes: whole.length }) + expect(shown.result).toBe(cut.slice(0, AI_TOOL_ERROR_PAYLOAD_MAX)) + }) + + it("still swaps an OpenInference schema for the arguments when both are cut", () => { + // The schema and its dual-written copy are cut to the same prefix, so + // they still compare equal and `input.value` still wins. + const schema = JSON.stringify({ + properties: { seat: { description: "d".repeat(AI_TOOL_ERROR_ATTRIBUTE_MAX) } }, + }) + const cut = schema.slice(0, AI_TOOL_ERROR_ATTRIBUTE_MAX) + expect( + payload( + { ...updateSeat(cut), "tool.parameters": cut }, + { "tool.parameters": schema.length, "gen_ai.tool.call.arguments": schema.length }, + ), + ).toMatchObject({ arguments: UPDATE_SEAT_ARGS, argumentsBytes: UPDATE_SEAT_ARGS.length }) }) }) diff --git a/packages/query-engine-integrations/src/ai/ai-tools.ts b/packages/query-engine-integrations/src/ai/ai-tools.ts index a177789aaf..dfbc3b656e 100644 --- a/packages/query-engine-integrations/src/ai/ai-tools.ts +++ b/packages/query-engine-integrations/src/ai/ai-tools.ts @@ -67,16 +67,22 @@ import * as CH from "@maple-dev/effect-clickhouse/expr" import * as T from "@maple-dev/effect-clickhouse/types" +import { compile } from "@maple-dev/effect-clickhouse/sql" import { from, fromQuery, inSubquery, param, unionAll, type CHUnionQuery } from "@maple-dev/effect-clickhouse" import { AI_TOOLS_BREAKDOWN_MAX, AI_TOOLS_OTHER_SERIES_KEY, type AiToolsPeriod } from "@maple/domain/http" -import type { AiGenAiField } from "@maple/domain/gen-ai" import { Array as Arr, Schema } from "effect" import type { CompiledQueryRowSchema } from "@maple-dev/effect-clickhouse" import { AiTraceIndex, TraceDetailSpans } from "@maple/query-engine/ch/tables" -import { finiteOrZero, isoBucket, leftUTF8 } from "@maple/query-engine/ch/format" +import { finiteOrZero, isoBucket } from "@maple/query-engine/ch/format" import { CHNumber } from "@maple/query-engine/ch/schema" -import { aiFieldSourceKeys } from "./ai-integrations" -import { SESSION_ORDER_SENTINEL, isSessionTraceCond, orderTuple, sessionKey } from "./ai-sessions" +import { aiToolCallPayload } from "./ai-integrations" +import { + SESSION_ORDER_SENTINEL, + aiSpanAttributes, + isSessionTraceCond, + orderTuple, + sessionKey, +} from "./ai-sessions" /** * The page's selection, as every read here takes it. @@ -627,6 +633,14 @@ export function aiToolsBreakdownsQuery(opts: AiToolsFilterOpts = {}) { * reports the true size beside it. */ export const AI_TOOL_ERROR_PAYLOAD_MAX = 4_000 +/** How much of each span attribute the payload read carries, in characters. A + * span's attributes hold whole files and message histories, and a hundred + * samples of them would leave the warehouse only to be cut to + * {@link AI_TOOL_ERROR_PAYLOAD_MAX}. Four times that cut, so the JSON behind a + * payload the modal can show whole still parses; a longer value arrives cut, + * decodes as its raw text, and is cut again for display. */ +export const AI_TOOL_ERROR_ATTRIBUTE_MAX = 16_384 + /** Error groups one breakdown returns, most failed calls first. */ export const AI_TOOL_ERRORS_LIMIT = 50 @@ -931,25 +945,23 @@ export interface AiToolErrorCallKey { readonly spanId: string } -export interface AiToolErrorPayloadsOutput { +export interface AiToolErrorPayloadsRow { readonly traceId: string readonly spanId: string readonly statusCode: string - /** Truncated to {@link AI_TOOL_ERROR_PAYLOAD_MAX}; `*Bytes` is the true size. */ - readonly arguments: string - readonly argumentsBytes: number - readonly result: string - readonly resultBytes: number + /** Each value cut to {@link AI_TOOL_ERROR_ATTRIBUTE_MAX} characters. */ + readonly spanAttributes: Record + /** The true size in bytes of every value that was cut, by key. */ + readonly cutAttributeBytes: Record } -export const aiToolErrorPayloadsRowSchema: CompiledQueryRowSchema = Schema.Struct({ +export const aiToolErrorPayloadsRowSchema: CompiledQueryRowSchema = Schema.Struct({ traceId: Schema.String, spanId: Schema.String, statusCode: Schema.String, - arguments: Schema.String, - argumentsBytes: CHNumber, - result: Schema.String, - resultBytes: CHNumber, + // A Map column arrives as a JSON object under FORMAT JSON. + spanAttributes: Schema.Record(Schema.String, Schema.String), + cutAttributeBytes: Schema.Record(Schema.String, CHNumber), }) /** A span row's `(TraceId, SpanId)`, for the payload read's tuple `IN`. Raw @@ -957,16 +969,19 @@ export const aiToolErrorPayloadsRowSchema: CompiledQueryRowSchema> -} +/** The attribute map with every value cut to {@link AI_TOOL_ERROR_ATTRIBUTE_MAX} + * characters. Raw because the DSL has no `mapApply`. */ +const cutAttributeValues = (attributes: CH.Expr>): CH.Expr> => + CH.rawExpr( + `mapApply((k, v) -> (k, leftUTF8(v, ${AI_TOOL_ERROR_ATTRIBUTE_MAX})), ${compile(attributes.toFragment())})`, + T.map(T.string, T.string), + ) -/** `coalesce(nullIf(a, ''), …, '')` — the first source key of a field that has - * a value, across every vendor dialect the integrations declare. */ -const spanField = ($: SpanAccessor, field: AiGenAiField): CH.Expr => - CH.coalesce( - ...aiFieldSourceKeys(field).map((key) => CH.nullIf(CH.mapGet($.SpanAttributes, key), "")), - CH.lit(""), +/** The byte size of each value {@link cutAttributeValues} cuts, by key. */ +const cutAttributeBytes = (attributes: CH.Expr>): CH.Expr> => + CH.rawExpr( + `mapApply((k, v) -> (k, length(v)), mapFilter((k, v) -> lengthUTF8(v) > ${AI_TOOL_ERROR_ATTRIBUTE_MAX}, ${compile(attributes.toFragment())}))`, + T.map(T.string, T.uint64), ) /** @@ -980,24 +995,25 @@ const spanField = ($: SpanAccessor, field: AiGenAiField): CH.Expr => * are exact and unpadded: the index copies `Timestamp` from the span verbatim, * so the call is inside its own bounds by construction, and a pad would only buy * partitions. + * + * It returns the span's attributes as the session page reads them, not the + * payloads: which attribute holds the arguments is a per-vendor decision the + * integrations make in TypeScript ({@link aiToolErrorPayload}), so the modal + * and the session page decode a tool span one way: an OpenInference tool shows + * its `input.value`, not the parameter schema its GenAI dual-write copied into + * `gen_ai.tool.call.arguments`. Each value is cut to + * {@link AI_TOOL_ERROR_ATTRIBUTE_MAX} in SQL, beside the true size of the ones + * that were, so a sample costs kilobytes on the wire rather than its whole span. */ export function aiToolErrorPayloadsQuery(calls: Arr.NonEmptyReadonlyArray) { return from(TraceDetailSpans) - .select(($) => { - const args = spanField($, "toolCallArguments") - const result = spanField($, "toolCallResult") - return { - traceId: $.TraceId, - spanId: $.SpanId, - statusCode: $.StatusCode, - // Characters, not bytes: `left` cuts mid-codepoint on any payload - // holding one. `length` stays byte-based — it reports a size. - arguments: leftUTF8(args, CH.lit(AI_TOOL_ERROR_PAYLOAD_MAX)), - argumentsBytes: CH.length_(args), - result: leftUTF8(result, CH.lit(AI_TOOL_ERROR_PAYLOAD_MAX)), - resultBytes: CH.length_(result), - } - }) + .select(($) => ({ + traceId: $.TraceId, + spanId: $.SpanId, + statusCode: $.StatusCode, + spanAttributes: cutAttributeValues(aiSpanAttributes($.SpanAttributes)), + cutAttributeBytes: cutAttributeBytes(aiSpanAttributes($.SpanAttributes)), + })) .where(($) => [ $.OrgId.eq(param.string("orgId")), $.Timestamp.gte(param.dateTimeString("sliceStart")), @@ -1012,6 +1028,52 @@ export function aiToolErrorPayloadsQuery(calls: Arr.NonEmptyReadonlyArray + value === undefined ? "" : typeof value === "string" ? value : JSON.stringify(value) + +/** Characters, not UTF-16 units: a cut never splits a codepoint. */ +const truncatePayload = (text: string): string => + text.length <= AI_TOOL_ERROR_PAYLOAD_MAX + ? text + : Array.from(text).slice(0, AI_TOOL_ERROR_PAYLOAD_MAX).join("") + +/** One {@link aiToolErrorPayloadsQuery} row as the modal shows it: the payloads + * the session page decodes for the same span, cut to + * {@link AI_TOOL_ERROR_PAYLOAD_MAX} characters beside their size in bytes. A + * payload decoded from a value the read cut no longer parses, so it is that + * value's text, and its size is the one the read reported for it. */ +export const aiToolErrorPayload = (row: AiToolErrorPayloadsRow): AiToolErrorPayloadsOutput => { + const payload = aiToolCallPayload(row.spanAttributes) + const args = payloadText(payload.arguments) + const result = payloadText(payload.result) + const bytes = (text: string): number => + Object.entries(row.cutAttributeBytes).find(([key]) => row.spanAttributes[key] === text)?.[1] ?? + utf8.encode(text).length + return { + traceId: row.traceId, + spanId: row.spanId, + statusCode: row.statusCode, + arguments: truncatePayload(args), + argumentsBytes: bytes(args), + result: truncatePayload(result), + resultBytes: bytes(result), + } +} + /** The bounds one {@link aiToolErrorPayloadsQuery} is read over: the earliest * and latest of the occurrences it is for, exactly as the index reported them. * No occurrences have no extent, which is why both take a non-empty list. */ diff --git a/packages/query-engine-integrations/src/ai/index.ts b/packages/query-engine-integrations/src/ai/index.ts index 70396d1425..3e026b040b 100644 --- a/packages/query-engine-integrations/src/ai/index.ts +++ b/packages/query-engine-integrations/src/ai/index.ts @@ -50,6 +50,7 @@ export { aiToolErrorBreakdownRowSchema, aiToolErrorOccurrencesQuery, aiToolErrorOccurrencesRowSchema, + aiToolErrorPayload, aiToolErrorPayloadSlice, aiToolErrorPayloadsQuery, aiToolErrorPayloadsRowSchema, @@ -81,6 +82,7 @@ export { type AiToolErrorOccurrencesOutput, type AiToolErrorPayloadSlice, type AiToolErrorPayloadsOutput, + type AiToolErrorPayloadsRow, type AiToolErrorSessionsOutput, type AiToolErrorsOpts, type AiToolErrorsOutput,