From 52ee0224a8d1c9861624794c6c5ed40edfa53459 Mon Sep 17 00:00:00 2001 From: JeremyFunk Date: Tue, 29 Sep 2026 20:14:34 +0200 Subject: [PATCH 1/2] fix(agent-sessions): decode tool-error payloads the way the session page does The tool-errors samples (web modal, get_agent_tool_error) read gen_ai.tool.call.arguments raw in SQL, so an OpenInference tool whose GenAI dual-write copied its parameter schema there showed the schema, and a tool without the dual-write showed nothing. The payload read now returns the span's projected attributes and runs them through the same integration decode as the session page (schema swapped for input.value, LangChain ToolMessage unwrapped), then truncates and sizes the result in TypeScript. --- .../mcp/tools/__tests__/agent-tools.test.ts | 26 ++-- .../routes/internal/ai-sessions.http.test.ts | 11 +- .../services/ai-sessions/ai-session-reads.ts | 6 +- .../warehouse/ai-tools.clickhouse.e2e.test.ts | 32 +++-- .../src/__sql_baseline__/integrations.sql | 5 +- .../src/ai/ai-integrations.ts | 45 +++++-- .../src/ai/ai-sessions.ts | 29 +++-- .../src/ai/ai-tools.test.ts | 115 +++++++++++++++--- .../src/ai/ai-tools.ts | 104 +++++++++------- .../query-engine-integrations/src/ai/index.ts | 2 + 10 files changed, 256 insertions(+), 119 deletions(-) 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 cd626fa774..bad9e2e71f 100644 --- a/packages/query-engine-integrations/src/__sql_baseline__/integrations.sql +++ b/packages/query-engine-integrations/src/__sql_baseline__/integrations.sql @@ -1152,10 +1152,7 @@ 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 + 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 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 eb689529ef..652a4def03 100644 --- a/packages/query-engine-integrations/src/ai/ai-sessions.ts +++ b/packages/query-engine-integrations/src/ai/ai-sessions.ts @@ -1337,6 +1337,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) => ({ @@ -1350,17 +1367,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 53ec155cec..4c7d42fbb8 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,7 @@ import { aiToolsTotalsQuery, AI_TOOLS_BREAKDOWN_LIMIT, AI_TOOLS_SERIES_MAX_KEYS, + AI_TOOL_ERROR_PAYLOAD_MAX, AI_TOOL_OCCURRENCES_LIMIT, type AiToolErrorCallKey, } from "./ai-tools" @@ -644,25 +646,102 @@ 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(") + }) +}) + +describe("aiToolErrorPayload", () => { + const payload = (spanAttributes: Record) => + aiToolErrorPayload({ traceId: "t1", spanId: "s1", statusCode: "Error", spanAttributes }) + + // `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( - decodeRows(compiled, [ - { - traceId: "t1", - spanId: "s1", - statusCode: "Error", - arguments: "{}", - argumentsBytes: "2", - result: "", - resultBytes: 0, - }, - ])[0], - ).toMatchObject({ argumentsBytes: 2, resultBytes: 0 }) + 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 }) }) }) diff --git a/packages/query-engine-integrations/src/ai/ai-tools.ts b/packages/query-engine-integrations/src/ai/ai-tools.ts index 38b7a61c4c..4c4078a8e5 100644 --- a/packages/query-engine-integrations/src/ai/ai-tools.ts +++ b/packages/query-engine-integrations/src/ai/ai-tools.ts @@ -69,14 +69,13 @@ import * as CH from "@maple-dev/effect-clickhouse/expr" import * as T from "@maple-dev/effect-clickhouse/types" 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, orderTuple, sessionKey } from "./ai-sessions" +import { aiToolCallPayload } from "./ai-integrations" +import { SESSION_ORDER_SENTINEL, aiSpanAttributes, orderTuple, sessionKey } from "./ai-sessions" /** * The page's selection, as every read here takes it. @@ -926,25 +925,19 @@ 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 + readonly spanAttributes: 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), }) /** A span row's `(TraceId, SpanId)`, for the payload read's tuple `IN`. Raw @@ -952,18 +945,6 @@ export const aiToolErrorPayloadsRowSchema: CompiledQueryRowSchema> -} - -/** `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 modal's right pane, step two: what each of those calls was called with and * what came back — the only facts on this page that live on the raw span. @@ -975,24 +956,22 @@ 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`. */ 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: aiSpanAttributes($.SpanAttributes), + })) .where(($) => [ $.OrgId.eq(param.string("orgId")), $.Timestamp.gte(param.dateTimeString("sliceStart")), @@ -1007,6 +986,47 @@ 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. */ +export const aiToolErrorPayload = (row: AiToolErrorPayloadsRow): AiToolErrorPayloadsOutput => { + const payload = aiToolCallPayload(row.spanAttributes) + const args = payloadText(payload.arguments) + const result = payloadText(payload.result) + return { + traceId: row.traceId, + spanId: row.spanId, + statusCode: row.statusCode, + arguments: truncatePayload(args), + argumentsBytes: utf8.encode(args).length, + result: truncatePayload(result), + resultBytes: utf8.encode(result).length, + } +} + /** 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, From 7525624f5957871946da2f6917567e90af5d71be Mon Sep 17 00:00:00 2001 From: JeremyFunk Date: Wed, 30 Sep 2026 02:42:19 +0200 Subject: [PATCH 2/2] fix(agent-tools): cut the payload read's attribute values in SQL --- .../src/__sql_baseline__/integrations.sql | 3 +- .../src/ai/ai-tools.test.ts | 64 ++++++++++++++++++- .../src/ai/ai-tools.ts | 46 +++++++++++-- 3 files changed, 105 insertions(+), 8 deletions(-) diff --git a/packages/query-engine-integrations/src/__sql_baseline__/integrations.sql b/packages/query-engine-integrations/src/__sql_baseline__/integrations.sql index 941d9e7591..eeee5326b7 100644 --- a/packages/query-engine-integrations/src/__sql_baseline__/integrations.sql +++ b/packages/query-engine-integrations/src/__sql_baseline__/integrations.sql @@ -1170,7 +1170,8 @@ SELECT TraceId AS traceId, SpanId AS spanId, StatusCode AS statusCode, - 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, 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-tools.test.ts b/packages/query-engine-integrations/src/ai/ai-tools.test.ts index 9fa59f29f2..f7f69cb2f5 100644 --- a/packages/query-engine-integrations/src/ai/ai-tools.test.ts +++ b/packages/query-engine-integrations/src/ai/ai-tools.test.ts @@ -24,6 +24,7 @@ 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, @@ -666,11 +667,42 @@ describe("aiToolErrorPayloadsQuery", () => { 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", + spanAttributes: { "input.value": "{}" }, + // UInt64 arrives quoted under FORMAT JSON. + cutAttributeBytes: { "output.value": "40000" }, + }, + ])[0]?.cutAttributeBytes, + ).toEqual({ "output.value": 40_000 }) + }) }) describe("aiToolErrorPayload", () => { - const payload = (spanAttributes: Record) => - aiToolErrorPayload({ traceId: "t1", spanId: "s1", statusCode: "Error", spanAttributes }) + 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 @@ -754,6 +786,34 @@ describe("aiToolErrorPayload", () => { 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 }) + }) }) describe("aiToolDescriptionQuery", () => { diff --git a/packages/query-engine-integrations/src/ai/ai-tools.ts b/packages/query-engine-integrations/src/ai/ai-tools.ts index 1b8ae52793..dfbc3b656e 100644 --- a/packages/query-engine-integrations/src/ai/ai-tools.ts +++ b/packages/query-engine-integrations/src/ai/ai-tools.ts @@ -67,6 +67,7 @@ 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 { Array as Arr, Schema } from "effect" @@ -632,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 @@ -940,7 +949,10 @@ export interface AiToolErrorPayloadsRow { readonly traceId: string readonly spanId: string readonly statusCode: string + /** 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({ @@ -949,6 +961,7 @@ export const aiToolErrorPayloadsRowSchema: CompiledQueryRowSchema>): CH.Expr> => + CH.rawExpr( + `mapApply((k, v) -> (k, leftUTF8(v, ${AI_TOOL_ERROR_ATTRIBUTE_MAX})), ${compile(attributes.toFragment())})`, + T.map(T.string, T.string), + ) + +/** 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), + ) + /** * The modal's right pane, step two: what each of those calls was called with and * what came back — the only facts on this page that live on the raw span. @@ -973,7 +1001,9 @@ const traceSpanKey = CH.rawExpr("(trace_detail_spans.TraceId, trace_detail_spans * 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`. + * `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) @@ -981,7 +1011,8 @@ export function aiToolErrorPayloadsQuery(calls: Arr.NonEmptyReadonlyArray [ $.OrgId.eq(param.string("orgId")), @@ -1022,19 +1053,24 @@ const truncatePayload = (text: string): string => /** 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. */ + * {@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: utf8.encode(args).length, + argumentsBytes: bytes(args), result: truncatePayload(result), - resultBytes: utf8.encode(result).length, + resultBytes: bytes(result), } }