Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 22 additions & 2 deletions packages/peer/src/body.ts
Original file line number Diff line number Diff line change
Expand Up @@ -38,15 +38,17 @@ export function toStandardBody(
if (rawContentType === undefined && bodyHint === ('event-stream' satisfies StandardBodyHint)) {
const eventStreamMessageQueue = new Queue<PeerEventStreamMessage>()
return {
resolveBody: async () => toAsyncIteratorObject(eventStreamMessageQueue, cleanup),
resolveBody: resolveStreamBodyOnce(() =>
toAsyncIteratorObject(eventStreamMessageQueue, cleanup),
),
eventStreamMessageQueue,
}
}

if (rawContentType !== undefined) {
const octetStreamMessageQueue = new Queue<PeerOctetStreamMessage>()
return {
resolveBody: async () => toOctetStream(octetStreamMessageQueue, cleanup),
resolveBody: resolveStreamBodyOnce(() => toOctetStream(octetStreamMessageQueue, cleanup)),
octetStreamMessageQueue,
}
}
Expand Down Expand Up @@ -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<StandardBody> {
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
Expand Down
45 changes: 45 additions & 0 deletions packages/peer/src/client.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<unknown>
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'),
)
})
})
})

Expand Down Expand Up @@ -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<Uint8Array>
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'),
)
})
})
})

Expand Down
99 changes: 99 additions & 0 deletions packages/peer/src/server.test.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import type { StandardLazyRequest, StandardResponse } from '@standard-server/core'
import type { Queue } from '@standard-server/shared'
import {
AbortError,
AsyncIteratorClass,
Expand Down Expand Up @@ -407,20 +408,51 @@ describe('serverPeer', () => {
await vi.waitFor(() => expect(handler).toHaveBeenCalled())
const request = handler.mock.calls[0]![0]
const iter = (await request.resolveBody()) as AsyncIterator<unknown>
const queue: Queue<unknown> = (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 })

// 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

// 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<unknown>
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)', () => {
Expand Down Expand Up @@ -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<Uint8Array>
const queue: Queue<unknown> = (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<Uint8Array>
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)', () => {
Expand Down
11 changes: 4 additions & 7 deletions packages/peer/src/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
Loading