Skip to content

Commit a0d7d76

Browse files
authored
fix(standard-server): settle sendStandardResponse when client disconnects before sending (#1738)
Fixes #1735 `handle()` hung forever when the client disconnected while the procedure was still running, because the response's `'close'` event had already fired before `sendStandardResponse()` started listening for it. Now: - `sendStandardResponse` (node, fastify, aws-lambda) settles immediately when the response is already ended and destroys the prepared body so event iterators are cleaned up. - `toAbortSignal` aborts immediately when the response is already closed, including http/2 client cancels. - Both detect an ended response via the new `isNodeResponseStreamEnded` util, which also handles `Http2ServerResponse` (it hides its stream state at runtime). Verified with real http/1 and http/2 servers: aborted requests now resolve instead of hanging.
1 parent ec6d86d commit a0d7d76

11 files changed

Lines changed: 532 additions & 17 deletions

File tree

packages/standard-server-aws-lambda/src/response.test.ts

Lines changed: 75 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -249,4 +249,79 @@ describe('sendStandardResponse', () => {
249249
await sendPromise
250250
})
251251
})
252+
253+
describe('response closed before sending', () => {
254+
it('resolves and destroys the body', async () => {
255+
let clean = false
256+
const res: StandardResponse = {
257+
body: (async function* () {
258+
try {
259+
yield 1
260+
}
261+
finally {
262+
clean = true
263+
}
264+
})(),
265+
headers: {},
266+
status: 200,
267+
}
268+
269+
const responseStream = new Stream.Writable({
270+
write(chunk, encoding, callback) {
271+
callback()
272+
},
273+
})
274+
275+
responseStream.destroy()
276+
277+
await vi.waitFor(() => {
278+
expect(responseStream.closed).toBe(true)
279+
})
280+
281+
await expect(sendStandardResponse(responseStream, res, { eventIteratorKeepAliveComment: 'test' })).resolves.toBeUndefined()
282+
283+
await vi.waitFor(() => {
284+
expect(clean).toBe(true)
285+
})
286+
287+
expect((globalThis as any).awslambda.HttpResponseStream.from).not.toHaveBeenCalled()
288+
})
289+
290+
it('rejects when response was destroyed with an error', async () => {
291+
let clean = false
292+
const res: StandardResponse = {
293+
body: (async function* () {
294+
try {
295+
yield 1
296+
}
297+
finally {
298+
clean = true
299+
}
300+
})(),
301+
headers: {},
302+
status: 200,
303+
}
304+
305+
const responseStream = new Stream.Writable({
306+
write(chunk, encoding, callback) {
307+
callback()
308+
},
309+
})
310+
311+
responseStream.once('error', () => {})
312+
responseStream.destroy(new Error('test'))
313+
314+
await vi.waitFor(() => {
315+
expect(responseStream.closed).toBe(true)
316+
})
317+
318+
await expect(sendStandardResponse(responseStream, res)).rejects.toThrow('test')
319+
320+
await vi.waitFor(() => {
321+
expect(clean).toBe(true)
322+
})
323+
324+
expect((globalThis as any).awslambda.HttpResponseStream.from).not.toHaveBeenCalled()
325+
})
326+
})
252327
})

packages/standard-server-aws-lambda/src/response.ts

