sequenceDiagram
participant Caller as 调用方
participant Agent as Agent
participant Models as 前文的模型调用封装
participant Bailian as 百炼模型
Caller->>Agent: 提交输入
Agent->>Models: 通过传入的函数调用模型
Models->>Bailian: 发送请求
Bailian-->>Models: 返回回复
Models-->>Agent: 传递模型回复事件
Agent-->>Caller: 回复过程与完整消息
Agent 位于调用方和模型调用封装之间。它负责组织这次处理,百炼配置、认证和协议适配仍由前文的模型调用封装完成。
- 把一次问答交给谁处理
@earendil-works/pi-agent-core 提供的 Agent 是一个保存运行状态、接收输入并推进执行的对象。调用方与它通过三个入口配合:
| 入口 | 谁调用谁 | 本章中的作用 |
|---|---|---|
| `agent.prompt(input)` | 调用方 → Agent | 提交输入,等待这次处理结束 |
| 创建 Agent 时传入的 `streamFn` | Agent → 模型调用函数 | 用模型、上下文和选项发起请求,返回模型事件流 |
| `agent.subscribe(listener)` | 调用方 → Agent(注册) | 注册接收函数,Agent 随后调用 `listener` 通知回复事件 |
streamFn 决定请求怎样发出,subscribe() 决定调用方怎样接收过程。Agent 使用前者取得模型事件流,再把处理后的事件交给后者。
- 用已有百炼配置完成一次流式问答
沿用第五章的 models 和 model,先注册接收事件的回调,再调用 prompt() 提交输入。完整百炼配置和断言放在 Lab 中:
<span>import</span> { <span>Agent</span> } <span>from</span> <span>"@earendil-works/pi-agent-core"</span>;
<span>import</span> <span>type</span> { <span>AssistantMessage</span> } <span>from</span> <span>"@earendil-works/pi-ai"</span>;
<span>const</span> agent = <span>new</span> <span>Agent</span>({
<span>initialState</span>: { model, <span>systemPrompt</span>: <span>"严格按用户要求回答。"</span>, <span>messages</span>: [] },
<span>streamFn</span>: <span>(<span>model, context, options</span>) =></span>
models.<span>streamSimple</span>(model, context, { ...options, <span>maxTokens</span>: <span>128</span> }),
});
<span>let</span> <span>reply</span>: <span>AssistantMessage</span> | <span>undefined</span>;
agent.<span>subscribe</span>(<span>(<span>event</span>) =></span> {
<span>if</span> (event.<span>type</span> === <span>"message_update"</span> && event.<span>assistantMessageEvent</span>.<span>type</span> === <span>"text_delta"</span>) {
process.<span>stdout</span>.<span>write</span>(event.<span>assistantMessageEvent</span>.<span>delta</span>);
}
<span>if</span> (event.<span>type</span> === <span>"message_end"</span> && event.<span>message</span>.<span>role</span> === <span>"assistant"</span>) {
reply = event.<span>message</span>;
}
});
<span>await</span> agent.<span>prompt</span>(<span>"请只回复:你好,Agent。"</span>);
<span>if</span> (!reply) <span>throw</span> <span>new</span> <span>Error</span>(<span>"没有收到完整回复"</span>);
<span>if</span> (reply.<span>stopReason</span> === <span>"error"</span> || reply.<span>stopReason</span> === <span>"aborted"</span>) {
<span>throw</span> <span>new</span> <span>Error</span>(reply.<span>errorMessage</span> ?? reply.<span>stopReason</span>);
}
<span>console</span>.<span>log</span>(<span>"\n完整回复:"</span>, reply.<span>content</span>);
initialState 设置当前模型、系统提示词和空消息列表。streamFn 使用前文的 models,通过 streamSimple() 接受 Agent 提供的通用调用选项;返回的仍是第五章介绍的模型事件流。
显示回调处理两种 Agent 事件:
message_update携带正在变化的回复,其中assistantMessageEvent是原始模型事件。遇到text_delta,继续追加新增文字。message_end携带已经结束的消息。用户输入也会产生这个事件,所以取得模型回复时还要判断role === "assistant"。
await agent.prompt() 等待这次处理及事件回调完成,返回值是 void。上面的 reply 来自结束事件;失败消息也能通过这个事件交付,因此仍要按第五章的方法检查 stopReason。
- 从提交输入追踪到模型调用,再回到调用方
这次问答沿下面的真实调用链执行:
Agent.prompt(input)
├─ normalizePromptInput(...) → 把文本整理为用户消息
└─ runPromptMessages(...)
└─ runWithLifecycle(回调) → 执行回调并管理运行状态
└─ runAgentLoop(...) → 准备本次上下文
└─ runLoop(...) → 推进执行
└─ streamAssistantResponse(...) → 调用模型并读取回复
└─ streamFunction(...) → 我们传入的 models.streamSimple(...)
下面保留与这条路径有关的源码,省略无关字段和分支。
prompt() 整理输入,并传入两个不同的回调
Agent.prompt() 先把输入变成消息,再开始执行:
<span>const</span> messages = <span>this</span>.<span>normalizePromptInput</span>(input, images);
<span>await</span> <span>this</span>.<span>runPromptMessages</span>(messages);
本例传入字符串。normalizePromptInput() 中的文本分支构造内容块,再返回一条用户消息:
<span>const</span> <span>content</span>: <span>Array</span><<span>TextContent</span> | <span>ImageContent</span>> = [{ <span>type</span>: <span>"text"</span>, <span>text</span>: input }];
<span>// ...</span>
<span>return</span> [{ <span>role</span>: <span>"user"</span>, content, <span>timestamp</span>: <span>Date</span>.<span>now</span>() }];
接下来,runPromptMessages() 的方法体把执行参数交给 runAgentLoop():
<span>await</span> <span>this</span>.<span>runWithLifecycle</span>(<span>async</span> (signal) => {
<span>await</span> <span>runAgentLoop</span>(
messages,
<span>this</span>.<span>createContextSnapshot</span>(),
<span>this</span>.<span>createLoopConfig</span>(options),
<span>(<span>event</span>) =></span> <span>this</span>.<span>processEvents</span>(event),
signal,
<span>this</span>.<span>streamFunction</span>,
);
});
createContextSnapshot() 提供当前系统提示词和消息列表的副本,本例起始消息列表为空;createLoopConfig() 提供当前模型和消息转换等调用配置。runWithLifecycle() 执行内部回调并管理运行状态,signal 是随调用传下去的中止信号。
这里最重要的是两个回调的方向:
| 传入的函数 | 内部参数名 | 用途 |
|---|---|---|
| `this.streamFunction` | `streamFn` | 发起模型请求 |
| `(event) => this.processEvents(event)` | `emit` | 把执行事件交回当前 Agent |
this.streamFunction 来自构造 Agent 时的赋值:
<span>this</span>.<span>streamFunction</span> = runtimeOptions.<span>streamFn</span> ?? <span>getDefaultStreamFn</span>();
runtimeOptions 来自创建 Agent 时传入的选项。本例明确提供了 streamFn,因此保存的就是第二节调用 models.streamSimple() 的函数。
runLoop() 将请求交给前文的模型调用封装
runAgentLoop() 接收到的新用户消息参数名是 prompts。它先准备本次上下文,再调用 runLoop(),继续传递同一个 emit 和 streamFn:
<span>const</span> <span>newMessages</span>: <span>AgentMessage</span>[] = [...prompts];
<span>const</span> <span>currentContext</span>: <span>AgentContext</span> = {
...context,
<span>messages</span>: [...context.<span>messages</span>, ...prompts],
};
<span>// ...</span>
<span>await</span> <span>runLoop</span>(currentContext, newMessages, config, signal, emit, streamFn ?? <span>getDefaultStreamFn</span>());
本例开始时 context.messages 为空,所以 currentContext.messages 只包含刚构造的那条用户消息。newMessages 收集本次产生的消息。runLoop() 接收模型调用函数时,参数名是 streamFunction;它再调用 streamAssistantResponse():
<span>// runLoop() 内部</span>
<span>let</span> currentContext = initialContext;
<span>let</span> config = initialConfig;
<span>// ...</span>
<span>const</span> message = <span>await</span> <span>streamAssistantResponse</span>(currentContext, config, signal, emit, streamFunction);
newMessages.<span>push</span>(message);
streamAssistantResponse() 把消息整理成前文的 Context,然后调用这个函数。以下保留请求消息的组装和实际调用位置:
<span>let</span> messages = context.<span>messages</span>;
<span>// ...</span>
<span>const</span> llmMessages = <span>await</span> config.<span>convertToLlm</span>(messages);
<span>const</span> <span>llmContext</span>: <span>Context</span> = {
<span>systemPrompt</span>: context.<span>systemPrompt</span>,
<span>messages</span>: llmMessages,
<span>// ...</span>
};
<span>// ...</span>
<span>const</span> response = <span>await</span> <span>streamFunction</span>(config.<span>model</span>, llmContext, {
...config,
<span>// ...</span>
});
config.convertToLlm 选出模型能接受的消息。本例使用默认转换,唯一一条用户消息会被保留。因此,streamFunction(...) 实际进入我们传入的回调,以这一条消息调用 models.streamSimple()。
models.streamSimple() 接受 Agent 的通用选项,复用第二、三章已经介绍的 Provider 选择、认证与百炼适配,返回模型事件流。Agent 与前文封装的衔接到这里就完成了。
模型事件如何变成调用方收到的 Agent 事件
streamAssistantResponse() 得到 response 后,在函数内部执行 for await (const event of response)。收到模型开始事件时建立临时回复;收到内容事件时,把原始事件放进 Agent 的 message_update。内容更新分支中的实际发送位置是:
partialMessage = event.<span>partial</span>;
<span>// ...</span>
<span>await</span> <span>emit</span>({
<span>type</span>: <span>"message_update"</span>,
<span>assistantMessageEvent</span>: event,
<span>message</span>: { ...partialMessage },
});
收到模型的 done 或 error 后,该函数调用第五章的 result() 取得完整消息,再发送 message_end。下面省略本次消息列表的更新:
<span>const</span> finalMessage = <span>await</span> response.<span>result</span>();
<span>// ...</span>
<span>await</span> <span>emit</span>({ <span>type</span>: <span>"message_end"</span>, <span>message</span>: finalMessage });
<span>return</span> finalMessage;
这两处 emit() 调用的都是最初传入的 (event) => this.processEvents(event)。所以事件沿下面的方向回到调用方:
streamAssistantResponse() 调用 emit(event)
→ Agent.processEvents(event) 更新状态
→ 调用 subscribe() 注册的 listener(event, signal)
listener 是调用方传给 subscribe() 的事件处理函数,第二节的显示回调就是其中一个。subscribe() 只把它存入 listeners,注册时并不执行:
<span>subscribe</span>(<span>listener</span>: <span>(<span>event: AgentEvent, signal: AbortSignal</span>) =></span> <span>Promise</span><<span>void</span>> | <span>void</span>): <span>() =></span> <span>void</span> {
<span>this</span>.<span>listeners</span>.<span>add</span>(listener);
<span>return</span> <span>() =></span> <span>this</span>.<span>listeners</span>.<span>delete</span>(listener);
}
真正调用 listener 的位置在 processEvents() 中。每次 emit(event) 把事件交给它后,它先更新 Agent 状态,再遍历 listeners,将同一个事件传给每个回调,并等待该回调执行完成:
<span>// 状态更新后,通知调用方</span>
<span>for</span> (<span>const</span> listener <span>of</span> <span>this</span>.<span>listeners</span>) {
<span>await</span> <span>listener</span>(event, signal);
}
这样,第二节的显示回调就会随着事件到达而执行:收到 message_update 时,检查 assistantMessageEvent 是否为 text_delta 并显示分片;收到助手消息的 message_end 时,从 message 取得完整回复并赋给 reply。这些字段正是前面的 streamAssistantResponse() 构造事件时放入的内容。
本例只请求一条文本回复。处理完成后,内部执行结束,外层 prompt() 的等待也随之结束;此时结束事件的回调已经把完整回复赋给 reply。
- 运行 Lab,核对一次请求的过程和结果
完整程序在 labs/06-agent-question-answer.ts。沿用第三章的 .env 配置,在项目根目录运行:
npm install
node labs/06-agent-question-answer.ts
依赖中新增了 @earendil-works/pi-agent-core。Lab 使用真实 Agent 和百炼服务,输入是:
请只回复:你好,Agent。
输出包含请求中的消息角色、逐次到达的分片和结束事件交付的完整回复。下面是一次成功运行的输出,分片边界会随运行变化:
model: qwen3.8-flash
input: 请只回复:你好,Agent。
request 1 roles: user
delta: "你好"
delta: ",Agent。"
complete reply: "你好,Agent。"
stopReason: stop
chapter 6 Agent question and answer passed
Lab 验证一次 prompt() 只发出一次模型请求,结束事件提供正常完成的指定问候语,并且显示的分片拼接后与完整回复一致。任何一项不满足,程序都会报错并以非零退出码结束。
- 本章小结
现在,调用方可以把一条输入交给 Agent。Agent 整理请求,通过传入的函数复用百炼调用,再把模型事件转为回复事件;调用方能够显示分片,并在这次处理结束时取得完整回复。
源码核对入口:agent.ts 包含输入入口、模型调用函数的接入和事件通知;agent-loop.ts 包含执行推进和模型回复处理;models.ts 包含复用 Provider 与认证的调用入口。
清晰拆解 Agent 与模型封装的职责边界,适合想理解 Agent 运行状态、事件流回调与流式问答落地路径的开发者按源码逐步对照。