Skip to content
Draft
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,5 @@
node_modules
.data
results.junit.xml
test-results
playwright-report
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
import * as Sentry from '@sentry/node';

Sentry.init({
environment: 'qa',
dsn: process.env.E2E_TEST_DSN,
tunnel: 'http://localhost:3031/', // proxy server
tracesSampleRate: 1.0,
});
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
{
"name": "node-pi-durable",
"version": "0.0.0",
"private": true,
"type": "module",
"scripts": {
"start": "node src/start.mjs",
"clean": "npx rimraf node_modules .data pnpm-lock.yaml",
"test:build": "pnpm install",
"test:build-latest": "pnpm install && pnpm add @earendil-works/pi-durable@latest @earendil-works/pi-ai@latest @earendil-works/chord@latest",
"test:assert": "OPENROUTER_API_KEY=$E2E_OPENROUTER_API_KEY playwright test"
},
"dependencies": {
"@earendil-works/chord": "1.0.0",
"@earendil-works/pi-ai": "1.0.0",
"@earendil-works/pi-durable": "1.0.0",
"@sentry/node": "file:../../packed/sentry-node-packed.tgz"
},
"devDependencies": {
"@playwright/test": "~1.63.0",
"@sentry-internal/test-utils": "link:../../../test-utils",
"@sentry/core": "file:../../packed/sentry-core-packed.tgz",
"@types/node": "24.x"
},
"engines": {
"node": "24.x"
},
"volta": {
"node": "24.15.0",
"extends": "../../package.json"
},
"sentryTest": {
"optional": true,
"optionalVariants": [
{
"build-command": "pnpm test:build-latest",
"label": "node-pi-durable (latest)"
}
]
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
import { getPlaywrightConfig } from '@sentry-internal/test-utils';

const config = getPlaywrightConfig(
{ startCommand: 'pnpm start' },
// Each test drives real OpenRouter tool-calling turns, and one restarts the server 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,162 @@
import { mkdirSync, writeFileSync } from 'node:fs';
import { createServer } from 'node:http';
import { fileURLToPath } from 'node:url';
import { BACKGROUND_CONTEXT } from '@earendil-works/chord/context';
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 {
AssistantEntry,
createRegistry,
defineExtension,
defineTool,
Harness,
section,
} from '@earendil-works/pi-durable';
import { openNodeSqliteStorage } from '@earendil-works/pi-durable/storage/sqlite/node';
import * as Sentry from '@sentry/node';

const context = BACKGROUND_CONTEXT;
const MODEL = { provider: 'openrouter', modelId: 'anthropic/claude-haiku-4.5' };
const CRASH_EXIT_CODE = 75;

if (!process.env.OPENROUTER_API_KEY) {
throw new Error('OPENROUTER_API_KEY is not set (E2E_OPENROUTER_API_KEY in the e2e-tests .env)');
}

const dataDir = fileURLToPath(new URL('../.data/', import.meta.url));
mkdirSync(dataDir, { recursive: true });

const models = createModels();
models.setProvider(openrouterProvider()); // reads OPENROUTER_API_KEY

const 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: 'flaky_step',
description: 'Runs one step of a job. Call this when the user asks to run the flaky step.',
parameters: Type.Object({ job: Type.String() }),
// Replay-safe, so pi-durable reruns the call after the crash instead of failing it.
replay: 'safe',
execute: async args => {
// The marker is created once, atomically; its presence means the first attempt already crashed.
try {
writeFileSync(`${dataDir}flaky-${args.job}`, 'crashed once', { flag: 'wx' });
} catch {
return { content: [{ type: 'text', text: `Step of job ${args.job} completed.` }] };
}
process.exit(CRASH_EXIT_CODE);
},
}),
defineTool({
name: 'delegate',
description: 'Delegate a self-contained task to a subagent and get its answer back.',
parameters: Type.Object({ task: Type.String() }),
replay: 'safe',
execute: async (args, api, toolContext) => {
const childId = await api.commit(async tx => {
const existing = (await tx.scanConversations({ ownerTaskId: api.taskId }, 1)).items[0];
if (existing !== undefined) return existing.id;
return (await tx.createConversation({ ownership: { kind: 'task', taskId: api.taskId } })).id;
}, toolContext);
const child = await api.conversation(childId, toolContext);
const request = { type: 'input', content: args.task, requestId: `delegate:${api.taskId}` };
const settled = await (await child.submit(request, toolContext)).wait(toolContext);
if (settled.status !== 'done') {
return { isError: true, content: [{ type: 'text', text: `Subagent ended: ${settled.reason}` }] };
}
const answer = await api.commit(tx => tx.entry(AssistantEntry, settled.answer), toolContext);
const text = (answer?.model?.[0]?.content ?? [])
.filter(part => part.type === 'text')
.map(part => part.text)
.join('');
return { content: [{ type: 'text', text }] };
},
}),
];

