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..46a27dd03b74 --- /dev/null +++ b/dev-packages/node-integration-tests/suites/tracing/pi-durable/instrument.mjs @@ -0,0 +1,9 @@ +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, +}); 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-tools.mjs b/dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-tools.mjs new file mode 100644 index 000000000000..630d1acdbe9c --- /dev/null +++ b/dev-packages/node-integration-tests/suites/tracing/pi-durable/scenario-tools.mjs @@ -0,0 +1,98 @@ +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 } 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: +// - the built-in `bash` tool runs a command that exits non-zero, which it reports by throwing; +// - `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('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: [ + 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..f39613afd345 --- /dev/null +++ b/dev-packages/node-integration-tests/suites/tracing/pi-durable/test.ts @@ -0,0 +1,317 @@ +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_SYSTEM_INSTRUCTIONS, + GEN_AI_TOOL_CALL_ARGUMENTS, + GEN_AI_TOOL_CALL_RESULT, + GEN_AI_TOOL_DEFINITIONS, + 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', + }, +}; + +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); + + 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.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); + }); + }, + 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 () => { + // 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. + await createRunner() + .ignore('span') + .expect({ + event: event => { + expect(event.exception?.values?.[0]?.value).toBe('Intentional pi-durable tool failure'); + }, + }) + .expect({ + event: event => { + expect(event.message).toBe('pi-durable tools done'); + }, + }) + .start() + .completed(); + }); + + 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({ + 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+/); + + // 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, + ); +}); 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 8b77b4f205e0..50ee792d4762 100644 --- a/packages/server-utils/src/ai/pi-ai/messages.ts +++ b/packages/server-utils/src/ai/pi-ai/messages.ts @@ -58,6 +58,80 @@ export function piAiAssistantMessageToGenAiMessage(message: unknown, finishReaso return { role: 'assistant', parts, ...(finishReason ? { finish_reason: finishReason } : {}) }; } +/** 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-durable/constants.ts b/packages/server-utils/src/ai/pi-durable/constants.ts new file mode 100644 index 000000000000..39138abdbbb7 --- /dev/null +++ b/packages/server-utils/src/ai/pi-durable/constants.ts @@ -0,0 +1,31 @@ +export const PI_DURABLE_INTEGRATION_NAME = 'PiDurable' as const; + +export const PI_DURABLE_MODULE_NAME = '@earendil-works/pi-durable'; + +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'; + +/** + * Name of pi-durable's built-in `CodingTools` extension (`read`, `write`, `edit`, `bash`). Its tools + * throw to report expected failures to the model, such as a command that exits non-zero. + */ +export const PI_CODING_TOOLS_EXTENSION = 'coding-tools'; + +/** + * 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..17c1e4d81752 --- /dev/null +++ b/packages/server-utils/src/ai/pi-durable/index.ts @@ -0,0 +1,213 @@ +import type { Scope } from '@sentry/core'; +import { + _INTERNAL_shouldSkipAiProviderWrapping, + _INTERNAL_skipAiProviderWrapping, + addNonEnumerableProperty, + captureException, + getCurrentScope, + getDefaultIsolationScope, + SPAN_STATUS_ERROR, + startNewTrace, + withActiveSpan, +} from '@sentry/core'; +import { ANTHROPIC_AI_INTEGRATION_NAME } from '../anthropic-ai/constants'; +import type { GenAiOptions } from '../core/utils'; +import { GOOGLE_GENAI_INTEGRATION_NAME } from '../google-genai/constants'; +import { OPENAI_INTEGRATION_NAME } from '../openai/constants'; +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 { + 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'); + +// 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 our `chat` span. +const SKIPPED_PROVIDERS = [OPENAI_INTEGRATION_NAME, ANTHROPIC_AI_INTEGRATION_NAME, GOOGLE_GENAI_INTEGRATION_NAME]; + +/** + * Return a copy of pi-durable `HarnessOptions` whose model requests, task phases and tool calls are + * traced. The caller's `models` and `registry` objects are wrapped, never modified. + * + * 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; + } + + // Applied on first request rather than here: the registry is reset per client, so a one-shot call + // would be undone by the next `init()`. + const skipProviders = (): void => { + if (!SKIPPED_PROVIDERS.every(provider => _INTERNAL_shouldSkipAiProviderWrapping(provider))) { + _INTERNAL_skipAiProviderWrapping(SKIPPED_PROVIDERS); + } + }; + + const { onReport } = harnessOptions; + const instrumented: T = { + ...harnessOptions, + ...(harnessOptions.models ? { models: instrumentPiModels(harnessOptions.models, options, skipProviders) } : {}), + ...(harnessOptions.registry ? { registry: instrumentRegistry(harnessOptions.registry, options) } : {}), + // pi-durable passes failures of extension code here, such as a throwing hook, and keeps going. + onReport: (error: unknown) => { + captureException(error, { mechanism: { handled: true, type: PI_DURABLE_ORIGIN } }); + onReport?.(error); + }, + }; + addNonEnumerableProperty(instrumented, INSTRUMENTED, true); + return instrumented; +} + +function instrumentRegistry(registry: PiRegistryReader, options: PiDurableOptions): PiRegistryReader { + const runs = createRuns(); + // 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; + const wrapPhase = + (phase: PiPhaseHandler): PiPhaseHandler => + (taskRecord, runtime, context) => + runPhase(definition.name, phase, taskRecord, runtime, context, runs, wrapTool); + const phases: Record = {}; + for (const [name, phase] of Object.entries(definition.phases)) { + phases[name] = typeof phase === 'function' ? wrapPhase(phase) : phase; + } + wrapped = { + ...task, + definition: { + ...definition, + phases, + ...(typeof definition.abort === 'function' ? { abort: wrapPhase(definition.abort) } : {}), + }, + }; + 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: PiPhaseHandler, + task: unknown, + runtime: PiTaskRuntime, + context: unknown, + 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(task, observed, context); + } 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..a51aa0a1cf1f --- /dev/null +++ b/packages/server-utils/src/ai/pi-durable/models.ts @@ -0,0 +1,171 @@ +import type { Span } from '@sentry/core'; +import { SEMANTIC_ATTRIBUTE_SENTRY_ORIGIN, SPAN_STATUS_ERROR, startSpanManual, stringify } from '@sentry/core'; +import { + 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, +} from '@sentry/conventions/attributes'; +import type { GenAiOptions } from '../core/utils'; +import { getGenAiSpanOp, resolveAIRecordingOptions } from '../core/utils'; +import { setUsageAttributes } from '../flue/utils'; +import { + piAiAssistantMessageToGenAiMessage, + piAiMessagesToGenAiMessages, + piAiSystemInstructions, + piAiToolDefinitions, +} from '../pi-ai/messages'; +import { PI_DURABLE_ORIGIN } from './constants'; +import type { PiAssistantMessage, PiContext, PiEventStream, PiModel, PiModels, PiStreamOptions } from './types'; + +type ModelCall = (model: PiModel, context: PiContext, options?: PiStreamOptions) => unknown; + +const STREAMING_METHODS = new Set(['stream', 'streamSimple']); +const COMPLETE_METHODS = new Set(['complete', 'completeSimple']); + +/** + * 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 streaming = STREAMING_METHODS.has(String(property)); + const instrumented = + streaming || COMPLETE_METHODS.has(String(property)) + ? (model: PiModel, context: PiContext, callOptions?: PiStreamOptions): unknown => { + onRequest(); + return traceModelRequest(model, context, callOptions, streaming, options, () => + original(model, context, callOptions), + ); + } + : original; + + wrapped.set(property, instrumented); + return instrumented; + }, + }); +} + +function traceModelRequest( + model: PiModel, + context: PiContext, + 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); + + 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 } : {}), + ...(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 ? 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 getRequestContentAttributes(context: PiContext): 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 = toFinishReason(message.stopReason); + if (finishReason) { + span.setAttribute(GEN_AI_RESPONSE_FINISH_REASONS, stringify([finishReason])); + } + + const failed = message.stopReason === 'error'; + setUsageAttributes(span, message.usage, failed); + + 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: message.errorMessage ?? 'internal_error' }); + } else if (message.stopReason === 'aborted') { + span.setStatus({ code: SPAN_STATUS_ERROR, message: 'cancelled' }); + } + span.end(); +} + +/** pi-ai names a tool-calling stop `toolUse`; the conventions call it `tool_call`. */ +function toFinishReason(stopReason: string | undefined): string | undefined { + return stopReason === 'toolUse' ? 'tool_call' : stopReason; +} 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..ce897f0e8c89 --- /dev/null +++ b/packages/server-utils/src/ai/pi-durable/runs.ts @@ -0,0 +1,208 @@ +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', 'withdrawn', '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; +} + +/** 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; + /** Resolved tools that come from pi-durable's built-in `coding-tools` extension. */ + builtInTools: WeakSet; +} + +/** 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(), + builtInTools: new WeakSet(), + }; +} + +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; + } + } + return result; + }; + + return runtime.commit(observedChange, context).then(() => { + if (ended) { + endRun(run, runs, toRunStatus(settlement)); + } + }); + }; + + 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..e92f5de48f80 --- /dev/null +++ b/packages/server-utils/src/ai/pi-durable/tools.ts @@ -0,0 +1,254 @@ +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_CODING_TOOLS_EXTENSION, + 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'; + +/** + * 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 => { + if (!resolved?.tools) { + return resolved; + } + markBuiltInTools(resolved, runs); + return { ...resolved, tools: resolved.tools.map(wrapTool) }; + }); + + 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 calls = runs.toolCalls.get(api?.taskId); + + 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) { + // Called outside a tool task phase, so no result entry follows. + 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, and the built-in tools throw to + // report expected failures to the model. + 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 (!runs.builtInTools.has(tool)) { + captureException(error, { mechanism: { handled: true, type: PI_DURABLE_ORIGIN } }); + } + } + if (!calls) { + span.end(call.endTimestamp); + } + throw error; + } + }, + ); + }; + + return new Proxy(tool, { + get(target, property) { + return property === 'execute' ? execute : Reflect.get(target, property, target); + }, + }); +} + +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); + }; + }, + }); +} + +/** + * Mark the resolved tools that come from the built-in `coding-tools` extension. A later + * extension's tool replaces an earlier one with the same name, so a tool comes from the last + * selected extension that has its name. A wrapper from another extension does not change that. + */ +function markBuiltInTools(agent: PiAgent, runs: PiRuns): void { + const providers = new Map(); + for (const extension of agent.extensions ?? []) { + for (const tool of extension.tools ?? []) { + providers.set(tool.name, extension.name); + } + } + for (const tool of agent.tools ?? []) { + if (providers.get(tool.name) === PI_CODING_TOOLS_EXTENSION) { + runs.builtInTools.add(tool); + } + } +} + +/** + * Record the conversations a tool call creates for itself, the subagent pattern: `api.commit()` with + * `tx.createConversation({ 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' || typeof value !== 'function') { + return value; + } + return async (createOptions: { ownership?: { kind?: string } }, ...rest: unknown[]) => { + const record = (await value(createOptions, ...rest)) as { id?: unknown } | undefined; + if (createOptions?.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 Reflect.get(target, property, target); + } + 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..9caea8d1660b --- /dev/null +++ b/packages/server-utils/src/ai/pi-durable/types.ts @@ -0,0 +1,155 @@ +/** + * 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. + */ + +export interface PiUsage { + input?: number; + output?: number; + cacheRead?: number; + cacheWrite?: number; + totalTokens?: number; + cost?: { + input?: number; + output?: number; + cacheRead?: number; + cacheWrite?: number; + total?: number; + }; +} + +export interface PiAssistantMessage { + role?: string; + content?: unknown[]; + provider?: string; + model?: string; + responseModel?: string; + responseId?: string; + usage?: PiUsage; + stopReason?: string; + errorMessage?: string; +} + +export interface PiModel { + id?: string; + provider?: string; + baseUrl?: string; +} + +export interface PiContext { + systemPrompt?: string; + messages?: unknown[]; + tools?: { name?: string; description?: string; parameters?: unknown }[]; +} + +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. */ +export interface PiModels { + stream?(model: PiModel, context: PiContext, options?: PiStreamOptions): PiEventStream; + streamSimple?(model: PiModel, context: PiContext, options?: PiStreamOptions): PiEventStream; + complete?(model: PiModel, context: PiContext, options?: PiStreamOptions): Promise; + completeSimple?(model: PiModel, context: PiContext, options?: PiStreamOptions): 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 PiExtension { + name?: string; + tools?: readonly { name: string }[]; +} + +export interface PiAgent { + /** The selected extensions, in order. */ + extensions?: readonly PiExtension[]; + 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; +} 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..b9f1bd841cb0 --- /dev/null +++ b/packages/server-utils/src/ai/pi-durable/utils.ts @@ -0,0 +1,17 @@ +import type { Scope } from '@sentry/core'; +import { 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. + */ +export function withCleanScopes(isolationScope: Scope, callback: () => T): T { + return withIsolationScope(isolationScope, () => withScope(getDefaultCurrentScope().clone(), callback)); +} diff --git a/packages/server-utils/src/index.ts b/packages/server-utils/src/index.ts index c162dadb1629..e775d13eceec 100644 --- a/packages/server-utils/src/index.ts +++ b/packages/server-utils/src/index.ts @@ -50,6 +50,8 @@ 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 { PiDurableOptions } from './ai/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..c6906e87774d --- /dev/null +++ b/packages/server-utils/src/integrations/pi-durable.ts @@ -0,0 +1,48 @@ +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 { instrumentPiDurableHarnessOptions } from '../ai/pi-durable'; +import { PI_DURABLE_INTEGRATION_NAME } from '../ai/pi-durable/constants'; +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[]; +} + +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 { + diagnosticsChannel + .tracingChannel(CHANNELS.PI_DURABLE_HARNESS_OPEN) + .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); + } + }); +} + +/** + * Traces `@earendil-works/pi-durable` Harnesses: 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. + */ +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..9bfa472d2e31 --- /dev/null +++ b/packages/server-utils/src/orchestrion/config/pi-durable.ts @@ -0,0 +1,24 @@ +import type { InstrumentationConfig } from '../apmTypes'; + +import { getModuleNames } from './module-names'; + +// `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: { + name: '@earendil-works/pi-durable', + versionRange: '>=1.0.0 <2.0.0', + filePath: 'dist/harness/harness.js', + }, + functionQuery: { methodName: 'open', kind: 'Async' as const }, + }, +] satisfies InstrumentationConfig[]; + +export const piDurableModuleNames = getModuleNames(piDurableConfig); + +export const piDurableChannels = { + PI_DURABLE_HARNESS_OPEN: 'orchestrion:@earendil-works/pi-durable:harnessOpen', +} 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..d498688f6266 --- /dev/null +++ b/packages/server-utils/test/ai/lib/tracing/pi-durable.test.ts @@ -0,0 +1,241 @@ +import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import type { Span } from '@sentry/core'; +import { + _INTERNAL_clearAiProviderSkips, + _INTERNAL_shouldSkipAiProviderWrapping, + getMainCarrier, + setCurrentClient, + spanToStaticSpanJSON, +} from '@sentry/core'; +import { instrumentPiDurableHarnessOptions } from '../../../../src/ai/pi-durable'; +import type { + PiCommitChange, + PiHarnessOptions, + PiLiveState, + PiModels, + PiRegistryReader, + PiTask, + PiTaskRuntime, +} from '../../../../src/ai/pi-durable/types'; +import { OPENAI_INTEGRATION_NAME } from '../../../../src/ai/openai/constants'; +import { getDefaultTestClientOptions, TestClient } from '../../../mocks/client'; + +const LIVE_DOC = { definition: { kind: 'pi.live' } }; + +describe('instrumentPiDurableHarnessOptions', () => { + let endedSpans: Span[]; + 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 = []; + client.on('spanEnd', span => endedSpans.push(span)); + }); + + afterEach(() => { + _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); + }); + + 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 = { + role: 'assistant', + content: [{ type: 'text', text: 'Paris.' }], + model: 'faux-model', + usage: { input: 10, output: 2, totalTokens: 12 }, + stopReason: 'stop', + }; + const models: PiModels = { streamSimple: () => ({ result: async () => message }) }; + const instrumented = instrumentPiDurableHarnessOptions({ models }); + + await instrumented.models!.streamSimple!( + { id: 'faux-model', provider: 'faux' }, + { 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('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!( + { id: 'faux-model', provider: 'faux' }, + { + 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."}]}]'); + }); + + it('captures the failures pi-durable reports and still hands them to the caller', async () => { + const reported: unknown[] = []; + const instrumented = instrumentPiDurableHarnessOptions({ 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); + }); + + describe('runs', () => { + /** A generation phase that commits each of `states` as the `pi.live` document, in order. */ + function generationPhase(states: (PiLiveState | 'settle-failed')[]): 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 (state === 'settle-failed') { + (tx as { settleSubmission: (...args: unknown[]) => void }).settleSubmission(5, { + status: 'unanswered', + reason: 'model_error', + }); + delete draft.run; + } else { + draft.run = state.run; + } + return undefined; + }, undefined); + } + }, + }, + }, + }; + } + + async function runPhase(task: PiTask): 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); + }, + }; + const registry: PiRegistryReader = { + snapshot: () => ({ task: () => task, tasks: () => [task] }), + subscribe: () => () => undefined, + }; + + const instrumented = instrumentPiDurableHarnessOptions({ registry }).registry!; + await instrumented.snapshot().task('pi.generation')!.definition.phases.request!(undefined, runtime, undefined); + } + + it('keeps the run open while run control moves to the next generation', async () => { + await runPhase(generationPhase([{ run: { inputs: [5] } }, { run: { inputs: [5, 6] } }])); + + expect(endedSpans.filter(span => spanToStaticSpanJSON(span).op === 'gen_ai.invoke_agent')).toHaveLength(0); + }); + + it('ends the run with the settlement of its first input', async () => { + await runPhase(generationPhase([{ run: { inputs: [5] } }, 'settle-failed'])); + + const runs = endedSpans.map(span => spanToStaticSpanJSON(span)).filter(span => span.op === 'gen_ai.invoke_agent'); + expect(runs).toHaveLength(1); + // The settlement reason is not a span status value, so the status reads as `internal_error`. + expect(runs[0]!.status).toBe('internal_error'); + expect(runs[0]!.data['gen_ai.conversation.id']).toMatch(/^[0-9a-f]{32}:1$/); + }); + + it('ends the run when a queued input starts the next run in the same commit', async () => { + await runPhase(generationPhase([{ run: { inputs: [5] } }, { run: { inputs: [7] } }])); + + const runs = endedSpans.map(span => spanToStaticSpanJSON(span)).filter(span => span.op === 'gen_ai.invoke_agent'); + expect(runs).toHaveLength(1); + expect(runs[0]!.status).toBe('ok'); + }); + }); +}); 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 f7a46b179266..42e04a4c1701 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 @@ -3,6 +3,8 @@ import { piAiAssistantMessageToGenAiMessage, piAiContentToString, piAiMessagesToGenAiMessages, + piAiSystemInstructions, + piAiToolDefinitions, } from '../../../../src/ai/pi-ai/messages'; describe('convert pi-ai messages to gen_ai messages', () => { @@ -100,6 +102,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([