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
35 changes: 35 additions & 0 deletions packages/peer/src/client.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<Uint8Array>()
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<Uint8Array>()
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<Uint8Array>({
Expand Down
21 changes: 9 additions & 12 deletions packages/peer/src/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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(() => {})
})
}

Expand Down Expand Up @@ -125,23 +125,20 @@ 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
await this.closeById(id, reason)

if (!state.requestSent) {
await this.closeById(id, reason)
} else if (!state.streamCancelled) {
await this.abortById(id, reason)
}
} finally {
if (untransmittedBody !== undefined) {
await cancelStandardBody(untransmittedBody, failure ?? request.signal?.reason).catch(
Expand Down
22 changes: 22 additions & 0 deletions packages/peer/tests/peer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<Uint8Array>()
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<void>()
const releaseEncode = promiseWithResolvers<void>()
Expand Down
Loading