Skip to content
2 changes: 2 additions & 0 deletions packages/core/src/shared-exports.ts
Original file line number Diff line number Diff line change
Expand Up @@ -230,6 +230,8 @@ export {
export { wrapToolsWithSpans, extractLLMFromParams, extractAgentNameFromParams } from './tracing/langgraph/utils';
export { LANGGRAPH_INTEGRATION_NAME } from './tracing/langgraph/constants';
export type { LangGraphOptions, LangGraphIntegration, CompiledGraph } from './tracing/langgraph/types';
export { instrumentWorkersAiClient } from './tracing/workers-ai';
export type { WorkersAiClient, WorkersAiOptions } from './tracing/workers-ai/types';
// eslint-disable-next-line typescript/no-deprecated
export type { OpenAiClient, OpenAiOptions, InstrumentedMethod } from './tracing/openai/types';
export type {
Expand Down
10 changes: 10 additions & 0 deletions packages/core/src/tracing/workers-ai/constants.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
/**
* The provider value for the `gen_ai.provider.name` attribute.
* @see https://developers.cloudflare.com/workers-ai/
*/
export const WORKERS_AI_PROVIDER_NAME = 'cloudflare.workers_ai';

/**
* The Sentry origin for spans created by the Workers AI instrumentation.
*/
export const WORKERS_AI_ORIGIN = 'auto.ai.cloudflare.workers_ai';
133 changes: 133 additions & 0 deletions packages/core/src/tracing/workers-ai/index.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,133 @@
import { SPAN_STATUS_ERROR } from '../../tracing';
import { startSpan, startSpanManual } from '../../tracing/trace';
import type { Span } from '../../types/span';
import { isObjectLike } from '../../utils/is';
import { resolveAIRecordingOptions, shouldEnableTruncation } from '../ai/utils';
import { instrumentWorkersAiStream } from './streaming';
import type { WorkersAiOptions } from './types';
import { addRequestAttributes, addResponseAttributes, extractRequestAttributes, getOperationName } from './utils';

// Adapted from /server-utils/src/vercel-ai/util.ts
// TODO(v11): Reuse this function once this gets moved to @sentry/server-utils
// Workers AI streaming responses are SSE byte streams, so we narrow to `Uint8Array`.
function isReadableStream(value: unknown): value is ReadableStream<Uint8Array> {
return (
isObjectLike(value) &&
typeof (value as { pipeThrough?: unknown }).pipeThrough === 'function' &&
typeof (value as { getReader?: unknown }).getReader === 'function'
);
}

/**
* Wrap the `run` method of the Workers AI binding with Sentry tracing.
*/
function instrumentRun(
originalRun: (...args: unknown[]) => Promise<unknown>,
context: unknown,
options: WorkersAiOptions & Required<Pick<WorkersAiOptions, 'recordInputs' | 'recordOutputs'>>,
): (...args: unknown[]) => Promise<unknown> {
return function instrumentedRun(...args: unknown[]): Promise<unknown> {
const [model, inputs, runOptions] = args as [unknown, unknown, Record<string, unknown> | undefined];

const operationName = getOperationName(inputs);
const requestAttributes = extractRequestAttributes(model, inputs, operationName);
const modelName = typeof model === 'string' ? model : 'unknown';

const isStreamRequested =
!!inputs && typeof inputs === 'object' && (inputs as { stream?: unknown }).stream === true;
const returnsRawResponse =
!!runOptions &&
typeof runOptions === 'object' &&
(runOptions.returnRawResponse === true || runOptions.websocket === true);

const spanConfig = {
name: `${operationName} ${modelName}`,
op: `gen_ai.${operationName}`,
attributes: requestAttributes,
};

if (isStreamRequested && !returnsRawResponse) {
return startSpanManual(spanConfig, (span: Span) => {
// `startSpanManual` does not auto-end the span, so we must end it on every exit path,
// including a synchronous throw from `run`.
const handleError = (error: unknown): never => {
span.setStatus({ code: SPAN_STATUS_ERROR, message: 'internal_error' });
span.end();
throw error;
};

let originalResult: Promise<unknown>;

try {
originalResult = originalRun.apply(context, args) as Promise<unknown>;
} catch (error) {
return handleError(error);
}

if (options.recordInputs) {
addRequestAttributes(span, inputs, operationName, shouldEnableTruncation(options.enableTruncation));
}

return originalResult.then(result => {
if (isReadableStream(result)) {
return instrumentWorkersAiStream(result, span, options.recordOutputs);
}

// The model did not actually return a stream — finalize the span eagerly.
addResponseAttributes(span, result, options.recordOutputs);
span.end();
return result;
}, handleError);
});
}

return startSpan(spanConfig, (span: Span) => {
const originalResult = originalRun.apply(context, args) as Promise<unknown>;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Bug: Synchronous errors in the non-streaming run method are not reported to Sentry as exceptions, only as a span status.
Severity: LOW

Suggested Fix

Wrap the originalRun.apply(context, args) call in a try...catch block, similar to the streaming path. In the catch block, call captureException to report the error to Sentry before re-throwing it.

Prompt for AI Agent
Review the code at the location below. A potential bug has been identified by an AI
agent. Verify if this is a real issue. If it is, propose a fix; if not, explain why it's
not valid.

Location: packages/core/src/tracing/workers-ai/index.ts#L90

Potential issue: In the non-streaming path for Workers AI tracing, if the original `run`
method throws a synchronous error, the error is not reported to Sentry via
`captureException`. The error is caught by `startSpan`'s internal `handleCallbackErrors`
function, which sets the span status to error but does not create a Sentry exception
event. This is inconsistent with the streaming path, which explicitly catches
synchronous errors and reports them. While synchronous throws from async methods are
unlikely, this represents a gap in error monitoring.


if (options.recordInputs) {
addRequestAttributes(span, inputs, operationName, shouldEnableTruncation(options.enableTruncation));
}

return originalResult.then(result => {
if (!returnsRawResponse) {
addResponseAttributes(span, result, options.recordOutputs);
}
return result;
});
});
};
}

/**
* Instrument a Cloudflare Workers AI binding (`env.AI`) with Sentry tracing.
*
* This wraps the binding's `run` method to create `gen_ai` spans following the
* Sentry AI Agents conventions. All other methods are passed through untouched.
*
* In `@sentry/cloudflare`, the `env.AI` binding is instrumented automatically —
* wrapping manually is only needed to pass custom options.
*
* @example
* ```javascript
* const ai = Sentry.instrumentWorkersAiClient(env.AI, { recordInputs: true, recordOutputs: true });
* const result = await ai.run('@cf/meta/llama-3.1-8b-instruct', { prompt: 'Hello' });
* ```
*/
export function instrumentWorkersAiClient<T extends object>(client: T, options?: WorkersAiOptions): T {
const resolvedOptions = resolveAIRecordingOptions(options);

const instrumented = new Proxy(client, {
get(target: object, prop: string | symbol, receiver: unknown): unknown {
Comment thread
isaacs marked this conversation as resolved.
const value = Reflect.get(target, prop, receiver);

if (prop === 'run' && typeof value === 'function') {
return instrumentRun(value as (...args: unknown[]) => Promise<unknown>, target, resolvedOptions);
}

// Bind passed-through functions to the original target to preserve `this` (e.g. private fields).
return typeof value === 'function' ? (value as (...args: unknown[]) => unknown).bind(target) : value;
},
}) as T;

return instrumented;
Comment thread
JPeer264 marked this conversation as resolved.
}
Comment thread
JPeer264 marked this conversation as resolved.
229 changes: 229 additions & 0 deletions packages/core/src/tracing/workers-ai/streaming.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,229 @@
import { SPAN_STATUS_ERROR } from '../../tracing';
import type { Span } from '../../types/span';
import { endStreamSpan, type StreamResponseState } from '../ai/utils';
import type { WorkersAiUsage } from './types';
import { setOutputMessagesAttribute } from './utils';

interface WorkersAiStreamingToolCall {
index?: number;
id?: string;
type?: string;
function?: { name?: string; arguments?: string };
// Some Workers AI models stream tool calls with the name/arguments at the top
// level of the tool-call object instead of nested under `function`.
name?: string;
arguments?: string;
}

interface WorkersAiStreamChunk {
// Native Workers AI streaming shape (`env.AI.run` with `stream: true`).
response?: unknown;
tool_calls?: unknown[];
// OpenAI-compatible streaming shape emitted for models routed through the
// OpenAI-compatible endpoint (e.g. via `workers-ai-provider`).
choices?: Array<{
delta?: { content?: unknown; tool_calls?: WorkersAiStreamingToolCall[] };
finish_reason?: unknown;
}>;
usage?: WorkersAiUsage & { prompt_tokens?: number; completion_tokens?: number; total_tokens?: number };
}

/**
* Accumulate a fragmented OpenAI-compatible tool call (delivered across multiple
* `choices[].delta.tool_calls` chunks) into the index-keyed accumulator.
*/
function accumulateStreamingToolCalls(
toolCalls: WorkersAiStreamingToolCall[],
accumulator: Record<number, WorkersAiStreamingToolCall>,
): void {
for (const toolCall of toolCalls) {
// Normalize both shapes: name/arguments nested under `function`, or at the top level.
const name = toolCall.function?.name ?? toolCall.name;
const args = toolCall.function?.arguments ?? toolCall.arguments;

// A tool call must carry at least a name or argument fragment to be meaningful.
if (name == null && args == null) {
continue;
}

const index = toolCall.index ?? 0;
const existing = accumulator[index];

if (!existing) {
accumulator[index] = {
index,
id: toolCall.id,
type: toolCall.type,
function: {
name,
arguments: args ?? '',
},
};
} else if (existing.function) {
if (name && !existing.function.name) {
existing.function.name = name;
}
if (args) {
existing.function.arguments = `${existing.function.arguments ?? ''}${args}`;
}
}
}
}

/**
* Parse a single SSE line (`data: {...}`) and accumulate its data into the streaming state.
*
* Handles both the native Workers AI shape (top-level `response`/`tool_calls`) and the
* OpenAI-compatible shape (`choices[].delta.content`/`choices[].delta.tool_calls`), because
* the same `run()` call transparently yields either format depending on the model.
*/
function processLine(
line: string,
state: StreamResponseState,
recordOutputs: boolean,
toolCallAccumulator: Record<number, WorkersAiStreamingToolCall>,
): void {
const trimmed = line.trim();
if (!trimmed.startsWith('data:')) {
return;
}

const data = trimmed.slice('data:'.length).trim();
if (!data || data === '[DONE]') {
return;
}

let parsed: WorkersAiStreamChunk;
try {
parsed = JSON.parse(data) as WorkersAiStreamChunk;
} catch {
return;
}

if (parsed.usage) {
if (typeof parsed.usage.prompt_tokens === 'number') {
state.promptTokens = parsed.usage.prompt_tokens;
}
if (typeof parsed.usage.completion_tokens === 'number') {
state.completionTokens = parsed.usage.completion_tokens;
}
if (typeof parsed.usage.total_tokens === 'number') {
state.totalTokens = parsed.usage.total_tokens;
}
}

if (recordOutputs && typeof parsed.response === 'string') {
state.responseTexts.push(parsed.response);
}

if (recordOutputs && Array.isArray(parsed.tool_calls) && parsed.tool_calls.length > 0) {
state.toolCalls.push(...parsed.tool_calls);
}

if (Array.isArray(parsed.choices)) {
for (const choice of parsed.choices) {
if (recordOutputs && typeof choice.delta?.content === 'string' && choice.delta.content) {
state.responseTexts.push(choice.delta.content);
}
if (recordOutputs && Array.isArray(choice.delta?.tool_calls)) {
accumulateStreamingToolCalls(choice.delta.tool_calls, toolCallAccumulator);
}
if (typeof choice.finish_reason === 'string') {
state.finishReasons.push(choice.finish_reason);
}
}
}
}

/**
* Wrap a Workers AI streaming response (a server-sent-events `ReadableStream`) so we can
* accumulate the response text and token usage while passing the original bytes through untouched.
*
* The span is ended once the consumer finishes reading (or cancels) the stream.
*/
export function instrumentWorkersAiStream(
stream: ReadableStream<Uint8Array>,
span: Span,
recordOutputs: boolean,
): ReadableStream<Uint8Array> {
const reader = stream.getReader();
const decoder = new TextDecoder();

const state: StreamResponseState = {
responseId: '',
responseModel: '',
finishReasons: [],
responseTexts: [],
toolCalls: [],
promptTokens: undefined,
completionTokens: undefined,
totalTokens: undefined,
};

// OpenAI-compatible tool calls arrive fragmented across chunks and are keyed by index;
// accumulate them here and flatten into `state.toolCalls` once the stream ends.
const toolCallAccumulator: Record<number, WorkersAiStreamingToolCall> = {};

let buffer = '';
let spanEnded = false;

const finish = (): void => {
if (spanEnded) {
return;
}
spanEnded = true;

if (recordOutputs) {
const accumulatedToolCalls = Object.values(toolCallAccumulator);
if (accumulatedToolCalls.length > 0) {
state.toolCalls.push(...accumulatedToolCalls);
}

// Set the authoritative `gen_ai.output.messages` alongside the deprecated response
// attributes `endStreamSpan` writes, so tool calls survive Relay's lossy migration.
setOutputMessagesAttribute(span, {
responseText: state.responseTexts.join(''),
toolCalls: state.toolCalls,
});
}

endStreamSpan(span, state, recordOutputs);
};

const flushBuffer = (isDone: boolean): void => {
const lines = buffer.split('\n');
// Keep the last (potentially incomplete) line in the buffer unless the stream is done.
buffer = isDone ? '' : (lines.pop() ?? '');
for (const line of lines) {
processLine(line, state, recordOutputs, toolCallAccumulator);
}
};

return new ReadableStream<Uint8Array>({
async pull(controller) {
try {
const { done, value } = await reader.read();

if (done) {
buffer += decoder.decode();
flushBuffer(true);
finish();
controller.close();
return;
}

buffer += decoder.decode(value, { stream: true });
flushBuffer(false);
controller.enqueue(value);
} catch (error) {
span.setStatus({ code: SPAN_STATUS_ERROR, message: 'internal_error' });
finish();
controller.error(error);
}
},
async cancel(reason) {
finish();
await reader.cancel(reason);
},
});
}
Loading
Loading