Node.js中Express流式x-ndjson测试阻塞与Promise未解析问题排查
我有一个使用Node v19.1.0的TypeScript库,该库包含一个监听流式服务器事件的函数。服务器提供/events路由,流式传输'application/x-ndjson'内容,类型可能为event/ping等(需定时发送ping维持连接)。我的observe函数解析流式数据,将有效事件传递给回调,调用者可通过中止函数按需终止流。
本地或CI运行测试时,会出现以下错误:
Warning: Test "observes events." generated asynchronous activity after the test ended. This activity created the error "AbortError: The operation was aborted." and would have caused the test to fail, but instead triggered an unhandledRejection event.
我用纯JavaScript简化了示例代码:
const assert = require('assert/strict'); const express = require('express'); const { it } = require('node:test'); it('observes events.', async () => { const expectedEvent = { type: 'event', payload: { metadata: { type: 'entity-created', commandId: 'commandId' } } }; const api = express(); const server = api .use(express.json()) .post('/events', (request, response) => { response.writeHead(200, { 'content-type': 'application/x-ndjson', }); const line = JSON.stringify(expectedEvent) + '\n'; response.write(line); }) .listen(3000); let stopObserving = () => { throw new Error('should never happen'); }; const actualEventPayload = await new Promise(async resolve => { stopObserving = await observeEvents(async newEvent => { resolve(newEvent); }); }); stopObserving(); server.closeAllConnections(); server.close(); assert.deepEqual(actualEventPayload, expectedEvent.payload); }); const observeEvents = async function (onReceivedFn) { const abortController = new AbortController(); const response = await fetch('http://localhost:3000/events', { method: 'POST', headers: { 'content-type': 'application/json' }, signal: abortController.signal, }); if (!response.ok) { throw new Error('error handling goes here - request failed'); } Promise.resolve().then(async () => { if (!response.body) { throw new Error('error handling goes here - missing response body'); } for await (const item of parseStream(response.body, abortController)) { switch (item.type) { case 'event': { await onReceivedFn(item.payload); break; } case 'ping': // Intentionally left blank break; case 'error': throw new Error('error handling goes here - stream failed'); default: throw new Error('error handling goes here - should never happen'); } } }); return () => { abortController.abort(); }; }; const parseLine = function () { return new TransformStream({ transform(chunk, controller) { try { const data = JSON.parse(chunk); // ... check if this is a valid line... controller.enqueue(data); } catch (error) { controller.error(error); } }, }); }; const splitLines = function () { let buffer = ''; return new TransformStream({ transform(chunk, controller) { buffer += chunk; const lines = buffer.split('\n'); for (let i = 0; i < lines.length - 1; i++) { controller.enqueue(lines[i]); } buffer = lines.at(-1) ?? ''; }, flush(controller) { if (buffer.length > 0) { controller.enqueue(buffer); } }, }); }; const parseStream = async function* (stream, abortController) { let streamReader; try { const pipedStream = stream .pipeThrough(new TextDecoderStream()) .pipeThrough(splitLines()) .pipeThrough(parseLine()); streamReader = pipedStream.getReader(); while (true) { const item = await streamReader.read(); if (item.done) { break; } yield item.value; } } finally { await streamReader?.cancel(); abortController.abort(); } };
运行node --test时测试无法结束,必须手动取消,问题卡在以下代码段:
const actualEventPayload = await new Promise(async resolve => { stopObserving = await observeEvents(async newEvent => { resolve(newEvent); }); });
我认为是Promise从未解析导致的,即使移除所有流解析代码,替换相关逻辑后问题仍存在,请问问题出在哪里或缺少什么?
核心问题1:服务器响应未结束,流长期挂起
测试服务器的/events路由仅调用response.write(line)发送数据,但未调用response.end()结束响应。这会导致客户端的fetch请求一直处于等待状态,流无法触发done信号,parseStream中的循环会无限卡住。
核心问题2:异步任务未处理错误,触发未捕获rejection
observeEvents中用Promise.resolve().then(async () => { ... })启动了独立异步任务,但未处理该任务的错误。当调用abortController.abort()时,流的read()方法会抛出AbortError,这个错误未被捕获,就会触发测试结束后的异步活动警告。
核心问题3:服务器关闭逻辑未等待完成
直接调用server.closeAllConnections()和server.close()但未等待关闭完成的回调,测试可能提前结束,但服务器连接或进程仍在运行,残留异步活动。
修复步骤
修改服务器路由,结束响应
在response.write(line)后添加response.end(),告知客户端流已结束:.post('/events', (request, response) => { response.writeHead(200, { 'content-type': 'application/x-ndjson', }); const line = JSON.stringify(expectedEvent) + '\n'; response.write(line); response.end(); // 新增:结束响应 })处理
observeEvents中的异步任务错误
移除独立的Promise.resolve().then()链,直接在函数内部处理流逻辑,并捕获错误:const observeEvents = async function (onReceivedFn) { const abortController = new AbortController(); const response = await fetch('http://localhost:3000/events', { method: 'POST', headers: { 'content-type': 'application/json' }, signal: abortController.signal, }); if (!response.ok) { throw new Error('error handling goes here - request failed'); } if (!response.body) { throw new Error('error handling goes here - missing response body'); } try { for await (const item of parseStream(response.body, abortController)) { switch (item.type) { case 'event': { await onReceivedFn(item.payload); break; } case 'ping': break; case 'error': throw new Error('error handling goes here - stream failed'); default: throw new Error('error handling goes here - should never happen'); } } } catch (err) { // 忽略AbortError,其他错误重新抛出 if (err.name !== 'AbortError') { throw err; } } return () => { abortController.abort(); }; };优化测试中的服务器关闭逻辑
改为异步等待服务器关闭完成,避免残留异步活动:stopObserving(); await new Promise((resolve) => { server.closeAllConnections(); server.close(resolve); // 等待服务器关闭完成 });简化Promise初始化逻辑
不需要在new Promise中使用async resolve,直接调整为:let stopObserving; const actualEventPayload = await new Promise(resolve => { observeEvents(async newEvent => { resolve(newEvent); }).then(fn => stopObserving = fn); });
修复后流程说明
- 服务器发送事件后立即结束响应,客户端流处理完数据后正常退出循环
observeEvents中的异步错误被捕获,不会触发未处理的rejection- 测试等待服务器完全关闭后再结束,彻底消除残留异步活动
内容的提问来源于stack exchange,提问作者baitendbidz

