Skip to content

Commit 4086ec4

Browse files
JPeer264claude
andcommitted
test(e2e): Add a node-pi-durable end-to-end application
Covers what the node integration test cannot: the packed SDK, real model requests through pi-ai's `@anthropic-ai/sdk` path, SQLite storage, and a run that resumes after the server crashes during a tool call. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
1 parent 1dad424 commit 4086ec4

9 files changed

Lines changed: 468 additions & 0 deletions

File tree

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
node_modules
2+
.data
3+
results.junit.xml
4+
test-results
5+
playwright-report
Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,8 @@
1+
import * as Sentry from '@sentry/node';
2+
3+
Sentry.init({
4+
environment: 'qa',
5+
dsn: process.env.E2E_TEST_DSN,
6+
tunnel: 'http://localhost:3031/', // proxy server
7+
tracesSampleRate: 1.0,
8+
});
Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,41 @@
1+
{
2+
"name": "node-pi-durable",
3+
"version": "0.0.0",
4+
"private": true,
5+
"type": "module",
6+
"scripts": {
7+
"start": "node src/start.mjs",
8+
"clean": "npx rimraf node_modules .data pnpm-lock.yaml",
9+
"test:build": "pnpm install",
10+
"test:build-latest": "pnpm install && pnpm add @earendil-works/pi-durable@latest @earendil-works/pi-ai@latest @earendil-works/chord@latest",
11+
"test:assert": "OPENROUTER_API_KEY=$E2E_OPENROUTER_API_KEY playwright test"
12+
},
13+
"dependencies": {
14+
"@earendil-works/chord": "1.0.0",
15+
"@earendil-works/pi-ai": "1.0.0",
16+
"@earendil-works/pi-durable": "1.0.0",
17+
"@sentry/node": "file:../../packed/sentry-node-packed.tgz"
18+
},
19+
"devDependencies": {
20+
"@playwright/test": "~1.56.0",
21+
"@sentry-internal/test-utils": "link:../../../test-utils",
22+
"@sentry/core": "file:../../packed/sentry-core-packed.tgz",
23+
"@types/node": "24.x"
24+
},
25+
"engines": {
26+
"node": "24.x"
27+
},
28+
"volta": {
29+
"node": "24.15.0",
30+
"extends": "../../package.json"
31+
},
32+
"sentryTest": {
33+
"optional": true,
34+
"optionalVariants": [
35+
{
36+
"build-command": "pnpm test:build-latest",
37+
"label": "node-pi-durable (latest)"
38+
}
39+
]
40+
}
41+
}
Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,10 @@
1+
import { getPlaywrightConfig } from '@sentry-internal/test-utils';
2+
3+
const config = getPlaywrightConfig(
4+
{ startCommand: 'pnpm start' },
5+
// Each test drives real OpenRouter tool-calling turns, and one restarts the server mid-run, which
6+
// does not fit the default 30s timeout when the provider is slow.
7+
{ timeout: 120_000 },
8+
);
9+
10+
export default config;
Lines changed: 157 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,157 @@
1+
import { existsSync, mkdirSync, writeFileSync } from 'node:fs';
2+
import { createServer } from 'node:http';
3+
import { fileURLToPath } from 'node:url';
4+
import { BACKGROUND_CONTEXT } from '@earendil-works/chord/context';
5+
import { Type } from '@earendil-works/pi-ai';
6+
import { createModels } from '@earendil-works/pi-ai/models';
7+
import { openrouterProvider } from '@earendil-works/pi-ai/providers/openrouter';
8+
import {
9+
AssistantEntry,
10+
createRegistry,
11+
defineExtension,
12+
defineTool,
13+
Harness,
14+
section,
15+
} from '@earendil-works/pi-durable';
16+
import { openNodeSqliteStorage } from '@earendil-works/pi-durable/storage/sqlite/node';
17+
import * as Sentry from '@sentry/node';
18+
19+
const context = BACKGROUND_CONTEXT;
20+
const MODEL = { provider: 'openrouter', modelId: 'anthropic/claude-haiku-4.5' };
21+
const CRASH_EXIT_CODE = 75;
22+
23+
const dataDir = fileURLToPath(new URL('../.data/', import.meta.url));
24+
mkdirSync(dataDir, { recursive: true });
25+
26+
const models = createModels();
27+
models.setProvider(openrouterProvider()); // reads OPENROUTER_API_KEY
28+
29+
const tools = [
30+
defineTool({
31+
name: 'get_weather',
32+
description: 'Get the current weather for a city.',
33+
parameters: Type.Object({ city: Type.String() }),
34+
// The manual span should nest under the SDK's `execute_tool` span.
35+
execute: async args =>
36+
Sentry.startSpan({ name: 'resolve-weather', attributes: { 'weather.city': args.city } }, () => ({
37+
content: [{ type: 'text', text: `It is 21 degrees and sunny in ${args.city}.` }],
38+
})),
39+
}),
40+
defineTool({
41+
name: 'fail_now',
42+
description: 'Always throws an error. Call this when the user asks to trigger a failure.',
43+
parameters: Type.Object({}),
44+
execute: async () => {
45+
throw new Error('Intentional pi-durable tool failure');
46+
},
47+
}),
48+
defineTool({
49+
name: 'flaky_step',
50+
description: 'Runs one step of a job. Call this when the user asks to run the flaky step.',
51+
parameters: Type.Object({ job: Type.String() }),
52+
// Replay-safe, so pi-durable reruns the call after the crash instead of failing it.
53+
replay: 'safe',
54+
execute: async args => {
55+
const marker = `${dataDir}flaky-${args.job}`;
56+
if (!existsSync(marker)) {
57+
writeFileSync(marker, 'crashed once');
58+
process.exit(CRASH_EXIT_CODE);
59+
}
60+
return { content: [{ type: 'text', text: `Step of job ${args.job} completed.` }] };
61+
},
62+
}),
63+
defineTool({
64+
name: 'delegate',
65+
description: 'Delegate a self-contained task to a subagent and get its answer back.',
66+
parameters: Type.Object({ task: Type.String() }),
67+
replay: 'safe',
68+
execute: async (args, api, toolContext) => {
69+
const childId = await api.commit(async tx => {
70+
const existing = (await tx.scanConversations({ ownerTaskId: api.taskId }, 1)).items[0];
71+
if (existing !== undefined) return existing.id;
72+
return (await tx.createConversation({ ownership: { kind: 'task', taskId: api.taskId } })).id;
73+
}, toolContext);
74+
const child = await api.conversation(childId, toolContext);
75+
const request = { type: 'input', content: args.task, requestId: `delegate:${api.taskId}` };
76+
const settled = await (await child.submit(request, toolContext)).wait(toolContext);
77+
if (settled.status !== 'done') {
78+
return { isError: true, content: [{ type: 'text', text: `Subagent ended: ${settled.reason}` }] };
79+
}
80+
const answer = await api.commit(tx => tx.entry(AssistantEntry, settled.answer), toolContext);
81+
const text = (answer?.model?.[0]?.content ?? [])
82+
.filter(part => part.type === 'text')
83+
.map(part => part.text)
84+
.join('');
85+
return { content: [{ type: 'text', text }] };
86+
},
87+
}),
88+
];
89+
90+
const registry = createRegistry();
91+
registry.install(
92+
defineExtension({
93+
name: 'e2e',
94+
sections: [
95+
section('preamble', () => 'You are a test assistant. Use the tools exactly as asked. Keep answers short.', {
96+
tag: false,
97+
}),
98+
],
99+
tools,
100+
}),
101+
);
102+
103+
const storage = await openNodeSqliteStorage(`${dataDir}pi.sqlite`);
104+
const harness = await Harness.open(storage, { models, registry }, context);
105+
// Continue whatever the previous process left unfinished, such as a run interrupted by `flaky_step`.
106+
harness.resume();
107+
108+
async function readJson(req) {
109+
let body = '';
110+
for await (const chunk of req) body += chunk;
111+
return body ? JSON.parse(body) : {};
112+
}
113+
114+
async function handle(req, res) {
115+
const url = new URL(req.url, 'http://localhost');
116+
const parts = url.pathname.split('/').filter(Boolean);
117+
118+
if (req.method === 'POST' && url.pathname === '/conversations') {
119+
const conversation = await harness.createConversation(
120+
{ ownership: { kind: 'ownerless' }, agent: { model: MODEL } },
121+
context,
122+
);
123+
return { id: conversation.id };
124+
}
125+
126+
// Returns as soon as the input is admitted. The run continues in the background, possibly in a
127+
// later process, so clients poll `GET /submissions/:id` for the outcome.
128+
if (req.method === 'POST' && parts[0] === 'conversations' && parts[2] === 'messages') {
129+
const conversation = await harness.conversation(Number(parts[1]), context);
130+
if (!conversation) return undefined;
131+
const { content } = await readJson(req);
132+
const submission = await conversation.submit({ type: 'input', content }, context);
133+
return { submissionId: submission.id };
134+
}
135+
136+
if (req.method === 'GET' && parts[0] === 'submissions') {
137+
const submission = await harness.submission(Number(parts[1]), context);
138+
if (!submission) return undefined;
139+
const record = await submission.status(context);
140+
return { status: record.status, reason: record.reason };
141+
}
142+
143+
return undefined;
144+
}
145+
146+
createServer((req, res) => {
147+
handle(req, res).then(
148+
result => {
149+
res.writeHead(result ? 200 : 404, { 'Content-Type': 'application/json' });
150+
res.end(JSON.stringify(result ?? { error: 'not found' }));
151+
},
152+
error => {
153+
res.writeHead(500, { 'Content-Type': 'application/json' });
154+
res.end(JSON.stringify({ error: String(error) }));
155+
},
156+
);
157+
}).listen(3030);
Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
1+
import { spawn } from 'node:child_process';
2+
import { rmSync } from 'node:fs';
3+
import { fileURLToPath } from 'node:url';
4+
5+
// The `flaky_step` tool exits the server mid-call with this code to simulate a crash. Restarting the
6+
// server on the same storage, as a process manager would, is what lets the run resume.
7+
const CRASH_EXIT_CODE = 75;
8+
9+
const appDir = fileURLToPath(new URL('..', import.meta.url));
10+
rmSync(new URL('../.data', import.meta.url), { recursive: true, force: true });
11+
12+
let server;
13+
14+
function startServer() {
15+
server = spawn(process.execPath, ['--import', './instrument.mjs', 'src/server.mjs'], {
16+
cwd: appDir,
17+
stdio: 'inherit',
18+
});
19+
server.on('exit', code => {
20+
if (code === CRASH_EXIT_CODE) {
21+
startServer();
22+
} else {
23+
process.exit(code ?? 1);
24+
}
25+
});
26+
}
27+
28+
for (const signal of ['SIGINT', 'SIGTERM']) {
29+
process.on(signal, () => {
30+
server.kill(signal);
31+
process.exit(0);
32+
});
33+
}
34+
35+
startServer();
Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
import { startEventProxyServer } from '@sentry-internal/test-utils';
2+
3+
startEventProxyServer({
4+
port: 3031,
5+
proxyServerName: 'node-pi-durable',
6+
});

0 commit comments

Comments
 (0)