const registry = createRegistry();
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,
}),
);

const storage = await openNodeSqliteStorage(`${dataDir}pi.sqlite`);
const harness = await Harness.open(storage, { models, registry }, context);
// Continue whatever the previous process left unfinished, such as a run interrupted by `flaky_step`.
harness.resume();

async function readJson(req) {
let body = '';
for await (const chunk of req) body += chunk;
return body ? JSON.parse(body) : {};
}

async function handle(req, res) {
const url = new URL(req.url, 'http://localhost');
const parts = url.pathname.split('/').filter(Boolean);

if (req.method === 'POST' && url.pathname === '/conversations') {
const conversation = await harness.createConversation(
{ ownership: { kind: 'ownerless' }, agent: { model: MODEL } },
context,
);
return { id: conversation.id };
}

// Returns as soon as the input is admitted. The run continues in the background, possibly in a
// later process, so clients poll `GET /submissions/:id` for the outcome.
if (req.method === 'POST' && parts[0] === 'conversations' && parts[2] === 'messages') {
const conversation = await harness.conversation(Number(parts[1]), context);
if (!conversation) return undefined;
const { content } = await readJson(req);
const submission = await conversation.submit({ type: 'input', content }, context);
return { submissionId: submission.id };
}

if (req.method === 'GET' && parts[0] === 'submissions') {
const submission = await harness.submission(Number(parts[1]), context);
if (!submission) return undefined;
const record = await submission.status(context);
return { status: record.status, reason: record.reason, detail: record.detail };
}

return undefined;
}

createServer((req, res) => {
handle(req, res).then(
result => {
res.writeHead(result ? 200 : 404, { 'Content-Type': 'application/json' });
res.end(JSON.stringify(result ?? { error: 'not found' }));
},
error => {
res.writeHead(500, { 'Content-Type': 'application/json' });
res.end(JSON.stringify({ error: String(error) }));
},
);
}).listen(3030);
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
import { spawn } from 'node:child_process';
import { rmSync } from 'node:fs';
import { fileURLToPath } from 'node:url';

// The `flaky_step` tool exits the server mid-call with this code to simulate a crash. Restarting the
// server on the same storage, as a process manager would, is what lets the run resume.
const CRASH_EXIT_CODE = 75;

const appDir = fileURLToPath(new URL('..', import.meta.url));
rmSync(new URL('../.data', import.meta.url), { recursive: true, force: true });

let server;

function startServer() {
server = spawn(process.execPath, ['--import', './instrument.mjs', 'src/server.mjs'], {
cwd: appDir,
stdio: 'inherit',
});
server.on('exit', code => {
if (code === CRASH_EXIT_CODE) {
startServer();
} else {
process.exit(code ?? 1);
}
});
}

for (const signal of ['SIGINT', 'SIGTERM']) {
process.on(signal, () => {
server.kill(signal);
process.exit(0);
});
}

startServer();
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: 'node-pi-durable',
});
Loading
Loading