From 4fc5960be4372b40f13a499727a3eb355e042392 Mon Sep 17 00:00:00 2001 From: JPeer264 Date: Fri, 2 Oct 2026 13:21:09 +0200 Subject: [PATCH] feat(server-utils): Add instrumentation for pi-durable pi-durable has no telemetry hooks, so `piDurableIntegration` wraps the `models` and `registry` passed to `Harness.open()`. Each run becomes its own `invoke_agent` trace with `chat` and `execute_tool` children, and a subagent run nests under the tool call that started it. Messages are mapped to the gen_ai conventions, and the system prompt and tools are read from pi-durable's positional system messages. Tool spans record the result the model receives. Failures pi-durable only reports, and task phases that throw, are captured. Throws of the built-in coding tools are not, because they report expected results to the model. On Cloudflare it also covers pi-durable in a Durable Object through the Agents SDK `PiHarness`, including requests that pi-ai sends through the Workers AI binding. Co-Authored-By: Claude Opus 5.5 Co-Authored-By: Claude Fable 5.1 --- .../node-suites/excludes.ts | 1 + .../suites/tracing/pi-durable/instrument.mjs | 12 + .../pi-durable/scenario-compaction.mjs | 38 + .../pi-durable/scenario-concurrent.mjs | 61 ++ .../pi-durable/scenario-interrupted.mjs | 114 +++ .../tracing/pi-durable/scenario-reports.mjs | 70 ++ .../tracing/pi-durable/scenario-subagent.mjs | 121 +++ .../tracing/pi-durable/scenario-tools.mjs | 110 +++ .../suites/tracing/pi-durable/scenario.mjs | 63 ++ .../suites/tracing/pi-durable/test.ts | 693 +++++++++++++++ packages/astro/src/index.server.ts | 1 + packages/aws-serverless/src/index.ts | 1 + packages/bun/src/index.ts | 1 + packages/deno/src/index.ts | 1 + packages/elysia/src/index.ts | 1 + packages/google-cloud-serverless/src/index.ts | 1 + packages/node/src/index.ts | 1 + .../server-utils/src/ai/pi-ai/messages.ts | 74 ++ .../server-utils/src/ai/pi-ai/providers.ts | 11 +- .../src/ai/pi-durable/constants.ts | 23 + .../server-utils/src/ai/pi-durable/index.ts | 227 +++++ .../server-utils/src/ai/pi-durable/models.ts | 201 +++++ .../server-utils/src/ai/pi-durable/runs.ts | 217 +++++ .../server-utils/src/ai/pi-durable/tools.ts | 250 ++++++ .../server-utils/src/ai/pi-durable/types.ts | 141 +++ .../server-utils/src/ai/pi-durable/utils.ts | 25 + packages/server-utils/src/index.ts | 1 + .../server-utils/src/integrations/index.ts | 2 + .../src/integrations/pi-durable.ts | 76 ++ .../server-utils/src/orchestrion/channels.ts | 2 + .../config/channel-integration-definitions.ts | 1 + .../src/orchestrion/config/index.ts | 2 + .../src/orchestrion/config/pi-durable.ts | 35 + .../test/ai/lib/tracing/pi-durable.test.ts | 840 ++++++++++++++++++ .../test/ai/lib/utils/pi-ai-messages.test.ts | 24 + 35 files changed, 3438 insertions(+), 4 deletions(-) create mode 100644 dev-packages/node-integration-tests/suites/tracing/pi-durable/instrument.mjs create mode 100644 dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-compaction.mjs create mode 100644 dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-concurrent.mjs create mode 100644 dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-interrupted.mjs create mode 100644 dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-reports.mjs create mode 100644 dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-subagent.mjs create mode 100644 dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-tools.mjs create mode 100644 dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario.mjs create mode 100644 dev-packages/node-integration-tests/suites/tracing/pi-durable/test.ts create mode 100644 packages/server-utils/src/ai/pi-durable/constants.ts create mode 100644 packages/server-utils/src/ai/pi-durable/index.ts create mode 100644 packages/server-utils/src/ai/pi-durable/models.ts create mode 100644 packages/server-utils/src/ai/pi-durable/runs.ts create mode 100644 packages/server-utils/src/ai/pi-durable/tools.ts create mode 100644 packages/server-utils/src/ai/pi-durable/types.ts create mode 100644 packages/server-utils/src/ai/pi-durable/utils.ts create mode 100644 packages/server-utils/src/integrations/pi-durable.ts create mode 100644 packages/server-utils/src/orchestrion/config/pi-durable.ts create mode 100644 packages/server-utils/test/ai/lib/tracing/pi-durable.test.ts diff --git a/dev-packages/bun-integration-tests/node-suites/excludes.ts b/dev-packages/bun-integration-tests/node-suites/excludes.ts index ba3fc6a6ef95..3b2d1257f8c2 100644 --- a/dev-packages/bun-integration-tests/node-suites/excludes.ts +++ b/dev-packages/bun-integration-tests/node-suites/excludes.ts @@ -123,6 +123,7 @@ export const NO_AUTO_INSTRUMENTATION = [ 'suites/tracing/openai/v6/test.ts', 'suites/tracing/openai/v7/test.ts', 'suites/tracing/orchestrion-lazy-registration/test.ts', + 'suites/tracing/pi-durable/test.ts', 'suites/tracing/postgres-streamed/test.ts', 'suites/tracing/postgres/test.ts', 'suites/tracing/postgresjs-streamed/test.ts', diff --git a/dev-packages/node-integration-tests/suites/tracing/pi-durable/instrument.mjs b/dev-packages/node-integration-tests/suites/tracing/pi-durable/instrument.mjs new file mode 100644 index 000000000000..f1f018ef471a --- /dev/null +++ b/dev-packages/node-integration-tests/suites/tracing/pi-durable/instrument.mjs @@ -0,0 +1,12 @@ +import * as Sentry from '@sentry/node'; +import { loggingTransport } from '@sentry-internal/node-integration-tests'; + +Sentry.init({ + dsn: 'https://public@dsn.ingest.sentry.io/1337', + release: '1.0', + tracesSampleRate: 1.0, + transport: loggingTransport, + ...(process.env.PI_DURABLE_RECORDING === 'off' + ? { dataCollection: { genAI: { inputs: false, outputs: false } } } + : {}), +}); diff --git a/dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-compaction.mjs b/dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-compaction.mjs new file mode 100644 index 000000000000..18c98f77f2ba --- /dev/null +++ b/dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-compaction.mjs @@ -0,0 +1,38 @@ +import { BACKGROUND_CONTEXT } from '@earendil-works/chord/context'; +import { createModels } from '@earendil-works/pi-ai/models'; +import { fauxAssistantMessage, fauxProvider } from '@earendil-works/pi-ai/providers/faux'; +import { createRegistry, Harness, MemoryStorage } from '@earendil-works/pi-durable'; + +const context = BACKGROUND_CONTEXT; + +// A context window so small that pi-durable compacts the conversation before it can answer the +// third question. The compaction request belongs to that run. +const faux = fauxProvider({ provider: 'faux', models: [{ id: 'tiny', contextWindow: 60 }] }); +const models = createModels(); +models.setProvider(faux.provider); +faux.setResponses([ + fauxAssistantMessage('First answer with quite a lot of detail so the transcript grows beyond the window.'), + fauxAssistantMessage('Second answer with even more detail so that the next request needs a compaction first.'), + fauxAssistantMessage('Summary of the conversation.'), + fauxAssistantMessage('Third answer.'), +]); + +const harness = await Harness.open( + new MemoryStorage(), + { + models, + registry: createRegistry(), + settings: { compaction: { enabled: true, reserveTokens: 0, backgroundTokens: 0, keepRecentTokens: 1 } }, + }, + context, +); +const root = await harness.root(context, { agent: { model: { provider: 'faux', modelId: 'tiny' } } }); +for (const content of [ + 'First question with some words in it.', + 'Second question with some more words.', + 'Third question.', +]) { + await (await root.submit({ type: 'input', content }, context)).wait(context); +} +await harness.waitForIdle(context); +await harness.close(context); diff --git a/dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-concurrent.mjs b/dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-concurrent.mjs new file mode 100644 index 000000000000..1fdf6b110bd2 --- /dev/null +++ b/dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-concurrent.mjs @@ -0,0 +1,61 @@ +import * as Sentry from '@sentry/node'; +import { BACKGROUND_CONTEXT } from '@earendil-works/chord/context'; +import { Type } from '@earendil-works/pi-ai'; +import { createModels } from '@earendil-works/pi-ai/models'; +import { fauxAssistantMessage, fauxProvider, fauxToolCall } from '@earendil-works/pi-ai/providers/faux'; +import { createRegistry, defineExtension, defineTool, Harness, MemoryStorage } from '@earendil-works/pi-durable'; + +const context = BACKGROUND_CONTEXT; + +// Two conversations run at the same time, each submitted from its own request with its own scope +// data. Each response is derived from the transcript, and the delays make the two runs interleave. +const delay = ms => new Promise(resolve => setTimeout(resolve, ms)); +const respond = async transcript => { + const lastUser = [...transcript.messages].reverse().find(message => message.role === 'user'); + const prompt = + typeof lastUser.content === 'string' ? lastUser.content : lastUser.content.map(part => part.text).join(''); + const afterTool = transcript.messages.some(message => message.role === 'toolResult'); + await delay(prompt.includes('A') ? 30 : 10); + return afterTool + ? fauxAssistantMessage(`Answer for ${prompt}`) + : fauxAssistantMessage(fauxToolCall('work', { tag: prompt }, { id: `call_${prompt.at(-1)}` }), { + stopReason: 'toolUse', + }); +}; +const faux = fauxProvider({ provider: 'faux', models: [{ id: 'faux-model' }] }); +const models = createModels(); +models.setProvider(faux.provider); +faux.setResponses([respond, respond, respond, respond]); + +const registry = createRegistry(); +registry.install( + defineExtension({ + name: 'app', + tools: [ + defineTool({ + name: 'work', + description: 'Works for a while, then fails.', + parameters: Type.Object({ tag: Type.String() }), + execute: async args => { + await delay(args.tag.includes('A') ? 10 : 30); + throw new Error(`work failed for ${args.tag}`); + }, + }), + ], + }), +); + +const harness = await Harness.open(new MemoryStorage(), { models, registry }, context); +const model = { provider: 'faux', modelId: 'faux-model' }; +const a = await harness.createConversation({ ownership: { kind: 'ownerless' }, agent: { model } }, context); +const b = await harness.createConversation({ ownership: { kind: 'ownerless' }, agent: { model } }, context); + +// Each request has its own isolation scope, as it would in a server, with data of its own. +const ask = (conversation, request) => + Sentry.withIsolationScope(async isolationScope => { + isolationScope.setTag('request', request); + isolationScope.setConversationId(`scope-of-${request}`); + await (await conversation.submit({ type: 'input', content: request }, context)).wait(context); + }); +await Promise.all([ask(a, 'request A'), ask(b, 'request B')]); +await harness.close(context); diff --git a/dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-interrupted.mjs b/dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-interrupted.mjs new file mode 100644 index 000000000000..f0efbea32483 --- /dev/null +++ b/dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-interrupted.mjs @@ -0,0 +1,114 @@ +import * as Sentry from '@sentry/node'; +import { BACKGROUND_CONTEXT } from '@earendil-works/chord/context'; +import { Type } from '@earendil-works/pi-ai'; +import { createModels } from '@earendil-works/pi-ai/models'; +import { fauxAssistantMessage, fauxProvider, fauxToolCall } from '@earendil-works/pi-ai/providers/faux'; +import { createRegistry, defineExtension, defineTool, Harness, MemoryStorage } from '@earendil-works/pi-durable'; + +const context = BACKGROUND_CONTEXT; + +// Six runs that do not end with an answer: +// 1. the model request fails (`stopReason` error), with retries off; +// 2. the provider throws, which faults the generation task; +// 3. the run is aborted while the `slow` tool waits; +// 4. the run is aborted while the model request is in flight; +// 5. the Harness closes while the `slow` tool waits; +// 6. in a second Harness, the model request fails and the Harness closes right after `wait()`. +const started = new Map(); +const toolStarted = job => new Promise(resolve => started.set(job, resolve)); + +const faux = fauxProvider({ provider: 'faux', models: [{ id: 'faux-model' }] }); +const realModels = createModels(); +realModels.setProvider(faux.provider); +// The faux provider turns a throwing response into an error message, so the throw of a broken +// provider is staged in front of pi-ai: the request after `breakNextRequest()` throws synchronously. +let brokenRequests = 0; +const breakNextRequest = () => brokenRequests++; +const models = new Proxy(realModels, { + get(target, property) { + const value = Reflect.get(target, property, target); + if (property === 'streamSimple') { + return (...args) => { + if (brokenRequests > 0) { + brokenRequests--; + throw new Error('provider broke'); + } + return value.apply(target, args); + }; + } + return typeof value === 'function' ? value.bind(target) : value; + }, +}); +faux.setResponses([ + fauxAssistantMessage('', { stopReason: 'error', errorMessage: 'provider exploded' }), + fauxAssistantMessage(fauxToolCall('slow', { job: 'abort' }, { id: 'call_abort' }), { stopReason: 'toolUse' }), + async (_transcript, options) => { + started.get('request')(); + await new Promise(resolve => options.signal.addEventListener('abort', resolve)); + return fauxAssistantMessage('partial', { stopReason: 'aborted' }); + }, + fauxAssistantMessage(fauxToolCall('slow', { job: 'close' }, { id: 'call_close' }), { stopReason: 'toolUse' }), + fauxAssistantMessage('', { stopReason: 'error', errorMessage: 'provider exploded again' }), +]); + +const registry = createRegistry(); +registry.install( + defineExtension({ + name: 'app', + tools: [ + defineTool({ + name: 'slow', + description: 'Waits until it is aborted.', + parameters: Type.Object({ job: Type.String() }), + execute: async (args, _api, toolContext) => { + started.get(args.job)(); + await new Promise((_, reject) => + toolContext.abortSignal.addEventListener('abort', () => reject(new Error('tool aborted'))), + ); + return { content: [] }; + }, + }), + ], + }), +); + +const harness = await Harness.open( + new MemoryStorage(), + { models, registry, settings: { retry: { enabled: false } } }, + context, +); +const root = await harness.root(context, { agent: { model: { provider: 'faux', modelId: 'faux-model' } } }); + +await (await root.submit({ type: 'input', content: 'Fail please.' }, context)).wait(context); +breakNextRequest(); +await (await root.submit({ type: 'input', content: 'Break please.' }, context)).wait(context); + +const abortRun = toolStarted('abort'); +const aborted = await root.submit({ type: 'input', content: 'Run slow, then abort.' }, context); +await abortRun; +await root.abort(context); +await aborted.wait(context); + +const requestRun = toolStarted('request'); +const abortedRequest = await root.submit({ type: 'input', content: 'Answer slowly, then abort.' }, context); +await requestRun; +await root.abort(context); +await abortedRequest.wait(context); + +const closeRun = toolStarted('close'); +await root.submit({ type: 'input', content: 'Run slow, then close.' }, context); +await closeRun; +await harness.close(context); + +const second = await Harness.open( + new MemoryStorage(), + { models, registry, settings: { retry: { enabled: false } } }, + context, +); +const secondRoot = await second.root(context, { agent: { model: { provider: 'faux', modelId: 'faux-model' } } }); +await (await secondRoot.submit({ type: 'input', content: 'Fail, then close.' }, context)).wait(context); +await second.close(context); + +// Everything captured during the runs is sent before this sentinel. +await Sentry.flush(2000); +Sentry.captureMessage('pi-durable interrupted done'); diff --git a/dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-reports.mjs b/dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-reports.mjs new file mode 100644 index 000000000000..50ff47c4f156 --- /dev/null +++ b/dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-reports.mjs @@ -0,0 +1,70 @@ +import { BACKGROUND_CONTEXT } from '@earendil-works/chord/context'; +import { createModels } from '@earendil-works/pi-ai/models'; +import { fauxAssistantMessage, fauxProvider } from '@earendil-works/pi-ai/providers/faux'; +import { + createRegistry, + defineExtension, + defineTask, + GenerationTask, + Harness, + hook, + MemoryStorage, +} from '@earendil-works/pi-durable'; + +const context = BACKGROUND_CONTEXT; + +// Failures pi-durable does not propagate: a hook that throws (pi-durable reports it and keeps the +// run going) and a durable task whose phase throws (the scheduler faults the task). A compaction +// started while the conversation is idle runs outside any run. +const faux = fauxProvider({ provider: 'faux', models: [{ id: 'faux-model' }] }); +const models = createModels(); +models.setProvider(faux.provider); +faux.setResponses([ + fauxAssistantMessage('First answer with some detail.'), + fauxAssistantMessage('Second answer with more detail.'), + fauxAssistantMessage('Summary of the conversation.'), +]); + +const Charge = defineTask({ + name: 'app.charge', + version: 1, + initial: () => ({ phase: 'charge' }), + phases: { + charge: async () => { + throw new Error('card declined'); + }, + }, + abort: async (_task, runtime, taskContext) => { + await runtime.commit(() => ({ status: 'terminal', outcome: { status: 'aborted' } }), taskContext); + }, +}); + +const registry = createRegistry(); +registry.install( + defineExtension({ + name: 'app', + tasks: [Charge], + hooks: [ + hook(GenerationTask, { + afterResponse: () => { + throw new Error('afterResponse hook failed'); + }, + }), + ], + }), +); + +const harness = await Harness.open( + new MemoryStorage(), + { models, registry, settings: { compaction: { keepRecentTokens: 1 } } }, + context, +); +const root = await harness.root(context, { agent: { model: { provider: 'faux', modelId: 'faux-model' } } }); +await (await root.submit({ type: 'input', content: 'First question.' }, context)).wait(context); +await (await root.submit({ type: 'input', content: 'Second question.' }, context)).wait(context); + +const charge = await root.commit(tx => tx.createTask(Charge, {}, { ownership: { kind: 'conversation' } }), context); +await harness.waitForTask(charge, context); + +await harness.waitForTask(await root.compact(undefined, context), context); +await harness.close(context); diff --git a/dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-subagent.mjs b/dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-subagent.mjs new file mode 100644 index 000000000000..da80b508f252 --- /dev/null +++ b/dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-subagent.mjs @@ -0,0 +1,121 @@ +import { BACKGROUND_CONTEXT } from '@earendil-works/chord/context'; +import { Type } from '@earendil-works/pi-ai'; +import { createModels } from '@earendil-works/pi-ai/models'; +import { fauxAssistantMessage, fauxProvider, fauxToolCall } from '@earendil-works/pi-ai/providers/faux'; +import { + AssistantEntry, + createRegistry, + defineExtension, + defineTool, + Harness, + MemoryStorage, +} from '@earendil-works/pi-durable'; + +const context = BACKGROUND_CONTEXT; + +// Three runs of the root conversation, each with one tool round: +// - `delegate` creates a conversation it owns and waits for its answer, the subagent pattern; +// - `fork_delegate` does the same with a fork of the current conversation; +// - `batch` runs the `inner` tool of its agent itself, passing its own api on. +// Each response is derived from the transcript, so the faux model answers every conversation. +const respond = transcript => { + const lastUser = [...transcript.messages].reverse().find(message => message.role === 'user'); + const text = + typeof lastUser.content === 'string' ? lastUser.content : lastUser.content.map(part => part.text).join(''); + const afterTool = transcript.messages.at(-1).role === 'toolResult'; + if (afterTool) { + return fauxAssistantMessage(`Done: ${text}`); + } + if (text === 'Delegate the weather question.') { + return fauxAssistantMessage(fauxToolCall('delegate', { task: 'Weather in Vienna?' }, { id: 'call_delegate' }), { + stopReason: 'toolUse', + }); + } + if (text === 'Fork a helper.') { + return fauxAssistantMessage(fauxToolCall('fork_delegate', {}, { id: 'call_fork' }), { stopReason: 'toolUse' }); + } + if (text === 'Run the batch.') { + return fauxAssistantMessage(fauxToolCall('batch', {}, { id: 'call_batch' }), { stopReason: 'toolUse' }); + } + return fauxAssistantMessage(`Answer: ${text}`); +}; +const faux = fauxProvider({ provider: 'faux', models: [{ id: 'faux-model' }] }); +const models = createModels(); +models.setProvider(faux.provider); +faux.setResponses(Array.from({ length: 12 }, () => respond)); + +const answerText = async (api, settled, toolContext) => { + const answer = await api.commit(tx => tx.entry(AssistantEntry, settled.answer), toolContext); + return (answer?.model?.[0]?.content ?? []) + .filter(part => part.type === 'text') + .map(part => part.text) + .join(''); +}; + +const registry = createRegistry(); +registry.install( + defineExtension({ + name: 'app', + tools: [ + defineTool({ + name: 'delegate', + description: 'Delegates a task to a subagent.', + parameters: Type.Object({ task: Type.String() }), + execute: async (args, api, toolContext) => { + const childId = await api.commit( + async tx => (await tx.createConversation({ ownership: { kind: 'task', taskId: api.taskId } })).id, + toolContext, + ); + const child = await api.conversation(childId, toolContext); + const settled = await ( + await child.submit({ type: 'input', content: args.task }, toolContext) + ).wait(toolContext); + return { content: [{ type: 'text', text: await answerText(api, settled, toolContext) }] }; + }, + }), + defineTool({ + name: 'fork_delegate', + description: 'Forks this conversation into a subagent.', + parameters: Type.Object({}), + execute: async (_args, api, toolContext) => { + const childId = await api.commit(async tx => { + const at = (await tx.scanEntries({ conversationId: api.conversationId }, 1)).items[0].id; + const fork = await tx.forkConversation(api.conversationId, at, { + ownership: { kind: 'task', taskId: api.taskId }, + }); + return fork.id; + }, toolContext); + const child = await api.conversation(childId, toolContext); + const settled = await ( + await child.submit({ type: 'input', content: 'Forked task.' }, toolContext) + ).wait(toolContext); + return { content: [{ type: 'text', text: await answerText(api, settled, toolContext) }] }; + }, + }), + defineTool({ + name: 'inner', + description: 'A plain tool.', + parameters: Type.Object({}), + execute: async () => ({ content: [{ type: 'text', text: 'INNER RESULT' }] }), + }), + defineTool({ + name: 'batch', + description: 'Runs the inner tool itself.', + parameters: Type.Object({}), + execute: async (_args, api, toolContext) => { + const inner = (await api.agent(toolContext)).tools.find(tool => tool.name === 'inner'); + await inner.execute({}, api, toolContext); + return { content: [{ type: 'text', text: 'OUTER RESULT' }] }; + }, + }), + ], + }), +); + +const harness = await Harness.open(new MemoryStorage(), { models, registry }, context); +const root = await harness.root(context, { agent: { model: { provider: 'faux', modelId: 'faux-model' } } }); +for (const content of ['Delegate the weather question.', 'Fork a helper.', 'Run the batch.']) { + await (await root.submit({ type: 'input', content }, context)).wait(context); +} +await harness.waitForIdle(context); +await harness.close(context); diff --git a/dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-tools.mjs b/dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-tools.mjs new file mode 100644 index 000000000000..fc42ebf2091f --- /dev/null +++ b/dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-tools.mjs @@ -0,0 +1,110 @@ +import * as Sentry from '@sentry/node'; +import { mkdtempSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { BACKGROUND_CONTEXT } from '@earendil-works/chord/context'; +import { Type } from '@earendil-works/pi-ai'; +import { createModels } from '@earendil-works/pi-ai/models'; +import { fauxAssistantMessage, fauxProvider, fauxToolCall } from '@earendil-works/pi-ai/providers/faux'; +import { + createRegistry, + defineExtension, + defineTool, + Harness, + hook, + MemoryStorage, + ToolTask, +} from '@earendil-works/pi-durable'; +import { NodeExecutionEnv } from '@earendil-works/pi-durable/env/node'; +import { CodingTools, createBashTool } from '@earendil-works/pi-durable/tools'; + +const context = BACKGROUND_CONTEXT; + +// One tool round, run one call at a time so the order of the error events is fixed: +// - `bash` is the built-in tool, registered by the app's own extension through its factory; it runs +// a command that exits non-zero, which it reports by throwing; +// - `read` is an app tool that replaces the built-in `read` by name, and throws; +// - `streamer` returns nothing, so its streamed output becomes the result; +// - `secret` returns a value an `afterTool` hook redacts before the model sees it; +// - `fail_now` is an app tool that throws. +const faux = fauxProvider({ provider: 'faux', models: [{ id: 'faux-model' }] }); +const models = createModels(); +models.setProvider(faux.provider); +faux.setResponses([ + fauxAssistantMessage( + [ + fauxToolCall('bash', { command: 'ls does-not-exist' }, { id: 'call_bash' }), + fauxToolCall('read', { path: 'notes.txt' }, { id: 'call_read' }), + fauxToolCall('streamer', {}, { id: 'call_streamer' }), + fauxToolCall('secret', {}, { id: 'call_secret' }), + fauxToolCall('fail_now', {}, { id: 'call_fail' }), + ], + { stopReason: 'toolUse' }, + ), + fauxAssistantMessage('Done.'), +]); + +const registry = createRegistry(); +registry.install(CodingTools); +registry.install( + defineExtension({ + name: 'app', + tools: [ + createBashTool(), + defineTool({ + name: 'read', + description: 'Reads notes.', + parameters: Type.Object({ path: Type.String() }), + execute: async () => { + throw new Error('app read failed'); + }, + }), + defineTool({ + name: 'streamer', + description: 'Streams its output.', + parameters: Type.Object({}), + execute: async (_args, api) => { + api.output('line 1\n'); + api.output('line 2\n'); + return {}; + }, + }), + defineTool({ + name: 'secret', + description: 'Returns a secret.', + parameters: Type.Object({}), + execute: async () => ({ content: [{ type: 'text', text: 'original secret value' }] }), + }), + defineTool({ + name: 'fail_now', + description: 'Always throws.', + parameters: Type.Object({}), + execute: async () => { + throw new Error('Intentional pi-durable tool failure'); + }, + }), + ], + hooks: [ + hook(ToolTask, { + afterTool: (call, result) => + call.name === 'secret' + ? { ...result, content: [{ type: 'text', text: 'redacted by afterTool' }] } + : undefined, + }), + ], + }), +); + +const cwd = mkdtempSync(join(tmpdir(), 'pi-durable-tools-')); +const harness = await Harness.open( + new MemoryStorage(), + { models, registry, env: () => new NodeExecutionEnv({ cwd }), settings: { toolExecution: 'sequential' } }, + context, +); +const root = await harness.root(context, { agent: { model: { provider: 'faux', modelId: 'faux-model' } } }); +await (await root.submit({ type: 'input', content: 'Run the tools.' }, context)).wait(context); +await harness.close(context); + +// Everything captured during the run is sent before this sentinel. +await Sentry.flush(2000); +Sentry.captureMessage('pi-durable tools done'); diff --git a/dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario.mjs b/dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario.mjs new file mode 100644 index 000000000000..42309d5ba2d4 --- /dev/null +++ b/dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario.mjs @@ -0,0 +1,63 @@ +import * as Sentry from '@sentry/node'; +import { BACKGROUND_CONTEXT } from '@earendil-works/chord/context'; +import { Type } from '@earendil-works/pi-ai'; +import { createModels } from '@earendil-works/pi-ai/models'; +import { fauxAssistantMessage, fauxProvider, fauxToolCall } from '@earendil-works/pi-ai/providers/faux'; +import { + createRegistry, + defineExtension, + defineTool, + Harness, + MemoryStorage, + section, +} from '@earendil-works/pi-durable'; + +const context = BACKGROUND_CONTEXT; + +// `pi-ai`'s faux provider scripts model responses in-process. The first run calls two tools and then +// answers, the second run answers directly. +const faux = fauxProvider({ provider: 'faux', models: [{ id: 'faux-model' }] }); +const models = createModels(); +models.setProvider(faux.provider); +faux.setResponses([ + fauxAssistantMessage( + [fauxToolCall('get_weather', { city: 'Berlin' }, { id: 'call_1' }), fauxToolCall('broken', {}, { id: 'call_2' })], + { stopReason: 'toolUse' }, + ), + fauxAssistantMessage('It is 21 degrees and sunny in Berlin.'), + fauxAssistantMessage('Goodbye.'), +]); + +const registry = createRegistry(); +registry.install( + defineExtension({ + name: 'weather', + sections: [section('preamble', () => 'You are a weather assistant.', { tag: false })], + tools: [ + defineTool({ + name: 'get_weather', + description: 'Get the current weather for a city.', + parameters: Type.Object({ city: Type.String() }), + execute: async args => ({ content: [{ type: 'text', text: `It is 21 degrees and sunny in ${args.city}.` }] }), + }), + defineTool({ + name: 'broken', + description: 'Always fails.', + parameters: Type.Object({}), + execute: async () => { + throw new Error('broken tool'); + }, + }), + ], + }), +); + +// The Harness is opened and driven inside a request span. Its runs must still start traces of their +// own: the scheduler runs them later, and their trace must not depend on which request woke it. +await Sentry.startSpan({ name: 'pi-durable-request', op: 'http.server' }, async () => { + const harness = await Harness.open(new MemoryStorage(), { models, registry }, context); + const root = await harness.root(context, { agent: { model: { provider: 'faux', modelId: 'faux-model' } } }); + await (await root.submit({ type: 'input', content: 'What is the weather in Berlin?' }, context)).wait(context); + await (await root.submit({ type: 'input', content: 'Thanks!' }, context)).wait(context); + await harness.close(context); +}); diff --git a/dev-packages/node-integration-tests/suites/tracing/pi-durable/test.ts b/dev-packages/node-integration-tests/suites/tracing/pi-durable/test.ts new file mode 100644 index 000000000000..f71f88804dc0 --- /dev/null +++ b/dev-packages/node-integration-tests/suites/tracing/pi-durable/test.ts @@ -0,0 +1,693 @@ +import { + GEN_AI_CONVERSATION_ID, + GEN_AI_INPUT_MESSAGES, + GEN_AI_OPERATION_NAME, + GEN_AI_OUTPUT_MESSAGES, + GEN_AI_PROVIDER_NAME, + GEN_AI_REQUEST_MODEL, + GEN_AI_RESPONSE_FINISH_REASONS, + GEN_AI_RESPONSE_STREAMING, + GEN_AI_SYSTEM_INSTRUCTIONS, + GEN_AI_TOOL_CALL_ARGUMENTS, + GEN_AI_TOOL_CALL_RESULT, + GEN_AI_TOOL_DEFINITIONS, + GEN_AI_TOOL_DESCRIPTION, + GEN_AI_TOOL_NAME, + GEN_AI_USAGE_INPUT_TOKENS, + GEN_AI_USAGE_OUTPUT_TOKENS, +} from '@sentry/conventions/attributes'; +import { afterAll, expect } from 'vitest'; +import { conditionalTest } from '../../../utils'; +import { cleanupChildProcesses, createEsmAndCjsTests } from '../../../utils/runner'; + +// The pi packages declare `engines.node >= 22.19`, so they can't live in the package's root +// `devDependencies` (that would break `yarn install` on the 20.19 CI matrix). Install them per-suite +// instead, guarded by the `min: 22` skip below. +const PI_DURABLE_DEPENDENCIES = { + additionalDependencies: { + '@earendil-works/pi-durable': '^1.0.0', + '@earendil-works/pi-ai': '^1.0.0', + '@earendil-works/chord': '^1.0.0', + }, +}; + +const CONTENT_ATTRIBUTES = [ + GEN_AI_INPUT_MESSAGES, + GEN_AI_OUTPUT_MESSAGES, + GEN_AI_SYSTEM_INSTRUCTIONS, + GEN_AI_TOOL_DEFINITIONS, + GEN_AI_TOOL_CALL_ARGUMENTS, + GEN_AI_TOOL_CALL_RESULT, +]; + +conditionalTest({ min: 22 })('pi-durable integration', () => { + afterAll(() => { + cleanupChildProcesses(); + }); + + createEsmAndCjsTests( + __dirname, + 'scenario.mjs', + 'instrument.mjs', + (createRunner, test, mode) => { + // The pi packages are ESM-only: their `exports` maps have no `require` condition. + if (mode === 'cjs') { + return; + } + + test('traces each run as its own invoke_agent trace with chat and execute_tool children', async () => { + let toolRunTraceId: string | undefined; + let brokenToolSpanId: string | undefined; + let errorTraceId: string | undefined; + let errorSpanId: string | undefined; + const conversationIds: unknown[] = []; + + await createRunner() + .unordered() + .expect({ + span: container => { + const spans = container.items; + const agent = spans.find(span => span.is_segment)!; + expect(agent.name).toBe('invoke_agent'); + expect(spans.filter(span => span.name === 'chat faux-model')).toHaveLength(2); + // One span per run, request and tool call: a tool wrapped twice would add a span. + expect( + spans.filter(span => span.attributes['sentry.origin']?.value === 'auto.ai.pi_durable'), + ).toHaveLength(5); + + expect(agent.attributes['sentry.op']?.value).toBe('gen_ai.invoke_agent'); + expect(agent.attributes['sentry.origin']?.value).toBe('auto.ai.pi_durable'); + expect(agent.attributes[GEN_AI_OPERATION_NAME]?.value).toBe('invoke_agent'); + const conversationId = agent.attributes[GEN_AI_CONVERSATION_ID]?.value; + // The root conversation's id is `1` in every pi-durable storage, so it is prefixed + // with a per-Harness id. + expect(conversationId).toMatch(/^[0-9a-f]{32}:1$/); + + const chats = spans + .filter(span => span.name === 'chat faux-model') + .sort((a, b) => a.start_timestamp - b.start_timestamp); + for (const chat of chats) { + expect(chat.parent_span_id).toBe(agent.span_id); + expect(chat.attributes['sentry.op']?.value).toBe('gen_ai.chat'); + expect(chat.attributes['sentry.origin']?.value).toBe('auto.ai.pi_durable'); + expect(chat.attributes[GEN_AI_PROVIDER_NAME]?.value).toBe('faux'); + expect(chat.attributes[GEN_AI_REQUEST_MODEL]?.value).toBe('faux-model'); + expect(chat.attributes[GEN_AI_CONVERSATION_ID]?.value).toBe(conversationId); + expect(chat.attributes[GEN_AI_USAGE_INPUT_TOKENS]?.value).toBeGreaterThan(0); + expect(chat.attributes[GEN_AI_USAGE_OUTPUT_TOKENS]?.value).toBeGreaterThan(0); + // pi-durable sends the prompt and the tools as positional system messages. + expect(chat.attributes[GEN_AI_SYSTEM_INSTRUCTIONS]?.value).toBe('You are a weather assistant.'); + expect( + JSON.parse(String(chat.attributes[GEN_AI_TOOL_DEFINITIONS]?.value)).map( + (tool: { name: string }) => tool.name, + ), + ).toEqual(['get_weather', 'broken']); + } + + // Mapped from pi-ai's shape to the conventions: the tool calls and tool results the + // second request carries are what Sentry's conversation view renders. + expect(JSON.parse(String(chats[0]!.attributes[GEN_AI_INPUT_MESSAGES]?.value))).toEqual([ + { role: 'user', parts: [{ type: 'text', content: 'What is the weather in Berlin?' }] }, + ]); + const answerInput = JSON.parse(String(chats[1]!.attributes[GEN_AI_INPUT_MESSAGES]?.value)); + expect(answerInput.map((message: { role: string }) => message.role)).toEqual([ + 'user', + 'assistant', + 'tool', + 'tool', + ]); + expect(answerInput[1].parts).toEqual([ + { type: 'tool_call', id: 'call_1', name: 'get_weather', arguments: '{"city":"Berlin"}' }, + { type: 'tool_call', id: 'call_2', name: 'broken', arguments: '{}' }, + ]); + expect(answerInput[2].parts).toEqual([ + { + type: 'tool_call_response', + id: 'call_1', + name: 'get_weather', + result: 'It is 21 degrees and sunny in Berlin.', + }, + ]); + expect(chats.map(chat => chat.attributes[GEN_AI_RESPONSE_FINISH_REASONS]?.value).sort()).toEqual([ + '["stop"]', + '["tool_call"]', + ]); + expect(chats.map(chat => chat.attributes[GEN_AI_OUTPUT_MESSAGES]?.value).join()).toContain( + 'It is 21 degrees and sunny in Berlin.', + ); + + const weather = spans.find(span => span.name === 'execute_tool get_weather')!; + expect(weather.parent_span_id).toBe(agent.span_id); + expect(weather.status).toBe('ok'); + expect(weather.attributes['sentry.op']?.value).toBe('gen_ai.execute_tool'); + expect(weather.attributes[GEN_AI_TOOL_NAME]?.value).toBe('get_weather'); + expect(weather.attributes[GEN_AI_TOOL_DESCRIPTION]?.value).toBe('Get the current weather for a city.'); + expect(weather.attributes['gen_ai.tool.call.id']?.value).toBe('call_1'); + expect(weather.attributes[GEN_AI_TOOL_CALL_ARGUMENTS]?.value).toBe('{"city":"Berlin"}'); + expect(weather.attributes[GEN_AI_TOOL_CALL_RESULT]?.value).toContain('sunny in Berlin'); + expect(weather.attributes[GEN_AI_CONVERSATION_ID]?.value).toBe(conversationId); + + const broken = spans.find(span => span.name === 'execute_tool broken')!; + expect(broken.parent_span_id).toBe(agent.span_id); + expect(broken.status).toBe('error'); + // The error result pi-durable gives the model, not the thrown error. + expect(broken.attributes[GEN_AI_TOOL_CALL_RESULT]?.value).toContain('[error] broken tool'); + + toolRunTraceId = agent.trace_id; + brokenToolSpanId = broken.span_id; + conversationIds.push(conversationId); + }, + }) + .expect({ + span: container => { + const spans = container.items; + const agent = spans.find(span => span.is_segment)!; + expect(agent.name).toBe('invoke_agent'); + expect(spans).toHaveLength(2); + + const chat = spans.find(span => span.name === 'chat faux-model')!; + expect(chat.parent_span_id).toBe(agent.span_id); + expect(chat.attributes[GEN_AI_OUTPUT_MESSAGES]?.value).toContain('Goodbye.'); + + conversationIds.push(agent.attributes[GEN_AI_CONVERSATION_ID]?.value); + }, + }) + .expect({ + span: container => { + // The request that opened the Harness and submitted both inputs holds no agent work. + expect(container.items).toHaveLength(1); + expect(container.items[0]!.name).toBe('pi-durable-request'); + }, + }) + .expect({ + event: event => { + const exception = event.exception?.values?.[0]; + expect(exception?.value).toBe('broken tool'); + expect(exception?.mechanism).toEqual({ type: 'auto.ai.pi_durable', handled: true }); + + errorTraceId = event.contexts?.trace?.trace_id; + errorSpanId = event.contexts?.trace?.span_id; + }, + }) + .start() + .completed(); + + // Both runs belong to the root conversation, so they share one conversation id. + expect(conversationIds).toHaveLength(2); + expect(conversationIds[0]).toBe(conversationIds[1]); + // The tool error is reported on the failing call's span, in the run's trace. + expect(errorTraceId).toBe(toolRunTraceId); + expect(errorSpanId).toBe(brokenToolSpanId); + }); + + test('records no message content when the client turns gen_ai recording off', async () => { + await createRunner() + .withEnv({ PI_DURABLE_RECORDING: 'off' }) + .unordered() + .expect({ + span: container => { + const spans = container.items; + expect(spans).toHaveLength(5); + for (const span of spans) { + for (const attribute of CONTENT_ATTRIBUTES) { + expect(span.attributes[attribute]).toBeUndefined(); + } + } + const weather = spans.find(span => span.name === 'execute_tool get_weather')!; + expect(weather.attributes[GEN_AI_TOOL_DESCRIPTION]?.value).toBe('Get the current weather for a city.'); + expect( + spans.find(span => span.name === 'chat faux-model')!.attributes[GEN_AI_USAGE_INPUT_TOKENS]?.value, + ).toBeGreaterThan(0); + }, + }) + .expect({ + span: container => { + expect(container.items).toHaveLength(2); + for (const span of container.items) { + for (const attribute of CONTENT_ATTRIBUTES) { + expect(span.attributes[attribute]).toBeUndefined(); + } + } + }, + }) + .expect({ + span: container => { + expect(container.items).toHaveLength(1); + expect(container.items[0]!.name).toBe('pi-durable-request'); + }, + }) + .expect({ + event: event => { + expect(event.exception?.values?.[0]?.value).toBe('broken tool'); + }, + }) + .start() + .completed(); + }); + }, + PI_DURABLE_DEPENDENCIES, + ); + + createEsmAndCjsTests( + __dirname, + 'scenario-tools.mjs', + 'instrument.mjs', + (createRunner, test, mode) => { + if (mode === 'cjs') { + return; + } + + test('reports app tool failures but not failures of the built-in coding tools', async () => { + const appErrors: unknown[] = []; + + // Ordered, with spans ignored: a captured failure of the built-in `bash` tool would be an + // extra event before the sentinel, which the scenario sends once everything else is sent. + // The two app failures are processed concurrently, so they may arrive in either order. + await createRunner() + .ignore('span') + .expect({ + event: event => { + expect(['app read failed', 'Intentional pi-durable tool failure']).toContain( + event.exception?.values?.[0]?.value, + ); + appErrors.push(event.exception?.values?.[0]?.value); + }, + }) + .expect({ + event: event => { + expect(['app read failed', 'Intentional pi-durable tool failure']).toContain( + event.exception?.values?.[0]?.value, + ); + appErrors.push(event.exception?.values?.[0]?.value); + }, + }) + .expect({ + event: event => { + expect(event.message).toBe('pi-durable tools done'); + }, + }) + .start() + .completed(); + + expect(appErrors.sort()).toEqual(['Intentional pi-durable tool failure', 'app read failed']); + }); + + test('records the tool result the model receives', async () => { + let errorSpanId: string | undefined; + let failingToolSpanId: string | undefined; + + await createRunner() + .unordered() + .expect({ + event: event => { + const exception = event.exception?.values?.[0]; + expect(exception?.value).toBe('Intentional pi-durable tool failure'); + expect(exception?.mechanism).toEqual({ type: 'auto.ai.pi_durable', handled: true }); + errorSpanId = event.contexts?.trace?.span_id; + }, + }) + .expect({ + event: event => { + expect(event.exception?.values?.[0]?.value).toBe('app read failed'); + }, + }) + .expect({ + span: container => { + const spans = container.items; + + const bash = spans.find(span => span.name === 'execute_tool bash')!; + expect(bash.status).toBe('error'); + expect(bash.attributes[GEN_AI_TOOL_CALL_RESULT]?.value).toMatch(/Command exited with code \d+/); + + const read = spans.find(span => span.name === 'execute_tool read')!; + expect(read.status).toBe('error'); + expect(read.attributes[GEN_AI_TOOL_CALL_RESULT]?.value).toContain('app read failed'); + + // The streamed output becomes the result, since the tool returns no content. + const streamer = spans.find(span => span.name === 'execute_tool streamer')!; + expect(streamer.status).toBe('ok'); + expect(streamer.attributes[GEN_AI_TOOL_CALL_RESULT]?.value).toBe('line 1\nline 2\n'); + + const secret = spans.find(span => span.name === 'execute_tool secret')!; + expect(secret.attributes[GEN_AI_TOOL_CALL_RESULT]?.value).toBe('redacted by afterTool'); + + const failNow = spans.find(span => span.name === 'execute_tool fail_now')!; + expect(failNow.status).toBe('error'); + failingToolSpanId = failNow.span_id; + }, + }) + .start() + .completed(); + + expect(errorSpanId).toBe(failingToolSpanId); + }); + }, + PI_DURABLE_DEPENDENCIES, + ); + + createEsmAndCjsTests( + __dirname, + 'scenario-reports.mjs', + 'instrument.mjs', + (createRunner, test, mode) => { + if (mode === 'cjs') { + return; + } + + test('captures failures pi-durable does not propagate, and tags idle compaction with its conversation', async () => { + let runConversationId: unknown; + let compactionConversationId: unknown; + + await createRunner() + .unordered() + .expect({ + event: event => { + const exception = event.exception?.values?.[0]; + expect(exception?.value).toBe('afterResponse hook failed'); + expect(exception?.mechanism).toEqual({ type: 'auto.ai.pi_durable', handled: true }); + }, + }) + .expect({ + event: event => { + const exception = event.exception?.values?.[0]; + expect(exception?.value).toBe('card declined'); + expect(exception?.mechanism).toEqual({ type: 'auto.ai.pi_durable', handled: false }); + }, + }) + .expect({ + span: container => { + const agent = container.items.find(span => span.is_segment)!; + expect(agent.name).toBe('invoke_agent'); + runConversationId = agent.attributes[GEN_AI_CONVERSATION_ID]?.value; + }, + }) + .expect({ + span: container => { + // A compaction started while idle belongs to no run, so its request is a trace of its own. + expect(container.items).toHaveLength(1); + const summary = container.items[0]!; + expect(summary.is_segment).toBe(true); + expect(summary.attributes['sentry.op']?.value).toBe('gen_ai.chat'); + compactionConversationId = summary.attributes[GEN_AI_CONVERSATION_ID]?.value; + }, + }) + .start() + .completed(); + + expect(runConversationId).toMatch(/^[0-9a-f]{32}:1$/); + expect(compactionConversationId).toBe(runConversationId); + }); + }, + PI_DURABLE_DEPENDENCIES, + ); + + createEsmAndCjsTests( + __dirname, + 'scenario-subagent.mjs', + 'instrument.mjs', + (createRunner, test, mode) => { + if (mode === 'cjs') { + return; + } + + test('nests subagent runs under the tool call that started them', async () => { + await createRunner() + .unordered() + .expect({ + span: container => { + const spans = container.items; + const parent = spans.find(span => span.is_segment)!; + const delegate = spans.find(span => span.name === 'execute_tool delegate')!; + expect(parent.name).toBe('invoke_agent'); + expect(delegate.parent_span_id).toBe(parent.span_id); + + // The subagent's run is a child of the delegating call, in the same trace, on its own + // conversation of the same Harness. + const child = spans.find(span => span.name === 'invoke_agent' && !span.is_segment)!; + expect(child.parent_span_id).toBe(delegate.span_id); + const parentConversation = String(parent.attributes[GEN_AI_CONVERSATION_ID]?.value); + const childConversation = String(child.attributes[GEN_AI_CONVERSATION_ID]?.value); + expect(childConversation).not.toBe(parentConversation); + expect(childConversation.split(':')[0]).toBe(parentConversation.split(':')[0]); + + const childChat = spans.find(span => span.parent_span_id === child.span_id)!; + expect(childChat.name).toBe('chat faux-model'); + expect(childChat.attributes[GEN_AI_CONVERSATION_ID]?.value).toBe(childConversation); + expect(delegate.attributes[GEN_AI_TOOL_CALL_RESULT]?.value).toBe('Answer: Weather in Vienna?'); + }, + }) + .expect({ + span: container => { + const spans = container.items; + const parent = spans.find(span => span.is_segment)!; + const fork = spans.find(span => span.name === 'execute_tool fork_delegate')!; + const child = spans.find(span => span.name === 'invoke_agent' && !span.is_segment)!; + expect(fork.parent_span_id).toBe(parent.span_id); + expect(child.parent_span_id).toBe(fork.span_id); + expect(fork.attributes[GEN_AI_TOOL_CALL_RESULT]?.value).toBe('Answer: Forked task.'); + }, + }) + .expect({ + span: container => { + const spans = container.items; + const batch = spans.find(span => span.name === 'execute_tool batch')!; + const inner = spans.find(span => span.name === 'execute_tool inner')!; + expect(batch.parent_span_id).toBe(spans.find(span => span.is_segment)!.span_id); + // A tool that runs another tool itself: the inner call ends with its own result. + expect(inner.parent_span_id).toBe(batch.span_id); + expect(inner.attributes[GEN_AI_TOOL_CALL_RESULT]?.value).toBe('INNER RESULT'); + expect(batch.attributes[GEN_AI_TOOL_CALL_RESULT]?.value).toBe('OUTER RESULT'); + }, + }) + .start() + .completed(); + }); + }, + PI_DURABLE_DEPENDENCIES, + ); + + createEsmAndCjsTests( + __dirname, + 'scenario-interrupted.mjs', + 'instrument.mjs', + (createRunner, test, mode) => { + if (mode === 'cjs') { + return; + } + + test('reports a provider that throws as a fault, and nothing for failed, aborted or closed runs', async () => { + await createRunner() + .ignore('span') + .expect({ + event: event => { + const exception = event.exception?.values?.[0]; + expect(exception?.value).toBe('provider broke'); + expect(exception?.mechanism).toEqual({ type: 'auto.ai.pi_durable', handled: false }); + }, + }) + .expect({ + event: event => { + expect(event.message).toBe('pi-durable interrupted done'); + }, + }) + .start() + .completed(); + }); + + test('ends every run that does not end with an answer', async () => { + await createRunner() + .unordered() + .expect({ + span: container => { + // 1. The model request failed: the run and the request are errors. + const chat = container.items.find(span => span.name === 'chat faux-model')!; + expect(chat.attributes[GEN_AI_INPUT_MESSAGES]?.value).toContain('Fail please.'); + expect(chat.attributes[GEN_AI_RESPONSE_FINISH_REASONS]?.value).toBe('["error"]'); + expect(chat.status).toBe('error'); + expect(chat.attributes['sentry.status.message']?.value).toBe('internal_error'); + const agent = container.items.find(span => span.is_segment)!; + expect(agent.name).toBe('invoke_agent'); + expect(agent.status).toBe('error'); + expect(agent.attributes['sentry.status.message']?.value).toBe('model_error'); + }, + }) + .expect({ + span: container => { + // 2. The provider threw: the generation faulted. + const chat = container.items.find(span => span.name === 'chat faux-model')!; + expect(chat.attributes[GEN_AI_RESPONSE_FINISH_REASONS]).toBeUndefined(); + expect(chat.attributes['sentry.status.message']?.value).toBe('internal_error'); + const agent = container.items.find(span => span.is_segment)!; + expect(agent.attributes['sentry.status.message']?.value).toBe('internal_error'); + }, + }) + .expect({ + span: container => { + // 3. Aborted while the tool ran: cancelled, which streamed spans report as ok. + const slow = container.items.find(span => span.name === 'execute_tool slow')!; + expect(slow.attributes[GEN_AI_TOOL_CALL_ARGUMENTS]?.value).toBe('{"job":"abort"}'); + expect(slow.status).toBe('ok'); + expect(container.items.find(span => span.is_segment)!.status).toBe('ok'); + }, + }) + .expect({ + span: container => { + // 4. Aborted while the model request ran. + const chat = container.items.find(span => span.name === 'chat faux-model')!; + expect(chat.attributes[GEN_AI_RESPONSE_FINISH_REASONS]?.value).toBe('["aborted"]'); + expect(container.items).toHaveLength(2); + }, + }) + .expect({ + span: container => { + // 5. The Harness closed while the tool ran: the run still ends, so the trace has a segment. + const slow = container.items.find(span => span.name === 'execute_tool slow')!; + expect(slow.attributes[GEN_AI_TOOL_CALL_ARGUMENTS]?.value).toBe('{"job":"close"}'); + const agent = container.items.find(span => span.is_segment)!; + expect(agent.name).toBe('invoke_agent'); + expect(slow.parent_span_id).toBe(agent.span_id); + }, + }) + .expect({ + span: container => { + // 6. The run failed, and the Harness closed right after `wait()`, before the commit that + // ended the run resolved: the run keeps its own status. + const chat = container.items.find(span => span.name === 'chat faux-model')!; + expect(chat.attributes[GEN_AI_INPUT_MESSAGES]?.value).toContain('Fail, then close.'); + const agent = container.items.find(span => span.is_segment)!; + expect(agent.status).toBe('error'); + expect(agent.attributes['sentry.status.message']?.value).toBe('model_error'); + }, + }) + .start() + .completed(); + }); + }, + PI_DURABLE_DEPENDENCIES, + ); + + createEsmAndCjsTests( + __dirname, + 'scenario-concurrent.mjs', + 'instrument.mjs', + (createRunner, test, mode) => { + if (mode === 'cjs') { + return; + } + + test('keeps concurrent runs apart and out of the scopes of the requests that submitted them', async () => { + const runTraceIds = new Set(); + const errorTraceIds = new Set(); + + await createRunner() + .unordered() + .expect({ + span: container => { + const spans = container.items; + const agent = spans.find(span => span.is_segment)!; + expect(agent.name).toBe('invoke_agent'); + const chats = spans.filter(span => span.name === 'chat faux-model'); + expect(JSON.parse(String(chats[0]!.attributes[GEN_AI_INPUT_MESSAGES]?.value))[0].parts[0].content).toBe( + 'request A', + ); + // The run's own conversation id, not the one the request set on its scope. + const conversationId = agent.attributes[GEN_AI_CONVERSATION_ID]?.value; + expect(conversationId).toMatch(/^[0-9a-f]{32}:\d+$/); + for (const span of spans) { + expect(span.attributes[GEN_AI_CONVERSATION_ID]?.value).toBe(conversationId); + expect(span.attributes['request']).toBeUndefined(); + } + expect( + spans.find(span => span.name === 'execute_tool work')!.attributes[GEN_AI_TOOL_CALL_ARGUMENTS]?.value, + ).toBe('{"tag":"request A"}'); + runTraceIds.add(agent.trace_id); + }, + }) + .expect({ + span: container => { + const spans = container.items; + const agent = spans.find(span => span.is_segment)!; + expect(agent.name).toBe('invoke_agent'); + const chats = spans.filter(span => span.name === 'chat faux-model'); + expect(JSON.parse(String(chats[0]!.attributes[GEN_AI_INPUT_MESSAGES]?.value))[0].parts[0].content).toBe( + 'request B', + ); + const conversationId = agent.attributes[GEN_AI_CONVERSATION_ID]?.value; + expect(conversationId).toMatch(/^[0-9a-f]{32}:\d+$/); + for (const span of spans) { + expect(span.attributes[GEN_AI_CONVERSATION_ID]?.value).toBe(conversationId); + expect(span.attributes['request']).toBeUndefined(); + } + expect( + spans.find(span => span.name === 'execute_tool work')!.attributes[GEN_AI_TOOL_CALL_ARGUMENTS]?.value, + ).toBe('{"tag":"request B"}'); + runTraceIds.add(agent.trace_id); + }, + }) + .expect({ + event: event => { + expect(event.exception?.values?.[0]?.value).toBe('work failed for request A'); + expect(event.tags?.request).toBeUndefined(); + errorTraceIds.add(String(event.contexts?.trace?.trace_id)); + }, + }) + .expect({ + event: event => { + expect(event.exception?.values?.[0]?.value).toBe('work failed for request B'); + expect(event.tags?.request).toBeUndefined(); + errorTraceIds.add(String(event.contexts?.trace?.trace_id)); + }, + }) + .start() + .completed(); + + // Each tool error belongs to its run's trace. + expect(errorTraceIds).toEqual(runTraceIds); + }); + }, + PI_DURABLE_DEPENDENCIES, + ); + + createEsmAndCjsTests( + __dirname, + 'scenario-compaction.mjs', + 'instrument.mjs', + (createRunner, test, mode) => { + if (mode === 'cjs') { + return; + } + + test('keeps a compaction that pi-durable starts for a run inside that run', async () => { + await createRunner() + .unordered() + .expect({ + span: container => { + const spans = container.items; + const summary = spans.find(span => + String(span.attributes[GEN_AI_OUTPUT_MESSAGES]?.value).includes('Summary'), + )!; + const agent = spans.find(span => span.is_segment)!; + expect(agent.name).toBe('invoke_agent'); + expect(summary.parent_span_id).toBe(agent.span_id); + // pi-durable completes the summary without streaming; the run's own request streams. + expect(summary.attributes[GEN_AI_RESPONSE_STREAMING]).toBeUndefined(); + const answer = spans.find(span => + String(span.attributes[GEN_AI_OUTPUT_MESSAGES]?.value).includes('Third answer'), + )!; + expect(answer.parent_span_id).toBe(agent.span_id); + expect(answer.attributes[GEN_AI_RESPONSE_STREAMING]?.value).toBe(true); + }, + }) + .expect({ + span: container => expect(container.items.find(span => span.is_segment)!.name).toBe('invoke_agent'), + }) + .expect({ + span: container => expect(container.items.find(span => span.is_segment)!.name).toBe('invoke_agent'), + }) + .start() + .completed(); + }); + }, + PI_DURABLE_DEPENDENCIES, + ); +}); diff --git a/packages/astro/src/index.server.ts b/packages/astro/src/index.server.ts index b67a20a788f8..cadfe7cdae59 100644 --- a/packages/astro/src/index.server.ts +++ b/packages/astro/src/index.server.ts @@ -102,6 +102,7 @@ export { langGraphIntegration, createFlueInstrumentation, mastraIntegration, + piDurableIntegration, mcpServerIntegration, SentryMastraExporter, parameterize, diff --git a/packages/aws-serverless/src/index.ts b/packages/aws-serverless/src/index.ts index 5d0c37614a3e..98d52745adb9 100644 --- a/packages/aws-serverless/src/index.ts +++ b/packages/aws-serverless/src/index.ts @@ -67,6 +67,7 @@ export { langChainIntegration, langGraphIntegration, mastraIntegration, + piDurableIntegration, mcpServerIntegration, SentryMastraExporter, createFlueInstrumentation, diff --git a/packages/bun/src/index.ts b/packages/bun/src/index.ts index d5b1d11ee4e5..d8860a5f2a54 100644 --- a/packages/bun/src/index.ts +++ b/packages/bun/src/index.ts @@ -89,6 +89,7 @@ export { langChainIntegration, langGraphIntegration, mastraIntegration, + piDurableIntegration, mcpServerIntegration, SentryMastraExporter, createFlueInstrumentation, diff --git a/packages/deno/src/index.ts b/packages/deno/src/index.ts index 8db992d91e6a..d77d424b9841 100644 --- a/packages/deno/src/index.ts +++ b/packages/deno/src/index.ts @@ -138,6 +138,7 @@ export { langChainIntegration, langGraphIntegration, mastraIntegration, + piDurableIntegration, mcpServerIntegration, SentryMastraExporter, createFlueInstrumentation, diff --git a/packages/elysia/src/index.ts b/packages/elysia/src/index.ts index cf35f6de4eb8..baea6e9f7f0e 100644 --- a/packages/elysia/src/index.ts +++ b/packages/elysia/src/index.ts @@ -69,6 +69,7 @@ export { langGraphIntegration, createFlueInstrumentation, mastraIntegration, + piDurableIntegration, mcpServerIntegration, SentryMastraExporter, modulesIntegration, diff --git a/packages/google-cloud-serverless/src/index.ts b/packages/google-cloud-serverless/src/index.ts index d789d83c91a8..293906ac87dc 100644 --- a/packages/google-cloud-serverless/src/index.ts +++ b/packages/google-cloud-serverless/src/index.ts @@ -67,6 +67,7 @@ export { langChainIntegration, langGraphIntegration, mastraIntegration, + piDurableIntegration, mcpServerIntegration, SentryMastraExporter, createFlueInstrumentation, diff --git a/packages/node/src/index.ts b/packages/node/src/index.ts index deed6c7c0b65..760ea6dd9f00 100644 --- a/packages/node/src/index.ts +++ b/packages/node/src/index.ts @@ -25,6 +25,7 @@ export { lruMemoizerIntegration, createFlueInstrumentation, mastraIntegration, + piDurableIntegration, mcpServerIntegration, SentryMastraExporter, mongoIntegration, diff --git a/packages/server-utils/src/ai/pi-ai/messages.ts b/packages/server-utils/src/ai/pi-ai/messages.ts index 5e3408d6f354..b34cfd22c5fb 100644 --- a/packages/server-utils/src/ai/pi-ai/messages.ts +++ b/packages/server-utils/src/ai/pi-ai/messages.ts @@ -66,6 +66,80 @@ export function piAiFinishReason(stopReason: unknown): string | undefined { return stopReason === 'toolUse' ? 'tool_call' : stopReason; } +/** The request fields of a pi-ai `Context` that carry the system prompt and the tools. */ +export interface PiAiContext { + systemPrompt?: string; + messages?: unknown; + tools?: unknown; +} + +/** + * The system prompt a pi-ai request sends. pi-durable sends it as positional system messages with + * `sections` instead of `systemPrompt`, so the messages are replayed like pi-ai's own + * `getCurrentSystemPrompt` does: `content` is appended, `sections` are patched by name. + */ +export function piAiSystemInstructions(context: PiAiContext): string | undefined { + const content: string[] = []; + const sections = new Map(); + for (const message of piAiSystemMessages(context)) { + const text = piAiContentText(message.content); + if (text) { + content.push(text); + } + if (isObjectLike(message.sections)) { + for (const [name, value] of Object.entries(message.sections)) { + if (value === null) { + sections.delete(name); + } else if (typeof value === 'string') { + sections.set(name, value); + } + } + } + } + return [content.join('\n\n'), ...sections.values()].filter(Boolean).join('\n\n') || undefined; +} + +/** The tools a pi-ai request offers, with `toolsAdded` and `toolsRemoved` replayed like pi-ai's `getCurrentTools`. */ +export function piAiToolDefinitions(context: PiAiContext): unknown[] { + const tools = new Map(); + for (const message of piAiSystemMessages(context)) { + for (const tool of Array.isArray(message.toolsRemoved) ? message.toolsRemoved : []) { + if (isObjectLike(tool) && typeof tool.name === 'string') { + tools.delete(tool.name); + } + } + for (const tool of Array.isArray(message.toolsAdded) ? message.toolsAdded : []) { + if (isObjectLike(tool) && typeof tool.name === 'string') { + tools.set(tool.name, tool); + } + } + } + return [...tools.values()]; +} + +/** The system messages pi-ai sends, led by the one it builds from `systemPrompt` and `tools`. */ +function piAiSystemMessages(context: PiAiContext): Record[] { + const hasTools = Array.isArray(context.tools) && context.tools.length > 0; + const leading = + context.systemPrompt || hasTools + ? [{ role: 'system', content: context.systemPrompt ?? '', toolsAdded: context.tools }] + : []; + const messages = Array.isArray(context.messages) ? context.messages : []; + return [...leading, ...messages].filter( + (message): message is Record => isObjectLike(message) && message.role === 'system', + ); +} + +function piAiContentText(content: unknown): string { + if (typeof content === 'string') { + return content; + } + return piAiContentToParts(content) + .filter(part => part.type === 'text') + .map(part => part.content) + .join('\n'); +} + /** Text-only content as its text, joined the way pi-ai joins it; anything else as its mapped parts. */ export function piAiContentToString(content: unknown): string | undefined { const parts = piAiContentToParts(content); diff --git a/packages/server-utils/src/ai/pi-ai/providers.ts b/packages/server-utils/src/ai/pi-ai/providers.ts index 524afe79d573..96bc6149f620 100644 --- a/packages/server-utils/src/ai/pi-ai/providers.ts +++ b/packages/server-utils/src/ai/pi-ai/providers.ts @@ -2,15 +2,18 @@ import { _INTERNAL_shouldSkipAiProviderWrapping, _INTERNAL_skipAiProviderWrappin import { ANTHROPIC_AI_INTEGRATION_NAME } from '../anthropic-ai/constants'; import { GOOGLE_GENAI_INTEGRATION_NAME } from '../google-genai/constants'; import { OPENAI_INTEGRATION_NAME } from '../openai/constants'; +import { WORKERS_AI_INTEGRATION_NAME } from '../workers-ai/constants'; -// pi-ai sends its requests through the `openai`, `@anthropic-ai/sdk` and `@google/genai` clients. -// Left alone, those integrations report the same request a second time beside the `chat` span of -// the framework that sent it. Bedrock requests go through `@aws-sdk/client-bedrock-runtime`, which -// `awsIntegration` still reports; that one has no skip yet. +// pi-ai sends its requests through the `openai`, `@anthropic-ai/sdk` and `@google/genai` clients, and +// on Cloudflare through the Workers AI binding (`createAI()` of `agents/models/pi-ai`). Left alone, +// those integrations report the same request a second time beside the `chat` span of the framework +// that sent it. Bedrock requests go through `@aws-sdk/client-bedrock-runtime`, which `awsIntegration` +// still reports; that one has no skip yet. const PI_AI_PROVIDER_INTEGRATIONS = [ OPENAI_INTEGRATION_NAME, ANTHROPIC_AI_INTEGRATION_NAME, GOOGLE_GENAI_INTEGRATION_NAME, + WORKERS_AI_INTEGRATION_NAME, ]; /** diff --git a/packages/server-utils/src/ai/pi-durable/constants.ts b/packages/server-utils/src/ai/pi-durable/constants.ts new file mode 100644 index 000000000000..f266a250c24a --- /dev/null +++ b/packages/server-utils/src/ai/pi-durable/constants.ts @@ -0,0 +1,23 @@ +export const PI_DURABLE_INTEGRATION_NAME = 'PiDurable' as const; + +export const PI_DURABLE_ORIGIN = 'auto.ai.pi_durable'; + +/** Task kinds of the pi-durable built-in tasks. */ +export const PI_TASK = { + GENERATION: 'pi.generation', + TOOL: 'pi.tool', + COMPACTION: 'pi.compaction', +} as const; + +/** Kind of the built-in document that holds a conversation's run control. */ +export const PI_LIVE_DOC_KIND = 'pi.live'; + +/** Kind of the transcript entry that holds the tool result the model receives. */ +export const PI_TOOL_RESULT_ENTRY_KIND = 'pi.tool-result'; + +/** + * Cap on open run spans. A run span is removed when the run settles, but a run the process never + * sees settle (a crash, or a settlement the scheduler writes outside a task phase) would otherwise + * stay in the map for the lifetime of the process. + */ +export const MAX_TRACKED_PI_RUNS = 1000; diff --git a/packages/server-utils/src/ai/pi-durable/index.ts b/packages/server-utils/src/ai/pi-durable/index.ts new file mode 100644 index 000000000000..e4a30393096e --- /dev/null +++ b/packages/server-utils/src/ai/pi-durable/index.ts @@ -0,0 +1,227 @@ +import type { Scope } from '@sentry/core'; +import { + addNonEnumerableProperty, + captureException, + getCurrentScope, + getDefaultIsolationScope, + SPAN_STATUS_ERROR, + startNewTrace, + withActiveSpan, +} from '@sentry/core'; +import type { GenAiOptions } from '../core/utils'; +import { skipPiAiProviderIntegrations } from '../pi-ai/providers'; +import { PI_DURABLE_ORIGIN, PI_TASK } from './constants'; +import { instrumentPiModels } from './models'; +import type { PiRuns } from './runs'; +import { createRuns, endRun, observeRunControl, startRun, toSentryConversationId } from './runs'; +import { instrumentTool, instrumentToolPhase } from './tools'; +import type { + PiHarness, + PiHarnessOptions, + PiPhaseHandler, + PiRegistryReader, + PiRegistrySnapshot, + PiTask, + PiTaskRuntime, + PiTool, +} from './types'; +import { bound, withCleanScopes } from './utils'; + +export type PiDurableOptions = GenAiOptions; + +const INSTRUMENTED = Symbol.for('sentry.pi-durable.instrumented'); +const RUNS = Symbol.for('sentry.pi-durable.runs'); + +/** + * Return pi-durable `HarnessOptions` whose model requests, task phases and tool calls are traced. + * The caller's `models` and `registry` objects are wrapped, never modified. The result has the + * caller's options as prototype: pi-durable reads `settings`, `env` and `conversationCreated` at + * every use, so they stay live, getters included. + * + * Spans: + * - `gen_ai.invoke_agent` for each run, a new trace per run, or a child of the tool call that + * created the conversation for a subagent; + * - `gen_ai.chat` for each model request, under its run; + * - `gen_ai.execute_tool` for each tool call, under its run. + * + * The scheduler starts task phases from whichever async context last woke it, often an unrelated + * request. Every phase therefore runs in its run's own scopes and trace, never in the ones it + * inherits. + */ +export function instrumentPiDurableHarnessOptions( + harnessOptions: T, + options: PiDurableOptions = {}, +): T { + if (!harnessOptions || (harnessOptions as Record)[INSTRUMENTED]) { + return harnessOptions; + } + + const runs = createRuns(); + const own = (value: unknown): PropertyDescriptor => ({ value, enumerable: true, configurable: true, writable: true }); + const instrumented: T = Object.create(harnessOptions, { + ...(harnessOptions.models + ? { models: own(instrumentPiModels(harnessOptions.models, options, skipPiAiProviderIntegrations)) } + : {}), + ...(harnessOptions.registry ? { registry: own(instrumentRegistry(harnessOptions.registry, options, runs)) } : {}), + // pi-durable passes failures of extension code here, such as a throwing hook, and keeps going. + onReport: own((error: unknown) => { + captureException(error, { mechanism: { handled: true, type: PI_DURABLE_ORIGIN } }); + harnessOptions.onReport?.(error); + }), + }); + addNonEnumerableProperty(instrumented, INSTRUMENTED, true); + addNonEnumerableProperty(instrumented, RUNS, runs); + return instrumented; +} + +/** + * End the open runs of `harness` when it closes. Closing stops the scheduler without settling the + * inputs in flight, so nothing else would end their `invoke_agent` spans. A run whose ending commit + * is still in flight ends with that commit's status. + */ +export function endRunsOnClose(harnessOptions: PiHarnessOptions, harness: unknown): void { + const runs = (harnessOptions as Record)[RUNS] as PiRuns | undefined; + const closable = harness as PiHarness | undefined; + if (!runs || typeof closable?.subscribeClose !== 'function') { + return; + } + closable.subscribeClose(() => { + for (const run of runs.active.values()) { + endRun(run, runs, run.pendingEnd ? run.pendingEnd.status : { code: SPAN_STATUS_ERROR, message: 'cancelled' }); + } + }); +} + +function instrumentRegistry(registry: PiRegistryReader, options: PiDurableOptions, runs: PiRuns): PiRegistryReader { + // Cached so identity stays stable: the scheduler hands a task over to a new definition whenever + // `snapshot().task(kind)` returns a different object than the one it runs. + const snapshots = new WeakMap(); + const tasks = new WeakMap(); + const tools = new WeakMap(); + + const wrapTool = (tool: PiTool): PiTool => { + let wrapped = tools.get(tool); + if (!wrapped) { + wrapped = instrumentTool(tool, runs, options); + tools.set(tool, wrapped); + } + return wrapped; + }; + + const wrapTask = (task: PiTask | undefined): PiTask | undefined => { + if (!task?.definition?.phases) { + return task; + } + let wrapped = tasks.get(task); + if (!wrapped) { + const { definition } = task; + // pi-durable calls a phase on the `phases` object and `abort` on the definition; keep that + // receiver, so handlers that use `this` work. + const wrapPhase = + (phase: PiPhaseHandler, receiver: unknown): PiPhaseHandler => + (taskRecord, runtime, context) => + runPhase( + definition.name, + observed => phase.call(receiver, taskRecord, observed, context), + runtime, + runs, + wrapTool, + ); + const phases: Record = {}; + for (const [name, phase] of Object.entries(definition.phases)) { + phases[name] = typeof phase === 'function' ? wrapPhase(phase, definition.phases) : phase; + } + wrapped = { + ...task, + definition: { + ...definition, + phases, + ...(typeof definition.abort === 'function' ? { abort: wrapPhase(definition.abort, definition) } : {}), + }, + }; + tasks.set(task, wrapped); + } + return wrapped; + }; + + const wrapSnapshot = (snapshot: PiRegistrySnapshot): PiRegistrySnapshot => { + let wrapped = snapshots.get(snapshot); + if (!wrapped) { + wrapped = new Proxy(snapshot, { + get(target, property) { + if (property === 'task') { + return (name: string) => wrapTask(target.task(name)); + } + if (property === 'tasks') { + return () => target.tasks().map(task => wrapTask(task)); + } + return bound(target, property); + }, + }); + snapshots.set(snapshot, wrapped); + } + return wrapped; + }; + + return new Proxy(registry, { + get(target, property) { + if (property === 'snapshot') { + return () => wrapSnapshot(target.snapshot()); + } + return bound(target, property); + }, + }); +} + +function runPhase( + kind: string, + phase: (runtime: PiTaskRuntime) => Promise, + runtime: PiTaskRuntime, + runs: PiRuns, + wrapTool: (tool: PiTool) => PiTool, +): Promise { + const { conversationId } = runtime; + const isRunWork = kind === PI_TASK.GENERATION || kind === PI_TASK.TOOL; + // A compaction joins the run it was started for. One started while idle has no run, and its + // `chat` span becomes the root of a trace of its own. Tasks of extensions never join a run. + const run = + isRunWork || kind === PI_TASK.COMPACTION + ? (runs.active.get(conversationId) ?? (isRunWork ? startRun(conversationId, runs) : undefined)) + : undefined; + + let observed = runtime; + let finishToolPhase: (() => void) | undefined; + if (run && kind === PI_TASK.GENERATION) { + observed = observeRunControl(runtime, run, runs); + } else if (kind === PI_TASK.TOOL) { + const toolPhase = instrumentToolPhase(runtime, runs, wrapTool); + observed = toolPhase.runtime; + finishToolPhase = toolPhase.finish; + } + + const runInScope = async (scope: Scope): Promise => { + scope.setConversationId(toSentryConversationId(runs, conversationId)); + try { + return await phase(observed); + } catch (error) { + // The scheduler faults a task whose phase throws, and a faulted generation ends its run. A + // throw after an abort mark, or once the Harness closes, is not a fault. + if (!runtime.signal?.aborted) { + captureException(error, { mechanism: { handled: false, type: PI_DURABLE_ORIGIN } }); + if (run && kind === PI_TASK.GENERATION) { + endRun(run, runs, { code: SPAN_STATUS_ERROR, message: 'internal_error' }); + } + } + throw error; + } finally { + finishToolPhase?.(); + } + }; + + if (!run) { + return withCleanScopes(getDefaultIsolationScope().clone(), () => + startNewTrace(() => runInScope(getCurrentScope())), + ); + } + return withCleanScopes(run.isolationScope, () => withActiveSpan(run.span, runInScope)); +} diff --git a/packages/server-utils/src/ai/pi-durable/models.ts b/packages/server-utils/src/ai/pi-durable/models.ts new file mode 100644 index 000000000000..ace332831479 --- /dev/null +++ b/packages/server-utils/src/ai/pi-durable/models.ts @@ -0,0 +1,201 @@ +import type { Span } from '@sentry/core'; +import { + getCurrentScope, + getIsolationScope, + isURLObjectRelative, + parseStringToURLObject, + SEMANTIC_ATTRIBUTE_SENTRY_ORIGIN, + SPAN_STATUS_ERROR, + startSpanManual, + stringify, +} from '@sentry/core'; +import { + GEN_AI_CONVERSATION_ID, + GEN_AI_INPUT_MESSAGES, + GEN_AI_OPERATION_NAME, + GEN_AI_OUTPUT_MESSAGES, + GEN_AI_PROVIDER_NAME, + GEN_AI_REQUEST_MAX_TOKENS, + GEN_AI_REQUEST_MODEL, + GEN_AI_REQUEST_REASONING_LEVEL, + GEN_AI_REQUEST_TEMPERATURE, + GEN_AI_RESPONSE_FINISH_REASONS, + GEN_AI_RESPONSE_ID, + GEN_AI_RESPONSE_MODEL, + GEN_AI_RESPONSE_STREAMING, + GEN_AI_SYSTEM_INSTRUCTIONS, + GEN_AI_TOOL_DEFINITIONS, + SERVER_ADDRESS, + SERVER_PORT, +} from '@sentry/conventions/attributes'; +import type { GenAiOptions } from '../core/utils'; +import { getGenAiSpanOp, resolveAIRecordingOptions } from '../core/utils'; +import type { PiAiContext } from '../pi-ai/messages'; +import { + piAiAssistantMessageToGenAiMessage, + piAiFinishReason, + piAiMessagesToGenAiMessages, + piAiSystemInstructions, + piAiToolDefinitions, +} from '../pi-ai/messages'; +import { setPiAiUsageAttributes } from '../pi-ai/usage'; +import { PI_DURABLE_ORIGIN } from './constants'; +import type { PiAssistantMessage, PiEventStream, PiModel, PiModels, PiStreamOptions } from './types'; + +type ModelCall = (model: PiModel, request: unknown, options?: unknown) => unknown; + +const STREAMING_METHODS = new Set(['stream', 'streamSimple', 'streamDeferred']); +const COMPLETE_METHODS = new Set(['complete', 'completeSimple', 'fetchDeferred']); +// These fetch the answer of a parked request; their second argument is a handle, not the request. +const DEFERRED_METHODS = new Set(['streamDeferred', 'fetchDeferred']); + +/** + * Wrap a pi-ai `Models` so every model request it sends gets a `gen_ai.chat` span. + * + * The proxy hands every other member through bound to the original, so a `Models` instance whose + * methods read `this` keeps working, and the caller's own object is never modified. + */ +export function instrumentPiModels(models: T, options: GenAiOptions, onRequest: () => void): T { + const wrapped = new Map(); + + return new Proxy(models, { + get(target, property) { + const value: unknown = Reflect.get(target, property, target); + if (typeof value !== 'function') { + return value; + } + + const cached = wrapped.get(property); + if (cached) { + return cached; + } + + const original = (value as ModelCall).bind(target); + const method = String(property); + const streaming = STREAMING_METHODS.has(method); + const deferred = DEFERRED_METHODS.has(method); + const instrumented = + streaming || COMPLETE_METHODS.has(method) + ? (model: PiModel, request: unknown, callOptions?: unknown): unknown => { + onRequest(); + return traceModelRequest( + model, + deferred ? undefined : (request as PiAiContext), + deferred ? undefined : (callOptions as PiStreamOptions | undefined), + streaming, + options, + () => original(model, request, callOptions), + ); + } + : original; + + wrapped.set(property, instrumented); + return instrumented; + }, + }); +} + +function traceModelRequest( + model: PiModel, + context: PiAiContext | undefined, + callOptions: PiStreamOptions | undefined, + streaming: boolean, + options: GenAiOptions, + send: () => unknown, +): unknown { + // Read per request: `dataCollection.genAI` belongs to the current client, which can change. + const { recordInputs, recordOutputs } = resolveAIRecordingOptions(options); + const conversationId = + getCurrentScope().getScopeData().conversationId ?? getIsolationScope().getScopeData().conversationId; + + return startSpanManual( + { + name: model.id ? `chat ${model.id}` : 'chat', + op: getGenAiSpanOp('chat'), + attributes: { + [SEMANTIC_ATTRIBUTE_SENTRY_ORIGIN]: PI_DURABLE_ORIGIN, + [GEN_AI_OPERATION_NAME]: 'chat', + ...(model.provider ? { [GEN_AI_PROVIDER_NAME]: model.provider } : {}), + ...(model.id ? { [GEN_AI_REQUEST_MODEL]: model.id } : {}), + ...(conversationId ? { [GEN_AI_CONVERSATION_ID]: conversationId } : {}), + ...getServerAttributes(model), + ...(streaming ? { [GEN_AI_RESPONSE_STREAMING]: true } : {}), + ...(callOptions?.temperature !== undefined ? { [GEN_AI_REQUEST_TEMPERATURE]: callOptions.temperature } : {}), + ...(callOptions?.maxTokens !== undefined ? { [GEN_AI_REQUEST_MAX_TOKENS]: callOptions.maxTokens } : {}), + ...(callOptions?.reasoning ? { [GEN_AI_REQUEST_REASONING_LEVEL]: callOptions.reasoning } : {}), + ...(recordInputs && context ? getRequestContentAttributes(context) : {}), + }, + }, + span => { + let result: unknown; + try { + result = send(); + } catch (error) { + span.setStatus({ code: SPAN_STATUS_ERROR, message: 'internal_error' }); + span.end(); + throw error; + } + + // pi-ai never rejects a request: a provider failure or an abort settles as a message with + // `stopReason` `error` or `aborted`. The rejection handler only covers a broken provider. + const settled = streaming ? (result as PiEventStream).result() : (result as Promise); + settled.then( + message => endChatSpan(span, message, recordOutputs), + () => { + span.setStatus({ code: SPAN_STATUS_ERROR, message: 'internal_error' }); + span.end(); + }, + ); + + return result; + }, + ); +} + +function getServerAttributes(model: PiModel): Record { + const url = model.baseUrl ? parseStringToURLObject(model.baseUrl) : undefined; + if (!url || isURLObjectRelative(url) || !url.hostname) { + return {}; + } + return { [SERVER_ADDRESS]: url.hostname, ...(url.port ? { [SERVER_PORT]: Number(url.port) } : {}) }; +} + +function getRequestContentAttributes(context: PiAiContext): Record { + const messages = piAiMessagesToGenAiMessages(context.messages); + const tools = piAiToolDefinitions(context); + return { + [GEN_AI_SYSTEM_INSTRUCTIONS]: piAiSystemInstructions(context), + [GEN_AI_INPUT_MESSAGES]: messages.length ? stringify(messages) : undefined, + [GEN_AI_TOOL_DEFINITIONS]: tools.length ? stringify(tools) : undefined, + }; +} + +function endChatSpan(span: Span, message: PiAssistantMessage, recordOutputs: boolean): void { + const responseModel = message.responseModel ?? message.model; + if (responseModel) { + span.setAttribute(GEN_AI_RESPONSE_MODEL, responseModel); + } + if (message.responseId) { + span.setAttribute(GEN_AI_RESPONSE_ID, message.responseId); + } + + const finishReason = piAiFinishReason(message.stopReason); + if (finishReason) { + span.setAttribute(GEN_AI_RESPONSE_FINISH_REASONS, stringify([finishReason])); + } + + const failed = message.stopReason === 'error'; + setPiAiUsageAttributes(span, message.usage, failed || message.stopReason === 'deferred'); + + const output = recordOutputs ? piAiAssistantMessageToGenAiMessage(message, finishReason) : undefined; + if (output) { + span.setAttribute(GEN_AI_OUTPUT_MESSAGES, stringify([output], String)); + } + + if (failed) { + span.setStatus({ code: SPAN_STATUS_ERROR, message: 'internal_error' }); + } else if (message.stopReason === 'aborted') { + span.setStatus({ code: SPAN_STATUS_ERROR, message: 'cancelled' }); + } + span.end(); +} diff --git a/packages/server-utils/src/ai/pi-durable/runs.ts b/packages/server-utils/src/ai/pi-durable/runs.ts new file mode 100644 index 000000000000..9067ab961600 --- /dev/null +++ b/packages/server-utils/src/ai/pi-durable/runs.ts @@ -0,0 +1,217 @@ +import type { Scope, Span, SpanStatus } from '@sentry/core'; +import { + debug, + getDefaultIsolationScope, + SEMANTIC_ATTRIBUTE_SENTRY_ORIGIN, + SPAN_STATUS_ERROR, + startInactiveSpan, + startNewTrace, + uuid4, +} from '@sentry/core'; +import { GEN_AI_CONVERSATION_ID, GEN_AI_OPERATION_NAME } from '@sentry/conventions/attributes'; +import { DEBUG_BUILD } from '../../debug-build'; +import { getGenAiSpanOp } from '../core/utils'; +import { MAX_TRACKED_PI_RUNS, PI_DURABLE_ORIGIN, PI_LIVE_DOC_KIND } from './constants'; +import type { PiCommitChange, PiDocToken, PiLiveState, PiSettlement, PiTaskRuntime } from './types'; +import { bound, withCleanScopes } from './utils'; + +// Unanswered reasons that mean the run was stopped, not that it failed. +const CANCELLED_REASONS = new Set(['aborted', 'reset']); + +/** + * One run: the work from an input to its answer. It spans several `pi.generation` tasks (run control + * moves to a new one after every tool round) and their `pi.tool` tasks, so it is tracked per + * conversation, which runs at most one run at a time. + */ +export interface PiRun { + conversationId: unknown; + span: Span; + isolationScope: Scope; + /** The run's first input submission, which identifies the run in `pi.live`. */ + firstInput?: unknown; + /** + * The status the run ends with, set while the commit that ends it is in flight. `harness.close()` + * can begin in that gap, because `wait()` resolves as soon as the commit settles the inputs. + */ + pendingEnd?: { status?: SpanStatus }; +} + +/** Open runs of one Harness, keyed by the pi-durable conversation id. */ +export interface PiRuns { + /** + * Random id of this Harness, prefixed to `gen_ai.conversation.id`. pi-durable numbers + * conversations per storage from 1, so the bare id would merge the root conversations of every + * Harness into one Sentry conversation. + */ + harnessId: string; + active: Map; + /** + * `execute_tool` spans of tool calls that created a conversation they own, keyed by that + * conversation. The child's first run becomes a child of the call, which is how a subagent's run + * joins the trace of the run that delegated to it. + */ + owners: Map; + /** + * Tool calls of running `pi.tool` phases, keyed by task id. Their spans end when the phase ends, + * after it has committed the result the model receives. + */ + toolCalls: Map; +} + +/** One `execute()` of a tool. */ +export interface PiToolCall { + span: Span; + callId?: string; + recordOutputs: boolean; + /** Set once the span has a status other than ok, so the committed result does not overwrite it. */ + failed?: boolean; + /** When `execute()` settled, used as the span's end time. */ + endTimestamp?: number; +} + +export function createRuns(): PiRuns { + return { + harnessId: uuid4(), + active: new Map(), + owners: new Map(), + toolCalls: new Map(), + }; +} + +export function toSentryConversationId(runs: PiRuns, conversationId: unknown): string { + return `${runs.harnessId}:${String(conversationId)}`; +} + +/** Start a run of `conversationId`: a new trace, or a child of the tool call that owns the conversation. */ +export function startRun(conversationId: unknown, runs: PiRuns): PiRun { + if (runs.active.size >= MAX_TRACKED_PI_RUNS) { + const oldest = runs.active.keys().next().value; + runs.active.get(oldest)?.span.end(); + runs.active.delete(oldest); + } + + const owner = runs.owners.get(conversationId); + runs.owners.delete(conversationId); + + const isolationScope = getDefaultIsolationScope().clone(); + const startRunSpan = (): Span => + startInactiveSpan({ + name: 'invoke_agent', + op: getGenAiSpanOp('invoke_agent'), + ...(owner ? { parentSpan: owner } : {}), + attributes: { + [SEMANTIC_ATTRIBUTE_SENTRY_ORIGIN]: PI_DURABLE_ORIGIN, + [GEN_AI_OPERATION_NAME]: 'invoke_agent', + [GEN_AI_CONVERSATION_ID]: toSentryConversationId(runs, conversationId), + }, + }); + const span = withCleanScopes(isolationScope, () => (owner ? startRunSpan() : startNewTrace(startRunSpan))); + + const run: PiRun = { conversationId, span, isolationScope }; + runs.active.set(conversationId, run); + DEBUG_BUILD && debug.log(`[pi-durable] run of conversation ${String(conversationId)} started`); + return run; +} + +export function endRun(run: PiRun, runs: PiRuns, status?: SpanStatus): void { + if (runs.active.get(run.conversationId) !== run) { + return; + } + runs.active.delete(run.conversationId); + DEBUG_BUILD && debug.log(`[pi-durable] run of conversation ${String(run.conversationId)} ended`, status ?? 'ok'); + if (status) { + run.span.setStatus(status); + } + run.span.end(); +} + +/** + * Watch the generation task's commits for the end of its run. + * + * pi-durable keeps run control in the conversation's `pi.live` document: `run` names the task that + * owns the run and the run's inputs, it moves to the next generation after every tool round, and it + * is removed, in the same commit that settles the inputs, when the run ends. Every commit that moves + * or ends a run edits `pi.live` through `tx.doc()`, so reading that draft after the commit's change + * function returns shows whether the run survived the commit. + */ +export function observeRunControl(runtime: PiTaskRuntime, run: PiRun, runs: PiRuns): PiTaskRuntime { + const commit = (change: PiCommitChange, context: unknown): Promise => { + let ended = false; + let settlement: PiSettlement | undefined; + + const observedChange: PiCommitChange = async (tx, current) => { + // A commit can run its change again, so only the attempt that committed counts. + ended = false; + settlement = undefined; + let live: Promise | undefined; + const observedTx = new Proxy(tx, { + get(target, property) { + const value = bound(target, property); + if (typeof value !== 'function') { + return value; + } + if (property === 'doc') { + return (token: PiDocToken, id: unknown, ...rest: unknown[]) => { + const draft = value(token, id, ...rest); + if (token?.definition?.kind === PI_LIVE_DOC_KIND && id === run.conversationId) { + live = Promise.resolve(draft); + } + return draft; + }; + } + if (property === 'settleSubmission') { + return (id: unknown, submissionSettlement: PiSettlement, ...rest: unknown[]) => { + if (run.firstInput === undefined || id === run.firstInput) { + settlement = submissionSettlement; + } + return value(id, submissionSettlement, ...rest); + }; + } + return value; + }, + }); + + const result = await change(observedTx, current); + const draft = (await live) as PiLiveState | undefined; + if (draft) { + const firstInput = draft.run?.inputs?.[0]; + if (firstInput === undefined) { + ended = true; + } else if (run.firstInput === undefined) { + run.firstInput = firstInput; + } else if (firstInput !== run.firstInput) { + // The run ended and a queued input started the next one in the same commit. + ended = true; + } + } + run.pendingEnd = ended ? { status: toRunStatus(settlement) } : undefined; + return result; + }; + + return runtime.commit(observedChange, context).then( + () => { + if (ended) { + endRun(run, runs, toRunStatus(settlement)); + } + }, + error => { + run.pendingEnd = undefined; + throw error; + }, + ); + }; + + return new Proxy(runtime, { + get(target, property) { + return property === 'commit' ? commit : bound(target, property); + }, + }); +} + +function toRunStatus(settlement: PiSettlement | undefined): SpanStatus | undefined { + if (settlement?.status !== 'unanswered') { + return undefined; + } + const reason = settlement.reason ?? 'internal_error'; + return { code: SPAN_STATUS_ERROR, message: CANCELLED_REASONS.has(reason) ? 'cancelled' : reason }; +} diff --git a/packages/server-utils/src/ai/pi-durable/tools.ts b/packages/server-utils/src/ai/pi-durable/tools.ts new file mode 100644 index 000000000000..2fbc287d4a2c --- /dev/null +++ b/packages/server-utils/src/ai/pi-durable/tools.ts @@ -0,0 +1,250 @@ +import type { Span } from '@sentry/core'; +import { + captureException, + isObjectLike, + SEMANTIC_ATTRIBUTE_SENTRY_ORIGIN, + SPAN_STATUS_ERROR, + startSpanManual, + stringify, + timestampInSeconds, +} from '@sentry/core'; +import { + GEN_AI_CONVERSATION_ID, + GEN_AI_OPERATION_NAME, + GEN_AI_TOOL_CALL_ARGUMENTS, + GEN_AI_TOOL_CALL_RESULT, + GEN_AI_TOOL_DESCRIPTION, + GEN_AI_TOOL_NAME, +} from '@sentry/conventions/attributes'; +import { GEN_AI_TOOL_CALL_ID_ATTRIBUTE } from '../core/gen-ai-attributes'; +import type { GenAiOptions } from '../core/utils'; +import { getGenAiSpanOp, resolveAIRecordingOptions } from '../core/utils'; +import { piAiContentToString } from '../pi-ai/messages'; +import { MAX_TRACKED_PI_RUNS, PI_DURABLE_ORIGIN, PI_TOOL_RESULT_ENTRY_KIND } from './constants'; +import type { PiRuns, PiToolCall } from './runs'; +import { toSentryConversationId } from './runs'; +import type { + PiAgent, + PiCommitChange, + PiEntryToken, + PiTaskRuntime, + PiTool, + PiToolExecutionApi, + PiToolResultMessage, +} from './types'; +import { bound } from './utils'; + +/** + * The tools pi-durable's `createBashTool`, `createReadTool`, `createEditTool` and `createWriteTool` + * factories return. They throw to report expected failures to the model, such as a command that + * exits non-zero, so their throws are not captured. Filled by the integration from the factories' + * channels, which also covers the tools an app registers in an extension of its own. + */ +const BUILT_IN_TOOLS = new WeakSet(); + +export function markBuiltInTool(tool: unknown): void { + if (isObjectLike(tool)) { + BUILT_IN_TOOLS.add(tool); + } +} + +/** + * Prepare the runtime of a `pi.tool` phase. The tool task resolves the called tool from its agent, + * so the runtime hands out traced tools, and it collects the result entries the phase commits. + * `finish` ends the phase's tool spans with the result the model receives, which pi-durable builds + * after `execute()` returns: from `api.output()`, through `afterTool` hooks and output limits. + */ +export function instrumentToolPhase( + runtime: PiTaskRuntime, + runs: PiRuns, + wrapTool: (tool: PiTool) => PiTool, +): { runtime: PiTaskRuntime; finish: () => void } { + const results = new Map(); + runs.toolCalls.set(runtime.taskId, []); + + const agent = (context: unknown): Promise => + runtime + .agent(context) + .then(resolved => (resolved?.tools ? { ...resolved, tools: resolved.tools.map(wrapTool) } : resolved)); + + const commit = (change: PiCommitChange, context: unknown): Promise => { + let committed: PiToolResultMessage[] = []; + const observedChange: PiCommitChange = (tx, current) => { + // A commit can run its change again, so only the attempt that committed counts. + committed = []; + return change(observeToolResultEntries(tx, committed), current); + }; + return runtime.commit(observedChange, context).then(() => { + for (const message of committed) { + results.set(message.toolCallId, message); + } + }); + }; + + return { + runtime: new Proxy(runtime, { + get(target, property) { + if (property === 'agent') { + return agent; + } + return property === 'commit' ? commit : bound(target, property); + }, + }), + finish: () => { + for (const call of runs.toolCalls.get(runtime.taskId) ?? []) { + finishToolCall(call, call.callId !== undefined ? results.get(call.callId) : undefined); + } + runs.toolCalls.delete(runtime.taskId); + }, + }; +} + +/** Wrap a tool so each call gets a `gen_ai.execute_tool` span. */ +export function instrumentTool(tool: PiTool, runs: PiRuns, options: GenAiOptions): PiTool { + const execute = (args: unknown, api: PiToolExecutionApi, context: unknown): Promise => { + const { recordInputs, recordOutputs } = resolveAIRecordingOptions(options); + const phaseCalls = runs.toolCalls.get(api?.taskId); + // A tool that runs another tool of its agent passes its own `api` on, so the inner call shares + // the call id. Only the outer call gets the committed result; the inner one ends with its own. + const calls = phaseCalls?.some(call => call.callId !== undefined && call.callId === api?.callId) + ? undefined + : phaseCalls; + + return startSpanManual( + { + name: `execute_tool ${tool.name}`, + op: getGenAiSpanOp('execute_tool'), + attributes: { + [SEMANTIC_ATTRIBUTE_SENTRY_ORIGIN]: PI_DURABLE_ORIGIN, + [GEN_AI_OPERATION_NAME]: 'execute_tool', + [GEN_AI_TOOL_NAME]: tool.name, + ...(tool.description ? { [GEN_AI_TOOL_DESCRIPTION]: tool.description } : {}), + ...(api?.callId ? { [GEN_AI_TOOL_CALL_ID_ATTRIBUTE]: api.callId } : {}), + ...(api?.conversationId !== undefined + ? { [GEN_AI_CONVERSATION_ID]: toSentryConversationId(runs, api.conversationId) } + : {}), + ...(recordInputs && args !== undefined ? { [GEN_AI_TOOL_CALL_ARGUMENTS]: stringify(args) } : {}), + }, + }, + async span => { + const call: PiToolCall = { span, callId: api?.callId, recordOutputs }; + calls?.push(call); + try { + const result = await tool.execute(args, trackOwnedConversations(api, span, runs), context); + call.endTimestamp = timestampInSeconds(); + if (!calls) { + // No result entry follows a call outside a tool task phase, nor a nested call. + finishToolCall(call, { content: result?.content, isError: result?.isError }); + } + return result; + } catch (error) { + call.endTimestamp = timestampInSeconds(); + call.failed = true; + // pi-durable turns the throw into an error result for the model, so nothing propagates to + // the global handlers. An aborted call is not a failure. + if ((context as { abortSignal?: AbortSignal } | undefined)?.abortSignal?.aborted) { + span.setStatus({ code: SPAN_STATUS_ERROR, message: 'cancelled' }); + } else { + span.setStatus({ code: SPAN_STATUS_ERROR, message: 'internal_error' }); + if (!BUILT_IN_TOOLS.has(tool)) { + captureException(error, { mechanism: { handled: true, type: PI_DURABLE_ORIGIN } }); + } + } + if (!calls) { + span.end(call.endTimestamp); + } + throw error; + } + }, + ); + }; + + // The proxy's target is an empty object with the tool as prototype, not the tool itself, so a + // frozen tool does not break the proxy invariants. Members are read bound to the tool, so a + // method that reads a private field works. + return new Proxy(Object.create(tool) as PiTool, { + get(_target, property) { + return property === 'execute' ? execute : bound(tool, property); + }, + }); +} + +function finishToolCall(call: PiToolCall, message: PiToolResultMessage | undefined): void { + if (message) { + const result = call.recordOutputs ? piAiContentToString(message.content) : undefined; + if (result) { + call.span.setAttribute(GEN_AI_TOOL_CALL_RESULT, result); + } + if (message.isError && !call.failed) { + call.span.setStatus({ code: SPAN_STATUS_ERROR, message: 'internal_error' }); + } + } + call.span.end(call.endTimestamp); +} + +function observeToolResultEntries(tx: object, committed: PiToolResultMessage[]): object { + return new Proxy(tx, { + get(target, property) { + const value = bound(target, property); + if (property !== 'appendEntry' || typeof value !== 'function') { + return value; + } + return ( + token: PiEntryToken, + conversationId: unknown, + draft: { model?: unknown[] } | undefined, + ...rest: unknown[] + ) => { + const message = token?.kind === PI_TOOL_RESULT_ENTRY_KIND ? draft?.model?.[0] : undefined; + if (isObjectLike(message)) { + committed.push(message); + } + return value(token, conversationId, draft, ...rest); + }; + }, + }); +} + +/** + * Record the conversations a tool call creates for itself, the subagent pattern: `api.commit()` with + * `tx.createConversation()` or `tx.forkConversation()` and `ownership: { kind: 'task', taskId: api.taskId }`. + */ +function trackOwnedConversations(api: PiToolExecutionApi, span: Span, runs: PiRuns): PiToolExecutionApi { + if (typeof api?.commit !== 'function') { + return api; + } + const commit = api.commit.bind(api); + + const trackCreatedConversations = (tx: object): object => + new Proxy(tx, { + get(target, property) { + const value = bound(target, property); + if ((property !== 'createConversation' && property !== 'forkConversation') || typeof value !== 'function') { + return value; + } + return async (...args: unknown[]) => { + const record = (await value(...args)) as { id?: unknown } | undefined; + const ownership = args.find( + (arg): arg is { ownership: { kind?: string } } => isObjectLike(arg) && isObjectLike(arg.ownership), + )?.ownership; + if (ownership?.kind === 'task' && record?.id !== undefined) { + if (runs.owners.size >= MAX_TRACKED_PI_RUNS) { + runs.owners.delete(runs.owners.keys().next().value); + } + runs.owners.set(record.id, span); + } + return record; + }; + }, + }); + + return new Proxy(api, { + get(target, property) { + if (property !== 'commit') { + return bound(target, property); + } + return (change: (tx: object) => unknown, context: unknown) => + commit((tx: object) => change(trackCreatedConversations(tx)), context); + }, + }); +} diff --git a/packages/server-utils/src/ai/pi-durable/types.ts b/packages/server-utils/src/ai/pi-durable/types.ts new file mode 100644 index 000000000000..6244b18af1e5 --- /dev/null +++ b/packages/server-utils/src/ai/pi-durable/types.ts @@ -0,0 +1,141 @@ +/** + * Structural subsets of the `@earendil-works/pi-durable` and `@earendil-works/pi-ai` types the + * instrumentation reads. The SDK does not depend on either package, so only the fields used here + * are declared, and every one of them is treated as possibly absent at runtime. + */ + +import type { PiAiContext } from '../pi-ai/messages'; +import type { PiAiUsage } from '../pi-ai/usage'; + +export interface PiAssistantMessage { + role?: string; + content?: unknown[]; + provider?: string; + model?: string; + responseModel?: string; + responseId?: string; + usage?: PiAiUsage; + stopReason?: string; + errorMessage?: string; +} + +export interface PiModel { + id?: string; + provider?: string; + baseUrl?: string; +} + +export interface PiStreamOptions { + temperature?: number; + maxTokens?: number; + reasoning?: string; +} + +export interface PiEventStream { + result(): Promise; +} + +/** + * The `Models` methods that send a model request. Everything else passes through untouched. The + * deferred methods fetch the answer of a request the provider parked, so they take a handle in + * place of the request context. + */ +export interface PiModels { + stream?(model: PiModel, context: PiAiContext, options?: PiStreamOptions): PiEventStream; + streamSimple?(model: PiModel, context: PiAiContext, options?: PiStreamOptions): PiEventStream; + complete?(model: PiModel, context: PiAiContext, options?: PiStreamOptions): Promise; + completeSimple?(model: PiModel, context: PiAiContext, options?: PiStreamOptions): Promise; + streamDeferred?(model: PiModel, handle: unknown, options?: unknown): PiEventStream; + fetchDeferred?(model: PiModel, handle: unknown, options?: unknown): Promise; +} + +export interface PiToolExecutionResult { + content?: unknown[]; + isError?: boolean; +} + +export interface PiToolExecutionApi { + taskId?: unknown; + conversationId?: unknown; + callId?: string; + commit?(change: (tx: object) => unknown, context: unknown): Promise; +} + +export interface PiTool { + name: string; + description?: string; + execute(args: unknown, api: PiToolExecutionApi, context: unknown): Promise; +} + +export interface PiAgent { + tools?: readonly PiTool[]; +} + +/** The tool result the model receives, as pi-durable commits it in a `pi.tool-result` entry. */ +export interface PiToolResultMessage { + toolCallId?: unknown; + content?: unknown[]; + isError?: boolean; +} + +export interface PiDocToken { + definition?: { kind?: string }; +} + +export interface PiEntryToken { + kind?: string; +} + +export interface PiSettlement { + status?: string; + reason?: string; +} + +/** The `pi.live` document; `run` is present exactly while the conversation is busy. */ +export interface PiLiveState { + run?: { inputs?: unknown[] }; +} + +export type PiCommitChange = (tx: object, current: unknown) => unknown; + +/** pi-durable ids (conversations, tasks, submissions) are numbers, typed as `unknown` here. */ +export interface PiTaskRuntime { + readonly taskId: unknown; + readonly conversationId: unknown; + readonly signal?: AbortSignal; + agent(context: unknown): Promise; + commit(change: PiCommitChange, context: unknown): Promise; +} + +export type PiPhaseHandler = (task: unknown, runtime: PiTaskRuntime, context: unknown) => Promise; + +export interface PiTaskDefinition { + name: string; + phases: Record; + abort?: PiPhaseHandler; +} + +export interface PiTask { + definition: PiTaskDefinition; +} + +export interface PiRegistrySnapshot { + task(name: string): PiTask | undefined; + tasks(): readonly PiTask[]; +} + +export interface PiRegistryReader { + snapshot(): PiRegistrySnapshot; + subscribe(listener: () => void): () => void; +} + +export interface PiHarnessOptions { + models?: PiModels; + registry?: PiRegistryReader; + onReport?: (error: unknown) => void; +} + +export interface PiHarness { + /** Calls `listener` when `close()` begins. */ + subscribeClose?(listener: () => void): () => void; +} diff --git a/packages/server-utils/src/ai/pi-durable/utils.ts b/packages/server-utils/src/ai/pi-durable/utils.ts new file mode 100644 index 000000000000..f7d1fc807728 --- /dev/null +++ b/packages/server-utils/src/ai/pi-durable/utils.ts @@ -0,0 +1,25 @@ +import type { Scope } from '@sentry/core'; +import { getClient, getDefaultCurrentScope, withIsolationScope, withScope } from '@sentry/core'; + +/** Read a member with `this` bound to the original, so getters and methods that use `this` keep working. */ +export function bound(target: object, property: PropertyKey): unknown { + const value: unknown = Reflect.get(target, property, target); + return typeof value === 'function' ? (value as (...args: unknown[]) => unknown).bind(target) : value; +} + +/** + * Run `callback` with `isolationScope` and a fresh copy of the default current scope. A forked + * scope would carry over the data of whatever async context last woke the scheduler, such as + * another conversation's id. Only the client of that context is kept: an SDK that initializes + * inside each request, such as `@sentry/cloudflare`, binds it to the request's scope alone. + */ +export function withCleanScopes(isolationScope: Scope, callback: () => T): T { + const client = getClient(); + return withIsolationScope(isolationScope, () => { + const scope = getDefaultCurrentScope().clone(); + if (client) { + scope.setClient(client); + } + return withScope(scope, callback); + }); +} diff --git a/packages/server-utils/src/index.ts b/packages/server-utils/src/index.ts index c162dadb1629..446812f8eed5 100644 --- a/packages/server-utils/src/index.ts +++ b/packages/server-utils/src/index.ts @@ -50,6 +50,7 @@ export { langChainIntegration } from './integrations/langchain'; export { langGraphIntegration } from './integrations/langgraph'; export { createFlueInstrumentation } from './ai/flue'; export { flueIntegration } from './integrations/flue'; +export { piDurableIntegration } from './integrations/pi-durable'; export type { FlueOptions } from './ai/flue'; export { mastraIntegration } from './integrations/mastra'; export { mcpServerIntegration } from './integrations/mcp-server'; diff --git a/packages/server-utils/src/integrations/index.ts b/packages/server-utils/src/integrations/index.ts index df8d488b9d28..a09d590e1121 100644 --- a/packages/server-utils/src/integrations/index.ts +++ b/packages/server-utils/src/integrations/index.ts @@ -15,6 +15,7 @@ import { mongooseIntegration } from './mongoose'; import { lruMemoizerIntegration } from './lru-memoizer'; import { langChainIntegration } from './langchain'; import { langGraphIntegration } from './langgraph'; +import { piDurableIntegration } from './pi-durable'; import { mastraIntegration } from './mastra'; import { mcpServerIntegration } from './mcp-server'; import { vercelAIIntegration } from './vercel-ai'; @@ -59,6 +60,7 @@ export function getTracingIntegrations(): Integration[] { langChainIntegration(), langGraphIntegration(), mastraIntegration(), + piDurableIntegration(), vercelAIIntegration(), openAIIntegration(), anthropicAIIntegration(), diff --git a/packages/server-utils/src/integrations/pi-durable.ts b/packages/server-utils/src/integrations/pi-durable.ts new file mode 100644 index 000000000000..94c91edb00c2 --- /dev/null +++ b/packages/server-utils/src/integrations/pi-durable.ts @@ -0,0 +1,76 @@ +import * as diagnosticsChannel from '../utils/diagnosticsChannel'; +import type { IntegrationFn } from '@sentry/core'; +import { debug, defineIntegration, isObjectLike } from '@sentry/core'; +import type { PiDurableOptions } from '../ai/pi-durable'; +import { endRunsOnClose, instrumentPiDurableHarnessOptions } from '../ai/pi-durable'; +import { PI_DURABLE_INTEGRATION_NAME } from '../ai/pi-durable/constants'; +import { markBuiltInTool } from '../ai/pi-durable/tools'; +import type { PiHarnessOptions } from '../ai/pi-durable/types'; +import { DEBUG_BUILD } from '../debug-build'; +import { CHANNELS } from '../orchestrion/channels'; +import { piDurableModuleNames } from '../orchestrion/config/pi-durable'; +import { invokeOrchestrionInstrumentation } from '../orchestrion/instrumentation'; + +interface HarnessOpenChannelContext { + arguments: unknown[]; + result?: unknown; +} + +interface CodingToolChannelContext { + result?: unknown; +} + +const _piDurableIntegration = ((options: PiDurableOptions = {}) => { + return { + name: PI_DURABLE_INTEGRATION_NAME, + setup(client) { + // Subscribing opens no span; spans open later, inside task phases. + invokeOrchestrionInstrumentation(client, piDurableModuleNames, instrumentPiDurable, [options], { + requiresTracingChannelBinding: false, + }); + }, + }; +}) satisfies IntegrationFn; + +function instrumentPiDurable(options: PiDurableOptions): void { + const harnessOpen = diagnosticsChannel.tracingChannel(CHANNELS.PI_DURABLE_HARNESS_OPEN); + + harnessOpen.start.subscribe(message => { + const args = (message as HarnessOpenChannelContext).arguments; + try { + if (args && isObjectLike(args[1])) { + args[1] = instrumentPiDurableHarnessOptions(args[1] as PiHarnessOptions, options); + } + } catch (error) { + DEBUG_BUILD && debug.error('[instrumentation:pi-durable] failed to instrument Harness options', error); + } + }); + + harnessOpen.asyncEnd.subscribe(message => { + const { arguments: args, result } = message as HarnessOpenChannelContext; + try { + if (args && isObjectLike(args[1])) { + endRunsOnClose(args[1], result); + } + } catch (error) { + DEBUG_BUILD && debug.error('[instrumentation:pi-durable] failed to observe Harness close', error); + } + }); + + diagnosticsChannel + .tracingChannel(CHANNELS.PI_DURABLE_CODING_TOOL) + .end.subscribe(message => markBuiltInTool((message as CodingToolChannelContext).result)); +} + +/** + * Diagnostics-channel-based integration for pi-durable (`@earendil-works/pi-durable` >= 1.0.0 < 2, + * whose API is experimental). Subscribes to the `orchestrion:@earendil-works/pi-durable:harnessOpen` + * channel injected into `Harness.open()` and to the `codingTool` channel injected into the built-in + * tool factories, so it requires the Sentry runtime hook or bundler plugin. + * + * Traces one `gen_ai.invoke_agent` span per run, with its model requests as `gen_ai.chat` and its + * tool calls as `gen_ai.execute_tool` children. From the first run on, the `openai`, + * `@anthropic-ai/sdk` and `@google/genai` integrations stop reporting requests in the whole process, + * because pi-ai sends its requests through those clients. + */ +export const piDurableIntegration = defineIntegration(_piDurableIntegration); diff --git a/packages/server-utils/src/orchestrion/channels.ts b/packages/server-utils/src/orchestrion/channels.ts index a6ab33d4d4c6..4211c8c88d5c 100644 --- a/packages/server-utils/src/orchestrion/channels.ts +++ b/packages/server-utils/src/orchestrion/channels.ts @@ -27,6 +27,7 @@ import { mysqlChannels } from './config/mysql'; import { nestjsChannels } from './config/nestjs'; import { openaiChannels } from './config/openai'; import { pgChannels } from './config/pg'; +import { piDurableChannels } from './config/pi-durable'; import { postgresJsChannels } from './config/postgres'; import { prismaChannels } from './config/prisma'; import { redisChannels } from './config/redis'; @@ -82,6 +83,7 @@ export const CHANNELS = { ...nestjsChannels, ...openaiChannels, ...pgChannels, + ...piDurableChannels, ...postgresJsChannels, ...prismaChannels, ...redisChannels, diff --git a/packages/server-utils/src/orchestrion/config/channel-integration-definitions.ts b/packages/server-utils/src/orchestrion/config/channel-integration-definitions.ts index 51d25ab28328..b12f44808610 100644 --- a/packages/server-utils/src/orchestrion/config/channel-integration-definitions.ts +++ b/packages/server-utils/src/orchestrion/config/channel-integration-definitions.ts @@ -48,6 +48,7 @@ export const CHANNEL_INTEGRATION_DEFINITIONS = [ { exportName: 'langGraphIntegration', modules: ['@langchain/langgraph'] }, { exportName: 'mastraIntegration', modules: ['@mastra/core'] }, { exportName: 'flueIntegration', modules: ['@flue/runtime'] }, + { exportName: 'piDurableIntegration', modules: ['@earendil-works/pi-durable'] }, { exportName: 'mcpServerIntegration', modules: ['@modelcontextprotocol/server', '@modelcontextprotocol/sdk'] }, { exportName: 'awsIntegration', modules: ['@aws-sdk/smithy-client', '@smithy/core', '@smithy/smithy-client'] }, { exportName: 'firebaseIntegration', modules: ['@firebase/firestore', 'firebase-functions'] }, diff --git a/packages/server-utils/src/orchestrion/config/index.ts b/packages/server-utils/src/orchestrion/config/index.ts index f8dded8a6a6d..2c6a7835db1a 100644 --- a/packages/server-utils/src/orchestrion/config/index.ts +++ b/packages/server-utils/src/orchestrion/config/index.ts @@ -23,6 +23,7 @@ import { langgraphConfig } from './langgraph'; import { lruMemoizerConfig } from './lru-memoizer'; import { flueConfig } from './flue'; import { mastraConfig } from './mastra'; +import { piDurableConfig } from './pi-durable'; import { mcpServerConfig } from './mcp-server'; import { mistralConfig } from './mistral'; import { mongodbConfig } from './mongodb'; @@ -88,6 +89,7 @@ export const SENTRY_INSTRUMENTATIONS: InstrumentationConfig[] = [ ...nestjsConfig, ...openaiConfig, ...pgConfig, + ...piDurableConfig, ...postgresJsConfig, ...prismaConfig, ...redisConfig, diff --git a/packages/server-utils/src/orchestrion/config/pi-durable.ts b/packages/server-utils/src/orchestrion/config/pi-durable.ts new file mode 100644 index 000000000000..c7dcaa902289 --- /dev/null +++ b/packages/server-utils/src/orchestrion/config/pi-durable.ts @@ -0,0 +1,35 @@ +import type { InstrumentationConfig } from '../apmTypes'; + +import { getModuleNames } from './module-names'; + +const module = { name: '@earendil-works/pi-durable', versionRange: '>=1.0.0 <2.0.0' }; + +// `Harness.open(storage, options, context)` is a method of the exported `Harness` object literal, the +// one entry point to a Harness. Its `options` carry the `models` and `registry` every model request +// and task phase goes through, so wrapping them at `start` covers the whole Harness. +export const piDurableConfig = [ + { + channelName: 'harnessOpen', + module: { ...module, filePath: 'dist/harness/harness.js' }, + functionQuery: { methodName: 'open', kind: 'Async' as const }, + }, + // The factories of the built-in coding tools. Their tools throw to report expected failures to the + // model, so the integration must know them wherever an app registers them. + ...[ + ['bash', 'createBashTool'], + ['read', 'createReadTool'], + ['edit', 'createEditTool'], + ['write', 'createWriteTool'], + ].map(([file, functionName]) => ({ + channelName: 'codingTool', + module: { ...module, filePath: `dist/tools/${file}.js` }, + functionQuery: { functionName: functionName as string, kind: 'Sync' as const }, + })), +] satisfies InstrumentationConfig[]; + +export const piDurableModuleNames = getModuleNames(piDurableConfig); + +export const piDurableChannels = { + PI_DURABLE_HARNESS_OPEN: 'orchestrion:@earendil-works/pi-durable:harnessOpen', + PI_DURABLE_CODING_TOOL: 'orchestrion:@earendil-works/pi-durable:codingTool', +} as const; diff --git a/packages/server-utils/test/ai/lib/tracing/pi-durable.test.ts b/packages/server-utils/test/ai/lib/tracing/pi-durable.test.ts new file mode 100644 index 000000000000..3a4a9acaea49 --- /dev/null +++ b/packages/server-utils/test/ai/lib/tracing/pi-durable.test.ts @@ -0,0 +1,840 @@ +import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import type { Event, Span } from '@sentry/core'; +import { + _INTERNAL_clearAiProviderSkips, + _INTERNAL_shouldSkipAiProviderWrapping, + _INTERNAL_skipAiProviderWrapping, + getCurrentScope, + getMainCarrier, + setCurrentClient, + spanToJSON, + spanToStaticSpanJSON, + withScope, +} from '@sentry/core'; +import { endRunsOnClose, instrumentPiDurableHarnessOptions } from '../../../../src/ai/pi-durable'; +import { MAX_TRACKED_PI_RUNS } from '../../../../src/ai/pi-durable/constants'; +import { createRuns, startRun } from '../../../../src/ai/pi-durable/runs'; +import { instrumentTool, markBuiltInTool } from '../../../../src/ai/pi-durable/tools'; +import type { + PiCommitChange, + PiEventStream, + PiHarnessOptions, + PiLiveState, + PiModels, + PiRegistryReader, + PiSettlement, + PiTask, + PiTaskRuntime, + PiTool, + PiToolExecutionApi, + PiToolExecutionResult, +} from '../../../../src/ai/pi-durable/types'; +import { ANTHROPIC_AI_INTEGRATION_NAME } from '../../../../src/ai/anthropic-ai/constants'; +import { GOOGLE_GENAI_INTEGRATION_NAME } from '../../../../src/ai/google-genai/constants'; +import { OPENAI_INTEGRATION_NAME } from '../../../../src/ai/openai/constants'; +import { instrumentWorkersAiClient } from '../../../../src/ai/workers-ai'; +import { getDefaultTestClientOptions, TestClient } from '../../../mocks/client'; + +const LIVE_DOC = { definition: { kind: 'pi.live' } }; +const MODEL = { id: 'faux-model', provider: 'faux' }; +const ANSWER = { role: 'assistant', content: [{ type: 'text', text: 'Paris.' }], stopReason: 'stop' }; + +describe('instrumentPiDurableHarnessOptions', () => { + let endedSpans: Span[]; + let events: Event[]; + let client: TestClient; + + beforeEach(() => { + _INTERNAL_clearAiProviderSkips(); + getMainCarrier().__SENTRY__ = undefined; + client = new TestClient( + getDefaultTestClientOptions({ dsn: 'https://public@dsn.ingest.sentry.io/1337', tracesSampleRate: 1 }), + ); + setCurrentClient(client); + client.init(); + + endedSpans = []; + events = []; + client.on('spanEnd', span => endedSpans.push(span)); + // The root spans of these tests are sent as transaction events; keep only the error events. + TestClient.sendEventCalled = event => { + if (event.type !== 'transaction') { + events.push(event); + } + }; + }); + + afterEach(() => { + TestClient.sendEventCalled = undefined; + _INTERNAL_clearAiProviderSkips(); + getMainCarrier().__SENTRY__ = undefined; + }); + + it('wraps the options once and leaves the caller objects untouched', () => { + const models = { completeSimple: async () => ({}) }; + const options: PiHarnessOptions = { models }; + + const instrumented = instrumentPiDurableHarnessOptions(options); + + expect(instrumented).not.toBe(options); + expect(instrumented.models).not.toBe(models); + expect(options.models).toBe(models); + expect(instrumentPiDurableHarnessOptions(instrumented)).toBe(instrumented); + }); + + // pi-durable reads `settings`, `env` and `conversationCreated` at every use, so getters and later + // assignments on the caller's options must reach it. + it('keeps the options the caller changes later live', () => { + let env = 'first'; + const options = { models: {}, settings: { toolExecution: 'sequential' } } as PiHarnessOptions & { + settings: unknown; + env?: string; + }; + Object.defineProperty(options, 'env', { get: () => env, enumerable: true }); + + const instrumented = instrumentPiDurableHarnessOptions(options); + options.settings = { toolExecution: 'parallel' }; + env = 'second'; + + expect(instrumented.settings).toEqual({ toolExecution: 'parallel' }); + expect(instrumented.env).toBe('second'); + }); + + it('hands members other than model requests through, bound to the original', () => { + class Models { + private _provider = 'faux'; + public getProvider(): string { + return this._provider; + } + } + + const instrumented = instrumentPiDurableHarnessOptions({ models: new Models() as PiModels }); + + expect((instrumented.models as unknown as Models).getProvider()).toBe('faux'); + expect(endedSpans).toHaveLength(0); + }); + + it('traces a streamed model request and skips the provider SDK integrations', async () => { + const message = { + ...ANSWER, + model: 'faux-model', + usage: { input: 10, output: 2, totalTokens: 12 }, + }; + const models: PiModels = { streamSimple: () => ({ result: async () => message }) }; + const instrumented = instrumentPiDurableHarnessOptions({ models }); + + await instrumented.models!.streamSimple!(MODEL, { + messages: [{ role: 'user', content: 'Capital of France?' }], + }).result(); + + expect(_INTERNAL_shouldSkipAiProviderWrapping(OPENAI_INTEGRATION_NAME)).toBe(true); + const chat = spanToStaticSpanJSON(endedSpans[0]!); + expect(chat.description).toBe('chat faux-model'); + expect(chat.op).toBe('gen_ai.chat'); + expect(chat.data).toMatchObject({ + 'gen_ai.provider.name': 'faux', + 'gen_ai.request.model': 'faux-model', + 'gen_ai.response.streaming': true, + 'gen_ai.response.finish_reasons': '["stop"]', + 'gen_ai.usage.input_tokens': 10, + 'gen_ai.usage.output_tokens': 2, + }); + expect(chat.data['gen_ai.input.messages']).toBe( + '[{"role":"user","parts":[{"type":"text","content":"Capital of France?"}]}]', + ); + expect(chat.data['gen_ai.output.messages']).toBe( + '[{"role":"assistant","parts":[{"type":"text","content":"Paris."}],"finish_reason":"stop"}]', + ); + }); + + it.each([ + ['stream', true], + ['streamSimple', true], + ['streamDeferred', true], + ['complete', false], + ['completeSimple', false], + ['fetchDeferred', false], + ])('traces %s with the request options and the response identity', async (method, streaming) => { + const message = { ...ANSWER, responseModel: 'faux-model-2026', responseId: 'resp_1' }; + const models = { + [method]: streaming ? () => ({ result: async () => message }) : async () => message, + } as PiModels; + const instrumented = instrumentPiDurableHarnessOptions({ models }); + + const result = (instrumented.models as Record unknown>)[method]!( + { ...MODEL, baseUrl: 'https://api.example.com:8443/v1' }, + { messages: [{ role: 'user', content: 'Hi.' }] }, + { temperature: 0.2, maxTokens: 50, reasoning: 'low' }, + ); + await (streaming ? (result as PiEventStream).result() : result); + + const chat = spanToJSON(endedSpans[0]!); + expect(chat.name).toBe('chat faux-model'); + expect(chat.attributes['gen_ai.response.streaming']).toBe(streaming ? true : undefined); + expect(chat.attributes['gen_ai.response.model']).toBe('faux-model-2026'); + expect(chat.attributes['gen_ai.response.id']).toBe('resp_1'); + expect(chat.attributes['server.address']).toBe('api.example.com'); + expect(chat.attributes['server.port']).toBe(8443); + // The deferred methods take a handle, not a request, so there is nothing to record. + const deferred = method.endsWith('Deferred'); + expect(chat.attributes['gen_ai.request.temperature']).toBe(deferred ? undefined : 0.2); + expect(chat.attributes['gen_ai.request.max_tokens']).toBe(deferred ? undefined : 50); + expect(chat.attributes['gen_ai.request.reasoning.level']).toBe(deferred ? undefined : 'low'); + expect(chat.attributes['gen_ai.input.messages']).toBe( + deferred ? undefined : '[{"role":"user","parts":[{"type":"text","content":"Hi."}]}]', + ); + }); + + it('sets the conversation id of the scope on the chat span itself', async () => { + const models: PiModels = { completeSimple: async () => ANSWER }; + const instrumented = instrumentPiDurableHarnessOptions({ models }); + + await withScope(async scope => { + scope.setConversationId('harness:7'); + await instrumented.models!.completeSimple!(MODEL, { messages: [] }); + }); + + expect(spanToJSON(endedSpans[0]!).attributes['gen_ai.conversation.id']).toBe('harness:7'); + }); + + it('counts the cached tokens of a response into its input tokens', async () => { + const message = { + ...ANSWER, + usage: { input: 37, output: 2, cacheRead: 52, cacheWrite: 38, reasoning: 1, totalTokens: 129 }, + }; + const models: PiModels = { completeSimple: async () => message }; + const instrumented = instrumentPiDurableHarnessOptions({ models }); + + await instrumented.models!.completeSimple!(MODEL, { messages: [] }); + + expect(spanToJSON(endedSpans[0]!).attributes).toMatchObject({ + 'gen_ai.usage.input_tokens': 127, + 'gen_ai.usage.cache_read.input_tokens': 52, + 'gen_ai.usage.cache_creation.input_tokens': 38, + 'gen_ai.usage.reasoning.output_tokens': 1, + 'gen_ai.usage.total_tokens': 129, + }); + }); + + // The provider parks the request and pi-durable fetches the answer later: only the fetch is billed. + it('records no usage for a request the provider parked, and the usage of its answer', async () => { + const parked = { + role: 'assistant', + content: [], + stopReason: 'deferred', + usage: { input: 0, output: 0, totalTokens: 0 }, + }; + const answer = { ...ANSWER, usage: { input: 8, output: 2, totalTokens: 10 } }; + const models: PiModels = { completeSimple: async () => parked, fetchDeferred: async () => answer }; + const instrumented = instrumentPiDurableHarnessOptions({ models }); + + await instrumented.models!.completeSimple!(MODEL, { messages: [] }); + await instrumented.models!.fetchDeferred!(MODEL, { id: 'handle_1' }); + + const [parkedChat, answerChat] = endedSpans.map(span => spanToJSON(span).attributes); + expect(parkedChat!['gen_ai.response.finish_reasons']).toBe('["deferred"]'); + expect(parkedChat!['gen_ai.usage.total_tokens']).toBeUndefined(); + expect(answerChat!['gen_ai.usage.total_tokens']).toBe(10); + }); + + it('marks a failed request with a fixed status message and maps a tool-calling stop', async () => { + const failed = { + role: 'assistant', + content: [], + stopReason: 'error', + errorMessage: '401 invalid x-api-key for user@example.com', + usage: { input: 0, output: 0, totalTokens: 0 }, + }; + const toolCall = { + role: 'assistant', + content: [{ type: 'toolCall', id: 'call_1', name: 'read', arguments: { path: 'a.txt' } }], + stopReason: 'toolUse', + }; + const models: PiModels = { completeSimple: async () => failed, complete: async () => toolCall }; + const instrumented = instrumentPiDurableHarnessOptions({ models }); + + await instrumented.models!.completeSimple!(MODEL, { messages: [] }); + await instrumented.models!.complete!(MODEL, { messages: [] }); + + const [chatFailed, chatToolCall] = endedSpans.map(span => spanToJSON(span)); + expect(chatFailed!.status).toBe('error'); + expect(chatFailed!.attributes['sentry.status.message']).toBe('internal_error'); + expect(JSON.stringify(chatFailed!.attributes)).not.toContain('x-api-key'); + expect(chatFailed!.attributes['gen_ai.usage.total_tokens']).toBeUndefined(); + expect(chatToolCall!.attributes['gen_ai.response.finish_reasons']).toBe('["tool_call"]'); + expect(chatToolCall!.attributes['gen_ai.output.messages']).toContain('"finish_reason":"tool_call"'); + }); + + it('marks a rejected or throwing request as errored and keeps the error', async () => { + const models: PiModels = { + completeSimple: async () => { + throw new Error('boom'); + }, + stream: () => { + throw new Error('sync boom'); + }, + }; + const instrumented = instrumentPiDurableHarnessOptions({ models }); + + await expect(instrumented.models!.completeSimple!(MODEL, { messages: [] })).rejects.toThrow('boom'); + expect(() => instrumented.models!.stream!(MODEL, { messages: [] })).toThrow('sync boom'); + + expect(endedSpans.map(span => spanToJSON(span).attributes['sentry.status.message'])).toEqual([ + 'internal_error', + 'internal_error', + ]); + }); + + it('reads the system prompt and the tools from the positional system messages', async () => { + const models: PiModels = { + streamSimple: () => ({ result: async () => ({ role: 'assistant', content: [], stopReason: 'stop' }) }), + }; + const instrumented = instrumentPiDurableHarnessOptions({ models }); + + // pi-durable never sets `systemPrompt` or `tools`: it sends prompt sections and tool changes as + // system messages, and each request replays them. + await instrumented.models!.streamSimple!(MODEL, { + messages: [ + { + role: 'system', + content: '', + sections: { preamble: 'You help.' }, + toolsAdded: [{ name: 'read' }, { name: 'bash' }], + }, + { role: 'user', content: 'Hi.' }, + { role: 'system', content: '', sections: { mode: 'Plan only.' }, toolsRemoved: [{ name: 'bash' }] }, + ], + }).result(); + + const chat = spanToStaticSpanJSON(endedSpans[0]!); + expect(chat.data['gen_ai.system_instructions']).toBe('You help.\n\nPlan only.'); + expect(chat.data['gen_ai.tool.definitions']).toBe('[{"name":"read"}]'); + expect(chat.data['gen_ai.input.messages']).toBe('[{"role":"user","parts":[{"type":"text","content":"Hi."}]}]'); + }); + + describe('provider skip', () => { + const models: PiModels = { completeSimple: async () => ANSWER }; + + it('skips the provider integrations on the first request, not when the options are wrapped', async () => { + const instrumented = instrumentPiDurableHarnessOptions({ models }); + expect(_INTERNAL_shouldSkipAiProviderWrapping(OPENAI_INTEGRATION_NAME)).toBe(false); + + await instrumented.models!.completeSimple!(MODEL, { messages: [] }); + + expect(_INTERNAL_shouldSkipAiProviderWrapping(OPENAI_INTEGRATION_NAME)).toBe(true); + expect(_INTERNAL_shouldSkipAiProviderWrapping(ANTHROPIC_AI_INTEGRATION_NAME)).toBe(true); + expect(_INTERNAL_shouldSkipAiProviderWrapping(GOOGLE_GENAI_INTEGRATION_NAME)).toBe(true); + }); + + // The registry is reset per client, so a one-shot call would be undone by the next `init()`. + it('applies the skip again after the registry is cleared', async () => { + const instrumented = instrumentPiDurableHarnessOptions({ models }); + await instrumented.models!.completeSimple!(MODEL, { messages: [] }); + _INTERNAL_clearAiProviderSkips(); + expect(_INTERNAL_shouldSkipAiProviderWrapping(OPENAI_INTEGRATION_NAME)).toBe(false); + + await instrumented.models!.completeSimple!(MODEL, { messages: [] }); + + expect(_INTERNAL_shouldSkipAiProviderWrapping(OPENAI_INTEGRATION_NAME)).toBe(true); + }); + + it('applies the skip to every provider when only some are already registered', async () => { + _INTERNAL_skipAiProviderWrapping([OPENAI_INTEGRATION_NAME]); + const instrumented = instrumentPiDurableHarnessOptions({ models }); + + await instrumented.models!.completeSimple!(MODEL, { messages: [] }); + + expect(_INTERNAL_shouldSkipAiProviderWrapping(ANTHROPIC_AI_INTEGRATION_NAME)).toBe(true); + expect(_INTERNAL_shouldSkipAiProviderWrapping(GOOGLE_GENAI_INTEGRATION_NAME)).toBe(true); + }); + + // `createAI()` of `agents/models/pi-ai` sends pi-ai requests through the Workers AI binding. + it('reports a request sent through the Workers AI binding once', async () => { + const ai = instrumentWorkersAiClient({ + run: async (_model: string, _inputs: unknown, _options: unknown) => new Response('{}'), + gateway: () => ({}), + toMarkdown: async () => [], + }); + const models: PiModels = { + completeSimple: async () => { + await ai.run('@cf/moonshotai/kimi-k2.7-code', { messages: [] }, { returnRawResponse: true }); + return ANSWER; + }, + }; + const instrumented = instrumentPiDurableHarnessOptions({ models }); + + await instrumented.models!.completeSimple!(MODEL, { messages: [] }); + + expect(endedSpans.map(span => spanToStaticSpanJSON(span).description)).toEqual(['chat faux-model']); + }); + }); + + describe('content recording', () => { + const tool: PiTool = { + name: 'get_weather', + description: 'Get the weather.', + execute: async () => ({ content: [{ type: 'text', text: 'sunny' }] }), + }; + const message = { + role: 'assistant', + content: [{ type: 'toolCall', id: 'call_1', name: 'get_weather', arguments: { city: 'Berlin' } }], + stopReason: 'toolUse', + }; + const context = { + messages: [ + { role: 'system', content: '', sections: { preamble: 'You help.' }, toolsAdded: [{ name: 'get_weather' }] }, + { role: 'user', content: 'Weather in Berlin?' }, + ], + }; + + it('records nothing when inputs and outputs are off', async () => { + const options = { recordInputs: false, recordOutputs: false }; + const instrumented = instrumentPiDurableHarnessOptions( + { models: { completeSimple: async () => message } as PiModels }, + options, + ); + + await instrumented.models!.completeSimple!(MODEL, context); + await instrumentTool(tool, createRuns(), options).execute({ city: 'Berlin' }, { callId: 'call_1' }, undefined); + + const [chat, execute] = endedSpans.map(span => spanToJSON(span).attributes); + expect(Object.keys(chat!)).not.toContain('gen_ai.input.messages'); + expect(Object.keys(chat!)).not.toContain('gen_ai.output.messages'); + expect(Object.keys(chat!)).not.toContain('gen_ai.system_instructions'); + expect(Object.keys(chat!)).not.toContain('gen_ai.tool.definitions'); + expect(Object.keys(execute!)).not.toContain('gen_ai.tool.call.arguments'); + expect(Object.keys(execute!)).not.toContain('gen_ai.tool.call.result'); + expect(execute!['gen_ai.tool.description']).toBe('Get the weather.'); + }); + + it('follows the data collection settings of the current client', async () => { + const strict = new TestClient( + getDefaultTestClientOptions({ + dsn: 'https://public@dsn.ingest.sentry.io/1337', + tracesSampleRate: 1, + dataCollection: { genAI: { inputs: false, outputs: false } }, + }), + ); + setCurrentClient(strict); + strict.init(); + strict.on('spanEnd', span => endedSpans.push(span)); + const models: PiModels = { completeSimple: async () => message }; + const instrumented = instrumentPiDurableHarnessOptions({ models }); + + await instrumented.models!.completeSimple!(MODEL, context); + await instrumentTool(tool, createRuns(), {}).execute({ city: 'Berlin' }, { callId: 'call_1' }, undefined); + + const recorded = endedSpans.flatMap(span => Object.keys(spanToJSON(span).attributes)); + expect(recorded).not.toContain('gen_ai.input.messages'); + expect(recorded).not.toContain('gen_ai.output.messages'); + expect(recorded).not.toContain('gen_ai.system_instructions'); + expect(recorded).not.toContain('gen_ai.tool.definitions'); + expect(recorded).not.toContain('gen_ai.tool.call.arguments'); + expect(recorded).not.toContain('gen_ai.tool.call.result'); + }); + }); + + it('captures the failures pi-durable reports and still hands them to the caller', async () => { + const reported: unknown[] = []; + const options: PiHarnessOptions = {}; + const instrumented = instrumentPiDurableHarnessOptions(options); + options.onReport = error => reported.push(error); + const error = new Error('afterResponse hook failed'); + + instrumented.onReport!(error); + await client.flush(); + + expect(reported).toEqual([error]); + expect(client.event?.exception?.values?.[0]?.value).toBe('afterResponse hook failed'); + expect(client.event?.exception?.values?.[0]?.mechanism).toEqual({ type: 'auto.ai.pi_durable', handled: true }); + }); + + it('returns the same wrapped task for the same definition in every snapshot', () => { + const generation: PiTask = { definition: { name: 'pi.generation', phases: { prepare: async () => undefined } } }; + const registry: PiRegistryReader = { + // A new snapshot object per call, as pi-durable publishes one per registry change. + snapshot: () => ({ task: () => generation, tasks: () => [generation] }), + subscribe: () => () => undefined, + }; + + const instrumented = instrumentPiDurableHarnessOptions({ registry }).registry!; + + const first = instrumented.snapshot().task('pi.generation'); + expect(first).not.toBe(generation); + expect(instrumented.snapshot().task('pi.generation')).toBe(first); + expect(instrumented.snapshot().tasks()[0]).toBe(first); + }); + + // pi-durable calls a phase on the `phases` object and `abort` on the definition. + it('calls task phases and abort handlers with their own receiver', async () => { + const definition = { + name: 'app.job', + phases: { + run(this: { helper(): Promise }) { + return this.helper(); + }, + helper: async () => 'helped', + }, + abort(this: { name: string }) { + return Promise.resolve(this.name); + }, + }; + const task = { definition } as unknown as PiTask; + const registry: PiRegistryReader = { + snapshot: () => ({ task: () => task, tasks: () => [task] }), + subscribe: () => () => undefined, + }; + const runtime: PiTaskRuntime = { + taskId: 1, + conversationId: 1, + agent: async () => ({}), + commit: async () => undefined, + }; + + const wrapped = instrumentPiDurableHarnessOptions({ registry }).registry!.snapshot().task('app.job')!; + + await expect(wrapped.definition.phases.run!(undefined, runtime, undefined)).resolves.toBe('helped'); + await expect(wrapped.definition.abort!(undefined, runtime, undefined)).resolves.toBe('app.job'); + }); + + describe('tools', () => { + it('runs a frozen tool and a tool that reads a private field', async () => { + const frozen: PiTool = Object.freeze({ + name: 'frozen', + execute: async () => ({ content: [{ type: 'text', text: 'ok' }] }), + }); + class PrefixTool { + #prefix = 'pre'; + public name = 'prefixed'; + public prepareArguments(args: string): string { + return `${this.#prefix}:${args}`; + } + public async execute(): Promise { + return { content: [{ type: 'text', text: this.#prefix }] }; + } + } + const runs = createRuns(); + + const wrappedFrozen = instrumentTool(frozen, runs, {}); + const wrappedPrefix = instrumentTool(new PrefixTool() as unknown as PiTool, runs, {}); + + expect(wrappedFrozen.name).toBe('frozen'); + await expect(wrappedFrozen.execute({}, {}, undefined)).resolves.toEqual({ + content: [{ type: 'text', text: 'ok' }], + }); + expect((wrappedPrefix as unknown as PrefixTool).prepareArguments('x')).toBe('pre:x'); + await expect(wrappedPrefix.execute({}, {}, undefined)).resolves.toEqual({ + content: [{ type: 'text', text: 'pre' }], + }); + expect(endedSpans.map(span => spanToJSON(span).name)).toEqual(['execute_tool frozen', 'execute_tool prefixed']); + }); + + it('ends a call outside a tool phase with its own result', async () => { + const runs = createRuns(); + const tool: PiTool = { + name: 'get_weather', + description: 'Get the weather.', + execute: async () => ({ content: [{ type: 'text', text: 'sunny' }] }), + }; + + await instrumentTool(tool, runs, {}).execute( + { city: 'Berlin' }, + { callId: 'call_1', conversationId: 4 }, + undefined, + ); + + const execute = spanToJSON(endedSpans[0]!); + expect(execute.name).toBe('execute_tool get_weather'); + expect(execute.status).toBe('ok'); + expect(execute.attributes).toMatchObject({ + 'sentry.op': 'gen_ai.execute_tool', + 'sentry.origin': 'auto.ai.pi_durable', + 'gen_ai.tool.name': 'get_weather', + 'gen_ai.tool.description': 'Get the weather.', + 'gen_ai.tool.call.id': 'call_1', + 'gen_ai.conversation.id': `${runs.harnessId}:4`, + 'gen_ai.tool.call.arguments': '{"city":"Berlin"}', + 'gen_ai.tool.call.result': 'sunny', + }); + }); + + it('marks error results and throws, and captures a throw unless the tool is built in', async () => { + const runs = createRuns(); + const erroring: PiTool = { + name: 'erroring', + execute: async () => ({ isError: true, content: [{ type: 'text', text: 'not found' }] }), + }; + const throwing: PiTool = { + name: 'throwing', + execute: async () => { + throw new Error('tool failed'); + }, + }; + const builtIn: PiTool = { + name: 'bash', + execute: async () => { + throw new Error('Command exited with code 1'); + }, + }; + markBuiltInTool(builtIn); + + await instrumentTool(erroring, runs, {}).execute({}, {}, undefined); + await expect(instrumentTool(throwing, runs, {}).execute({}, {}, undefined)).rejects.toThrow('tool failed'); + await expect(instrumentTool(builtIn, runs, {}).execute({}, {}, undefined)).rejects.toThrow('code 1'); + await client.flush(); + + const [errored, thrown, bash] = endedSpans.map(span => spanToJSON(span)); + expect(errored!.attributes['sentry.status.message']).toBe('internal_error'); + expect(errored!.attributes['gen_ai.tool.call.result']).toBe('not found'); + expect(thrown!.attributes['sentry.status.message']).toBe('internal_error'); + expect(bash!.attributes['sentry.status.message']).toBe('internal_error'); + expect(events.map(event => event.exception?.values?.[0]?.value)).toEqual(['tool failed']); + expect(events[0]?.exception?.values?.[0]?.mechanism).toEqual({ type: 'auto.ai.pi_durable', handled: true }); + }); + + it('marks an aborted call as cancelled and captures nothing', async () => { + const controller = new AbortController(); + const tool: PiTool = { + name: 'slow', + execute: async () => { + controller.abort(); + throw new Error('tool aborted'); + }, + }; + + await expect( + instrumentTool(tool, createRuns(), {}).execute({}, {}, { abortSignal: controller.signal }), + ).rejects.toThrow('tool aborted'); + await client.flush(); + + expect(spanToStaticSpanJSON(endedSpans[0]!).status).toBe('cancelled'); + expect(events).toHaveLength(0); + }); + + it('ends a nested tool call with its own result and links a forked conversation to the call', async () => { + const runs = createRuns(); + runs.toolCalls.set(7, []); + const inner = instrumentTool( + { name: 'inner', execute: async () => ({ content: [{ type: 'text', text: 'INNER' }] }) }, + runs, + {}, + ); + const outer = instrumentTool( + { + name: 'outer', + execute: async (_args, api) => { + await inner.execute({}, api, undefined); + await api.commit!( + tx => + (tx as { forkConversation: (...args: unknown[]) => unknown }).forkConversation(1, 3, { + ownership: { kind: 'task', taskId: 7 }, + }), + undefined, + ); + return { content: [{ type: 'text', text: 'OUTER' }] }; + }, + }, + runs, + {}, + ); + const api: PiToolExecutionApi = { + taskId: 7, + callId: 'call_1', + commit: async change => change({ forkConversation: async () => ({ id: 42 }) }), + }; + + await outer.execute({}, api, undefined); + + expect(endedSpans.map(span => spanToJSON(span).name)).toEqual(['execute_tool inner']); + expect(spanToJSON(endedSpans[0]!).attributes['gen_ai.tool.call.result']).toBe('INNER'); + // The outer call waits for the result entry its phase commits. + const [outerCall] = runs.toolCalls.get(7)!; + expect(outerCall!.span.spanContext().spanId).toBe(spanToJSON(endedSpans[0]!).parent_span_id); + expect(spanToJSON(startRun(42, runs).span).parent_span_id).toBe(outerCall!.span.spanContext().spanId); + }); + }); + + describe('runs', () => { + /** A generation phase that commits each of `states` as the `pi.live` document, in order. */ + function generationPhase(states: (PiLiveState | { settle: PiSettlement })[]): PiTask { + return { + definition: { + name: 'pi.generation', + phases: { + request: async (_task, runtime) => { + for (const state of states) { + await runtime.commit(async tx => { + const draft = (await (tx as { doc: (...args: unknown[]) => Promise }).doc(LIVE_DOC, 1))!; + if ('settle' in state) { + (tx as { settleSubmission: (...args: unknown[]) => void }).settleSubmission(5, state.settle); + delete draft.run; + } else { + draft.run = state.run; + } + return undefined; + }, undefined); + } + }, + }, + }, + }; + } + + function harnessFor(task: PiTask): PiHarnessOptions { + const registry: PiRegistryReader = { + snapshot: () => ({ task: () => task, tasks: () => [task] }), + subscribe: () => () => undefined, + }; + return instrumentPiDurableHarnessOptions({ registry }); + } + + async function runGeneration(harness: PiHarnessOptions): Promise { + const live: PiLiveState = {}; + const tx = { doc: async () => live, settleSubmission: () => undefined }; + const runtime: PiTaskRuntime = { + taskId: 3, + conversationId: 1, + agent: async () => ({}), + commit: async (change: PiCommitChange) => { + await change(tx, undefined); + }, + }; + await harness.registry!.snapshot().task('pi.generation')!.definition.phases.request!( + undefined, + runtime, + undefined, + ); + } + + function runSpans(): ReturnType[] { + return endedSpans + .map(span => spanToJSON(span)) + .filter(span => span.attributes['sentry.op'] === 'gen_ai.invoke_agent'); + } + + it('keeps the run open while run control moves to the next generation', async () => { + await runGeneration(harnessFor(generationPhase([{ run: { inputs: [5] } }, { run: { inputs: [5, 6] } }]))); + + expect(runSpans()).toHaveLength(0); + }); + + it.each([ + ['model_error', 'error', 'model_error'], + ['no_model', 'error', 'no_model'], + ['aborted', 'ok', undefined], + ['reset', 'ok', undefined], + [undefined, 'error', 'internal_error'], + ])('ends an unanswered run with reason %s as %s', async (reason, status, message) => { + await runGeneration( + harnessFor(generationPhase([{ run: { inputs: [5] } }, { settle: { status: 'unanswered', reason } }])), + ); + + const runs = runSpans(); + expect(runs).toHaveLength(1); + expect(runs[0]!.status).toBe(status); + expect(runs[0]!.attributes['sentry.status.message']).toBe(message); + expect(runs[0]!.attributes['gen_ai.conversation.id']).toMatch(/^[0-9a-f]{32}:1$/); + }); + + // `@sentry/cloudflare` initializes the SDK inside each request, so the client is bound to the + // scope of that request and never to the default scope. + it('traces a run when only the scope that starts the phase has the client', async () => { + getCurrentScope().setClient(undefined); + const harness = harnessFor(generationPhase([{ run: { inputs: [5] } }, { settle: { status: 'done' } }])); + + await withScope(async scope => { + scope.setClient(client); + await runGeneration(harness); + }); + + expect(runSpans()).toHaveLength(1); + }); + + it('ends the run when a queued input starts the next run in the same commit', async () => { + await runGeneration(harnessFor(generationPhase([{ run: { inputs: [5] } }, { run: { inputs: [7] } }]))); + + const runs = runSpans(); + expect(runs).toHaveLength(1); + expect(runs[0]!.status).toBe('ok'); + }); + + it('gives every Harness its own conversation id prefix', async () => { + const phase = generationPhase([{ run: { inputs: [5] } }, { settle: { status: 'done' } }]); + const first = harnessFor(phase); + + await runGeneration(first); + await runGeneration(first); + await runGeneration(harnessFor(phase)); + + const ids = runSpans().map(run => String(run.attributes['gen_ai.conversation.id'])); + expect(ids).toHaveLength(3); + expect(ids[0]).toBe(ids[1]); + expect(ids[2]).not.toBe(ids[0]); + expect(ids.every(id => id.endsWith(':1'))).toBe(true); + }); + + it('ends the open runs when the Harness closes', async () => { + const harness = harnessFor(generationPhase([{ run: { inputs: [5] } }])); + let onClose: (() => void) | undefined; + endRunsOnClose(harness, { subscribeClose: (listener: () => void) => ((onClose = listener), () => undefined) }); + + await runGeneration(harness); + expect(runSpans()).toHaveLength(0); + + onClose!(); + + expect(runSpans()).toHaveLength(1); + expect(spanToStaticSpanJSON(endedSpans[0]!).status).toBe('cancelled'); + }); + + // `wait()` resolves once the ending commit settles the inputs, before `commit()` resolves, so an + // app that closes right after `wait()` closes inside that gap. + it('ends a run with its own settlement when the Harness closes before the ending commit resolves', async () => { + const harness = harnessFor( + generationPhase([{ run: { inputs: [5] } }, { settle: { status: 'unanswered', reason: 'model_error' } }]), + ); + let onClose: (() => void) | undefined; + endRunsOnClose(harness, { subscribeClose: (listener: () => void) => ((onClose = listener), () => undefined) }); + const live: PiLiveState = {}; + const tx = { doc: async () => live, settleSubmission: () => undefined }; + let commits = 0; + const runtime: PiTaskRuntime = { + taskId: 3, + conversationId: 1, + agent: async () => ({}), + commit: async (change: PiCommitChange) => { + await change(tx, undefined); + commits++; + if (commits === 2) { + onClose!(); + } + }, + }; + + await harness.registry!.snapshot().task('pi.generation')!.definition.phases.request!( + undefined, + runtime, + undefined, + ); + + const runs = runSpans(); + expect(runs).toHaveLength(1); + expect(runs[0]!.status).toBe('error'); + expect(runs[0]!.attributes['sentry.status.message']).toBe('model_error'); + }); + + // A run the process never sees settle would otherwise stay in the map for the process lifetime. + it('ends the oldest run span when the run tracker overflows', () => { + const runs = createRuns(); + + for (let conversation = 0; conversation <= MAX_TRACKED_PI_RUNS; conversation++) { + startRun(conversation, runs); + } + + expect(runs.active.size).toBe(MAX_TRACKED_PI_RUNS); + expect(runs.active.has(0)).toBe(false); + expect(endedSpans).toHaveLength(1); + expect(spanToJSON(endedSpans[0]!).attributes['gen_ai.conversation.id']).toBe(`${runs.harnessId}:0`); + }); + }); +}); diff --git a/packages/server-utils/test/ai/lib/utils/pi-ai-messages.test.ts b/packages/server-utils/test/ai/lib/utils/pi-ai-messages.test.ts index d75389725cf5..f9b2d3b31dc6 100644 --- a/packages/server-utils/test/ai/lib/utils/pi-ai-messages.test.ts +++ b/packages/server-utils/test/ai/lib/utils/pi-ai-messages.test.ts @@ -4,6 +4,8 @@ import { piAiContentToString, piAiFinishReason, piAiMessagesToGenAiMessages, + piAiSystemInstructions, + piAiToolDefinitions, } from '../../../../src/ai/pi-ai/messages'; describe('convert pi-ai messages to gen_ai messages', () => { @@ -127,6 +129,28 @@ describe('convert pi-ai messages to gen_ai messages', () => { expect(piAiAssistantMessageToGenAiMessage({ role: 'assistant', content: [] }, 'stop')).toBeUndefined(); }); + it('replays positional system messages into the current system prompt and tools', () => { + const context = { + systemPrompt: 'Base prompt.', + tools: [{ name: 'read' }], + messages: [ + { + role: 'system', + content: '', + sections: { preamble: 'You help.', mode: 'Edit.' }, + toolsAdded: [{ name: 'bash' }], + }, + { role: 'user', content: 'Hi.' }, + { role: 'system', content: '', sections: { mode: null }, toolsRemoved: [{ name: 'read' }] }, + { role: 'system', content: 'Appended.', sections: { mode: 'Plan only.' } }, + ], + }; + + expect(piAiSystemInstructions(context)).toBe('Base prompt.\n\nAppended.\n\nYou help.\n\nPlan only.'); + expect(piAiToolDefinitions(context)).toStrictEqual([{ name: 'bash' }]); + expect(piAiSystemInstructions({ messages: [{ role: 'user', content: 'Hi.' }] })).toBeUndefined(); + }); + it('renders text-only content as text and other content as mapped parts', () => { expect( piAiContentToString([