https://earendil.com/posts/pi-durable/

Pi 发布 1.0 的同时,推出了 Pi Durable。这篇文章让我们简单快速了解它的原理与正确使用姿势。

什么是 Pi Durable?

Pi 1.0 一般运行在你终端中,如果遇到进程挂了,只能人工介入手动恢复继续。

而 Pi Durable 提供了一种 harness 框架,可以让 agent 跑在任何地方,被多个用户各种客户端同时操作,支持无限的长对话,以及从故障中自动恢复(很持久)。

什么是 Harness?

每个人对 Harness 的定义不同,本文的定义:

  • Harness = 存储 + 运行机制(与模型交互),它提供了供模型调用的 tools & 运行环境。
  • Agent = 模型 + 配置,例如 thinking level,可用的 tools 等

注意 tools 可以运行在任何环境,例如你的笔记本,远程虚拟机,或一个运行在内存的沙箱,不需要 Harness 跑在同一台机器上。

p.s. Pi Durable 和 Pi 一样极简,仅包含15,000 行代码,方便 Agent 直接阅读和理解它。

Long runs anywhere

目标:我们希望 agent 可以在任意地方持久运行。

关于存储,Pi Durable 原生支持内存、SQLite 和 JSONL。接口也非常轻量,方便接入任意实现(注意同一份存储同时只由一个进程管理,其他客户端通过连接这个进程来访问)。

export interface Storage {
    // 原子写入:会话、消息、任务、请求、文档等
    commit(writes: readonly StorageWrite[], context: Context): Promise<Seq>;

    // 分配 ID
    mintId<I extends Id<string>>(): Promise<I>;

    // 读取会话
    conversation(
        id: ConversationId,
        context: Context,
    ): Promise<ConversationRecord | undefined>;

    // 其余接口省略:
    // entry / scanEntries              消息记录
    // task / scanTasks                 任务状态
    // submission / scanSubmissions     提交的请求
    // submissionByRequest              按 requestId 查重
    // document / scanDocuments         应用状态