Lines changed: 17 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -12,10 +12,25 @@ export function sendStandardResponse(
1212
options: SendStandardResponseOptions = {},
1313
): Promise<void> {
1414
return new Promise((resolve, reject) => {
15+
const [body, standardHeaders] = toLambdaBody(standardResponse.body, standardResponse.headers, options)
16+
17+
if (responseStream.closed || responseStream.destroyed) {
18+
if (typeof body === 'object' && !body.closed) {
19+
body.destroy()
20+
}
21+
22+
if (responseStream.errored) {
23+
reject(responseStream.errored)
24+
}
25+
else {
26+
resolve()
27+
}
28+
29+
return
30+
}
31+
1532
responseStream.once('error', reject)
1633
responseStream.once('close', resolve)
17-
18-
const [body, standardHeaders] = toLambdaBody(standardResponse.body, standardResponse.headers, options)
1934
const [headers, setCookies] = toLambdaHeaders(standardHeaders)
2035

2136
// awslambda is global aws lambda global object

packages/standard-server-fastify/src/response.test.ts

Lines changed: 82 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
import { Readable } from 'node:stream'
1+
import { Readable, Writable } from 'node:stream'
22
import FastifyCookie from '@fastify/cookie'
33
import * as StandardServerNode from '@orpc/standard-server-node'
44
import Fastify from 'fastify'
@@ -303,4 +303,85 @@ describe('sendStandardResponse', () => {
303303
expect(res.text).toEqual(JSON.stringify({ foo: 'bar' }))
304304
})
305305
})
306+
307+
describe('response closed before sending', () => {
308+
it('resolves and destroys the body', async () => {
309+
let clean = false
310+
const body = (async function* () {
311+
try {
312+
yield 1
313+
}
314+
finally {
315+
clean = true
316+
}
317+
})()
318+
319+
const raw = new Writable({
320+
write(chunk, encoding, callback) {
321+
callback()
322+
},
323+
})
324+
325+
raw.destroy()
326+
327+
const reply = {
328+
raw,
329+
status: vi.fn(),
330+
headers: vi.fn(),
331+
send: vi.fn(),
332+
} as any
333+
334+
await expect(sendStandardResponse(reply, {
335+
status: 200,
336+
headers: {},
337+
body,
338+
})).resolves.toBeUndefined()
339+
340+
await vi.waitFor(() => {
341+
expect(clean).toBe(true)
342+
})
343+
344+
expect(reply.send).not.toHaveBeenCalled()
345+
})
346+
347+
it('rejects when response was destroyed with an error', async ({ onTestFinished }) => {
348+
let clean = false
349+
const body = (async function* () {
350+
try {
351+
yield 1
352+
}
353+
finally {
354+
clean = true
355+
}
356+
})()
357+
358+
const raw = new Writable({
359+
write(chunk, encoding, callback) {
360+
callback()
361+
},
362+
})
363+
364+
raw.once('error', () => { })
365+
raw.destroy(new Error('test'))
366+
367+
const reply = {
368+
raw,
369+
status: vi.fn(),
370+
headers: vi.fn(),
371+
send: vi.fn(),
372+
} as any
373+
374+
await expect(sendStandardResponse(reply, {
375+
status: 200,
376+
headers: {},
377+
body,
378+
})).rejects.toThrow('test')
379+
380+
await vi.waitFor(() => {
381+
expect(clean).toBe(true)
382+
})
383+
384+
expect(reply.send).not.toHaveBeenCalled()
385+
})
386+
})
306387
})

packages/standard-server-fastify/src/response.ts

Lines changed: 19 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
import type { StandardHeaders, StandardResponse } from '@orpc/standard-server'
22
import type { ToNodeHttpBodyOptions } from '@orpc/standard-server-node'
33
import type { FastifyReply } from 'fastify'
4-
import { toNodeHttpBody, toNodeHttpHeaders } from '@orpc/standard-server-node'
4+
import { isNodeResponseStreamEnded, toNodeHttpBody, toNodeHttpHeaders } from '@orpc/standard-server-node'
55

66
export interface SendStandardResponseOptions extends ToNodeHttpBodyOptions { }
77

