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 正式化

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 洞察

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 正式化

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 洞察

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 正式化

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 洞察

讨论 / Comments

评论托管在本仓库的 GitHub Discussions, 需 GitHub 账号。