Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
{
"name": "cloudflare-pi-durable",
"version": "0.0.0",
"private": true,
"scripts": {
"dev": "vite dev",
"build": "vite build",
"preview": "vite preview --port 38789",
"typecheck": "tsc --noEmit",
"test:build": "pnpm install && pnpm build",
"test:build-latest": "pnpm install && pnpm add agents@latest @earendil-works/pi-durable@latest @earendil-works/pi-ai@latest && pnpm build",
"test:assert": "playwright test"
},
"dependencies": {
"@earendil-works/pi-ai": "1.0.0",
"@earendil-works/pi-durable": "1.0.0",
"@sentry/cloudflare": "file:../../packed/sentry-cloudflare-packed.tgz",
"agents": "0.26.0"
},
"devDependencies": {
"@cloudflare/vite-plugin": "1.57.2",
"@cloudflare/workers-types": "^4.20260426.0",
"@playwright/test": "~1.63.0",
"@sentry-internal/test-utils": "link:../../../test-utils",
"typescript": "^5.5.2",
"vite": "7.3.5",
"wrangler": "^4.136.2"
},
"volta": {
"extends": "../../package.json"
},
"sentryTest": {
"optional": true,
"optionalVariants": [
{
"build-command": "pnpm test:build-latest",
"label": "cloudflare-pi-durable (latest)"
}
]
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
import { getPlaywrightConfig } from '@sentry-internal/test-utils';

const config = getPlaywrightConfig(
{
startCommand: 'pnpm preview',
port: 38789,
},
// Each test drives real OpenRouter tool-calling turns, and one resets the Durable Object mid-run,
// which does not fit the default 30s timeout when the provider is slow.
{ timeout: 120_000 },
);

export default config;
Original file line number Diff line number Diff line change
@@ -0,0 +1,96 @@
import { Type } from '@earendil-works/pi-ai';
import { createModels } from '@earendil-works/pi-ai/models';
import { openrouterProvider } from '@earendil-works/pi-ai/providers/openrouter';
import { createRegistry, defineExtension, defineTool, Harness, section } from '@earendil-works/pi-durable';
import * as Sentry from '@sentry/cloudflare';
import { Agent } from 'agents';
import { PiHarness } from 'agents/harness/pi';

/**
* pi-durable hosted by the Agents SDK `PiHarness`: pi keeps its state in this object's SQLite
* database, and the harness reopens pi when the object restarts. Nothing wires Sentry into
* pi-durable by hand, the Vite plugin injects the channel into `Harness.open()`.
*/
export class Assistant extends Agent<Env> {
public registry = createRegistry();

public harness = new PiHarness({
harness: ({ storage, context }) => {
// pi-ai reads provider keys from `process.env` by default, so the Worker hands over its secret.
const models = createModels({
authContext: {
env: async name => (name === 'OPENROUTER_API_KEY' ? this.env.E2E_OPENROUTER_API_KEY : undefined),
fileExists: async () => false,
},
});
models.setProvider(openrouterProvider());
return Harness.open(storage, { models, registry: this.registry }, context);
},
defaults: { model: { provider: 'openrouter', id: 'anthropic/claude-haiku-4.5' } },
});

public constructor(ctx: DurableObjectState, env: Env) {
super(ctx, env);
this.registry.install(
defineExtension({
name: 'e2e',
sections: [
section('preamble', () => 'You are a test assistant. Use the tools exactly as asked. Keep answers short.', {
tag: false,
}),
],
tools: [
defineTool({
name: 'get_weather',
description: 'Get the current weather for a city.',
parameters: Type.Object({ city: Type.String() }),
// The manual span should nest under the SDK's `execute_tool` span.
execute: async args =>
Sentry.startSpan({ name: 'resolve-weather', attributes: { 'weather.city': args.city } }, () => ({
content: [{ type: 'text', text: `It is 21 degrees and sunny in ${args.city}.` }],
})),
}),
defineTool({
name: 'fail_now',
description: 'Always throws an error. Call this when the user asks to trigger a failure.',
parameters: Type.Object({}),
execute: async () => {
throw new Error('Intentional pi-durable tool failure');
},
}),
defineTool({
name: 'crash_once',
description: 'Runs one step of a job. Call this when the user asks for crash_once.',
parameters: Type.Object({ job: Type.String() }),
// Replay-safe, so pi-durable reruns the call after the reset instead of failing it.
replay: 'safe',
execute: async args => {
const marker = `crashed:${args.job}`;
if (await this.ctx.storage.get(marker)) {
return { content: [{ type: 'text', text: `Step of job ${args.job} completed.` }] };
}
await this.ctx.storage.put(marker, true);
await this.ctx.storage.sync();
// Drops the object mid-call, like an eviction. The next request starts it again.
this.ctx.abort('crash_once resets the Durable Object');
throw new Error('The Durable Object did not reset');
},
}),
],
}),
);
this.lifecycle.use(this.harness);
}

public async onRequest(request: Request): Promise<Response> {
if (request.method === 'GET') {
const operationId = new URL(request.url).searchParams.get('operation') ?? '';
const { status, reason, text } = await this.harness.wait(operationId);
return Response.json({ status, reason, text });
}

const { message, operationId } = await request.json<{ message: string; operationId?: string }>();
const { status, reason, text } = await this.harness.prompt(message, { operationId });
return Response.json({ status, reason, text });
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
declare namespace Cloudflare {
interface Env {
E2E_TEST_DSN: string;
E2E_OPENROUTER_API_KEY: string;
Assistant: DurableObjectNamespace;
}
}

interface Env extends Cloudflare.Env {}
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
import { routeAgentRequest } from 'agents';

export { Assistant } from './assistant';

export default {
async fetch(request, env): Promise<Response> {
return (await routeAgentRequest(request, env)) ?? new Response('Not found', { status: 404 });
},
} satisfies ExportedHandler<Env>;
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
import { defineCloudflareOptions } from '@sentry/cloudflare';

// The Sentry Vite plugin picks this file up by convention, next to the worker entry named in
// wrangler's `main`, and hands its default export to `withSentry`.
export default defineCloudflareOptions((env: Env) => ({
dsn: env.E2E_TEST_DSN,
environment: 'qa',
tunnel: 'http://localhost:3031/',
tracesSampleRate: 1.0,
}));
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
import { startEventProxyServer } from '@sentry-internal/test-utils';

startEventProxyServer({
port: 3031,
proxyServerName: 'cloudflare-pi-durable',
});
Original file line number Diff line number Diff line change
@@ -0,0 +1,133 @@
import { expect, test } from '@playwright/test';
import { collectStreamedSpans, getSpanOp, waitForError } from '@sentry-internal/test-utils';
import { newAgentId, prompt, submitPrompt, waitForOperation } from './utils';

const APP = 'cloudflare-pi-durable';

test('traces a PiHarness prompt as invoke_agent with chat, execute_tool and provider spans', async ({ baseURL }) => {
// A run is its own trace, and its `invoke_agent` segment ends last. Only this test calls get_weather.
const spansPromise = collectStreamedSpans(
APP,
spansOfTrace =>
spansOfTrace.some(span => span.is_segment && getSpanOp(span) === 'gen_ai.invoke_agent') &&
spansOfTrace.some(span => span.attributes['gen_ai.tool.name']?.value === 'get_weather'),
);

const answer = await prompt(
baseURL!,
newAgentId('weather'),
'What is the weather in Vienna? Use the get_weather tool.',
);
expect(answer.status, answer.reason).toBe('done');

const spans = await spansPromise;
const agent = spans.find(span => span.is_segment)!;
const chats = spans.filter(span => getSpanOp(span) === 'gen_ai.chat');
// The manual span of the tool is the one way to pick the call that ran, should the model call twice.
const manualSpan = spans.find(span => span.name === 'resolve-weather')!;
const tool = spans.find(span => span.span_id === manualSpan.parent_span_id)!;
const providerCalls = spans.filter(span => getSpanOp(span) === 'http.client');

expect(agent.attributes['sentry.origin']?.value).toBe('auto.ai.pi_durable');
// The PiHarness root session is pi conversation 1, in a Harness of its own per Durable Object.
expect(agent.attributes['gen_ai.conversation.id']?.value).toMatch(/^[0-9a-f]{32}:1$/);
// The request that submitted the prompt is a trace of its own: the scheduler runs the work.
expect(spans.some(span => getSpanOp(span) === 'http.server')).toBe(false);
// pi-durable keeps its state in the `pi_` tables of `PiHarness`, whose statements start no span.
expect(spans.filter(span => getSpanOp(span) === 'db.query' && /\bpi_/.test(String(span.name)))).toEqual([]);
// pi-ai sends the requests through `@anthropic-ai/sdk`, whose own integration must stay out so
// each request is reported once.
expect(spans.filter(span => String(span.attributes['sentry.origin']?.value).startsWith('auto.ai.'))).toEqual(
spans.filter(span => span.attributes['sentry.origin']?.value === 'auto.ai.pi_durable'),
);

// One tool-calling response, then the answer.
expect(chats.length).toBeGreaterThanOrEqual(2);
for (const chat of chats) {
expect(chat.parent_span_id).toBe(agent.span_id);
expect(chat.attributes['gen_ai.provider.name']?.value).toBe('openrouter');
expect(chat.attributes['gen_ai.request.model']?.value).toBe('anthropic/claude-haiku-4.5');
expect(chat.attributes['gen_ai.conversation.id']?.value).toBe(agent.attributes['gen_ai.conversation.id']?.value);
// pi-durable retries a provider error inside the run; such a request has no usage and no HTTP span
// of its own to assert on.
if (chat.status === 'ok') {
expect(typeof chat.attributes['gen_ai.usage.input_tokens']?.value).toBe('number');
expect(typeof chat.attributes['gen_ai.usage.output_tokens']?.value).toBe('number');
expect(providerCalls.some(providerCall => providerCall.parent_span_id === chat.span_id)).toBe(true);
}
}

expect(tool.attributes['gen_ai.tool.name']?.value).toBe('get_weather');
expect(tool.parent_span_id).toBe(agent.span_id);
expect(tool.status).toBe('ok');
expect(tool.attributes['gen_ai.tool.call.arguments']?.value).toContain('Vienna');
expect(tool.attributes['gen_ai.tool.call.result']?.value).toContain('21 degrees and sunny in');

// The provider's HTTP calls nest inside the `chat` span that sent them.
expect(providerCalls.length).toBeGreaterThan(0);
for (const providerCall of providerCalls) {
expect(chats.map(chat => chat.span_id)).toContain(providerCall.parent_span_id);
expect(providerCall.attributes['server.address']?.value).toBe('openrouter.ai');
}
});

test('reports a throwing tool as an error on its execute_tool span', async ({ baseURL }) => {
const errorPromise = waitForError(
APP,
event => event.exception?.values?.[0]?.value === 'Intentional pi-durable tool failure',
);
const spansPromise = collectStreamedSpans(
APP,
spansOfTrace =>
spansOfTrace.some(span => span.is_segment && getSpanOp(span) === 'gen_ai.invoke_agent') &&
spansOfTrace.some(span => span.attributes['gen_ai.tool.name']?.value === 'fail_now'),
);

const answer = await prompt(baseURL!, newAgentId('fail'), 'Call the fail_now tool, then tell me what happened.');
expect(answer.status, answer.reason).toBe('done');

const [error, spans] = await Promise.all([errorPromise, spansPromise]);
// The error is reported on the span of the call that threw.
const tool = spans.find(span => span.span_id === error.contexts?.trace?.span_id)!;

expect(tool.attributes['gen_ai.tool.name']?.value).toBe('fail_now');
expect(tool.status).toBe('error');
expect(tool.trace_id).toBe(error.contexts?.trace?.trace_id);
expect(error.exception?.values?.[0]?.mechanism).toEqual({ type: 'auto.ai.pi_durable', handled: true });
});

test('resumes a run in a new trace after the Durable Object resets during a tool call', async ({ baseURL }) => {
const agentId = newAgentId('reset');
const operationId = crypto.randomUUID();
// The tool span of the reset object never ends, so a finished `crash_once` span can only come
// from the rerun after the restart.
const spansPromise = collectStreamedSpans(
APP,
spansOfTrace =>
spansOfTrace.some(span => span.is_segment && getSpanOp(span) === 'gen_ai.invoke_agent') &&
spansOfTrace.some(span => span.attributes['gen_ai.tool.name']?.value === 'crash_once' && span.status === 'ok'),
);

// The reset fails the request that submitted the prompt, but pi has stored the input already.
await submitPrompt(
baseURL!,
agentId,
'Run the crash_once tool for job "nightly", then reply with the word DONE.',
operationId,
);
const answer = await waitForOperation(baseURL!, agentId, operationId);
expect(answer.status, answer.reason).toBe('done');

const spans = await spansPromise;
const agent = spans.find(span => span.is_segment)!;
const tool = spans.find(span => span.attributes['gen_ai.tool.name']?.value === 'crash_once')!;
const chats = spans.filter(span => getSpanOp(span) === 'gen_ai.chat');
const finalAnswer = chats.find(span => span.attributes['gen_ai.response.finish_reasons']?.value === '["stop"]');

expect(getSpanOp(agent)).toBe('gen_ai.invoke_agent');
expect(tool.parent_span_id).toBe(agent.span_id);
expect(finalAnswer?.parent_span_id).toBe(agent.span_id);
// The request that called the tool ran before the reset, so this trace starts with the rerun of the
// tool. A run that never reset would start with that request.
expect(chats.every(chat => chat.start_timestamp >= tool.start_timestamp)).toBe(true);
});
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
import { expect } from '@playwright/test';

export type Operation = { status: string; reason?: string; text?: string };

/**
* A Durable Object name nothing has used yet. The name is the `:name` path segment
* `routeAgentRequest` maps to an object, so each test gets its own SQLite database and pi Harness.
*/
export function newAgentId(prefix: string): string {
return `${prefix}-${Date.now()}-${Math.random().toString(36).slice(2, 8)}`;
}

/** Submit a prompt to the root session. The response arrives once pi settles it. */
export function submitPrompt(
baseURL: string,
agentId: string,
message: string,
operationId?: string,
): Promise<Response> {
return fetch(`${baseURL}/agents/assistant/${agentId}`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ message, operationId }),
});
}

export async function prompt(baseURL: string, agentId: string, message: string): Promise<Operation> {
const res = await submitPrompt(baseURL, agentId, message);
expect(res.status).toBe(200);
return (await res.json()) as Operation;
}

/** Wait until pi settles an operation. The request starts the object again if it was reset. */
export async function waitForOperation(baseURL: string, agentId: string, operationId: string): Promise<Operation> {
const res = await fetch(`${baseURL}/agents/assistant/${agentId}?operation=${operationId}`);
expect(res.status).toBe(200);
return (await res.json()) as Operation;
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
{
"compilerOptions": {
"target": "es2022",
"lib": ["es2022"],
"module": "es2022",
"moduleResolution": "Bundler",
"resolveJsonModule": true,
"allowJs": true,
"checkJs": false,
"noEmit": true,
"isolatedModules": true,
"allowSyntheticDefaultImports": true,
"forceConsistentCasingInFileNames": true,
"strict": true,
"skipLibCheck": true,
"types": ["@cloudflare/workers-types"]
},
"exclude": ["tests"],
"include": ["src/**/*.ts"]
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
import { cloudflare } from '@cloudflare/vite-plugin';
import { sentryCloudflareVitePlugin } from '@sentry/cloudflare/vite';
import { defineConfig } from 'vite';

export default defineConfig({
plugins: [cloudflare(), sentryCloudflareVitePlugin()],
});
Loading
Loading