From 5120d2d99a3a515210374cf6ca6c9613097b77a1 Mon Sep 17 00:00:00 2001 From: Claude Date: Sat, 3 Oct 2026 12:16:21 +0000 Subject: [PATCH 1/3] fix(peer): cancel the server request when the client fails after sending it An exception thrown in transmitRequest after the request message went out (e.g. a request body ReadableStream the caller already locked, which makes the OctetStreamTransmitter constructor throw) was handled with closeById, which sends no cancel. The server kept the request open, its handler blocked on the body, until the peer closed. Once the request message is sent, failures now go through abortById so the server receives a cancel. A failed cancel delivery is swallowed so it cannot surface as an unhandled rejection. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_019Pu97o66FJwoRA8ZmYjFWo --- packages/peer/src/client.test.ts | 35 ++++++++++++++++++++++++++++++++ packages/peer/src/client.ts | 9 +++++++- packages/peer/tests/peer.test.ts | 22 ++++++++++++++++++++ 3 files changed, 65 insertions(+), 1 deletion(-) diff --git a/packages/peer/src/client.test.ts b/packages/peer/src/client.test.ts index c64e44c..d9b1732 100644 --- a/packages/peer/src/client.test.ts +++ b/packages/peer/src/client.test.ts @@ -911,6 +911,41 @@ describe('clientPeer', () => { await promise }) + it('sends cancel message and rejects request when the body cannot be read after the request is sent', async () => { + const stream = new ReadableStream() + stream.getReader() + + await expect( + peer.request(makeRequest({ method: 'POST', headers: {}, body: stream })), + ).rejects.toThrow(TypeError) + + // the server already received the request, so it must be told to drop it + const id = (send.mock.calls[0]![0] as PeerRequestMessage).id + expect(send.mock.calls.map(([m]) => m)).toEqual([ + expect.objectContaining({ id, kind: 'request' }), + { id, kind: 'cancel' }, + ]) + }) + + it('silently ignores transport failures when cancelling after the body cannot be read', async () => { + send.mockImplementation(async (message) => { + if (message.kind === 'cancel') { + throw new Error('transport down') + } + }) + + const stream = new ReadableStream() + stream.getReader() + + await expect( + peer.request(makeRequest({ method: 'POST', headers: {}, body: stream })), + ).rejects.toThrow(TypeError) + + // let the failed cancel delivery settle; it must not surface anywhere + await sleep(1) + expect(send.mock.calls.map(([m]) => m.kind)).toEqual(['request', 'cancel']) + }) + it('stops transmitting the octet-stream request body when a full response arrives', async () => { const cancel = vi.fn() const stream = new ReadableStream({ diff --git a/packages/peer/src/client.ts b/packages/peer/src/client.ts index 96bafde..7014dd7 100644 --- a/packages/peer/src/client.ts +++ b/packages/peer/src/client.ts @@ -141,7 +141,14 @@ export class ClientPeer { } } catch (reason) { failure = reason - await this.closeById(id, reason) + + if (state.requestSent) { + // the server already holds the request, so it must be cancelled there as well; + // a failed cancel delivery must not surface as an unhandled rejection + await this.abortById(id, reason).catch(() => {}) + } else { + await this.closeById(id, reason) + } } finally { if (untransmittedBody !== undefined) { await cancelStandardBody(untransmittedBody, failure ?? request.signal?.reason).catch( diff --git a/packages/peer/tests/peer.test.ts b/packages/peer/tests/peer.test.ts index d4f1ef7..c7d4f78 100644 --- a/packages/peer/tests/peer.test.ts +++ b/packages/peer/tests/peer.test.ts @@ -380,6 +380,28 @@ describe('peer integration (client <-> server over encoded wire)', () => { await vi.waitFor(() => expect(serverSignal!.aborted).toBe(true)) }) + it('releases the server request when the client fails to stream a body after sending the request', async () => { + let serverSignal: AbortSignal | undefined + + const { client, server } = connect(async (request) => { + serverSignal = request.signal + // blocks until the upload ends or the request is cancelled + await ((await request.resolveBody()) as ReadableStream).getReader().read() + return { status: 200, headers: {} } + }) + + // a body the caller already locked cannot be streamed once the request message is out + const body = new ReadableStream() + body.getReader() + + await expect( + client.request({ url: '/upload', method: 'POST', headers: {}, body }), + ).rejects.toThrow(TypeError) + + await vi.waitFor(() => expect(serverSignal?.aborted).toBe(true)) + expect((server as any).requests.size).toBe(0) + }) + it('propagates a client abort fired while the request message is still being sent', async () => { const encodeStarted = promiseWithResolvers() const releaseEncode = promiseWithResolvers() From 5675adf4ec16a8cf02134ca733a73f25c0c42f8a Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 4 Oct 2026 03:15:40 +0000 Subject: [PATCH 2/3] refactor(peer): handle request body transmit failures in one place The outer catch in transmitRequest already cancels the server request once the request message is sent, so the per-transmitter .catch blocks are redundant. Their guard (skip the abort when the transmitter was cleared) moves to the outer catch as a streamCancelled check: a stream/cancel from the server is the only way a transmitter is cleared while the request stays open, so a failing transmitter after it is expected and the request keeps waiting for its response. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_019Pu97o66FJwoRA8ZmYjFWo --- packages/peer/src/client.ts | 20 +++++++------------- 1 file changed, 7 insertions(+), 13 deletions(-) diff --git a/packages/peer/src/client.ts b/packages/peer/src/client.ts index 7014dd7..149be06 100644 --- a/packages/peer/src/client.ts +++ b/packages/peer/src/client.ts @@ -125,30 +125,24 @@ export class ClientPeer { if (isAsyncIteratorObject(request.body)) { const transmitter = new EventStreamTransmitter(request.body, id, this.send) state.eventStreamTransmitter = transmitter - await transmitter.transmit().catch((error) => { - if (state.eventStreamTransmitter) { - return this.abortById(id, error) - } - }) + await transmitter.transmit() } else if (request.body instanceof ReadableStream) { const transmitter = new OctetStreamTransmitter(request.body, id, this.send) state.octetStreamTransmitter = transmitter - await transmitter.transmit().catch((error) => { - if (state.octetStreamTransmitter) { - return this.abortById(id, error) - } - }) + await transmitter.transmit() } } catch (reason) { failure = reason - if (state.requestSent) { + if (!state.requestSent) { + await this.closeById(id, reason) + } else if (!state.streamCancelled) { // the server already holds the request, so it must be cancelled there as well; // a failed cancel delivery must not surface as an unhandled rejection await this.abortById(id, reason).catch(() => {}) - } else { - await this.closeById(id, reason) } + // otherwise the server stopped the upload itself, so a failing transmitter is expected + // and the request stays open for its response } finally { if (untransmittedBody !== undefined) { await cancelStandardBody(untransmittedBody, failure ?? request.signal?.reason).catch( From 7151006c250c66051c55d74fdd1d48765d1e53d7 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 4 Oct 2026 03:20:35 +0000 Subject: [PATCH 3/3] refactor(peer): swallow transmitRequest failures at its call site transmitRequest is fire-and-forget and settles the request itself before anything can fail, so a rejection from it (e.g. a failed cancel delivery) has no observer. Catch it once where it is started instead of inside the outer catch, and drop comments whose behavior the tests already cover. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_019Pu97o66FJwoRA8ZmYjFWo --- packages/peer/src/client.ts | 8 ++------ 1 file changed, 2 insertions(+), 6 deletions(-) diff --git a/packages/peer/src/client.ts b/packages/peer/src/client.ts index 149be06..5a42a33 100644 --- a/packages/peer/src/client.ts +++ b/packages/peer/src/client.ts @@ -66,7 +66,7 @@ export class ClientPeer { state.removeAbortListener = () => signal.removeEventListener('abort', abortListener) } - void this.transmitRequest(id, state, request) + void this.transmitRequest(id, state, request).catch(() => {}) }) } @@ -137,12 +137,8 @@ export class ClientPeer { if (!state.requestSent) { await this.closeById(id, reason) } else if (!state.streamCancelled) { - // the server already holds the request, so it must be cancelled there as well; - // a failed cancel delivery must not surface as an unhandled rejection - await this.abortById(id, reason).catch(() => {}) + await this.abortById(id, reason) } - // otherwise the server stopped the upload itself, so a failing transmitter is expected - // and the request stays open for its response } finally { if (untransmittedBody !== undefined) { await cancelStandardBody(untransmittedBody, failure ?? request.signal?.reason).catch(