From 404e484785dbb447d35fe07501911d0082893e87 Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 7 Oct 2026 03:07:09 +0000 Subject: [PATCH 1/4] fix(peer): reject resolving a streamed body more than once A streamed peer body (event or octet stream) has a single message queue, but every `resolveBody()` call built a new consumer of it. Two consumers split the messages between them, and on the server the second one could hang forever: the success cleanup dropped the queue without closing it, so its pending pull never settled. - Throw `TypeError('Failed to read body: body stream already read')` on a second `resolveBody()` for streamed bodies, matching the fetch and node adapters. - Close the request body queues in the server's success/error cleanup instead of only unsetting them, so nothing is left waiting on them. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_012g4ovPb64wru9czTYsrwGe --- packages/peer/README.md | 2 + packages/peer/src/body.ts | 24 +++++++- packages/peer/src/client.test.ts | 45 +++++++++++++++ packages/peer/src/server.test.ts | 99 ++++++++++++++++++++++++++++++++ packages/peer/src/server.ts | 9 ++- 5 files changed, 174 insertions(+), 5 deletions(-) diff --git a/packages/peer/README.md b/packages/peer/README.md index 248d187..9e96e57 100644 --- a/packages/peer/README.md +++ b/packages/peer/README.md @@ -128,6 +128,8 @@ const payload = await response.resolveBody() Unlike the HTTP adapters, `resolveBody(hint?)` ignores the `hint` argument in this adapter. HTTP adapters receive the body as a raw byte stream and must decide how to parse it, so a hint can steer that decision. The peer protocol instead encodes the body in structured form at send time: JSON values travel as JSON, binary payloads travel as binary, event and octet streams flow as dedicated stream messages, and markers in the message distinguish the ambiguous cases such as `form-data` vs. `file`. By the time a message arrives, there are no raw bytes left to reinterpret — the body always resolves to exactly the representation the sender had, so a hint has nothing to override. +As with the HTTP adapters, a streamed body (an event or octet stream) can be resolved only once: calling `resolveBody()` again throws a `TypeError`, because the stream messages can only be delivered to a single consumer. + ## Codec helpers Use `encodePeerMessage()` and `decodePeerMessage()` to bridge between the peer protocol and your underlying transport. 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..a89041a 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 closed, so nothing can be left waiting on it + await expect(queue.pull()).rejects.toThrow('Queue was closed.') 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 closed, so nothing can be left waiting on it + await expect(queue.pull()).rejects.toThrow('Queue was closed.') + + 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..668507c 100644 --- a/packages/peer/src/server.ts +++ b/packages/peer/src/server.ts @@ -68,15 +68,18 @@ 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. + * Close 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) { + if (kind === 'cancelled') { state.eventStreamMessageQueue?.abort(error) state.octetStreamMessageQueue?.abort(error) + } else { + state.eventStreamMessageQueue?.close(error) + state.octetStreamMessageQueue?.close(error) } state.eventStreamMessageQueue = undefined From 52c889bf91543d5862c0451c8d20a93c816993e2 Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 7 Oct 2026 03:13:06 +0000 Subject: [PATCH 2/4] refactor(peer): always abort request body queues in server cleanup Once the request body is finished, its single consumer is done, so closing and aborting the queue behave the same. Abort unconditionally instead of branching on the cleanup kind; this also releases late buffered messages. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_012g4ovPb64wru9czTYsrwGe --- packages/peer/src/server.test.ts | 4 ++-- packages/peer/src/server.ts | 12 +++--------- 2 files changed, 5 insertions(+), 11 deletions(-) diff --git a/packages/peer/src/server.test.ts b/packages/peer/src/server.test.ts index a89041a..347b9b3 100644 --- a/packages/peer/src/server.test.ts +++ b/packages/peer/src/server.test.ts @@ -417,7 +417,7 @@ describe('serverPeer', () => { await peer.message(makeEventStreamMessage('1', 'late'), vi.fn()) expect((peer as any).requests.get('1').eventStreamMessageQueue).toBeUndefined() // the dropped queue is closed, so nothing can be left waiting on it - await expect(queue.pull()).rejects.toThrow('Queue was closed.') + await expect(queue.pull()).rejects.toThrow('Queue was aborted.') box.resolve(jsonResponse()) await promise @@ -725,7 +725,7 @@ describe('serverPeer', () => { await peer.message(makeOctetStreamMessage('1', true, new Uint8Array([2])), vi.fn()) expect((peer as any).requests.get('1').octetStreamMessageQueue).toBeUndefined() // the dropped queue is closed, so nothing can be left waiting on it - await expect(queue.pull()).rejects.toThrow('Queue was closed.') + await expect(queue.pull()).rejects.toThrow('Queue was aborted.') box.resolve(jsonResponse()) await promise diff --git a/packages/peer/src/server.ts b/packages/peer/src/server.ts index 668507c..648217d 100644 --- a/packages/peer/src/server.ts +++ b/packages/peer/src/server.ts @@ -68,20 +68,14 @@ export class ServerPeer { const decoded = toStandardBody(message, async ({ kind, error }) => { /** * The request body is finished (fully read, errored, or cancelled). - * Close the queues so nothing is left waiting on them, then drop them + * 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') { - state.eventStreamMessageQueue?.abort(error) - state.octetStreamMessageQueue?.abort(error) - } else { - state.eventStreamMessageQueue?.close(error) - state.octetStreamMessageQueue?.close(error) - } - + state.eventStreamMessageQueue?.abort(error) + state.octetStreamMessageQueue?.abort(error) state.eventStreamMessageQueue = undefined state.octetStreamMessageQueue = undefined From c67569bfc806555412704edfdea21355a956ba57 Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 7 Oct 2026 03:33:33 +0000 Subject: [PATCH 3/4] test(peer): say the dropped queue is aborted, not closed The server cleanup now aborts the request body queues, and close and abort are distinct Queue operations, so the test comments should match. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_012g4ovPb64wru9czTYsrwGe --- packages/peer/src/server.test.ts | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/packages/peer/src/server.test.ts b/packages/peer/src/server.test.ts index 347b9b3..10e7d70 100644 --- a/packages/peer/src/server.test.ts +++ b/packages/peer/src/server.test.ts @@ -416,7 +416,7 @@ 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 closed, so nothing can be left waiting on it + // 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()) @@ -724,7 +724,7 @@ describe('serverPeer', () => { // 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 closed, so nothing can be left waiting on it + // 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()) From 1b59e607f27f1a440a7fc2114dd7f88d0eb326a5 Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 8 Oct 2026 01:54:33 +0000 Subject: [PATCH 4/4] docs(peer): drop the read-once note from the README Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_012g4ovPb64wru9czTYsrwGe --- packages/peer/README.md | 2 -- 1 file changed, 2 deletions(-) diff --git a/packages/peer/README.md b/packages/peer/README.md index 9e96e57..248d187 100644 --- a/packages/peer/README.md +++ b/packages/peer/README.md @@ -128,8 +128,6 @@ const payload = await response.resolveBody() Unlike the HTTP adapters, `resolveBody(hint?)` ignores the `hint` argument in this adapter. HTTP adapters receive the body as a raw byte stream and must decide how to parse it, so a hint can steer that decision. The peer protocol instead encodes the body in structured form at send time: JSON values travel as JSON, binary payloads travel as binary, event and octet streams flow as dedicated stream messages, and markers in the message distinguish the ambiguous cases such as `form-data` vs. `file`. By the time a message arrives, there are no raw bytes left to reinterpret — the body always resolves to exactly the representation the sender had, so a hint has nothing to override. -As with the HTTP adapters, a streamed body (an event or octet stream) can be resolved only once: calling `resolveBody()` again throws a `TypeError`, because the stream messages can only be delivered to a single consumer. - ## Codec helpers Use `encodePeerMessage()` and `decodePeerMessage()` to bridge between the peer protocol and your underlying transport.