    close(context: Context): Promise<void>;

以 SQLite 为例,完整记录持久化保存在磁盘中,harness 只把当前所需的数据放在内存中,例如活跃对话、正在运行的任务、待提交的任务。一个对话天然对应一个模型的 context window,值得注意的是 Pi Durable 会在快到窗口上限时,提前总结上下文,然后丝滑继续,所以即使上万条消息也不会有内存问题。

下面是一个最小例子,通过 env 动态决定执行环境:

import { BACKGROUND_CONTEXT } from "@earendil-works/chord/context";
import { createModels } from "@earendil-works/pi-ai/models";
import { openaiProvider } from "@earendil-works/pi-ai/providers/openai";
import { createRegistry, Harness } from "@earendil-works/pi-durable";
import { NodeExecutionEnv } from "@earendil-works/pi-durable/env/node";
import {
    openNodeSqliteStorage,
} from "@earendil-works/pi-durable/storage/sqlite/node";
import { CodingTools } from "@earendil-works/pi-durable/tools";

const context = BACKGROUND_CONTEXT; // every call takes a context for cancellation
const models = createModels();
models.setProvider(openaiProvider());

const registry = createRegistry();
registry.install(CodingTools); // read, write, edit, bash

const env = ({ cwd }: { cwd?: string }) =>
    new NodeExecutionEnv({ cwd: cwd ?? process.cwd() });
const harness = await Harness.open(
    await openNodeSqliteStorage("./agent.sqlite"),
    { models, registry, env },
    context,
);
// The root conversation: created on first use, and the same one after every
// restart.
const root = await harness.root(context, {
    agent: {
        model: { provider: "openai", modelId: "gpt-6.1-sol" },
        cwd: "/work/repo",
    },
});

故障自动恢复

在 Pi Durable 中,一次运行被分为多个步骤,分别对应持久化的 checkpoint,所以在进程中断重启后,可以“断点续传”:

const job = {
    type: "input",
    content: "Fix the flaky login test",
    requestId: "job-42",
} as const;
await root.submit(job, context);
// The process dies here, in the middle of a tool call.

// A new process opens the same storage.
const harness = await Harness.open(
    await openNodeSqliteStorage("./agent.sqlite"),
    { models, registry, env },
    context,
);
harness.resume(); // continue the interrupted run
const root = await harness.root(context);
// the same submission, answered
const settled = await (await root.submit(job, context)).wait(context);

同时多个对话

我们希望 harness 可以并行运行多个对话,相互不阻塞。Pi Durable 支持从任意节点分叉,新对话可以访问分叉点之前的所有历史。

想象在 slack 主频道中,每个子 thread 就是一个小分叉 ✌️ —— 同时也对应一个 agent(自定义模型,thinking level,tools,额外的指令,以及运行环境等等),例如 reviewer 可以用更便宜的模型和只读的 tools。

const channel = await harness.root(context);
const question = await channel.submit(
    { type: "input", content: "@agent why did the deploy fail?" },
    context,
);
const answered = await question.wait(context);

// Someone replies to the agent's answer in a thread. Every conversation names
// its owner, which decides what an abort reaches (more on that under Tasks).
// The thread has none.
const thread = await channel.fork(
    answered.answer!,
    { ownership: { kind: "ownerless" } },
    context,
);

// Both conversations work at the same time.
const inThread = await thread.submit(
    { type: "input", content: "@agent can we roll it back?" },
    context,
);
const inChannel = await channel.submit(
    { type: "input", content: "@agent who is on call today?" },
    context,
);
await Promise.all([inThread.wait(context), inChannel.wait(context)]);

插件

我们希望 agent 的任何能力都可以作为插件,Pi Durable 中的插件包含:

defineExtension({
    name: "project-context",
    sections: [...], // 提示词
    tools: [...],    // 工具
    hooks: [...],    // 执行钩子
    tasks: [...],    // 持久化任务
});

System prompt sections

每次请求发出前,都会动态获取最新的 system prompt(注意对应的变化也记录在历史 session 中,所以即使进程重启后,模型看到上下文的和重启前一模一样):

import { defineExtension, section } from "@earendil-works/pi-durable";

const ProjectContext = defineExtension({
    name: "project-context",
    sections: [
        // Read from the conversation's execution environment. The files can be
        // loaded and watched in the background; every request renders the
        // latest state.
        section("agents_md", (input) => agentsMd.latest(input.env)),
        section("skills", (input) => skills.latest(input.env)),
    ],
});

Tools

从下面的代码可以看到,每个 tools 定义了是否允许在进程崩溃后重新执行。每个 thread 也可以选择各自的 tools:

