From 04eee2ef11d3234e30f5a807b1da12189c2f1bac Mon Sep 17 00:00:00 2001 From: Artem Zakharchenko Date: Thu, 7 May 2026 21:57:15 +0200 Subject: [PATCH 1/8] fix(sse): invoke `finalize` on response stream end --- src/core/handlers/RequestHandler.ts | 2 +- src/core/sse.ts | 135 +++++++++++++++-------- src/mockServiceWorker.js | 6 +- test/browser/sse-api/sse.finally.test.ts | 127 +++++++++++++++++++++ 4 files changed, 220 insertions(+), 50 deletions(-) create mode 100644 test/browser/sse-api/sse.finally.test.ts diff --git a/src/core/handlers/RequestHandler.ts b/src/core/handlers/RequestHandler.ts index f52eb0a47..d8580abbe 100644 --- a/src/core/handlers/RequestHandler.ts +++ b/src/core/handlers/RequestHandler.ts @@ -505,7 +505,7 @@ export abstract class RequestHandler< this.scheduledCleanups.set(requestId, cleanups) } - private async exhaustCleanups( + protected async exhaustCleanups( cleanups: Array<() => MaybePromise>, ): Promise { const errors: Array = [] diff --git a/src/core/sse.ts b/src/core/sse.ts index a07b378f2..c012f0811 100644 --- a/src/core/sse.ts +++ b/src/core/sse.ts @@ -1,6 +1,6 @@ import { invariant } from 'outvariant' import { DeferredPromise } from '@open-draft/deferred-promise' -import { Emitter } from 'strict-event-emitter' +import { Emitter, TypedEvent } from 'rettime' import type { ResponseResolver } from './handlers/RequestHandler' import { HttpHandler, @@ -14,6 +14,7 @@ import { getTimestamp } from './utils/logging/getTimestamp' import { devUtils } from './utils/internal/devUtils' import { colors } from './ws/utils/attachWebSocketLogger' import { toPublicUrl } from './utils/request/toPublicUrl' +import type { MaybePromise } from './typeUtils' type EventMapConstraint = { message?: unknown @@ -140,7 +141,14 @@ class ServerSentEventHandler< return matches } - async log(_args: { request: Request; response: Response }): Promise { + async log(args: { request: Request; response: Response }): Promise { + /** + * @note Cancel the response stream because it's not needed for logging. + * Otherwise, this cloned response remains unconsumed and its original + * doesn't propagate stream cancelations at all. + */ + args.response.body?.cancel() + /** * @note Skip the default `this.log()` logic so that when this handler is logged * upon handling the request, nothing is printed (we log SSE requests early). @@ -155,16 +163,14 @@ class ServerSentEventHandler< const publicUrl = toPublicUrl(request.url) /* eslint-disable no-console */ - emitter.on('message', (payload) => { + emitter.on('message', ({ data }) => { console.groupCollapsed( - devUtils.formatMessage( - `${getTimestamp()} SSE %s %c⇣%c ${payload.event}`, - ), + devUtils.formatMessage(`${getTimestamp()} SSE %s %c⇣%c ${data.event}`), publicUrl, `color:${colors.mocked}`, 'color:inherit', ) - console.log(payload.frames) + console.log(data.frames) console.groupEnd() }) @@ -191,6 +197,18 @@ class ServerSentEventHandler< }) /* eslint-enable no-console */ } + + protected async exhaustCleanups( + cleanups: Array<() => MaybePromise>, + ): Promise { + const onClose = () => { + this.#emitter.removeListener('error', onClose) + this.#emitter.removeListener('close', onClose) + void super.exhaustCleanups(cleanups) + } + + this.#emitter.once('error', onClose).once('close', onClose) + } } type Values = T[keyof T] @@ -216,9 +234,9 @@ type ToEventDiscriminatedUnion = Values<{ }> type ServerSentEventClientEventMap = { - message: [payload: EventStreamMessage] - error: [] - close: [] + message: TypedEvent + error: TypedEvent + close: TypedEvent } const kClientEmitter = Symbol.for('kClientEmitter') @@ -229,11 +247,13 @@ class ServerSentEventClient< private [kClientEmitter]?: Emitter #encoder: TextEncoder - #writer: WritableStreamDefaultWriter + #controller: ReadableStreamDefaultController + #closed: DeferredPromise - constructor(writable: WritableStream) { + constructor(controller: ReadableStreamDefaultController) { this.#encoder = new TextEncoder() - this.#writer = writable.getWriter() + this.#controller = controller + this.#closed = new DeferredPromise() } /** @@ -289,37 +309,50 @@ class ServerSentEventClient< * error. */ public error(): void { - this.#writer.abort().catch((error) => { - console.error(error) - devUtils.error( - 'Failed to abort server-side EventSource. Please see the original error above.', - ) - }) - this[kClientEmitter]?.emit('error') + if (this.#closed.state !== 'pending') { + return + } + + this.#controller.error() + this.#closed.resolve() + this[kClientEmitter]?.emit(new TypedEvent('error')) } /** * Closes the underlying `EventSource`, closing the connection. */ public close(): void { - this.#writer.close().catch((error) => { + if (this.#closed.state !== 'pending') { + return + } + + try { + this.#controller.close() + this.#closed.resolve() + } catch { + // + } + + this[kClientEmitter]?.emit(new TypedEvent('close')) + } + + #enqueue(chunk: Uint8Array): void { + if (this.#closed.state !== 'pending') { + return + } + + try { + this.#controller.enqueue(chunk) + } catch (error) { console.error(error) devUtils.error( - 'Failed to close server-side EventSource. Please see the original error above.', + 'Failed to write to server-side EventSource. Please see the original error above.', ) - }) - this[kClientEmitter]?.emit('close') + } } #sendRetry(retry: number): void { - this.#writer - .write(this.#encoder.encode(`retry:${retry}\n\n`)) - .catch((error) => { - console.error(error) - devUtils.error( - 'Failed to send a retry packet to server-side EventSource. Please see the original error above.', - ) - }) + this.#enqueue(this.#encoder.encode(`retry:${retry}\n\n`)) } #sendMessage(message: { @@ -350,21 +383,18 @@ class ServerSentEventClient< frames.push('', '') - this.#writer - .write(this.#encoder.encode(frames.join('\n'))) - .catch((error) => { - console.error(error) - devUtils.error( - 'Failed to send a message to server-side EventSource. Please see the original error above.', - ) - }) + this.#enqueue(this.#encoder.encode(frames.join('\n'))) - this[kClientEmitter]?.emit('message', { - id: message.id, - event: message.event?.toString() || 'message', - data: message.data, - frames, - }) + this[kClientEmitter]?.emit( + new TypedEvent('message', { + data: { + id: message.id, + event: message.event?.toString() || 'message', + data: message.data, + frames, + }, + }), + ) } } @@ -991,9 +1021,18 @@ function createEventStream( request.url, ) - const { readable, writable } = new TransformStream() + let controller!: ReadableStreamDefaultController + + const readable = new ReadableStream({ + start(defaultController) { + controller = defaultController + }, + cancel() { + client.close() + }, + }) - const client = new ServerSentEventClient(writable) + const client = new ServerSentEventClient(controller) const server = new ServerSentEventServer({ request, client, diff --git a/src/mockServiceWorker.js b/src/mockServiceWorker.js index d02b80689..d2d616db6 100644 --- a/src/mockServiceWorker.js +++ b/src/mockServiceWorker.js @@ -134,7 +134,11 @@ async function handleRequest(event, requestId, requestInterceptedAt) { // Send back the response clone for the "response:*" life-cycle events. // Ensure MSW is active and ready to handle the message, otherwise // this message will pend indefinitely. - if (client && activeClientIds.has(client.id)) { + if ( + client && + activeClientIds.has(client.id) && + response.headers.get('content-type') !== 'text/event-stream' + ) { const serializedRequest = await serializeRequest(requestCloneForEvents) // Clone the response so both the client and the library could consume it. diff --git a/test/browser/sse-api/sse.finally.test.ts b/test/browser/sse-api/sse.finally.test.ts new file mode 100644 index 000000000..c45e8c992 --- /dev/null +++ b/test/browser/sse-api/sse.finally.test.ts @@ -0,0 +1,127 @@ +import type { sse } from 'msw' +import type { setupWorker } from 'msw/browser' +import { DeferredPromise } from '@open-draft/deferred-promise' +import { test, expect } from '../playwright.extend' + +declare namespace window { + export const msw: { + setupWorker: typeof setupWorker + sse: typeof sse + } +} + +const EXAMPLE_URL = new URL('./sse.mocks.ts', import.meta.url) + +test('runs cleanup after the event source is closed by the client', async ({ + loadExample, + page, +}) => { + await loadExample(EXAMPLE_URL, { + skipActivation: true, + }) + + const finalizedAt = new DeferredPromise() + await page.exposeFunction('notifyFinalized', () => { + finalizedAt.resolve(Date.now()) + }) + + await page.evaluate(async () => { + const { setupWorker, sse } = window.msw + + const worker = setupWorker( + sse('http://localhost/stream', ({ finalize }) => { + finalize(() => window.notifyFinalized()) + }), + ) + await worker.start() + }) + + const closedAt = await page.evaluate(() => { + const source = new EventSource('http://localhost/stream') + + return new Promise((resolve) => { + source.addEventListener('open', () => { + source.close() + resolve(Date.now()) + }) + }) + }) + + await expect(finalizedAt).resolves.toBeGreaterThanOrEqual(closedAt) +}) + +test('runs cleanup after the event source is closed by the handler', async ({ + loadExample, + page, +}) => { + await loadExample(EXAMPLE_URL, { + skipActivation: true, + }) + + const finalizedAt = new DeferredPromise() + await page.exposeFunction('notifyFinalized', () => { + finalizedAt.resolve(Date.now()) + }) + + await page.evaluate(async () => { + const { setupWorker, sse } = window.msw + + const worker = setupWorker( + sse('http://localhost/stream', ({ client, finalize }) => { + setTimeout(() => client.close(), 250) + finalize(() => window.notifyFinalized()) + }), + ) + await worker.start() + }) + + const closedAt = await page.evaluate(() => { + const source = new EventSource('http://localhost/stream') + + return new Promise((resolve) => { + source.addEventListener('open', () => { + resolve(Date.now()) + }) + }) + }) + + await expect(finalizedAt).resolves.toBeGreaterThanOrEqual(closedAt) +}) + +test('runs cleanup after the event source is errored by the handler', async ({ + loadExample, + page, +}) => { + await loadExample(EXAMPLE_URL, { + skipActivation: true, + }) + + const finalizedAt = new DeferredPromise() + await page.exposeFunction('notifyFinalized', () => { + finalizedAt.resolve(Date.now()) + }) + + await page.evaluate(async () => { + const { setupWorker, sse } = window.msw + + const worker = setupWorker( + sse('http://localhost/stream', ({ client, finalize }) => { + setTimeout(() => client.error(), 250) + finalize(() => window.notifyFinalized()) + }), + ) + await worker.start() + }) + + const closedAt = await page.evaluate(() => { + const source = new EventSource('http://localhost/stream') + + return new Promise((resolve) => { + source.addEventListener('open', () => { + resolve(Date.now()) + }) + }) + }) + + await expect(finalizedAt).resolves.toBeGreaterThanOrEqual(closedAt) +}) From 82e48799292174a0b679bceda14921bcb4ecc062 Mon Sep 17 00:00:00 2001 From: Artem Zakharchenko Date: Tue, 12 May 2026 12:10:47 +0200 Subject: [PATCH 2/8] fix: permissive `text/event-stream` check --- src/mockServiceWorker.js | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/src/mockServiceWorker.js b/src/mockServiceWorker.js index d2d616db6..0772b04f6 100644 --- a/src/mockServiceWorker.js +++ b/src/mockServiceWorker.js @@ -137,7 +137,10 @@ async function handleRequest(event, requestId, requestInterceptedAt) { if ( client && activeClientIds.has(client.id) && - response.headers.get('content-type') !== 'text/event-stream' + !response.headers + .get('content-type') + ?.toLowerCase() + .startsWith('text/event-stream') ) { const serializedRequest = await serializeRequest(requestCloneForEvents) From f9138a9a82ff2604fac6764830cc7ed1f26cd757 Mon Sep 17 00:00:00 2001 From: Artem Zakharchenko Date: Tue, 7 Jul 2026 18:15:25 +0200 Subject: [PATCH 3/8] fix(RequestHandler): implement `finalize` --- src/core/handlers/RequestHandler.ts | 122 ++++++++++++++---- src/core/sse.ts | 13 -- src/core/utils/HttpResponse/decorators.ts | 25 ++++ .../observe-response-body-stream.test.ts | 64 +++++++++ .../internal/observe-response-body-stream.ts | 72 +++++++++++ 5 files changed, 260 insertions(+), 36 deletions(-) create mode 100644 src/core/utils/internal/observe-response-body-stream.test.ts create mode 100644 src/core/utils/internal/observe-response-body-stream.ts diff --git a/src/core/handlers/RequestHandler.ts b/src/core/handlers/RequestHandler.ts index d8580abbe..969e360bd 100644 --- a/src/core/handlers/RequestHandler.ts +++ b/src/core/handlers/RequestHandler.ts @@ -15,6 +15,7 @@ import { import type { GraphQLRequestBody } from './GraphQLHandler' import { devUtils } from '../utils/internal/devUtils' import { getRawSetCookie } from '../utils/HttpResponse/decorators' +import { observeResponseBodyStream } from '../utils/internal/observe-response-body-stream' export type DefaultRequestMultipartBody = Record< string, @@ -31,12 +32,7 @@ export type DefaultBodyType = | undefined export type JsonBodyType = - | Record - | string - | number - | boolean - | null - | undefined + Record | string | number | boolean | null | undefined export interface RequestHandlerDefaultInfo { header: string @@ -94,6 +90,10 @@ export type ResponseResolverInfo< /** * Schedule a callback to run after this response resolver completes. * Handy for cleaning up the side effects introduced in the resolver. + * + * For responses with a `ReadableStream` body (including `sse()` handlers), + * the callback runs once the response stream settles: it is read to + * completion, errored, or canceled, or the request is aborted. * @example * sse('/', ({ client, finalize }) => { * const interval = setInterval(() => client.send({ data: 'ping' })) @@ -423,9 +423,13 @@ export abstract class RequestHandler< } if (!isIterable(result)) { - // Otherwise, run the cleanup immediately if it's a plain response resolver. - await this.runScheduledCleanups(info.requestId) - return result + // Otherwise, run the cleanups if it's a plain response resolver + // (deferred until the response stream settles for streamed responses). + return this.complete({ + request: info.request, + requestId: info.requestId, + response: result, + }) } /** @@ -465,11 +469,13 @@ export abstract class RequestHandler< // only after it's been completely exhausted. this.isUsed = true - await this.runScheduledCleanups(info.requestId) - // Clone the previously stored response so it can be read // when receiving it repeatedly from the "done" generator. - return this.resolverIteratorResult?.clone() + return this.complete({ + request: info.request, + requestId: info.requestId, + response: this.resolverIteratorResult?.clone(), + }) } return nextResponse @@ -505,7 +511,7 @@ export abstract class RequestHandler< this.scheduledCleanups.set(requestId, cleanups) } - protected async exhaustCleanups( + private async exhaustCleanups( cleanups: Array<() => MaybePromise>, ): Promise { const errors: Array = [] @@ -532,28 +538,98 @@ export abstract class RequestHandler< } } - private async runScheduledCleanups(requestId: string): Promise { + /** + * Remove and return the cleanups scheduled for the given request + * (or the pending iterator cleanups for generator resolvers). + */ + private takeScheduledCleanups( + requestId: string, + ): Array<() => MaybePromise> | undefined { if ( this.resolverIterator && this.resolverIteratorCleanups != null && this.resolverIteratorCleanups.length > 0 ) { - try { - await this.exhaustCleanups(this.resolverIteratorCleanups) - } finally { - this.resolverIteratorCleanups = undefined - } - return + const cleanups = this.resolverIteratorCleanups + this.resolverIteratorCleanups = undefined + return cleanups } const cleanups = this.scheduledCleanups.get(requestId) - if (!cleanups || cleanups.length == 0) { - return + if (!cleanups || cleanups.length === 0) { + return undefined } - await this.exhaustCleanups(cleanups) this.scheduledCleanups.delete(requestId) + return cleanups + } + + private async runScheduledCleanups(requestId: string): Promise { + const cleanups = this.takeScheduledCleanups(requestId) + + if (cleanups) { + await this.exhaustCleanups(cleanups) + } + } + + /** + * Conclude the response resolution for the given request. + * Runs the scheduled cleanups immediately for responses without a + * `ReadableStream` body. For streamed responses, returns an observed + * copy of the response and defers the cleanups until its body settles + * (is read to completion, errored, or canceled) or the request is + * aborted, whichever comes first. + */ + private async complete(args: { + request: StrictRequest + requestId: string + response: Response | undefined | void + }): Promise { + const cleanups = this.takeScheduledCleanups(args.requestId) + + if (!cleanups) { + return args.response + } + + const observedResponse = args.response + ? observeResponseBodyStream(args.response) + : null + + if (!observedResponse) { + await this.exhaustCleanups(cleanups) + return args.response + } + + const listenerController = new AbortController() + + const runCleanupsOnce = (): void => { + if (listenerController.signal.aborted) { + return + } + + listenerController.abort() + void this.exhaustCleanups(cleanups) + } + + void observedResponse.settled.then(runCleanupsOnce) + + /** + * @note Also run the cleanups when the request is aborted. + * Stream cancellation does not always propagate to the observed + * response (e.g. if an unconsumed clone of the response exists), + * while the request abort reliably means the response is unused. + */ + if (args.request.signal.aborted) { + runCleanupsOnce() + } else { + args.request.signal.addEventListener('abort', runCleanupsOnce, { + once: true, + signal: listenerController.signal, + }) + } + + return observedResponse.response } } diff --git a/src/core/sse.ts b/src/core/sse.ts index c012f0811..33cefcc4e 100644 --- a/src/core/sse.ts +++ b/src/core/sse.ts @@ -14,7 +14,6 @@ import { getTimestamp } from './utils/logging/getTimestamp' import { devUtils } from './utils/internal/devUtils' import { colors } from './ws/utils/attachWebSocketLogger' import { toPublicUrl } from './utils/request/toPublicUrl' -import type { MaybePromise } from './typeUtils' type EventMapConstraint = { message?: unknown @@ -197,18 +196,6 @@ class ServerSentEventHandler< }) /* eslint-enable no-console */ } - - protected async exhaustCleanups( - cleanups: Array<() => MaybePromise>, - ): Promise { - const onClose = () => { - this.#emitter.removeListener('error', onClose) - this.#emitter.removeListener('close', onClose) - void super.exhaustCleanups(cleanups) - } - - this.#emitter.once('error', onClose).once('close', onClose) - } } type Values = T[keyof T] diff --git a/src/core/utils/HttpResponse/decorators.ts b/src/core/utils/HttpResponse/decorators.ts index 0025ff2c1..9988a0d2d 100644 --- a/src/core/utils/HttpResponse/decorators.ts +++ b/src/core/utils/HttpResponse/decorators.ts @@ -59,3 +59,28 @@ export function decorateResponse( export function getRawSetCookie(response: Response): string | undefined { return Reflect.get(response, kSetCookie) } + +/** + * Copy the instance-level response decorations, like the mocked + * response type or the raw "Set-Cookie" header record, from one + * response instance to another. + */ +export function copyResponseDecorations( + source: Response, + target: Response, +): void { + const typeDescriptor = Object.getOwnPropertyDescriptor(source, 'type') + + if (typeDescriptor) { + Object.defineProperty(target, 'type', typeDescriptor) + } + + const setCookieDescriptor = Object.getOwnPropertyDescriptor( + source, + kSetCookie, + ) + + if (setCookieDescriptor) { + Object.defineProperty(target, kSetCookie, setCookieDescriptor) + } +} diff --git a/src/core/utils/internal/observe-response-body-stream.test.ts b/src/core/utils/internal/observe-response-body-stream.test.ts new file mode 100644 index 000000000..ac2c6770d --- /dev/null +++ b/src/core/utils/internal/observe-response-body-stream.test.ts @@ -0,0 +1,64 @@ +import { observeResponseBodyStream } from './observe-response-body-stream' + +it('returns null for a response without a body', () => { + expect(observeResponseBodyStream(new Response(null))).toBeNull() +}) + +it('returns null for a response whose body is already used', async () => { + const response = new Response('hello') + await response.text() + + expect(observeResponseBodyStream(response)).toBeNull() +}) + +it('returns null for a response whose body is locked', () => { + const response = new Response('hello') + response.body?.getReader() + + expect(observeResponseBodyStream(response)).toBeNull() +}) + +it('resolves "settled" when the response body is read to completion', async () => { + const observed = observeResponseBodyStream(new Response('hello')) + expect(observed).not.toBeNull() + + await expect(observed?.response.text()).resolves.toBe('hello') + await expect(observed?.settled).resolves.toBeUndefined() +}) + +it('resolves "settled" when the response body errors', async () => { + const stream = new ReadableStream({ + start(controller) { + controller.error(new Error('stream error')) + }, + }) + const observed = observeResponseBodyStream(new Response(stream)) + expect(observed).not.toBeNull() + + await expect(observed?.response.text()).rejects.toThrow('stream error') + await expect(observed?.settled).resolves.toBeUndefined() +}) + +it('resolves "settled" when the response body is canceled', async () => { + const stream = new ReadableStream({ + pull(controller) { + controller.enqueue(new TextEncoder().encode('ping')) + }, + }) + const observed = observeResponseBodyStream(new Response(stream)) + const reader = observed?.response.body?.getReader() + expect(reader).toBeDefined() + + await reader?.read() + await reader?.cancel(new Error('client disconnected')) + + await expect(observed?.settled).resolves.toBeUndefined() +}) + +it('resolves "settled" when the response body is consumed via piping', async () => { + const observed = observeResponseBodyStream(new Response('hello')) + expect(observed).not.toBeNull() + + await observed?.response.body?.pipeTo(new WritableStream()) + await expect(observed?.settled).resolves.toBeUndefined() +}) diff --git a/src/core/utils/internal/observe-response-body-stream.ts b/src/core/utils/internal/observe-response-body-stream.ts new file mode 100644 index 000000000..a2baeef17 --- /dev/null +++ b/src/core/utils/internal/observe-response-body-stream.ts @@ -0,0 +1,72 @@ +import { DeferredPromise } from '@open-draft/deferred-promise' +import { FetchResponse } from '@mswjs/interceptors' +import { copyResponseDecorations } from '../HttpResponse/decorators' + +export interface ObservedResponse { + response: Response + /** + * A promise that resolves once the response body stream has settled: + * it was closed, errored, or canceled, and will not be used anymore. + */ + settled: Promise +} + +/** + * The `cancel` transformer callback is missing from the TypeScript + * DOM types. It is invoked when the readable side of the transform + * stream is canceled by the consumer or its writable side is aborted + * (e.g. when the source stream errors). + * @see https://streams.spec.whatwg.org/#transformer-api + */ +interface TransformerWithCancel extends Transformer< + Input, + Output +> { + cancel?: (reason: unknown) => void | PromiseLike +} + +/** + * Observe the `ReadableStream` body of the given response. + * Returns a copy of that response whose body reports when it has + * settled (was read to completion, errored, or canceled by the consumer). + * Returns `null` for responses whose body cannot be observed + * (no body, already used, or locked). + */ +export function observeResponseBodyStream( + response: Response, +): ObservedResponse | null { + if (response.body == null || response.bodyUsed || response.body.locked) { + return null + } + + const settled = new DeferredPromise() + const settle = (): void => { + settled.resolve() + } + + /** + * @note Reconstruct the response because the body of an existing + * response cannot be replaced. Use `FetchResponse` to support + * non-standard response status codes (e.g. 101). + */ + const observedResponse = new FetchResponse( + response.body.pipeThrough( + new TransformStream({ + flush: settle, + cancel: settle, + } as TransformerWithCancel), + ), + { + status: response.status, + statusText: response.statusText, + headers: response.headers, + }, + ) + + copyResponseDecorations(response, observedResponse) + + return { + response: observedResponse, + settled, + } +} From 50f7d17c2ed2a9294a0adf522f269fe84bc26d6a Mon Sep 17 00:00:00 2001 From: Artem Zakharchenko Date: Tue, 7 Jul 2026 18:22:05 +0200 Subject: [PATCH 4/8] chore: copy all response own properties --- src/core/utils/HttpResponse/decorators.ts | 32 ++++++++----------- .../internal/observe-response-body-stream.ts | 4 +-- 2 files changed, 16 insertions(+), 20 deletions(-) diff --git a/src/core/utils/HttpResponse/decorators.ts b/src/core/utils/HttpResponse/decorators.ts index 9988a0d2d..fb60f286c 100644 --- a/src/core/utils/HttpResponse/decorators.ts +++ b/src/core/utils/HttpResponse/decorators.ts @@ -1,5 +1,5 @@ import statuses from '../../../shims/statuses' -import type { HttpResponseInit } from '../../HttpResponse' +import type { HttpResponse, HttpResponseInit } from '../../HttpResponse' const { message } = statuses @@ -27,7 +27,7 @@ export function normalizeResponseInit( } export function decorateResponse( - response: Response, + response: HttpResponse, init: HttpResponseDecoratedInit, ): Response { // Allow mocking the response type. @@ -61,26 +61,22 @@ export function getRawSetCookie(response: Response): string | undefined { } /** - * Copy the instance-level response decorations, like the mocked - * response type or the raw "Set-Cookie" header record, from one - * response instance to another. + * Copy the given response own properties, like internal symbols, + * onto another response. Used for faithful internal copying of responses. */ -export function copyResponseDecorations( +export function copyResponseOwnProperties( source: Response, target: Response, ): void { - const typeDescriptor = Object.getOwnPropertyDescriptor(source, 'type') + for (const propertyName of Reflect.ownKeys(source)) { + const descriptor = Object.getOwnPropertyDescriptor(source, propertyName) + const existingDescriptor = Object.getOwnPropertyDescriptor( + target, + propertyName, + ) - if (typeDescriptor) { - Object.defineProperty(target, 'type', typeDescriptor) - } - - const setCookieDescriptor = Object.getOwnPropertyDescriptor( - source, - kSetCookie, - ) - - if (setCookieDescriptor) { - Object.defineProperty(target, kSetCookie, setCookieDescriptor) + if (descriptor && existingDescriptor == null) { + Object.defineProperty(target, propertyName, descriptor) + } } } diff --git a/src/core/utils/internal/observe-response-body-stream.ts b/src/core/utils/internal/observe-response-body-stream.ts index a2baeef17..771a5911f 100644 --- a/src/core/utils/internal/observe-response-body-stream.ts +++ b/src/core/utils/internal/observe-response-body-stream.ts @@ -1,6 +1,6 @@ import { DeferredPromise } from '@open-draft/deferred-promise' import { FetchResponse } from '@mswjs/interceptors' -import { copyResponseDecorations } from '../HttpResponse/decorators' +import { copyResponseOwnProperties } from '../HttpResponse/decorators' export interface ObservedResponse { response: Response @@ -63,7 +63,7 @@ export function observeResponseBodyStream( }, ) - copyResponseDecorations(response, observedResponse) + copyResponseOwnProperties(response, observedResponse) return { response: observedResponse, From d8146efdd99a1d13a86e3b33e437e49297431a4f Mon Sep 17 00:00:00 2001 From: Artem Zakharchenko Date: Tue, 7 Jul 2026 18:29:05 +0200 Subject: [PATCH 5/8] fix(RequestHandler): apply finalize only if accessed and called --- src/core/handlers/HttpHandler.test.ts | 89 +++++++++++++++++++++++++++ src/core/handlers/RequestHandler.ts | 41 +++++++++--- src/core/sse.ts | 16 +++-- test/node/msw-api/finalize.test.ts | 2 +- 4 files changed, 132 insertions(+), 16 deletions(-) diff --git a/src/core/handlers/HttpHandler.test.ts b/src/core/handlers/HttpHandler.test.ts index 70f018423..8a14cc366 100644 --- a/src/core/handlers/HttpHandler.test.ts +++ b/src/core/handlers/HttpHandler.test.ts @@ -252,3 +252,92 @@ describe('run', () => { await expect(run()).resolves.toBe('complete') }) }) + +describe('finalize', () => { + it('returns the exact response instance if the resolver never accesses finalize', async () => { + const response = new HttpResponse( + new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode('hello')) + controller.close() + }, + }), + ) + const handler = new HttpHandler('GET', '/resource', () => { + return response + }) + + const result = await handler.run({ + request: new Request(new URL('/resource', location.href)), + requestId: createRequestId(), + }) + + expect(result?.response).toBe(response) + }) + + it('returns the exact response instance if finalize is accessed but never called', async () => { + const response = new HttpResponse( + new ReadableStream({ + start(controller) { + controller.close() + }, + }), + ) + const handler = new HttpHandler('GET', '/resource', ({ finalize }) => { + expect(finalize).toBeInstanceOf(Function) + return response + }) + + const result = await handler.run({ + request: new Request(new URL('/resource', location.href)), + requestId: createRequestId(), + }) + + expect(result?.response).toBe(response) + }) + + it('defers the cleanup until the response stream settles', async () => { + const cleanup = vi.fn() + const response = new HttpResponse( + new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode('hello')) + controller.close() + }, + }), + ) + const handler = new HttpHandler('GET', '/resource', ({ finalize }) => { + finalize(cleanup) + return response + }) + + const result = await handler.run({ + request: new Request(new URL('/resource', location.href)), + requestId: createRequestId(), + }) + + // The observed response is a copy of the original response. + expect(result?.response).not.toBe(response) + expect(cleanup).not.toHaveBeenCalled() + + await expect(result?.response?.text()).resolves.toBe('hello') + await vi.waitFor(() => { + expect(cleanup).toHaveBeenCalledOnce() + }) + }) + + it('runs the cleanup immediately for responses without a body', async () => { + const cleanup = vi.fn() + const handler = new HttpHandler('GET', '/resource', ({ finalize }) => { + finalize(cleanup) + return new HttpResponse(null) + }) + + await handler.run({ + request: new Request(new URL('/resource', location.href)), + requestId: createRequestId(), + }) + + expect(cleanup).toHaveBeenCalledOnce() + }) +}) diff --git a/src/core/handlers/RequestHandler.ts b/src/core/handlers/RequestHandler.ts index 969e360bd..53d12f77f 100644 --- a/src/core/handlers/RequestHandler.ts +++ b/src/core/handlers/RequestHandler.ts @@ -355,20 +355,41 @@ export abstract class RequestHandler< const listenerController = new AbortController() - args.request.signal.addEventListener( - 'abort', - () => this.runScheduledCleanups(args.requestId), - { - once: true, - signal: listenerController.signal, - }, - ) + /** + * @note Initialize the `finalize` machinery lazily, on the first + * access of the `finalize` property by the resolver. If the resolver + * never accesses it, the handler behaves as if `finalize` never + * existed: no abort listeners, no scheduled cleanups, and no response + * body stream observation (see `this.complete()`). + */ + let finalizeFunction: ResponseResolverFinalizeFunction | undefined + + const getFinalize = (): ResponseResolverFinalizeFunction => { + if (finalizeFunction == null) { + // Run any scheduled cleanups if the request gets aborted + // while the resolver is still executing. + args.request.signal.addEventListener( + 'abort', + () => this.runScheduledCleanups(args.requestId), + { + once: true, + signal: listenerController.signal, + }, + ) + + finalizeFunction = (callback) => { + this.scheduleCleanup(args.requestId, callback) + } + } + + return finalizeFunction + } const mockedResponsePromise = ( executeResolver({ ...resolverExtras, - finalize: (callback) => { - this.scheduleCleanup(args.requestId, callback) + get finalize(): ResponseResolverFinalizeFunction { + return getFinalize() }, requestId: args.requestId, request: args.request, diff --git a/src/core/sse.ts b/src/core/sse.ts index 33cefcc4e..8eceb9a11 100644 --- a/src/core/sse.ts +++ b/src/core/sse.ts @@ -96,11 +96,17 @@ class ServerSentEventHandler< client[kClientEmitter] = this.#emitter - await resolver({ - ...info, - client, - server, - }) + /** + * @note Extend the resolver info via assignment instead of a spread. + * Spreading the info would invoke its lazy `finalize` getter, + * initializing the finalize machinery for resolvers that never use it. + */ + await resolver( + Object.assign(info, { + client, + server, + }), + ) return response }) diff --git a/test/node/msw-api/finalize.test.ts b/test/node/msw-api/finalize.test.ts index 1566f61a1..a8feee5aa 100644 --- a/test/node/msw-api/finalize.test.ts +++ b/test/node/msw-api/finalize.test.ts @@ -1,7 +1,7 @@ // @vitest-environment node +import { setTimeout } from 'node:timers/promises' import { http, passthrough } from 'msw' import { setupServer } from 'msw/node' -import { setTimeout } from 'node:timers/promises' const server = setupServer() From be43d68a74a5c5b279186a97f59f2f8c1f804206 Mon Sep 17 00:00:00 2001 From: Artem Zakharchenko Date: Tue, 7 Jul 2026 18:51:03 +0200 Subject: [PATCH 6/8] fix: support chromium in response stream observing --- .../internal/observe-response-body-stream.ts | 65 +++++++------- ...e.finally.test.ts => sse.finalize.test.ts} | 87 +++++++++++++++++++ 2 files changed, 122 insertions(+), 30 deletions(-) rename test/browser/sse-api/{sse.finally.test.ts => sse.finalize.test.ts} (59%) diff --git a/src/core/utils/internal/observe-response-body-stream.ts b/src/core/utils/internal/observe-response-body-stream.ts index 771a5911f..2f63fac4f 100644 --- a/src/core/utils/internal/observe-response-body-stream.ts +++ b/src/core/utils/internal/observe-response-body-stream.ts @@ -11,20 +11,6 @@ export interface ObservedResponse { settled: Promise } -/** - * The `cancel` transformer callback is missing from the TypeScript - * DOM types. It is invoked when the readable side of the transform - * stream is canceled by the consumer or its writable side is aborted - * (e.g. when the source stream errors). - * @see https://streams.spec.whatwg.org/#transformer-api - */ -interface TransformerWithCancel extends Transformer< - Input, - Output -> { - cancel?: (reason: unknown) => void | PromiseLike -} - /** * Observe the `ReadableStream` body of the given response. * Returns a copy of that response whose body reports when it has @@ -40,28 +26,47 @@ export function observeResponseBodyStream( } const settled = new DeferredPromise() - const settle = (): void => { - settled.resolve() - } + const reader = response.body.getReader() + + /** + * @note Relay the body through a manual underlying source instead of + * `.pipeThrough(new TransformStream({ flush, cancel }))`. The `cancel` + * transformer callback is not implemented in Chromium, which loses + * the stream error/cancelation signals there entirely. + */ + const observedStream = new ReadableStream({ + async pull(controller) { + try { + const readResult = await reader.read() + + if (readResult.done) { + settled.resolve() + controller.close() + return + } + + controller.enqueue(readResult.value) + } catch (error) { + settled.resolve() + throw error + } + }, + async cancel(reason) { + settled.resolve() + await reader.cancel(reason) + }, + }) /** * @note Reconstruct the response because the body of an existing * response cannot be replaced. Use `FetchResponse` to support * non-standard response status codes (e.g. 101). */ - const observedResponse = new FetchResponse( - response.body.pipeThrough( - new TransformStream({ - flush: settle, - cancel: settle, - } as TransformerWithCancel), - ), - { - status: response.status, - statusText: response.statusText, - headers: response.headers, - }, - ) + const observedResponse = new FetchResponse(observedStream, { + status: response.status, + statusText: response.statusText, + headers: response.headers, + }) copyResponseOwnProperties(response, observedResponse) diff --git a/test/browser/sse-api/sse.finally.test.ts b/test/browser/sse-api/sse.finalize.test.ts similarity index 59% rename from test/browser/sse-api/sse.finally.test.ts rename to test/browser/sse-api/sse.finalize.test.ts index c45e8c992..ab9fc0cc8 100644 --- a/test/browser/sse-api/sse.finally.test.ts +++ b/test/browser/sse-api/sse.finalize.test.ts @@ -125,3 +125,90 @@ test('runs cleanup after the event source is errored by the handler', async ({ await expect(finalizedAt).resolves.toBeGreaterThanOrEqual(closedAt) }) + +test('runs independent cleanups for parallel event sources', async ({ + loadExample, + page, +}) => { + await loadExample(EXAMPLE_URL, { + skipActivation: true, + }) + + const finalizedIds: Array = [] + await page.exposeFunction('notifyFinalized', (id: string) => { + finalizedIds.push(id) + }) + + await page.evaluate(async () => { + const { setupWorker, sse } = window.msw + + const worker = setupWorker( + sse('http://localhost/stream', ({ request, finalize }) => { + const url = new URL(request.url) + finalize(() => window.notifyFinalized(url.searchParams.get('id'))) + }), + ) + await worker.start() + }) + + await page.evaluate(() => { + const sourceOne = new EventSource('http://localhost/stream?id=one') + const sourceTwo = new EventSource('http://localhost/stream?id=two') + Object.assign(window, { sourceOne, sourceTwo }) + + return Promise.all([ + new Promise((resolve) => { + sourceOne.addEventListener('open', resolve) + }), + new Promise((resolve) => { + sourceTwo.addEventListener('open', resolve) + }), + ]) + }) + + // Closing one event source must only run the cleanups + // scheduled for that connection. + await page.evaluate(() => { + window.sourceOne.close() + }) + + await expect.poll(() => finalizedIds).toEqual(['one']) + + await page.evaluate(() => { + window.sourceTwo.close() + }) + + await expect.poll(() => finalizedIds).toEqual(['one', 'two']) +}) + +test('runs cleanup when the handler closes the event source synchronously', async ({ + loadExample, + page, +}) => { + await loadExample(EXAMPLE_URL, { + skipActivation: true, + }) + + const finalized = new DeferredPromise() + await page.exposeFunction('notifyFinalized', () => { + finalized.resolve() + }) + + await page.evaluate(async () => { + const { setupWorker, sse } = window.msw + + const worker = setupWorker( + sse('http://localhost/stream', ({ client, finalize }) => { + finalize(() => window.notifyFinalized()) + client.close() + }), + ) + await worker.start() + }) + + await page.evaluate(() => { + new EventSource('http://localhost/stream') + }) + + await expect(finalized).resolves.toBeUndefined() +}) From a1da3e89ef2a0b0e3812b04542c2a13c792e54dc Mon Sep 17 00:00:00 2001 From: Artem Zakharchenko Date: Tue, 7 Jul 2026 19:00:20 +0200 Subject: [PATCH 7/8] fix: skip event-stream body cloning, preserve response events --- src/mockServiceWorker.js | 31 ++++++++++++++++++------------- 1 file changed, 18 insertions(+), 13 deletions(-) diff --git a/src/mockServiceWorker.js b/src/mockServiceWorker.js index 0772b04f6..035c844a5 100644 --- a/src/mockServiceWorker.js +++ b/src/mockServiceWorker.js @@ -134,18 +134,21 @@ async function handleRequest(event, requestId, requestInterceptedAt) { // Send back the response clone for the "response:*" life-cycle events. // Ensure MSW is active and ready to handle the message, otherwise // this message will pend indefinitely. - if ( - client && - activeClientIds.has(client.id) && - !response.headers + if (client && activeClientIds.has(client.id)) { + const serializedRequest = await serializeRequest(requestCloneForEvents) + + // Omit the body of server-sent event stream responses. + // Cloning such responses would prevent client-side stream cancelations + // from reaching the original stream (a teed stream only cancels its + // source once both of its branches cancel) and would buffer the + // entire stream into the unconsumed clone indefinitely. + const isEventStreamResponse = response.headers .get('content-type') ?.toLowerCase() .startsWith('text/event-stream') - ) { - const serializedRequest = await serializeRequest(requestCloneForEvents) // Clone the response so both the client and the library could consume it. - const responseClone = response.clone() + const responseClone = isEventStreamResponse ? null : response.clone() sendToClient( client, @@ -158,15 +161,17 @@ async function handleRequest(event, requestId, requestInterceptedAt) { ...serializedRequest, }, response: { - type: responseClone.type, - status: responseClone.status, - statusText: responseClone.statusText, - headers: Object.fromEntries(responseClone.headers.entries()), - body: responseClone.body, + type: response.type, + status: response.status, + statusText: response.statusText, + headers: Object.fromEntries(response.headers.entries()), + body: responseClone ? responseClone.body : null, }, }, }, - responseClone.body ? [serializedRequest.body, responseClone.body] : [], + responseClone && responseClone.body + ? [serializedRequest.body, responseClone.body] + : [], ) } From 170959c7488de46e4af1a64720d5408f676e5154 Mon Sep 17 00:00:00 2001 From: Artem Zakharchenko Date: Tue, 7 Jul 2026 19:24:01 +0200 Subject: [PATCH 8/8] fix(RequestHandler): improve abort listener support --- src/core/handlers/RequestHandler.ts | 29 ++++++++++---- test/node/msw-api/finalize.test.ts | 60 +++++++++++++++++++++++++++++ 2 files changed, 81 insertions(+), 8 deletions(-) diff --git a/src/core/handlers/RequestHandler.ts b/src/core/handlers/RequestHandler.ts index 53d12f77f..aaeb261ac 100644 --- a/src/core/handlers/RequestHandler.ts +++ b/src/core/handlers/RequestHandler.ts @@ -368,17 +368,30 @@ export abstract class RequestHandler< if (finalizeFunction == null) { // Run any scheduled cleanups if the request gets aborted // while the resolver is still executing. - args.request.signal.addEventListener( - 'abort', - () => this.runScheduledCleanups(args.requestId), - { - once: true, - signal: listenerController.signal, - }, - ) + if (!args.request.signal.aborted) { + args.request.signal.addEventListener( + 'abort', + () => this.runScheduledCleanups(args.requestId), + { + once: true, + signal: listenerController.signal, + }, + ) + } finalizeFunction = (callback) => { this.scheduleCleanup(args.requestId, callback) + + /** + * @note Run the cleanup immediately if the request has already + * been aborted. The "abort" listener above never fires for an + * already-aborted signal (and fires at most once), while + * long-lived resolvers (streams, generators) may never settle + * to run the cleanups on completion. + */ + if (args.request.signal.aborted) { + void this.runScheduledCleanups(args.requestId) + } } } diff --git a/test/node/msw-api/finalize.test.ts b/test/node/msw-api/finalize.test.ts index a8feee5aa..5bf4a3c16 100644 --- a/test/node/msw-api/finalize.test.ts +++ b/test/node/msw-api/finalize.test.ts @@ -149,6 +149,66 @@ it('runs after the request has been aborted', async () => { expect(cleanup).toHaveBeenCalledOnce() }) +it('runs immediately when scheduled after the request has been aborted', async () => { + const cleanup = vi.fn() + + server.use( + http.get('http://localhost/resource', async ({ finalize }) => { + // Await past the point where the request gets aborted so the + // first `finalize` access happens on an already-aborted signal. + await setTimeout(250) + finalize(cleanup) + + // Simulate a long-lived resolver that never settles. + await new Promise(() => {}) + }), + ) + + const controller = new AbortController() + const responsePromise = fetch('http://localhost/resource', { + signal: controller.signal, + }) + await setTimeout(100) + controller.abort() + + await expect(responsePromise).rejects.toThrow() + await vi.waitFor(() => { + expect(cleanup).toHaveBeenCalledOnce() + }) +}) + +it('runs cleanups scheduled after the abort listener has fired', async () => { + const cleanupBeforeAbort = vi.fn() + const cleanupAfterAbort = vi.fn() + + server.use( + http.get('http://localhost/resource', async ({ finalize }) => { + finalize(cleanupBeforeAbort) + + // Await past the point where the request gets aborted + // (the "abort" listener fires and runs `cleanupBeforeAbort`). + await setTimeout(250) + finalize(cleanupAfterAbort) + + // Simulate a long-lived resolver that never settles. + await new Promise(() => {}) + }), + ) + + const controller = new AbortController() + const responsePromise = fetch('http://localhost/resource', { + signal: controller.signal, + }) + await setTimeout(100) + controller.abort() + + await expect(responsePromise).rejects.toThrow() + await vi.waitFor(() => { + expect(cleanupBeforeAbort).toHaveBeenCalledOnce() + expect(cleanupAfterAbort).toHaveBeenCalledOnce() + }) +}) + it('runs once the generator resolver is exhausted', async () => { const cleanup = vi.fn()