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
6 changes: 5 additions & 1 deletion examples/next/channels/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,9 @@
"deploy": "vite build && wrangler deploy",
"types": "wrangler types env.d.ts --include-runtime false",
"tui": "node src/ai-sdk-tui.ts",
"tui2": "sh -c 'tsx ../../../packages/agents/src/cli/index.ts tui \"${AGENT_ORIGIN:-http://localhost:5173}/channels/${1:-ai-sdk}/${2:-default}\" --as \"terminal-$$\"' tui2"
"tui2": "sh -c 'tsx ../../../packages/agents/src/cli/index.ts tui \"${AGENT_ORIGIN:-http://localhost:5173}/channels/${1:-ai-sdk}/${2:-default}\" --as \"terminal-$$\"' tui2",
"typecheck": "tsc --noEmit",
"test": "vitest run --config src/tests/vitest.config.ts"
},
"dependencies": {
"@cloudflare/kumo": "^2.6.0",
Expand All @@ -28,6 +30,7 @@
"devDependencies": {
"@ai-sdk/tui": "1.0.119",
"@cloudflare/vite-plugin": "1.62.3",
"@cloudflare/vitest-pool-workers": "^0.19.1",
"@cloudflare/workers-types": "^5.20260729.1",
"@tailwindcss/vite": "^4",
"@types/node": "^26.0.1",
Expand All @@ -37,6 +40,7 @@
"tailwindcss": "^4.3.2",
"typescript": "^6.0.3",
"vite": "^8.1.0",
"vitest": "4.1.11",
"wrangler": "^4.145.0"
}
}
55 changes: 43 additions & 12 deletions examples/next/channels/src/pi/channels-harness.ts
Original file line number Diff line number Diff line change
Expand Up @@ -168,9 +168,10 @@ class PiChannelsSession implements HarnessSession {
// pi's tools run on the server, and it has no approvals.
if (!("parts" in input)) throw new Error("pi takes no tool answers");
// pi has no participants, so `from` is dropped.
const content = toUserInput(input);
const operationId = options.operationId ?? crypto.randomUUID();
this.ids.record(this.id, operationId, input.messageId ?? operationId);
return this.session.submit(toUserInput(input), {
return this.session.submit(content, {
operationId,
...(options.whenBusy && { whenBusy: options.whenBusy })
});
Expand Down Expand Up @@ -231,9 +232,28 @@ class PiChannelsSession implements HarnessSession {

/** The active transcript, with user messages under the caller's ids. */
async transcript(): Promise<TranscriptMessage[]> {
const messages = projectEntries(await this.session.messages());
const users = messages.filter((m) => m.role === "user").map((m) => m.id);
const ids = await this.ids.resolve(this.id, users);
const entries = await this.session.messages();
const messages = projectEntries(entries);
const users = new Set(
messages.filter((m) => m.role === "user").map((m) => m.id)
);
// A fork inherits its parent's entries, which keep the parent's
// conversation id: resolve each against the session that placed it.
const bySession = new Map<string, string[]>();
for (const entry of entries) {
if (!users.has(String(entry.id))) continue;
const session = String(entry.conversationId);
bySession.set(session, [
...(bySession.get(session) ?? []),
String(entry.id)
]);
}
const ids = new Map<string, string>();
for (const [session, entryIds] of bySession) {
for (const [entry, id] of await this.ids.resolve(session, entryIds)) {
ids.set(entry, id);
}
}
return toTranscript(messages, (entryId) => ids.get(entryId));
}

Expand Down Expand Up @@ -500,16 +520,27 @@ class RunChunks {
}
}

/**
* The caller's parts as pi user content. pi takes text and base64 images;
* an inline text file becomes text. Anything else throws, so the submit
* fails rather than pi answering an incomplete prompt.
*/
function toUserInput(input: HarnessInput): UserInput {
const parts = input.parts.flatMap(
(part): Exclude<UserInput, string>[number][] => {
if (part.type === "text") return [{ type: "text", text: part.text }];
const data = /^data:([^;,]+);base64,(.*)$/.exec(part.url);
return data && part.mediaType.startsWith("image/")
? [{ type: "image", mimeType: data[1], data: data[2] }]
: [];
const parts = input.parts.map((part): Exclude<UserInput, string>[number] => {
if (part.type === "text") return { type: "text", text: part.text };
const data = /^data:([^;,]+)(?:;[^;,]*)*;base64,(.*)$/.exec(part.url);
if (data && part.mediaType.startsWith("image/")) {
return { type: "image", mimeType: part.mediaType, data: data[2] };
}
if (data && part.mediaType.startsWith("text/")) {
const bytes = Uint8Array.from(atob(data[2]), (c) => c.charCodeAt(0));
return { type: "text", text: new TextDecoder().decode(bytes) };
}
);
const name = part.filename ? ` (${part.filename})` : "";
throw new Error(
`pi cannot take a ${part.mediaType} attachment${name}${data ? "" : " by URL"}; it takes text, inline text files, and inline images`
);
});
return parts.length === 1 && parts[0].type === "text" ? parts[0].text : parts;
}

Expand Down
86 changes: 86 additions & 0 deletions examples/next/channels/src/tests/channels-harness.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,86 @@
import { env } from "cloudflare:workers";
import { evictDurableObject } from "cloudflare:test";
import { describe, expect, it } from "vitest";

const stub = () => env.PI_CHANNELS_TEST.getByName(crypto.randomUUID());
const PNG = "iVBORw0KGgo=";

describe("piChannelsHarness input", () => {
it("passes text and inline images to pi", async () => {
const result = await stub().send(
[
{ type: "text", text: "look" },
{
type: "file",
mediaType: "image/png",
url: `data:image/png;base64,${PNG}`
}
],
"m1"
);
expect(JSON.parse(result)).toMatchObject({
status: "done",
text: "echo: look [image]"
});
});

it("passes inline text files to pi as text", async () => {
const result = await stub().send(
[
{
type: "file",
mediaType: "text/plain",
url: `data:text/plain;base64,${btoa("notes")}`
}
],
"m1"
);
expect(JSON.parse(result)).toMatchObject({ text: "echo: notes" });
});

it("rejects attachments pi cannot take instead of dropping them", async () => {
const pdf = await stub().send(
[
{ type: "text", text: "summarize" },
{
type: "file",
mediaType: "application/pdf",
url: "data:application/pdf;base64,JVBERi0="
}
],
"m1"
);
expect(pdf).toMatch(/^rejected: .*application\/pdf/);
const remote = await stub().send(
[
{
type: "file",
mediaType: "image/png",
url: "https://example.com/a.png"
}
],
"m1"
);
expect(remote).toMatch(/^rejected: /);
});
});

describe("piChannelsHarness forks", () => {
it("keeps caller message ids on inherited user messages", async () => {
const agent = stub();
const root = await agent.rootId();
await agent.send([{ type: "text", text: "one" }], "first", root);
const fork = await agent.fork(root);
await agent.send([{ type: "text", text: "two" }], "second", fork);
expect(await agent.userIds(fork)).toEqual(["first", "second"]);
});

it("keeps them after eviction", async () => {
const agent = stub();
const root = await agent.rootId();
await agent.send([{ type: "text", text: "one" }], "first", root);
const fork = await agent.fork(root);
await evictDurableObject(agent);
expect(await agent.userIds(fork)).toEqual(["first"]);
});
});
11 changes: 11 additions & 0 deletions examples/next/channels/src/tests/cloudflare-test.d.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
import type { PiChannelsTestObject } from "./worker";

declare global {
namespace Cloudflare {
interface Env {
PI_CHANNELS_TEST: DurableObjectNamespace<PiChannelsTestObject>;
}
}
}

export {};
28 changes: 28 additions & 0 deletions examples/next/channels/src/tests/vitest.config.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
import path from "node:path";
import { cloudflareTest } from "@cloudflare/vitest-pool-workers";
import { defineConfig } from "vitest/config";

const testsDir = import.meta.dirname;

export default defineConfig({
plugins: [
cloudflareTest({
wrangler: { configPath: path.join(testsDir, "wrangler.jsonc") }
})
],
resolve: {
// The adapter imports from ../harnesses/pi; resolve pi from one place.
dedupe: [
"@earendil-works/chord",
"@earendil-works/pi-ai",
"@earendil-works/pi-durable",
"@earendil-works/pi-telemetry"
]
},
test: {
name: "next-channels",
include: [path.join(testsDir, "**/*.test.ts")],
testTimeout: 30_000,
hookTimeout: 30_000
}
});
94 changes: 94 additions & 0 deletions examples/next/channels/src/tests/worker.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,94 @@
import {
fauxAssistantMessage,
fauxProvider,
fauxText,
type Message,
type TranscriptContext
} from "@earendil-works/pi-ai";
import { createModels } from "@earendil-works/pi-ai/models";
import { createRegistry, Harness } from "@earendil-works/pi-durable";
import { DurableObject } from "cloudflare:workers";
import type { InputPart } from "agents/experimental/channels";
import { PiHarness } from "agents/harness/pi";
import { Lifecycle } from "agents/lifecycle";
import { piChannelsHarness } from "../pi/channels-harness";

function textOf(content: Message["content"] | undefined): string {
if (content === undefined) return "";
if (typeof content === "string") return content;
return content
.map((part) => (part.type === "text" ? part.text : `[${part.type}]`))
.join(" ");
}

/** Echoes the last user message, so tests can see what pi was given. */
function script(context: TranscriptContext) {
const last = context.messages.filter((m) => m.role === "user").at(-1);
return fauxAssistantMessage([fauxText(`echo: ${textOf(last?.content)}`)]);
}

/** `piChannelsHarness` over a real `PiHarness` with pi-ai's faux provider. */
export class PiChannelsTestObject extends DurableObject<Env> {
readonly #faux = fauxProvider();
readonly pi = new PiHarness({
harness: async ({ storage, context }) => {
const models = createModels();
models.setProvider(this.#faux.provider);
return Harness.open(
storage,
{
models,
registry: createRegistry(),
settings: {
retry: { enabled: false, maxRetries: 0, baseDelayMs: 0 }
},
onReport: (error) => console.warn("pi report", error)
},
context
);
},
defaults: { model: this.#faux.getModel() }
});
readonly harness = piChannelsHarness(this.pi, { kv: this.ctx.storage.kv });
readonly lifecycle = Lifecycle.install(this).use(this.pi);

constructor(ctx: DurableObjectState, env: Env) {
super(ctx, env);
this.#faux.setResponses(Array.from({ length: 200 }, () => script));
}

/** Submit and wait; the result as JSON, or the submit error's message. */
async send(
parts: InputPart[],
messageId: string,
session?: string
): Promise<string> {
const target = this.harness.session(session);
let receipt;
try {
receipt = await target.submit({ parts, messageId });
} catch (error) {
return `rejected: ${error instanceof Error ? error.message : String(error)}`;
}
return JSON.stringify(await target.wait(receipt.operationId));
}

async fork(from: string): Promise<string> {
return (await this.harness.sessions.fork(from)).id;
}

rootId(): string {
return this.harness.session().id;
}

/** User message ids in a session's transcript, as a client sees them. */
async userIds(session?: string): Promise<string[]> {
const watch = await this.harness.session(session).watch();
await watch.stop();
return watch.state.messages
.filter((m) => m.role === "user")
.map((m) => m.id);
}
}

export default { fetch: () => new Response("Not found", { status: 404 }) };
14 changes: 14 additions & 0 deletions examples/next/channels/src/tests/wrangler.jsonc
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
{
"$schema": "../../node_modules/wrangler/config-schema.json",
"main": "worker.ts",
"compatibility_date": "2026-06-11",
"compatibility_flags": ["nodejs_compat"],
"durable_objects": {
"bindings": [
{ "name": "PI_CHANNELS_TEST", "class_name": "PiChannelsTestObject" }
]
},
"migrations": [
{ "tag": "v1", "new_sqlite_classes": ["PiChannelsTestObject"] }
]
}
10 changes: 9 additions & 1 deletion examples/next/channels/tsconfig.json
Original file line number Diff line number Diff line change
@@ -1,3 +1,11 @@
{
"extends": "agents/tsconfig"
"extends": "agents/tsconfig",
"compilerOptions": {
"types": [
"node",
"@cloudflare/workers-types",
"@cloudflare/vitest-pool-workers/types",
"vite/client"
]
}
}
21 changes: 21 additions & 0 deletions pnpm-lock.yaml

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading