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
Expand Up @@ -123,6 +123,7 @@ export const NO_AUTO_INSTRUMENTATION = [
'suites/tracing/openai/v6/test.ts',
'suites/tracing/openai/v7/test.ts',
'suites/tracing/orchestrion-lazy-registration/test.ts',
'suites/tracing/pi-durable/test.ts',
'suites/tracing/postgres-streamed/test.ts',
'suites/tracing/postgres/test.ts',
'suites/tracing/postgresjs-streamed/test.ts',
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
import * as Sentry from '@sentry/node';
import { loggingTransport } from '@sentry-internal/node-integration-tests';

Sentry.init({
dsn: 'https://public@dsn.ingest.sentry.io/1337',
release: '1.0',
tracesSampleRate: 1.0,
transport: loggingTransport,
});
Original file line number Diff line number Diff line change
@@ -0,0 +1,70 @@
import { BACKGROUND_CONTEXT } from '@earendil-works/chord/context';
import { createModels } from '@earendil-works/pi-ai/models';
import { fauxAssistantMessage, fauxProvider } from '@earendil-works/pi-ai/providers/faux';
import {
createRegistry,
defineExtension,
defineTask,
GenerationTask,
Harness,
hook,
MemoryStorage,
} from '@earendil-works/pi-durable';

const context = BACKGROUND_CONTEXT;

// Failures pi-durable does not propagate: a hook that throws (pi-durable reports it and keeps the
// run going) and a durable task whose phase throws (the scheduler faults the task). A compaction
// started while the conversation is idle runs outside any run.
const faux = fauxProvider({ provider: 'faux', models: [{ id: 'faux-model' }] });
const models = createModels();
models.setProvider(faux.provider);
faux.setResponses([
fauxAssistantMessage('First answer with some detail.'),
fauxAssistantMessage('Second answer with more detail.'),
fauxAssistantMessage('Summary of the conversation.'),
]);

const Charge = defineTask({
name: 'app.charge',
version: 1,
initial: () => ({ phase: 'charge' }),
phases: {
charge: async () => {
throw new Error('card declined');
},
},
abort: async (_task, runtime, taskContext) => {
await runtime.commit(() => ({ status: 'terminal', outcome: { status: 'aborted' } }), taskContext);
},
});

const registry = createRegistry();
registry.install(
defineExtension({
name: 'app',
tasks: [Charge],
hooks: [
hook(GenerationTask, {
afterResponse: () => {
throw new Error('afterResponse hook failed');
},
}),
],
}),
);

const harness = await Harness.open(
new MemoryStorage(),
{ models, registry, settings: { compaction: { keepRecentTokens: 1 } } },
context,
);
const root = await harness.root(context, { agent: { model: { provider: 'faux', modelId: 'faux-model' } } });
await (await root.submit({ type: 'input', content: 'First question.' }, context)).wait(context);
await (await root.submit({ type: 'input', content: 'Second question.' }, context)).wait(context);

const charge = await root.commit(tx => tx.createTask(Charge, {}, { ownership: { kind: 'conversation' } }), context);
await harness.waitForTask(charge, context);

await harness.waitForTask(await root.compact(undefined, context), context);
await harness.close(context);
Original file line number Diff line number Diff line change
@@ -0,0 +1,98 @@
import * as Sentry from '@sentry/node';
import { mkdtempSync } from 'node:fs';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
import { BACKGROUND_CONTEXT } from '@earendil-works/chord/context';
import { Type } from '@earendil-works/pi-ai';
import { createModels } from '@earendil-works/pi-ai/models';
import { fauxAssistantMessage, fauxProvider, fauxToolCall } from '@earendil-works/pi-ai/providers/faux';
import {
createRegistry,
defineExtension,
defineTool,
Harness,
hook,
MemoryStorage,
ToolTask,
} from '@earendil-works/pi-durable';
import { NodeExecutionEnv } from '@earendil-works/pi-durable/env/node';
import { CodingTools } from '@earendil-works/pi-durable/tools';

const context = BACKGROUND_CONTEXT;

// One tool round, run one call at a time so the order of the error events is fixed:
// - the built-in `bash` tool runs a command that exits non-zero, which it reports by throwing;
// - `streamer` returns nothing, so its streamed output becomes the result;
// - `secret` returns a value an `afterTool` hook redacts before the model sees it;
// - `fail_now` is an app tool that throws.
const faux = fauxProvider({ provider: 'faux', models: [{ id: 'faux-model' }] });
const models = createModels();
models.setProvider(faux.provider);
faux.setResponses([
fauxAssistantMessage(
[
fauxToolCall('bash', { command: 'ls does-not-exist' }, { id: 'call_bash' }),
fauxToolCall('streamer', {}, { id: 'call_streamer' }),
fauxToolCall('secret', {}, { id: 'call_secret' }),
fauxToolCall('fail_now', {}, { id: 'call_fail' }),
],
{ stopReason: 'toolUse' },
),
fauxAssistantMessage('Done.'),
]);

const registry = createRegistry();
registry.install(CodingTools);
registry.install(
defineExtension({
name: 'app',
tools: [
defineTool({
name: 'streamer',
description: 'Streams its output.',
parameters: Type.Object({}),
execute: async (_args, api) => {
api.output('line 1\n');
api.output('line 2\n');
return {};
},
}),
defineTool({
name: 'secret',
description: 'Returns a secret.',
parameters: Type.Object({}),
execute: async () => ({ content: [{ type: 'text', text: 'original secret value' }] }),
}),
defineTool({
name: 'fail_now',
description: 'Always throws.',
parameters: Type.Object({}),
execute: async () => {
throw new Error('Intentional pi-durable tool failure');
},
}),
],
hooks: [
hook(ToolTask, {
afterTool: (call, result) =>
call.name === 'secret'
? { ...result, content: [{ type: 'text', text: 'redacted by afterTool' }] }
: undefined,
}),
],
}),
);

const cwd = mkdtempSync(join(tmpdir(), 'pi-durable-tools-'));
const harness = await Harness.open(
new MemoryStorage(),
{ models, registry, env: () => new NodeExecutionEnv({ cwd }), settings: { toolExecution: 'sequential' } },
context,
);
const root = await harness.root(context, { agent: { model: { provider: 'faux', modelId: 'faux-model' } } });
await (await root.submit({ type: 'input', content: 'Run the tools.' }, context)).wait(context);
await harness.close(context);

// Everything captured during the run is sent before this sentinel.
await Sentry.flush(2000);
Sentry.captureMessage('pi-durable tools done');
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
import * as Sentry from '@sentry/node';
import { BACKGROUND_CONTEXT } from '@earendil-works/chord/context';
import { Type } from '@earendil-works/pi-ai';
import { createModels } from '@earendil-works/pi-ai/models';
import { fauxAssistantMessage, fauxProvider, fauxToolCall } from '@earendil-works/pi-ai/providers/faux';
import {
createRegistry,
defineExtension,
defineTool,
Harness,
MemoryStorage,
section,
} from '@earendil-works/pi-durable';

const context = BACKGROUND_CONTEXT;

// `pi-ai`'s faux provider scripts model responses in-process. The first run calls two tools and then
// answers, the second run answers directly.
const faux = fauxProvider({ provider: 'faux', models: [{ id: 'faux-model' }] });
const models = createModels();
models.setProvider(faux.provider);
faux.setResponses([
fauxAssistantMessage(
[fauxToolCall('get_weather', { city: 'Berlin' }, { id: 'call_1' }), fauxToolCall('broken', {}, { id: 'call_2' })],
{ stopReason: 'toolUse' },
),
fauxAssistantMessage('It is 21 degrees and sunny in Berlin.'),
fauxAssistantMessage('Goodbye.'),
]);

const registry = createRegistry();
registry.install(
defineExtension({
name: 'weather',
sections: [section('preamble', () => 'You are a weather assistant.', { tag: false })],
tools: [
defineTool({
name: 'get_weather',
description: 'Get the current weather for a city.',
parameters: Type.Object({ city: Type.String() }),
execute: async args => ({ content: [{ type: 'text', text: `It is 21 degrees and sunny in ${args.city}.` }] }),
}),
defineTool({
name: 'broken',
description: 'Always fails.',
parameters: Type.Object({}),
execute: async () => {
throw new Error('broken tool');
},
}),
],
}),
);

// The Harness is opened and driven inside a request span. Its runs must still start traces of their
// own: the scheduler runs them later, and their trace must not depend on which request woke it.
await Sentry.startSpan({ name: 'pi-durable-request', op: 'http.server' }, async () => {
const harness = await Harness.open(new MemoryStorage(), { models, registry }, context);
const root = await harness.root(context, { agent: { model: { provider: 'faux', modelId: 'faux-model' } } });
await (await root.submit({ type: 'input', content: 'What is the weather in Berlin?' }, context)).wait(context);
await (await root.submit({ type: 'input', content: 'Thanks!' }, context)).wait(context);
await harness.close(context);
});
Loading
Loading