  • replay: "safe":恢复后可以重试。
  • replay: "unsafe"(默认):恢复后不重试,将中断的信息交给模型处理。
import { Type } from "@earendil-works/pi-ai";
import { defineTool } from "@earendil-works/pi-durable";

const searchIssues = defineTool({
    name: "search_issues",
    description: "Search the issue tracker",
    parameters: Type.Object({ query: Type.String() }),
    replay: "safe", // only reads, so a rerun after a crash is fine
    execute: async (args, api) => {
        // streamed to every client watching
        api.output(`searching for ${args.query}\n`);
        return {
            content: [{ type: "text", text: await tracker.search(args.query) }],
        };
    },
});

const deploy = defineTool({
    name: "deploy",
    description: "Deploy a version to production",
    parameters: Type.Object({ version: Type.String() }),
    // No replay: a deploy interrupted by a crash is reported to the model,
    // never repeated.
    execute: async (args) => ({
        content: [{ type: "text", text: await ci.deploy(args.version) }],
    }),
});

registry.install(defineExtension({ name: "ops", tools: [searchIssues, deploy] }));

// The thread may search, but not deploy.
await thread.configure({ tools: { remove: [deploy] } }, context);

tool 在执行时,可以直接调用 harness 的接口,例如用自定义模型开始新的任务,甚至一个新的 session,甚至可以与其他 session 进行交互。这种灵活的设计,让 subagent 的实现变得非常简单和自然。

p.s. 由 tool 创建的 session 本质上与主 session 没有区别,都在存储中持久化,所以也可以从故障中自动恢复继续。

import type { AssistantMessage } from "@earendil-works/pi-ai";
import { AssistantEntry, configure } from "@earendil-works/pi-durable";

const triage = defineTool({
    name: "triage",
    description: "Label an incoming issue as bug, feature, or question",
    parameters: Type.Object({ issue: Type.String() }),
    // a rerun after a crash finds the same subagent and the same submission
    replay: "safe",
    execute: async (args, api, context) => {
        const child = await api.commit(async (tx) => {
            const existing = (
                await tx.scanConversations({ ownerTaskId: api.taskId }, 1)
            ).items[0];
            if (existing !== undefined) return existing.id;
            // Owned by this call, so aborting the call aborts the subagent.
            const created = await tx.createConversation({
                ownership: { kind: "task", taskId: api.taskId },
            });
            // It starts as a copy of this conversation's agent. Make it a small
            // model without tools.
            await configure(tx, created.id, {
                model: { provider: "openai", modelId: "gpt-6-luna" },
                tools: [],
                instructions: "Answer with one word: bug, feature, or question.",
            });
            return created.id;
        }, context);
        // lets a UI show the subagent under the call
        await api.details({ conversationId: child }, context);
        const subagent = await api.conversation(child, context);
        const request = {
            type: "input",
            content: args.issue,
            requestId: `triage:${api.taskId}`,
        } as const;
        const settled = await (
            await subagent!.submit(request, context)
        ).wait(context);
        // The answer is an entry in the subagent's transcript. Read it and take
        // its text.
        const entry = await api.commit(
            (tx) => tx.entry(AssistantEntry, settled.answer!),
            context,
        );
        const message = entry?.model?.[0] as AssistantMessage;
        const text = message.content
            .flatMap((content) => (content.type === "text" ? [content.text] : []))
            .join("");
        return { content: [{ type: "text", text }] };
    },
});

插件还可以覆盖其他插件的 tools 实现,以 bash 为例:

import { wrapTool } from "@earendil-works/pi-durable";
import { createBashTool } from "@earendil-works/pi-durable/tools";

// Times every bash call, whichever bash the conversation ends up with.
const Timing = defineExtension({
    name: "timing",
    wraps: [
        wrapTool(createBashTool(), (bash) => ({
            ...bash,
            execute: async (args, api, context) => {
                const start = Date.now();
                try {
                    return await bash.execute(args, api, context);
                } finally {
                    metrics.record("bash", Date.now() - start);
                }
            },
        })),
    ],
});

Hooks

例如在每次 deploy tool 执行前,通过 askInSlack 审批:

import { hook, ToolTask } from "@earendil-works/pi-durable";

const Approval = defineExtension({
    name: "approval",
    hooks: [
        hook(ToolTask, {
            beforeTool: async (call, api, context) => {
                if (call.name !== "deploy") return undefined;
                // After a restart, the hook finds the stored answer instead of
                // asking again.
                let approved = await api.memo<boolean>(
                    "approval:deploy",
                    context,
                );
                approved ??= await api.memo(
                    "approval:deploy",
                    await askInSlack(call),
                    context,
                );
                return approved
                    ? undefined
                    : { block: "Nobody approved the deploy." };
            },
        }),
    ],
});

多个 extension 可以对同个 hook 进行注入,最终串行执行。注意 beforeTool 比较特别,可以中断串行执行。

Tasks

harness 本身就是通过 原生的 task 执行对话交互的,例如1)模型请求 2)tool 调用 3)上下文压缩。Extension 同理,值得一提的是 checkpoint 机制。

下面完整例子的调用关系:

