Skip to content

Commit 1dad424

Browse files
JPeer264claude
andcommitted
feat(server-utils): Add instrumentation for pi-durable
pi-durable has no telemetry hooks, so `piDurableIntegration` wraps the `models` and `registry` passed to `Harness.open()`. Each run becomes its own `invoke_agent` trace with `chat` and `execute_tool` children, and a subagent run nests under the tool call that started it. Messages are mapped to the gen_ai conventions, and the system prompt and tools are read from pi-durable's positional system messages. Tool spans record the result the model receives. Failures pi-durable only reports, and faults of any task, are captured. Failures of the built-in coding tools are not, because they report expected results to the model. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
1 parent 3a5576b commit 1dad424

29 files changed

Lines changed: 2039 additions & 0 deletions

File tree

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,9 @@
1+
import * as Sentry from '@sentry/node';
2+
import { loggingTransport } from '@sentry-internal/node-integration-tests';
3+
4+
Sentry.init({
5+
dsn: 'https://public@dsn.ingest.sentry.io/1337',
6+
release: '1.0',
7+
tracesSampleRate: 1.0,
8+
transport: loggingTransport,
9+
});
Lines changed: 73 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,73 @@
1+
import * as Sentry from '@sentry/node';
2+
import { BACKGROUND_CONTEXT } from '@earendil-works/chord/context';
3+
import { createModels } from '@earendil-works/pi-ai/models';
4+
import { fauxAssistantMessage, fauxProvider } from '@earendil-works/pi-ai/providers/faux';
5+
import {
6+
createRegistry,
7+
defineExtension,
8+
defineTask,
9+
GenerationTask,
10+
Harness,
11+
hook,
12+
MemoryStorage,
13+
} from '@earendil-works/pi-durable';
14+
15+
const context = BACKGROUND_CONTEXT;
16+
17+
// Failures pi-durable does not propagate: a hook that throws (pi-durable reports it and keeps the
18+
// run going) and a durable task whose phase throws (the scheduler faults the task). A compaction
19+
// started while the conversation is idle runs outside any run.
20+
const faux = fauxProvider({ provider: 'faux', models: [{ id: 'faux-model' }] });
21+
const models = createModels();
22+
models.setProvider(faux.provider);
23+
faux.setResponses([
24+
fauxAssistantMessage('First answer with some detail.'),
25+
fauxAssistantMessage('Second answer with more detail.'),
26+
fauxAssistantMessage('Summary of the conversation.'),
27+
]);
28+
29+
const Charge = defineTask({
30+
name: 'app.charge',
31+
version: 1,
32+
initial: () => ({ phase: 'charge' }),
33+
phases: {
34+
charge: async () => {
35+
throw new Error('card declined');
36+
},
37+
},
38+
abort: async (_task, runtime, taskContext) => {
39+
await runtime.commit(() => ({ status: 'terminal', outcome: { status: 'aborted' } }), taskContext);
40+
},
41+
});
42+
43+
const registry = createRegistry();
44+
registry.install(
45+
defineExtension({
46+
name: 'app',
47+
tasks: [Charge],
48+
hooks: [
49+
hook(GenerationTask, {
50+
afterResponse: () => {
51+
throw new Error('afterResponse hook failed');
52+
},
53+
}),
54+
],
55+
}),
56+
);
57+
58+
const harness = await Harness.open(
59+
new MemoryStorage(),
60+
{ models, registry, settings: { compaction: { keepRecentTokens: 1 } } },
61+
context,
62+
);
63+
const root = await harness.root(context, { agent: { model: { provider: 'faux', modelId: 'faux-model' } } });
64+
await (await root.submit({ type: 'input', content: 'First question.' }, context)).wait(context);
65+
await (await root.submit({ type: 'input', content: 'Second question.' }, context)).wait(context);
66+
67+
const charge = await root.commit(tx => tx.createTask(Charge, {}, { ownership: { kind: 'conversation' } }), context);
68+
await harness.waitForTask(charge, context);
69+
70+
await harness.waitForTask(await root.compact(undefined, context), context);
71+
await harness.close(context);
72+
73+
await Sentry.flush(2000);
Lines changed: 99 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,99 @@
1+
import * as Sentry from '@sentry/node';
2+
import { mkdtempSync } from 'node:fs';
3+
import { tmpdir } from 'node:os';
4+
import { join } from 'node:path';
5+
import { BACKGROUND_CONTEXT } from '@earendil-works/chord/context';
6+
import { Type } from '@earendil-works/pi-ai';
7+
import { createModels } from '@earendil-works/pi-ai/models';
8+
import { fauxAssistantMessage, fauxProvider, fauxToolCall } from '@earendil-works/pi-ai/providers/faux';
9+
import {
10+
createRegistry,
11+
defineExtension,
12+
defineTool,
13+
Harness,
14+
hook,
15+
MemoryStorage,
16+
ToolTask,
17+
} from '@earendil-works/pi-durable';
18+
import { NodeExecutionEnv } from '@earendil-works/pi-durable/env/node';
19+
import { CodingTools } from '@earendil-works/pi-durable/tools';
20+
21+
const context = BACKGROUND_CONTEXT;
22+
23+
// One tool round, run one call at a time so the order of the error events is fixed:
24+
// - the built-in `bash` tool runs a command that exits non-zero, which it reports by throwing;
25+
// - `streamer` returns nothing, so its streamed output becomes the result;
26+
// - `secret` returns a value an `afterTool` hook redacts before the model sees it;
27+
// - `fail_now` is an app tool that throws.
28+
const faux = fauxProvider({ provider: 'faux', models: [{ id: 'faux-model' }] });
29+
const models = createModels();
30+
models.setProvider(faux.provider);
31+
faux.setResponses([
32+
fauxAssistantMessage(
33+
[
34+
fauxToolCall('bash', { command: 'ls does-not-exist' }, { id: 'call_bash' }),
35+
fauxToolCall('streamer', {}, { id: 'call_streamer' }),
36+
fauxToolCall('secret', {}, { id: 'call_secret' }),
37+
fauxToolCall('fail_now', {}, { id: 'call_fail' }),
38+
],
39+
{ stopReason: 'toolUse' },
40+
),
41+
fauxAssistantMessage('Done.'),
42+
]);
43+
44+
const registry = createRegistry();
45+
registry.install(CodingTools);
46+
registry.install(
47+
defineExtension({
48+
name: 'app',
49+
tools: [
50+
defineTool({
51+
name: 'streamer',
52+
description: 'Streams its output.',
53+
parameters: Type.Object({}),
54+
execute: async (_args, api) => {
55+
api.output('line 1\n');
56+
api.output('line 2\n');
57+
return {};
58+
},
59+
}),
60+
defineTool({
61+
name: 'secret',
62+
description: 'Returns a secret.',
63+
parameters: Type.Object({}),
64+
execute: async () => ({ content: [{ type: 'text', text: 'original secret value' }] }),
65+
}),
66+
defineTool({
67+
name: 'fail_now',
68+
description: 'Always throws.',
69+
parameters: Type.Object({}),
70+
execute: async () => {
71+
throw new Error('Intentional pi-durable tool failure');
72+
},
73+
}),
74+
],
75+
hooks: [
76+
hook(ToolTask, {
77+
afterTool: (call, result) =>
78+
call.name === 'secret'
79+
? { ...result, content: [{ type: 'text', text: 'redacted by afterTool' }] }
80+
: undefined,
81+
}),
82+
],
83+
}),
84+
);
85+
86+
const cwd = mkdtempSync(join(tmpdir(), 'pi-durable-tools-'));
87+
const harness = await Harness.open(
88+
new MemoryStorage(),
89+
{ models, registry, env: () => new NodeExecutionEnv({ cwd }), settings: { toolExecution: 'sequential' } },
90+
context,
91+
);
92+
const root = await harness.root(context, { agent: { model: { provider: 'faux', modelId: 'faux-model' } } });
93+
await (await root.submit({ type: 'input', content: 'Run the tools.' }, context)).wait(context);
94+
await harness.close(context);
95+
96+
// Everything captured during the run is sent before this sentinel.
97+
await Sentry.flush(2000);
98+
Sentry.captureMessage('pi-durable tools done');
99+
await Sentry.flush(2000);
Lines changed: 65 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,65 @@
1+
import * as Sentry from '@sentry/node';
2+
import { BACKGROUND_CONTEXT } from '@earendil-works/chord/context';
3+
import { Type } from '@earendil-works/pi-ai';
4+
import { createModels } from '@earendil-works/pi-ai/models';
5+
import { fauxAssistantMessage, fauxProvider, fauxToolCall } from '@earendil-works/pi-ai/providers/faux';
6+
import {
7+
createRegistry,
8+
defineExtension,
9+
defineTool,
10+
Harness,
11+
MemoryStorage,
12+
section,
13+
} from '@earendil-works/pi-durable';
14+
15+
const context = BACKGROUND_CONTEXT;
16+
17+
// `pi-ai`'s faux provider scripts model responses in-process. The first run calls two tools and then
18+
// answers, the second run answers directly.
19+
const faux = fauxProvider({ provider: 'faux', models: [{ id: 'faux-model' }] });
20+
const models = createModels();
21+
models.setProvider(faux.provider);
22+
faux.setResponses([
23+
fauxAssistantMessage(
24+
[fauxToolCall('get_weather', { city: 'Berlin' }, { id: 'call_1' }), fauxToolCall('broken', {}, { id: 'call_2' })],
25+
{ stopReason: 'toolUse' },
26+
),
27+
fauxAssistantMessage('It is 21 degrees and sunny in Berlin.'),
28+
fauxAssistantMessage('Goodbye.'),
29+
]);
30+
31+
const registry = createRegistry();
32+
registry.install(
33+
defineExtension({
34+
name: 'weather',
35+
sections: [section('preamble', () => 'You are a weather assistant.', { tag: false })],
36+
tools: [
37+
defineTool({
38+
name: 'get_weather',
39+
description: 'Get the current weather for a city.',
40+
parameters: Type.Object({ city: Type.String() }),
41+
execute: async args => ({ content: [{ type: 'text', text: `It is 21 degrees and sunny in ${args.city}.` }] }),
42+
}),
43+
defineTool({
44+
name: 'broken',
45+
description: 'Always fails.',
46+
parameters: Type.Object({}),
47+
execute: async () => {
48+
throw new Error('broken tool');
49+
},
50+
}),
51+
],
52+
}),
53+
);
54+
55+
// The Harness is opened and driven inside a request span. Its runs must still start traces of their
56+
// own: the scheduler runs them later, and their trace must not depend on which request woke it.
57+
await Sentry.startSpan({ name: 'pi-durable-request', op: 'http.server' }, async () => {
58+
const harness = await Harness.open(new MemoryStorage(), { models, registry }, context);
59+
const root = await harness.root(context, { agent: { model: { provider: 'faux', modelId: 'faux-model' } } });
60+
await (await root.submit({ type: 'input', content: 'What is the weather in Berlin?' }, context)).wait(context);
61+
await (await root.submit({ type: 'input', content: 'Thanks!' }, context)).wait(context);
62+
await harness.close(context);
63+
});
64+
65+
await Sentry.flush(2000);

0 commit comments

Comments
 (0)