TS·07 异步与流式:Promise、async/await、与 async generator
这是系列里最"pi"的一章,也是理解整个 agent 框架数据流的钥匙。pi 是一个 AI coding-agent:它向 LLM 发一个请求,模型不是一次性把答案吐回来,而是一个 token 一个 token 地流式返回;pi 要一边接收、一边解析、一边把它变成一串类型化的事件(text_delta、toolcall_start……)往上层推。承载这整套机制的语言特性,就是本章的三块拼图:Promise(将来才有的值)、async/await(取那个值)、以及 async generator(边算边一个个吐值)。前两块是异步编程的地基,第三块是把地基垒成"流"的那道拱。
如果你写过 Python 的 asyncio,这三块都有直接对照:Promise ≈ awaitable/coroutine,async/await 就是 async def/await,async function* + for await...of 就是异步生成器 async def + yield 配 async for。差异在细节里,我们逐个拆。读完你会看懂 pi 是怎么把一条网络 SSE 流,变成 agent 主循环里那个 for await (const event of response) 的。
有一点值得先建立心智:JS/TS 是单线程的,并发靠事件循环而非多线程。所谓"异步"不是真的并行计算,而是"遇到 IO 就让出控制权,别的任务先跑,IO 好了再回来"。这一点和 Python asyncio 完全一致——它们的并发模型是同一套(协作式、单事件循环),所以你在 Python 里对 await 的所有直觉都能平移过来。真正需要重新学的只是 JS 特有的写法:Promise 这个对象、.then 链、以及取消要手动接 AbortSignal。三章走完,这些就都补齐了。
1. Promise 与 async/await
1.1 直觉
Promise<T> 是"将来会有一个 T 类型的值"的占位符——一个可能还没兑现的凭证。异步函数不能立刻给你结果(结果要等 IO、等网络),于是先给你一个 Promise,值到了再兑现。
关键要分清两件独立的事:一个函数是不是 async,和它返不返回 Promise。async 只是语法糖——它让你能在函数体里用 await,并自动把返回值包进 Promise。但一个普通(非 async)函数完全可以直接 return 一个已经存在的 Promise。对调用方来说,两者拿到的都是 Promise<T>,await 起来没区别。
和 Python 协程还有一处微妙区别:Python 里 foo()(foo 是 async def)返回一个协程对象但不会开始执行,直到你 await 它或交给事件循环;而 JS 里 foo()(foo 是 async)会立即开始执行函数体,一路跑到第一个 await 才让出——返回的 Promise 代表"已经在跑、将来会有结果"。这个"急切启动"的差异,决定了下面的写法为什么合理。
Python 对照:
async函数 ≈async def;await expr就是 Python 的await expr;Promise<T>≈ 一个 awaitable(协程/Future)。区别是 JS 只有一个内置事件循环、且自动运行——你不用像 Python 那样asyncio.run(...)手动启动。
1.2 最小 demo
// 教学示例 — 非生产代码
// 情形 A:async 函数,自动返回 Promise<string>
async function fetchName(id: number): Promise<string> {
// await 会"暂停"本函数直到 Promise 兑现,拿到里面的值
const res = await fakeDb(id); // res: string
return res.toUpperCase(); // 返回值被自动包成 Promise<string>
}
// 情形 B:普通函数,不是 async,却直接返回一个 Promise
function fetchNameEager(id: number): Promise<string> {
return fakeDb(id); // 直接把已有的 Promise 透传出去
}
// 两者调用方写法完全一样:
async function main() {
const a = await fetchName(1); // a: string
const b = await fetchNameEager(2); // b: string
console.log(a, b);
}
function fakeDb(id: number): Promise<string> {
return Promise.resolve("row-" + id);
}
fetchName 用了 async/await,fetchNameEager 没有——但它们的返回类型都是 Promise<string>,main 里 await 的用法一模一样。这就是"是 async"与"返回 Promise"两回事的直观演示。
1.3 正式化
- 声明:
async function f(): Promise<T>或箭头const f = async (): Promise<T> => {...};类方法写async method(): Promise<T>。 - 返回类型:
async函数的返回类型总是Promise<某>。写Promise<string>;若不返回值,写Promise<void>。你return x(x 是T),TS 自动认成Promise<T>;你return一个Promise<T>,TS 不会嵌套成Promise<Promise<T>>,会自动摊平。 await:只能出现在async函数体内(或模块顶层)。await p中,若p是Promise<T>,表达式的值就是T;若p不是 Promise,await原样返回它。await会让出事件循环,让别的任务先跑。- 不 await 的后果:调用一个
async函数但不await,你拿到的是没兑现的 Promise 本身,不是里面的值——常见 bug 来源(相当于 Python 里忘了await一个协程)。 - 错误:Promise 兑现失败会"reject";
await一个 rejected 的 Promise 会抛异常,用普通try/catch捕获即可(对照 Python:被 await 的协程抛异常同样在 await 处抛出)。 - 不用 async 也能消费 Promise:除了
await,还有更老的写法p.then(v => ...)注册"兑现后回调"、p.catch(e => ...)注册"失败回调"。async/await只是这套回调的语法糖,让异步代码读起来像同步代码。pi 里绝大多数地方用await,但理解.then有助于读懂库代码。 - 顶层 await:在 ES 模块(见类与模块)的最外层可以直接写
await,不必包一层async函数;但在普通函数体内await仍要求该函数是async。
1.4 代码引用
先看 pi 里一个"标准 async 函数返回 Promise"的例子——一个查询当前时间的工具:
pi/packages/agent/test/utils/get-current-time.ts:L6-L18 — async 函数,返回类型显式标注 Promise<...>,函数体内 return 一个普通对象,被自动包成 Promise
export async function getCurrentTime(timezone?: string): Promise<GetCurrentTimeResult> {
const date = new Date();
if (timezone) {
try {
const timeStr = date.toLocaleString("en-US", {
timeZone: timezone,
dateStyle: "full",
timeStyle: "long",
});
return {
content: [{ type: "text", text: timeStr }],
details: { utcTimestamp: date.getTime() },
};
} catch (_e) {
throw new Error(`Invalid timezone: ${timezone}. Current UTC time: ${date.toISOString()}`);
}
}
再看返回 Promise 但不是 async 的对照——EventStream 的 result():
pi/packages/ai/src/utils/event-stream.ts:L8-L18,L64-L66 — 构造时先造好一个 Promise 存起来;result() 是普通方法(无 async),直接把这个 Promise 透传出去
private finalResultPromise: Promise<R>;
private resolveFinalResult!: (result: R) => void;
// ...
constructor(isComplete: (event: T) => boolean, extractResult: (event: T) => R) {
this.isComplete = isComplete;
this.extractResult = extractResult;
this.finalResultPromise = new Promise((resolve) => {
this.resolveFinalResult = resolve;
});
}
// ...
result(): Promise<R> {
return this.finalResultPromise;
}
对照本节语法点:getCurrentTime 有 async,函数体里 return 一个对象字面量,返回类型标注成 Promise<GetCurrentTimeResult>——async 帮你把对象包进 Promise。而 result() 没有 async 关键字,却也返回 Promise<R>:它在构造函数里用 new Promise((resolve) => {...}) 手工造了一个 Promise 存进 finalResultPromise,result() 只是把这个"将来才会 resolve 的凭证"原样递出去——谁 await stream.result(),谁就会一直挂到别处调用 resolveFinalResult(...) 那一刻。这正是"是 async"和"返回 Promise"两件事的真实对比。
1.5 洞察
async是给函数体的糖,不是给调用方的契约:调用方只看返回类型Promise<T>,不关心里面有没有async。手工构造并透传 Promise(如result())是完全合法且常见的写法,用于"值稍后由别的代码 resolve"的场景。- 忘了
await不会报类型错:const x = getCurrentTime()里x是Promise<...>而不是结果本身,后续当普通对象用就会出错。养成"调 async 函数必 await(或显式.then/存 Promise)"的习惯。一个特别隐蔽的坑:在if (getCurrentTime())这类布尔上下文里,Promise 对象恒为真,判断永远成立——Python 里 await 一个协程的返回同样需要显式,但 JS 这里连运行时都不报错,尤其要警惕。 new Promise((resolve) => ...)里的resolve可以存起来晚点再调——这就是 pi 把"事件流结束时才有最终结果"接到一个 Promise 上的手法,下一章的 AbortSignal 也用同一招。!断言的用意:上面resolveFinalResult!的感叹号是"确定赋值断言",告诉 TS"这个字段虽然构造函数体里没直接赋值,但我保证它会被赋值(在new Promise的回调里)",绕过严格初始化检查。这是把 resolve 存成字段时的常见配套写法。
2. 并发与取消:Promise.all、AbortSignal
2.1 直觉
有了单个 Promise,自然要问两件事:怎么同时等一堆、怎么中途喊停。
Promise.all([p1, p2, ...]) 接一个 Promise 数组,返回一个新 Promise,在全部兑现后一次性给你结果数组——这些 Promise 是并发跑的,总耗时约等于最慢那个,而不是逐个相加。这里的"并发"仍是单线程事件循环下的交错执行:三个网络请求可以同时在途,谁的响应先回来先处理,但没有多线程。AbortSignal 则是取消的标准协议:一个 AbortController 持有 signal,调用 controller.abort() 时,所有拿到这个 signal 的异步操作都能收到"该停了"的通知,提前退出。对 agent 尤其重要——用户按 Ctrl-C 时,正在途中的 LLM 请求和工具调用都要能被同一个 signal 一起叫停。
Python 对照:
Promise.all(...)≈asyncio.gather(*coros);AbortSignal/AbortController≈ 取消一个 task 后协程收到的asyncio.CancelledError。差别是 JS 的取消不是异常自动传播,而是你要主动监听signal的"abort"事件、自己决定怎么中止。
2.2 最小 demo
// 教学示例 — 非生产代码
// Promise.all:并发等三个,拿到有序结果数组
async function loadAll(): Promise<string[]> {
const results = await Promise.all([
fakeDb(1),
fakeDb(2),
fakeDb(3),
]);
return results; // results: string[],顺序与入参一致
}
// AbortSignal:把"取消"接到一个 Promise 上
function wait(ms: number, signal?: AbortSignal): Promise<void> {
return new Promise((resolve, reject) => {
if (signal?.aborted) return reject(new Error("aborted"));
const id = setTimeout(resolve, ms);
signal?.addEventListener("abort", () => {
clearTimeout(id); // 清理副作用
reject(new Error("aborted"));
}, { once: true });
});
}
async function demo() {
const c = new AbortController();
setTimeout(() => c.abort(), 50); // 50ms 后喊停
await wait(1000, c.signal); // 会在 50ms 时 reject
}
function fakeDb(id: number): Promise<string> {
return Promise.resolve("row-" + id);
}
loadAll 三个查询并发跑;wait 演示了取消的惯用法:在 new Promise 里同时挂一个 setTimeout 和一个 signal 的 "abort" 监听,谁先到谁决定这个 Promise 是 resolve 还是 reject。
2.3 正式化
Promise.all(iterable):入参是 Promise(或普通值)的数组,返回Promise<T[]>。任一 reject 会让整体立刻 reject(短路)。要"全部跑完不管成败"用Promise.allSettled。结果数组顺序对应入参顺序,与谁先兑现无关。.map(...)造数组:常见搭配是Promise.all(items.map(x => asyncFn(x)))——把一组输入映射成一组 Promise 再一起等。.map见类与模块一章的数组方法。AbortController/AbortSignal:const c = new AbortController()后c.signal是只读信号,c.abort()触发。检查用signal.aborted(布尔);监听用signal.addEventListener("abort", handler, { once: true })——{ once: true }表示触发一次后自动移除监听,省得手动清理。signal?.:?.是可选链(见窄化),signal可能是undefined(没传取消能力),signal?.addEventListener(...)在signal为空时整体求值为undefined而不报错。这让"可选的取消能力"写起来很干净:传了 signal 就接上取消,没传就当无事发生。- 相关工具:
Promise.race([...])返回最先兑现/失败的那个,常用来做超时("请求 vs 定时器,谁先到");Promise.any([...])返回最先成功的一个。它们和Promise.all组成一套并发原语,对照 Python 的asyncio.wait(..., return_when=FIRST_COMPLETED)。
2.4 代码引用
pi 的 agent 主循环里,一批工具调用是并行执行的——用 Promise.all + .map:
pi/packages/agent/src/agent-loop.ts:L542-L544 — Promise.all 并发跑完一批工具调用;每个 entry 可能是函数(要调用)或已算好的结果
const orderedFinalizedCalls = await Promise.all(
finalizedCalls.map((entry) => (typeof entry === "function" ? entry() : Promise.resolve(entry))),
);
而 read 工具演示了"把取消接到一个 Promise 上"的标准写法:
pi/packages/coding-agent/src/core/tools/read.ts:L223-L234 — new Promise 里先查 signal.aborted,再挂一个 abort 监听,取消时 reject 整个 Promise
return new Promise<{ content: (TextContent | ImageContent)[]; details: ReadToolDetails | undefined }>(
(resolve, reject) => {
if (signal?.aborted) {
reject(new Error("Operation aborted"));
return;
}
let aborted = false;
const onAbort = () => {
aborted = true;
reject(new Error("Operation aborted"));
};
signal?.addEventListener("abort", onAbort, { once: true });
对照本节语法点:Promise.all(finalizedCalls.map(...)) 正是"一组输入 → 一组 Promise → 一起等"的惯用式,await 后 orderedFinalizedCalls 是保持原顺序的结果数组(变量名里的 ordered 就是在强调这点)。read 里则把上一章 result() 的手法用在取消上:进入 new Promise 先用 signal?.aborted 挡掉"已经取消了"的情形,再 addEventListener("abort", onAbort, { once: true }) 把未来的取消翻译成对这个 Promise 的 reject——于是外层 await 这个 Promise 的代码会以抛异常的形式感知到取消。
2.5 洞察
Promise.all是短路的:一个失败就整体 reject,其余 Promise 仍在后台跑完(不会自动取消)。要真正停掉它们,得配合把同一个signal传给每个操作。- 取消在 JS 里是"协作式"的:
abort()不会强行杀死正在跑的代码,只是发信号;每个异步操作必须自己监听signal并主动退出——这和 PythonCancelledError会在 await 点自动抛出略有不同,JS 要你手写监听。 { once: true }省心:一次性监听触发后自动解绑,避免内存泄漏;若不用它,记得在操作正常结束时removeEventListener手动清理(pi 的read工具在正常路径结束时确实调了removeEventListener,双保险)。- 一个 signal 贯穿全链:pi 把同一个
AbortSignal从 CLI 顶层一路透传给 LLM 请求、每个工具调用、乃至上一章的 SSE 解析生成器,于是一次abort()能同时叫停所有在途异步操作——这是 agent 能"干净地被 Ctrl-C 打断"的工程前提。
3. async generator + for await...of:pi 如何流式吐 LLM
3.1 直觉
前两章的 Promise 是"一个将来的值"。但 LLM 流式输出是"一连串陆续到来的值":token 一片片来,工具调用一段段来。承载它的是异步生成器 async function*:函数体里每 yield 一次就"吐"出一个值给消费方,吐完再继续算下一个;消费方用 for await (const x of stream) 一个个接,await 保证"没有新值就挂起等待"(即背压——生产快了消费方自然等,消费慢了生产方也不会淹没它)。
一个类还可以 implements AsyncIterable<T>,声明"我能被 for await...of 消费",只要提供一个 [Symbol.asyncIterator]() 方法。pi 正是用这套机制,把一条底层网络 SSE 流,层层包成"一串类型化事件",最后在 agent 主循环里 for await + switch 消费——和 T3 的可辨识联合 switch 无缝接上。
为什么非用生成器不可?因为 LLM 流式响应的本质是"值的时间维度":你事先不知道总共会有多少个 token、多少段工具调用,它们随时间陆续到达。用一个普通的 Promise<全部结果> 就得等全部结束才能拿到第一个字,交互体验很差;而异步生成器让你"来一个处理一个",UI 能实时上屏。这正是 agent 类应用把生成器当一等公民的原因。
Python 对照:
async function*≈ Python 里async def函数体内用yield(异步生成器);yield x就是 Python 的yield x;for await (const x of it)≈async for x in it;implements AsyncIterable<T>≈ 实现__aiter__/__anext__协议。
3.2 最小 demo
// 教学示例 — 非生产代码
// async function*:边算边 yield 一串值(注意 function 后的星号 *)
async function* countUp(n: number): AsyncGenerator<number> {
for (let i = 0; i < n; i++) {
await tick(); // 模拟等待下一片数据(IO/网络)
yield i; // 吐出一个值,函数在此暂停,直到消费方要下一个
}
}
// 消费端:for await...of 一个个接
async function consume() {
for await (const x of countUp(3)) {
console.log(x); // 依次打印 0, 1, 2
}
}
// 一个类也能被 for await...of 消费:实现 AsyncIterable
class Ticker implements AsyncIterable<number> {
async *[Symbol.asyncIterator](): AsyncIterator<number> {
yield 1;
yield 2;
}
}
// for await (const t of new Ticker()) { ... }
function tick(): Promise<void> {
return new Promise((r) => setTimeout(r, 10));
}
countUp 是异步生成器:function* 的星号标记"这是生成器",async 让它能在 yield 之间 await。consume 用 for await...of 逐个取。Ticker 展示了另一条路:类通过实现 [Symbol.asyncIterator]() 让自己可被 for await...of 迭代——这正是 pi 的 EventStream 用的形态。
3.3 正式化
- 声明异步生成器:
async function* g(): AsyncGenerator<T> {...}或方法形式async *methodName(): AsyncIterator<T> {...}。星号*紧跟function(或方法名前),缺了它就是普通 async 函数。 yield:每yield v产出一个T,函数暂停;下次被要值时从暂停处继续。生成器体内可自由混用await(等 IO)和yield(产值)。- 返回类型:
AsyncGenerator<T>或更宽的AsyncIterable<T>/AsyncIterator<T>。三者关系:AsyncIterable<T>是"有[Symbol.asyncIterator]()方法"的东西;AsyncGenerator<T>是它的一个具体实现。 for await...of:for await (const x of src)中src只需是AsyncIterable<T>(生成器、或实现了该协议的类实例都行),x的类型是T。循环体每轮隐式await下一个值;src迭代结束(生成器return或耗尽)时循环退出。只能在async函数内使用。[Symbol.asyncIterator]:一个"计算属性名"——用内置符号Symbol.asyncIterator作方法名,是"我可被异步迭代"的协议约定。类里写async *[Symbol.asyncIterator]() {...}即同时满足implements AsyncIterable<T>。方括号[...]表示"用括号里表达式的值当键名",因为Symbol.asyncIterator不是普通字符串而是一个唯一符号,必须这么写。对照 Python:这相当于给类实现__aiter__/__anext__双下方法。- 提前结束:
for await...of循环里break或抛异常会让底层生成器收到"提前终止"信号,触发它的finally块做清理(如关闭网络 reader)。这也是 pi 能在取消时干净收尾的原因之一。 - 同步版对照:去掉
async/await就是同步生成器function*+for...of,产出的是立刻可得的一串值(如遍历一棵树)。异步版的唯一区别是每次取值可能要等 IO。
3.4 代码引用
先看 pi 的 EventStream——它 implements AsyncIterable<T>,手写异步迭代器,内建背压:
pi/packages/ai/src/utils/event-stream.ts:L4,L50-L62 — 类声明 implements AsyncIterable<T>;async * 方法既 yield 已排队的值,队列空又没结束时 await 一个 Promise 挂起等待(背压)
export class EventStream<T, R = T> implements AsyncIterable<T> {
// ...
async *[Symbol.asyncIterator](): AsyncIterator<T> {
while (true) {
if (this.queue.length > 0) {
yield this.queue.shift()!;
} else if (this.done) {
return;
} else {
const result = await new Promise<IteratorResult<T>>((resolve) => this.waiting.push(resolve));
if (result.done) return;
yield result.value;
}
}
}
再看数据的源头:pi 解析底层网络流(ReadableStream),用 async function* 边解析边 yield 事件:
pi/packages/ai/src/api/anthropic-messages.ts:L384-L412 — async function* 读取 ReadableStream,逐块 decode、逐行解析,每解析出一个 SSE 事件就 yield
async function* iterateSseMessages(
body: ReadableStream<Uint8Array>,
signal?: AbortSignal,
): AsyncGenerator<ServerSentEvent> {
const reader = body.getReader();
const decoder = new TextDecoder();
const state: SseDecoderState = { event: null, data: [], raw: [] };
let buffer = "";
try {
while (true) {
if (signal?.aborted) {
throw new Error("Request was aborted");
}
const { value, done } = await reader.read();
if (done) {
break;
}
buffer += decoder.decode(value, { stream: true });
let consumed = consumeLine(buffer);
while (consumed) {
buffer = consumed.rest;
const event = decodeSseLine(consumed.line, state);
if (event) {
yield event;
}
}
最后,这串事件在 agent 主循环里被 for await + switch 消费——和 T3 的可辨识联合收口:
pi/packages/agent/src/agent-loop.ts:L319-L327 — for await 逐个取事件,switch (event.type) 按可辨识联合的判别字段分派
for await (const event of response) {
switch (event.type) {
case "start":
partialMessage = event.partial;
context.messages.push(partialMessage);
addedPartial = true;
await emit({ type: "message_start", message: { ...partialMessage } });
break;
对照本节语法点:三段拼成一条完整数据流。iterateSseMessages 是源头生成器——async function* 里 await reader.read() 等下一块网络数据(前两章的 Promise/await),yield event 把解析出的事件一个个吐出;它同时把 signal?.aborted 检查嵌进循环,顺手接上了第 2 章的取消。EventStream 则展示了类形态的同一能力:implements AsyncIterable<T> + async *[Symbol.asyncIterator](),并在队列空且未结束时 await new Promise(...) 挂起——把第 1 章"存起 resolve 晚点再调"的手法用作背压,消费方要值而暂时没有,就静静等着 push 触发。到了 agent-loop,for await (const event of response) 只管逐个接,switch (event.type) 按判别字段 type 分派(见可辨识联合)。至此:网络流 → 生成器 yield 事件 → for await 消费 → switch 分派,整条链闭合。
3.5 洞察
- 星号别漏:
async function*(有*)是异步生成器,yield一串值;async function(无*)是普通异步函数,return一个值。漏了*,yield会直接语法报错。 for await...of天然带背压与顺序:每轮隐式await下一个值,生产方快了消费方自然按自己节奏取,不用手写缓冲队列——这是把"流式 + 反压"写得像普通for一样简单的关键。- 生成器 vs 类实现二选一:能用
async function*就用它(最简);当你需要额外方法(如EventStream.result()那样"流之外还要一个最终值")或要被implements AsyncIterable<T>约束时,才落到"手写[Symbol.asyncIterator]()类"这一层。pi 两种都用:底层解析用async function*,需要"边流边留最终结果"的中间层用类——按需求选形态,别一律套类。 - 对照 Python 迁移:整套机制心智模型等同
async for x in agen,几乎一一对应;主要新概念只是[Symbol.asyncIterator]这个"用符号当方法名"的协议写法,以及取消要手动接signal。 - 这就是读懂 pi 事件流的钥匙:从
iterateSseMessages解析网络字节,到EventStream带背压地转发,到agent-loop的for await+switch分派——三层都是同一套"异步可迭代"抽象。抓住"生成器 yield 一串类型化事件、消费方 for await 逐个 switch"这条主线,pi 里几乎所有流式代码都是它的变体。
讨论 / Comments
评论托管在本仓库的 GitHub Discussions, 需 GitHub 账号。