- [tool] checkout
    - [task] shop.checkout
        - [phase] pay
            - [task] shop.payment
                - [phase] charge
        - [phase] decide 
完整代码
import { defineTask, type TaskId } from "@earendil-works/pi-durable";

const Payment = defineTask<{ card: string }, { phase: "charge" }, string>({
    name: "shop.payment",
    version: 1,
    initial: () => ({ phase: "charge" }),
    phases: {
        charge: async (task, runtime, context) => {
            // The key makes the charge idempotent: if a crash reruns this
            // phase, the card is only charged once.
            const charge = await bank.charge(
                task.input.card,
                `payment-${task.id}`,
            );
            await runtime.commit(
                () => ({
                    status: "terminal",
                    outcome: charge.ok
                        ? { status: "completed", result: charge.receipt }
                        : { status: "failed", error: { message: charge.error } },
                }),
                context,
            );
        },
    },
    // Another payment failed, or the checkout was cancelled: undo this one.
    abort: async (task, runtime, context) => {
        await bank.refund(`payment-${task.id}`);
        await runtime.commit(
            () => ({ status: "terminal", outcome: { status: "aborted" } }),
            context,
        );
    },
});

type CheckoutState =
    | { phase: "pay" }
    | { phase: "decide"; payments: TaskId<string>[] };
const Checkout = defineTask<{ cards: string[] }, CheckoutState, string>({
    name: "shop.checkout",
    version: 1,
    initial: () => ({ phase: "pay" }),
    phases: {
        pay: async (task, runtime, context) => {
            await runtime.commit(async (tx) => {
                const payments: TaskId<string>[] = [];
                for (const card of task.input.cards) {
                    payments.push(
                        await tx.createTask(Payment, { card }, {
                            ownership: { kind: "task", taskId: task.id },
                        }),
                    );
                }
                // Run no code until every payment is done. The first failed
                // payment aborts the others.
                return {
                    status: "waiting",
                    checkpoint: { phase: "decide", payments },
                    on: payments,
                    policy: "failFast",
                };
            }, context);
        },
        decide: async (task, runtime, context) => {
            const outcomes = await runtime.outcomes(
                task.state.checkpoint.payments,
                context,
            );
            const paid = outcomes.every(
                (outcome) => outcome.status === "completed",
            );
            await runtime.commit(
                () => ({
                    status: "terminal",
                    outcome: paid
                        ? { status: "completed", result: "Order placed." }
                        : {
                            status: "failed",
                            error: { message: "A payment failed." },
                        },
                }),
                context,
            );
        },
    },
    abort: (_task, runtime, context) =>
        runtime.commit(
            () => ({ status: "terminal", outcome: { status: "aborted" } }),
            context,
        ),
});

// The agent starts a checkout with a tool.
const checkout = defineTool({
    name: "checkout",
    description: "Pay for the cart, split across several cards",
    parameters: Type.Object({ cards: Type.Array(Type.String()) }),
    execute: async (args, api, context) => {
        // Owned by this call: aborting the call aborts the checkout and refunds
        // its payments.
        const owner = {
            ownership: { kind: "task", taskId: api.taskId },
        } as const;
        const id = await api.createTask(
            Checkout,
            { cards: args.cards },
            owner,
            context,
        );
        const { outcome } = (await api.waitForTask(id, context)).state;
        const text =
            outcome.status === "completed" ? outcome.result : outcome.status;
        return { content: [{ type: "text", text }] };
    },
});

registry.install(defineExtension({
    name: "shop",
    tools: [checkout],
    tasks: [Payment, Checkout],
}));

