diff --git a/packages/peer/src/body.ts b/packages/peer/src/body.ts index 6ec59c7..e42a0d6 100644 --- a/packages/peer/src/body.ts +++ b/packages/peer/src/body.ts @@ -38,7 +38,9 @@ export function toStandardBody( if (rawContentType === undefined && bodyHint === ('event-stream' satisfies StandardBodyHint)) { const eventStreamMessageQueue = new Queue() return { - resolveBody: async () => toAsyncIteratorObject(eventStreamMessageQueue, cleanup), + resolveBody: resolveStreamBodyOnce(() => + toAsyncIteratorObject(eventStreamMessageQueue, cleanup), + ), eventStreamMessageQueue, } } @@ -46,7 +48,7 @@ export function toStandardBody( if (rawContentType !== undefined) { const octetStreamMessageQueue = new Queue() return { - resolveBody: async () => toOctetStream(octetStreamMessageQueue, cleanup), + resolveBody: resolveStreamBodyOnce(() => toOctetStream(octetStreamMessageQueue, cleanup)), octetStreamMessageQueue, } } @@ -107,6 +109,24 @@ export function toStandardBody( return { resolveBody } } +/** + * A streamed body has a single queue, so a second consumer would split its messages + * with the first one. Like the fetch and node adapters, reject reading it more than once. + */ +function resolveStreamBodyOnce(resolve: () => StandardBody): () => Promise { + let resolved = false + + return async () => { + if (resolved) { + // native fetch error use TypeError + throw new TypeError('Failed to read body: body stream already read') + } + + resolved = true + return resolve() + } +} + export interface EncodedAtomicStandardBody { jsonBody: unknown headers: StandardHeaders diff --git a/packages/peer/src/client.test.ts b/packages/peer/src/client.test.ts index d9b1732..1baa5b9 100644 --- a/packages/peer/src/client.test.ts +++ b/packages/peer/src/client.test.ts @@ -797,6 +797,28 @@ describe('clientPeer', () => { await expect(iter.next()).resolves.toEqual({ done: false, value: 'evt1' }) await expect(iter.next()).rejects.toThrow('Server canceled the request') }) + + it('throws on resolving the event-stream response body multiple times', async () => { + const { id, promise } = await requestAndGetId() + await peer.message(makeStreamingResponse(id, 'event-stream')) + + const response = await promise + const iter = (await response.resolveBody()) as AsyncIterator + await expect(response.resolveBody()).rejects.toThrow( + new TypeError('Failed to read body: body stream already read'), + ) + + // the first iterator still receives every event + await peer.message(makeEventStreamMessage(id, 0)) + await peer.message(makeEventStreamMessage(id, 1)) + await peer.message(makeEventStreamMessage(id, 2, 'close')) + await expect(iter.next()).resolves.toEqual({ done: false, value: 0 }) + await expect(iter.next()).resolves.toEqual({ done: false, value: 1 }) + await expect(iter.next()).resolves.toEqual({ done: true, value: 2 }) + await expect(response.resolveBody()).rejects.toThrow( + new TypeError('Failed to read body: body stream already read'), + ) + }) }) }) @@ -1132,6 +1154,29 @@ describe('clientPeer', () => { await expect(reader.read()).rejects.toThrow('Server canceled the request') }) + + it('throws on resolving the octet-stream response body multiple times', async () => { + const { id, promise } = await requestAndGetId() + await peer.message(makeStreamingResponse(id, 'octet-stream')) + + const response = await promise + const body = (await response.resolveBody()) as ReadableStream + await expect(response.resolveBody()).rejects.toThrow( + new TypeError('Failed to read body: body stream already read'), + ) + + // the first stream still receives every chunk + await peer.message(makeOctetStreamMessage(id, false, new Uint8Array([1]))) + await peer.message(makeOctetStreamMessage(id, false, new Uint8Array([2]))) + await peer.message(makeOctetStreamMessage(id, true)) + const reader = body.getReader() + expect(await reader.read()).toEqual({ done: false, value: new Uint8Array([1]) }) + expect(await reader.read()).toEqual({ done: false, value: new Uint8Array([2]) }) + expect(await reader.read()).toEqual({ done: true, value: undefined }) + await expect(response.resolveBody()).rejects.toThrow( + new TypeError('Failed to read body: body stream already read'), + ) + }) }) }) diff --git a/packages/peer/src/server.test.ts b/packages/peer/src/server.test.ts index c383187..10e7d70 100644 --- a/packages/peer/src/server.test.ts +++ b/packages/peer/src/server.test.ts @@ -1,4 +1,5 @@ import type { StandardLazyRequest, StandardResponse } from '@standard-server/core' +import type { Queue } from '@standard-server/shared' import { AbortError, AsyncIteratorClass, @@ -407,6 +408,7 @@ describe('serverPeer', () => { await vi.waitFor(() => expect(handler).toHaveBeenCalled()) const request = handler.mock.calls[0]![0] const iter = (await request.resolveBody()) as AsyncIterator + const queue: Queue = (peer as any).requests.get('1').eventStreamMessageQueue await peer.message(makeEventStreamMessage('1', 'bye', 'close'), vi.fn()) await expect(iter.next()).resolves.toEqual({ value: 'bye', done: true }) @@ -414,6 +416,8 @@ describe('serverPeer', () => { // late messages are ignored instead of buffered forever await peer.message(makeEventStreamMessage('1', 'late'), vi.fn()) expect((peer as any).requests.get('1').eventStreamMessageQueue).toBeUndefined() + // the dropped queue is aborted, so nothing can be left waiting on it + await expect(queue.pull()).rejects.toThrow('Queue was aborted.') box.resolve(jsonResponse()) await promise @@ -421,6 +425,34 @@ describe('serverPeer', () => { // the fully consumed body does not trigger a stream/cancel on close expect(send.mock.calls.map(([m]) => m.kind)).toEqual(['response']) }) + + it('throws on resolving the event-stream request body multiple times', async () => { + const { handler, box } = deferredHandler() + const msg = makeRequestMessage({ headers: { 'standard-server': 'event-stream' } }) + const promise = peer.message(msg, handler) + + await vi.waitFor(() => expect(handler).toHaveBeenCalled()) + const request = handler.mock.calls[0]![0] + const iter = (await request.resolveBody()) as AsyncIterator + await expect(request.resolveBody()).rejects.toThrow( + new TypeError('Failed to read body: body stream already read'), + ) + + // the first iterator still receives every event + await peer.message(makeEventStreamMessage('1', 0), vi.fn()) + await peer.message(makeEventStreamMessage('1', 1), vi.fn()) + await peer.message(makeEventStreamMessage('1', 2, 'close'), vi.fn()) + await expect(iter.next()).resolves.toEqual({ value: 0, done: false }) + await expect(iter.next()).resolves.toEqual({ value: 1, done: false }) + await expect(iter.next()).resolves.toEqual({ value: 2, done: true }) + await expect(request.resolveBody()).rejects.toThrow( + new TypeError('Failed to read body: body stream already read'), + ) + + box.resolve(jsonResponse()) + await promise + expect(send.mock.calls.map(([m]) => m.kind)).toEqual(['response']) + }) }) describe('response body (outgoing)', () => { @@ -668,6 +700,73 @@ describe('serverPeer', () => { vi.fn(), ) }) + + it('stops accepting octet-stream messages once the body is fully consumed', async () => { + const { handler, box } = deferredHandler() + const msg = makeRequestMessage({ + headers: { + 'standard-server': 'octet-stream', + 'content-type': 'application/octet-stream', + }, + }) + const promise = peer.message(msg, handler) + + await vi.waitFor(() => expect(handler).toHaveBeenCalled()) + const request = handler.mock.calls[0]![0] + const body = (await request.resolveBody()) as ReadableStream + const queue: Queue = (peer as any).requests.get('1').octetStreamMessageQueue + + await peer.message(makeOctetStreamMessage('1', true, new Uint8Array([1])), vi.fn()) + const reader = body.getReader() + expect(await reader.read()).toEqual({ value: new Uint8Array([1]), done: false }) + expect(await reader.read()).toEqual({ value: undefined, done: true }) + + // late messages are ignored instead of buffered forever + await peer.message(makeOctetStreamMessage('1', true, new Uint8Array([2])), vi.fn()) + expect((peer as any).requests.get('1').octetStreamMessageQueue).toBeUndefined() + // the dropped queue is aborted, so nothing can be left waiting on it + await expect(queue.pull()).rejects.toThrow('Queue was aborted.') + + box.resolve(jsonResponse()) + await promise + + // the fully consumed body does not trigger a stream/cancel on close + expect(send.mock.calls.map(([m]) => m.kind)).toEqual(['response']) + }) + + it('throws on resolving the octet-stream request body multiple times', async () => { + const { handler, box } = deferredHandler() + const msg = makeRequestMessage({ + headers: { + 'standard-server': 'octet-stream', + 'content-type': 'application/octet-stream', + }, + }) + const promise = peer.message(msg, handler) + + await vi.waitFor(() => expect(handler).toHaveBeenCalled()) + const request = handler.mock.calls[0]![0] + const body = (await request.resolveBody()) as ReadableStream + await expect(request.resolveBody()).rejects.toThrow( + new TypeError('Failed to read body: body stream already read'), + ) + + // the first stream still receives every chunk + await peer.message(makeOctetStreamMessage('1', false, new Uint8Array([1])), vi.fn()) + await peer.message(makeOctetStreamMessage('1', false, new Uint8Array([2])), vi.fn()) + await peer.message(makeOctetStreamMessage('1', true), vi.fn()) + const reader = body.getReader() + expect(await reader.read()).toEqual({ value: new Uint8Array([1]), done: false }) + expect(await reader.read()).toEqual({ value: new Uint8Array([2]), done: false }) + expect(await reader.read()).toEqual({ value: undefined, done: true }) + await expect(request.resolveBody()).rejects.toThrow( + new TypeError('Failed to read body: body stream already read'), + ) + + box.resolve(jsonResponse()) + await promise + expect(send.mock.calls.map(([m]) => m.kind)).toEqual(['response']) + }) }) describe('response body (outgoing)', () => { diff --git a/packages/peer/src/server.ts b/packages/peer/src/server.ts index ea7acfa..648217d 100644 --- a/packages/peer/src/server.ts +++ b/packages/peer/src/server.ts @@ -68,17 +68,14 @@ export class ServerPeer { const decoded = toStandardBody(message, async ({ kind, error }) => { /** * The request body is finished (fully read, errored, or cancelled). - * Drop the queues so late stream messages are ignored instead of - * buffered forever. + * Abort the queues so nothing is left waiting on them, then drop them + * so late stream messages are ignored instead of buffered forever. */ const streamActive = state.eventStreamMessageQueue !== undefined || state.octetStreamMessageQueue !== undefined - if (kind === 'cancelled' && streamActive) { - state.eventStreamMessageQueue?.abort(error) - state.octetStreamMessageQueue?.abort(error) - } - + state.eventStreamMessageQueue?.abort(error) + state.octetStreamMessageQueue?.abort(error) state.eventStreamMessageQueue = undefined state.octetStreamMessageQueue = undefined