@@ -11,13 +11,28 @@ export function sendStandardResponse(
1111
options: SendStandardResponseOptions = {},
1212
): Promise<void> {
1313
return new Promise((resolve, reject) => {
14-
reply.raw.once('error', reject)
15-
reply.raw.once('close', resolve)
16-
1714
const resHeaders: StandardHeaders = { ...standardResponse.headers }
1815

1916
const resBody = toNodeHttpBody(standardResponse.body, resHeaders, options)
2017

18+
if (isNodeResponseStreamEnded(reply.raw)) {
19+
if (typeof resBody === 'object' && !resBody.closed) {
20+
resBody.destroy()
21+
}
22+
23+
if (reply.raw.errored) {
24+
reject(reply.raw.errored)
25+
}
26+
else {
27+
resolve()
28+
}
29+
30+
return
31+
}
32+
33+
reply.raw.once('error', reject)
34+
reply.raw.once('close', resolve)
35+
2136
reply.status(standardResponse.status)
2237
// Fastify treats undefined headers as empty string, so remember to use toNodeHttpHeaders
2338
// to filter out undefined headers

packages/standard-server-node/src/index.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,3 +7,4 @@ export * from './response'
77
export * from './signal'
88
export * from './types'
99
export * from './url'
10+
export * from './utils'

packages/standard-server-node/src/response.test.ts

Lines changed: 79 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -302,4 +302,83 @@ describe('sendStandardResponse', () => {
302302
await sendPromise
303303
})
304304
})
305+
306+
describe('response closed before sending', () => {
307+
it('resolves and destroys the body', async () => {
308+
let clean = false
309+
const res: StandardResponse = {
310+
body: (async function* () {
311+
try {
312+
yield 1
313+
}
314+
finally {
315+
clean = true
316+
}
317+
})(),
318+
headers: {},
319+
status: 200,
320+
}
321+
322+
const responseStream = new Stream.Writable({
323+
write(chunk, encoding, callback) {
324+
callback()
325+
},
326+
})
327+
328+
;(responseStream as any).writeHead = vi.fn()
329+
330+
responseStream.destroy()
331+
332+
await vi.waitFor(() => {
333+
expect(responseStream.closed).toBe(true)
334+
})
335+
336+
await expect(sendStandardResponse(responseStream as any, res)).resolves.toBeUndefined()
337+
338+
await vi.waitFor(() => {
339+
expect(clean).toBe(true)
340+
})
341+
342+
expect((responseStream as any).writeHead).not.toHaveBeenCalled()
343+
})
344+
345+
it('rejects when response was destroyed with an error', async () => {
346+
let clean = false
347+
const res: StandardResponse = {
348+
body: (async function* () {
349+
try {
350+
yield 1
351+
}
352+
finally {
353+
clean = true
354+
}
355+
})(),
356+
headers: {},
357+
status: 200,
358+
}
359+
360+
const responseStream = new Stream.Writable({
361+
write(chunk, encoding, callback) {
362+
callback()
363+
},
364+
})
365+
366+
;(responseStream as any).writeHead = vi.fn()
367+
368+
responseStream.once('error', () => {})
369+
responseStream.destroy(new Error('test'))
370+
371+
await vi.waitFor(() => {
372+
expect(responseStream.closed).toBe(true)
373+
})
374+
375+
await expect(sendStandardResponse(responseStream as any, res)).rejects.toThrow('test')
376+
377+
await vi.waitFor(() => {
378+
expect(clean).toBe(true)
379+
})
380+
381+
expect((responseStream as any).writeHead).not.toHaveBeenCalled()
382+
})
383+
})
305384
})

packages/standard-server-node/src/response.ts

Lines changed: 19 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ import type { ToNodeHttpBodyOptions } from './body'
33
import type { NodeHttpResponse } from './types'
44
import { toNodeHttpBody } from './body'
55
import { toNodeHttpHeaders } from './headers'
6+
import { isNodeResponseStreamEnded } from './utils'
67

78
export interface SendStandardResponseOptions extends ToNodeHttpBodyOptions {}
89

@@ -12,13 +13,28 @@ export function sendStandardResponse(
1213
options: SendStandardResponseOptions = {},
1314
): Promise<void> {
1415
return new Promise((resolve, reject) => {
15-
res.once('error', reject)
16-
res.once('close', resolve)
17-
1816
const resHeaders: StandardHeaders = { ...standardResponse.headers }
1917

2018
const resBody = toNodeHttpBody(standardResponse.body, resHeaders, options)
2119

20+
if (isNodeResponseStreamEnded(res)) {
21+
if (typeof resBody === 'object' && !resBody.closed) {
22+
resBody.destroy()
23+
}
24+
25+
if (res.errored) {
26+
reject(res.errored)
27+
}
28+
else {
29+
resolve()
30+
}
31+
32+
return
33+
}
34+
35+
res.once('error', reject)
36+
res.once('close', resolve)
37+
2238
// Node.js throws an error when a header is undefined, so remember to use toNodeHttpHeaders
2339
// to filter out undefined headers
2440
res.writeHead(standardResponse.status, toNodeHttpHeaders(resHeaders))

0 commit comments

Comments
 (0)