假如进程在 charge 阶段崩溃(task 尚未提交完成状态),恢复后会根据当前的 checkpoint(i.e. phase)重新执行 charge,并通过 taskid 做幂等:

   phases: {
       charge: async (task, runtime, context) => {
           // 根据 checkpoint 恢复
           const charge = await bank.charge(
               task.input.card,
               `payment-${task.id}`,
           );

           // 持久化提交任务的终态
           await runtime.commit(
               () => ({
                   status: "terminal",
               // ...

所有的 taks 形成一个 ownership 树:假如用户取消 checkout,各个 payment 子节点会执行对应的 abort 定义,只有当子节点清理完成,任务结束后,父节点才结束。

const owner = {
    ownership: { kind: "task", taskId: api.taskId },
} as const;

task 默认在前台执行,可通过 { background: true } 参数开启后台任务(不会影响对话的空闲状态)。

上下文压缩

上面提到过 compact 在后台提前运行(本身也是一个 task),而不打断对话。当然你也可以随时手动调用触发。

const harness = await Harness.open(storage, {
    models,
    registry,
    settings: {
        compaction: {
            // past contextWindow - reserveTokens, the next request waits for a
            // summary
            reserveTokens: 16384,
            // this far before that, a summary starts in the background
            backgroundTokens: 32768,
        },
    },
}, context);

// Manual, also while the agent is working.
await root.compact("Keep the names of the failing tests", context);

应用状态

应用状态期望也和对话一样持久化,例如 todo list,plan,sandbox 等等。这些信息保存在 typed JSON 中(称为 Documents):

import { defineDoc } from "@earendil-works/pi-durable";

const Todos = defineDoc<{ items: string[] }>({
    kind: "app.todos",
    version: 1,
    scope: "conversation",
    history: "rewindable",
    fork: "asOf", // a fork starts with the todos its parent had at the fork entry
    initial: () => ({ items: [] }),
});

const Todo = defineExtension({
    name: "todo",
    tools: [
        defineTool({
            name: "todo",
            description: "Add an item to your todo list",
            parameters: Type.Object({ item: Type.String() }),
            execute: async (args, api, context) => {
                await api.commit(async (tx) => {
                    const todos = await tx.doc(Todos, api.conversationId);
                    todos.items.push(args.item);
                }, context);
                const text = `Added ${args.item}`;
                return { content: [{ type: "text", text }] };
            },
        }),
    ],
    // The model sees the list before every request.
    sections: [
        section("todos", async (input, context) => {
            const todos = await input.read.snapshot(
                Todos,
                input.conversationId,
                context,
            );
            return todos?.items.join("\n") || undefined;
        }),
    ],
});

// A UI subscribes to the committed value.
const todos = await harness.documentState(Todos, channel.id, context);
todos?.subscribe((value) => renderTodos(value?.items ?? []));

动态可插拔

在 agent 运行中,支持动态热加载或替换插件(在下一次 tool 调用时生效):

// The extension's file changed on disk.
// same name "ops": replaces the installed one
registry.install(await loadExtension("./ops.ts"));

多人互动

有点类似 tmux,每个客户端可以 attach 到同一个对话,获取最新的状态(committed state)并持续同步更新(通过 thread.watch() 获取每一条新的 commit)。

// A second client joins the thread while the agent is working.
const view = await thread.viewState(context);
render(view.value);
view.subscribe((value) => render(value));

// And steers it. The message joins the running work after the current tool
// calls.
await thread.submit(
    { type: "input", content: "Check the staging logs first", whenBusy: "steer" },
    context,
);

Try it

心动不如行动,将 agent 指向 Pi Durable 的 README,开始构建你的 harness 应用吧!

碎碎念

原本希望通过让 agent 学习我的写作风格,使用「说人话」skill,自动生成这篇博客。第一印象 效果颇为不错,结构合理,文字也没有太多的 AI slop 味道。

然而人成为了整个流程的瓶颈,因为我只能从总结压缩后的内容摄取 10% 的知识。故还是选择人工花费数小时,慢慢翻译文章,通过输出消化知识。就像看一本书,作者可能就一句话的观点,却通过 300 页的例子和“废话”反复让读者理解接受。

所以即使随着 AI 模型的快速进化,个人并不觉得工作会被 AI 迅速取代。人类自身的局限性,导致了极低的理解能力,沟通效率等,最终产生了无数就业岗位。