文章
合集Pi Agent 源码阅读第 9 / 9 篇

rpc-mode:RPC 模式

概述

rpc-mode.ts 有 817 行,是三个模式中最复杂的。它实现了一个无头(headless)运行模式:通过 JSON 协议在 stdin/stdout 上通信,用于嵌入其他应用程序。

典型使用场景:

  • VS Code 插件调用 Agent 的代码编辑能力
  • 其他 CLI 工具通过管道发送 JSON-RPC 请求
  • 自动化脚本调用 Agent 的文件操作能力

触发条件:

pi --mode rpc

协议格式

输入(stdin):JSON 命令

每行一个 JSON 对象,包含 type 字段和可选的 id 用于关联响应:

{"type":"prompt","message":"hello","id":"cmd-1"}
{"type":"set_model","provider":"anthropic","modelId":"claude-3-5-sonnet"}
{"type":"get_state"}

输出(stdout):JSON 响应 + 事件流

{"type":"response","command":"prompt","success":true,"id":"cmd-1"}
{"type":"message","message":{"role":"user","content":"hello"}}
{"type":"model_change","provider":"anthropic","modelId":"claude-3-5-sonnet"}
{"type":"message","message":{"role":"assistant","content":"你好!..."}}

核心函数 runRpcMode()

export async function runRpcMode(runtimeHost: AgentSessionRuntime): Promise<never> {

返回类型是 Promise<never>,这个函数永远不会返回,它会一直运行直到进程被杀死。

初始化

takeOverStdout();
let session = runtimeHost.session;
let unsubscribe: (() => void) | undefined;
let unsubscribeBackpressure: (() => void) | undefined;

接管 stdout,初始化状态变量。

输出辅助函数

const output = (obj: RpcResponse | RpcExtensionUIRequest | object) => {
    writeRawStdout(serializeJsonLine(obj));
};

const success = <T extends RpcCommand["type"]>(
    id: string | undefined,
    command: T,
    data?: object | null,
): RpcResponse => {
    if (data === undefined) {
        return { id, type: "response", command, success: true } as RpcResponse;
    }
    return { id, type: "response", command, success: true, data } as RpcResponse;
};

const error = (id: string | undefined, command: string, message: string): RpcResponse => {
    return { id, type: "response", command, success: false, error: message };
};

三个辅助函数:

  • output(),把对象序列化成 JSON 行输出
  • success(),构建成功响应
  • error(),构建错误响应

命令处理 handleCommand()

这是 RPC 模式的核心,处理所有输入命令:

const handleCommand = async (command: RpcCommand): Promise<RpcResponse | undefined> => {
    const id = command.id;

    switch (command.type) {
        case "prompt": {...}
        case "steer": {...}
        case "follow_up": {...}
        case "abort": {...}
        case "new_session": {...}
        case "get_state": {...}
        case "set_model": {...}
        case "cycle_model": {...}
        case "get_available_models": {...}
        case "set_thinking_level": {...}
        case "cycle_thinking_level": {...}
        case "compact": {...}
        case "set_auto_compaction": {...}
        case "bash": {...}
        case "abort_bash": {...}
        case "switch_session": {...}
        case "fork": {...}
        case "clone": {...}
        case "get_entries": {...}
        case "get_tree": {...}
        case "get_messages": {...}
        case "get_commands": {...}
        default: {
            return error(id, unknownCommand.type, `Unknown command: ${unknownCommand.type}`);
        }
    }
};

命令分类

分类命令作用
提示prompt、steer、follow_up、abort发送消息给 AI
状态get_state获取当前会话状态
模型set_model、cycle_model、get_available_models管理 AI 模型
思考set_thinking_level、cycle_thinking_level管理思考级别
压缩compact、set_auto_compaction上下文压缩
Bashbash、abort_bash执行 bash 命令
会话switch_session、fork、clone、get_entries、get_tree会话管理
消息get_messages获取消息历史
命令get_commands获取可用命令列表

prompt 命令的特殊处理

case "prompt": {
    let preflightSucceeded = false;
    void session
        .prompt(command.message, {
            images: command.images,
            streamingBehavior: command.streamingBehavior,
            source: "rpc",
            preflightResult: (didSucceed) => {
                if (didSucceeded) {
                    preflightSucceeded = true;
                    output(success(id, "prompt"));
                }
            },
        })
        .catch((e) => {
            if (!preflightSucceeded) {
                output(error(id, "prompt", e.message));
            }
        });
    return undefined;
}

prompt 命令是异步的:

  1. 立即返回 undefined(不发送响应)
  2. 在后台执行 session.prompt()
  3. 预检成功后发送成功响应
  4. 如果失败且预检未成功,发送错误响应

输入处理

const handleInputLine = async (line: string) => {
    let parsed: unknown;
    try {
        parsed = JSON.parse(line);
    } catch (parseError: unknown) {
        output(error(undefined, "parse", `Failed to parse command: ...`));
        await waitForRawStdoutBackpressure();
        return;
    }

    // 处理扩展 UI 响应
    if (typeof parsed === "object" && parsed !== null && "type" in parsed && parsed.type === "extension_ui_response") {
        const response = parsed as RpcExtensionUIResponse;
        const pending = pendingExtensionRequests.get(response.id);
        if (pending) {
            pendingExtensionRequests.delete(response.id);
            pending.resolve(response);
        }
        return;
    }

    const command = parsed as RpcCommand;
    try {
        const response = await handleCommand(command);
        if (response) {
            output(response);
            await waitForRawStdoutBackpressure();
        }
        await checkShutdownRequested();
    } catch (commandError: unknown) {
        output(error(command.id, command.type, ...));
        await waitForRawStdoutBackpressure();
    }
};

处理流程:

读取一行 JSON
  │
  ├─ 解析失败? → 输出错误响应
  │
  ├─ 是扩展 UI 响应? → 匹配 pending 请求,resolve Promise
  │
  └─ 是命令? → handleCommand() → 输出响应

扩展 UI 上下文

RPC 模式下,扩展可以通过 JSON 协议请求 UI 操作:

const createExtensionUIContext = (): ExtensionUIContext => ({
    select: (title, options, opts) =>
        createDialogPromise(opts, undefined, { method: "select", title, options }, (r) => ...),
    confirm: (title, message, opts) =>
        createDialogPromise(opts, false, { method: "confirm", title, message }, (r) => ...),
    input: (title, placeholder, opts) =>
        createDialogPromise(opts, undefined, { method: "input", title, placeholder }, (r) => ...),
    notify(message, type) {...},
    setStatus(key, text) {...},
    setTitle(title) {...},
    // ... 更多方法
});

dialog 的工作流程

扩展调用 ui.select("选择模型", [...])
  │
  ├─ 生成唯一 ID
  ├─ 创建 Promise,存入 pendingExtensionRequests
  ├─ 输出 extension_ui_request 到 stdout
  │    └─ {"type":"extension_ui_request","id":"xxx","method":"select","title":"选择模型","options":[...]}
  │
  └─ 等待客户端响应...

客户端收到请求,显示选择框,用户选择后:
  │
  ├─ 输出 extension_ui_response 到 stdin
  │    └─ {"type":"extension_ui_response","id":"xxx","value":"gpt-4o"}
  │
  └─ handleInputLine() 匹配 pending 请求,resolve Promise

扩展的 select() 调用返回 "gpt-4o"

信号处理与关闭

async function shutdown(exitCode = 0, signal?: NodeJS.Signals): Promise<never> {
    if (shuttingDown) {
        process.exit(exitCode);
    }
    shuttingDown = true;
    for (const cleanup of signalCleanupHandlers) {
        cleanup();
    }
    unsubscribe?.();
    unsubscribeBackpressure?.();
    await runtimeHost.dispose();
    detachInput();
    process.stdin.pause();
    if (signal !== "SIGTERM") {
        await flushRawStdout();
    }
    process.exit(exitCode);
}

关闭流程:

  1. 防止重复关闭(shuttingDown 标志)
  2. 移除信号处理器
  3. 取消事件订阅
  4. 销毁运行时
  5. 断开 stdin 输入
  6. 暂停 stdin
  7. 刷新 stdout(SIGTERM 除外)
  8. 退出进程

进程保活

// Keep process alive forever
return new Promise(() => {});

RPC 模式永远不会自己退出,一直等待输入。只有收到 SIGTERM/SIGHUP 或 stdin 关闭时才退出。

会话绑定(rebindSession)

const rebindSession = async (): Promise<void> => {
    session = runtimeHost.session;
    await session.bindExtensions({
        uiContext: createExtensionUIContext(),
        mode: "rpc",
        commandContextActions: {
            waitForIdle: () => session.waitForIdle(),
            newSession: async (options) => runtimeHost.newSession(options),
            fork: async (entryId, forkOptions) => {...},
            navigateTree: async (targetId, options) => {...},
            switchSession: async (sessionPath, options) => {...},
            reload: async () => {...},
        },
        shutdownHandler: () => {
            shutdownRequested = true;
        },
        onError: (err) => {
            output({ type: "extension_error", ... });
        },
    });

    unsubscribe?.();
    unsubscribeBackpressure?.();
    unsubscribe = session.subscribe((event) => {
        output(toJsonEvent(event));
        if (event.type === "agent_settled") {
            void checkShutdownRequested();
        }
    });
    unsubscribeBackpressure = session.agent.subscribe(async () => {
        await waitForRawStdoutBackpressure();
    });
};

与 print 模式的区别:

  • RPC 模式传递了 uiContext,扩展可以请求 UI 操作
  • RPC 模式注册了 shutdownHandler,扩展可以请求关闭进程
  • 事件订阅中检查 agent_settled 事件,触发关闭检查

RPC 模式 vs Print 模式

维度RPC 模式Print 模式
通信方式stdin/stdout JSON 协议直接输出文本/JSON
交互性双向通信,可发送多个命令单次执行
会话管理支持切换、fork、克隆无(或加载已有会话)
扩展 UI支持对话框、通知等不支持
进程生命周期永远运行,直到被杀死执行完就退出
使用场景IDE 插件、程序化调用脚本、管道
返回类型Promise<never>Promise<number>

总结

rpc-mode.ts 实现了一个完整的 RPC 服务器:

  1. 初始化,接管 stdout,设置输出辅助函数
  2. 绑定,绑定扩展,订阅事件
  3. 命令循环,读取 JSON 命令,执行,返回响应
  4. 扩展 UI,支持对话框、通知等 UI 操作
  5. 关闭,信号处理,资源清理

没有 UI、没有交互界面,但提供了完整的程序化接口。适合嵌入其他应用程序。