From eb87aab082356ce8ec7d243b7b6e4dea79d68b9f Mon Sep 17 00:00:00 2001 From: Elliott Johnson Date: Wed, 19 Aug 2026 13:45:08 -0600 Subject: [PATCH 1/5] tweaks --- .../src/runtime/server/remote-functions.js | 291 +++++++++--------- .../runtime/server/remote-functions.spec.js | 72 +++++ packages/kit/src/runtime/server/utils.js | 34 -- packages/kit/src/runtime/server/utils.spec.js | 76 ----- 4 files changed, 223 insertions(+), 250 deletions(-) create mode 100644 packages/kit/src/runtime/server/remote-functions.spec.js delete mode 100644 packages/kit/src/runtime/server/utils.spec.js diff --git a/packages/kit/src/runtime/server/remote-functions.js b/packages/kit/src/runtime/server/remote-functions.js index 04c7761d0422..f6cfabf37e29 100644 --- a/packages/kit/src/runtime/server/remote-functions.js +++ b/packages/kit/src/runtime/server/remote-functions.js @@ -1,7 +1,7 @@ /** @import { RequestEvent, SSRManifest } from '@sveltejs/kit' */ /** @import { RemoteForm } from '$app/server' */ /** @import { ActionResult } from '$app/forms' */ -/** @import { RemoteFormInternals, RemoteFunctionData, RemoteFunctionResponse, RemoteInternals, RemoteQueryLiveInternals, RequestState, SSROptions } from 'types' */ +/** @import { RemoteFormInternals, RemoteFunctionData, RemoteFunctionResponse, RemoteInternals, RequestState, SSROptions } from 'types' */ import { error } from '@sveltejs/kit'; import { Redirect, SvelteKitError } from '@sveltejs/kit/internal'; @@ -15,7 +15,8 @@ import { normalize_error } from '../../utils/error.js'; import { check_incorrect_fail_use, get_action_location } from './page/actions.js'; import { DEV } from 'esm-env'; import { deserialize_binary_form } from '../form-utils.js'; -import { stream_from_iterator, with_version_header } from './utils.js'; +import { text_encoder } from '../utils.js'; +import { with_version_header } from './utils.js'; /** * How long (in milliseconds) to wait after the last message was sent before @@ -24,80 +25,15 @@ import { stream_from_iterator, with_version_header } from './utils.js'; */ const KEEP_ALIVE_INTERVAL = 30_000; -/** @type {typeof handle_remote_call_internal} */ -export async function handle_remote_call(event, state, options, manifest, id) { - return record_span({ - name: 'sveltekit.remote.call', - attributes: { - 'sveltekit.remote.call.id': id - }, - fn: async (current) => { - const traced_event = merge_tracing(event, current); - const response = await with_request_store({ event: traced_event, state }, () => - handle_remote_call_internal(traced_event, state, options, manifest, id) - ); - return with_version_header(response); - } - }); -} - /** - * Looks a remote function up in the manifest by its request id. - * @param {SSRManifest} manifest - * @param {string} id - */ -async function resolve_remote_function(manifest, id) { - const [hash, name, additional_args] = id.split('/'); - const remotes = manifest._.remotes; - - if (!Object.hasOwn(remotes, hash)) error(404); - - const module = await remotes[hash](); - const fn = Object.hasOwn(module.default, name) ? module.default[name] : undefined; - - if (!fn) error(404); - - return { fn, internals: /** @type {RemoteInternals} */ (fn.__), additional_args }; -} - -/** - * @param {RemoteFunctionData} data - * @param {HeadersInit | undefined} headers - */ -function result_response(data, headers) { - return Response.json( - /** @type {RemoteFunctionResponse} */ ({ - type: 'result', - data: stringify(data) - }), - { headers } - ); -} - -/** - * Handles a `query.live` call: runs the generator and streams its values as - * server-sent events. * @param {RequestEvent} event * @param {RequestState} state * @param {SSROptions} options - * @param {RemoteQueryLiveInternals} internals + * @param {import('types').RemoteQueryLiveInternals} internals + * @param {any} arg */ -function handle_live_query(event, state, options, internals) { - if (event.request.method !== 'GET') { - throw new SvelteKitError( - 405, - 'Method Not Allowed', - `\`query.live\` functions must be invoked via GET request, not ${event.request.method}` - ); - } - - const payload = /** @type {string} */ (new URL(event.request.url).searchParams.get('payload')); - - // aborted whenever the stream is torn down, so unlike the request signal it - // also fires on response teardown, which the generator could otherwise - // never observe +export function create_live_query_response(event, state, options, internals, arg) { const cancellation = new AbortController(); - const live_event = { ...event, request: new Request(event.request, { @@ -105,82 +41,91 @@ function handle_live_query(event, state, options, internals) { }) }; - const generator = internals.run(live_event, state, parse_remote_arg(payload)); - - /** @param {any} payload */ - const frame = (payload) => 'data: ' + JSON.stringify(payload) + '\n\n'; + const generator = internals.run(live_event, state, arg); + let open = true; + let pulling = false; + /** @type {ReadableStreamDefaultController} */ + let stream_controller; + /** @type {ReturnType | undefined} */ + let keep_alive; /** @type {string | undefined} */ let result; - // everything the stream sends, as a generator of SSE strings — it holds no - // reference to the stream controller, so it cannot touch a dead one - async function* frames() { - /** @type {Promise> | null} */ - let pending = null; - let settled = false; - /** @type {() => void} */ - let wake = () => {}; - const settle = () => { - settled = true; - wake(); - }; + function schedule_keep_alive() { + clearTimeout(keep_alive); + keep_alive = setTimeout(() => { + if (!open) return; + stream_controller.enqueue(text_encoder.encode(': keep-alive\n\n')); + schedule_keep_alive(); + }, KEEP_ALIVE_INTERVAL); + } - try { - while (true) { - if (!pending) { - settled = false; - // one reaction per next() call, so an idle stream doesn't - // accumulate one per keep-alive tick - pending = generator.next(); - pending.then(settle, settle); - } + /** @param {any} data */ + function send(data) { + if (!open) return; + stream_controller.enqueue(text_encoder.encode('data: ' + JSON.stringify(data) + '\n\n')); + schedule_keep_alive(); + } - if (!settled) { - await new Promise((resolve) => { - const timer = setTimeout(resolve, KEEP_ALIVE_INTERVAL); - wake = () => { - clearTimeout(timer); - resolve(undefined); - }; - }); - } + /** @param {boolean} cancelled */ + function teardown(cancelled) { + if (!open) return; + open = false; + clearTimeout(keep_alive); + cancellation.abort(); + if (!cancelled) stream_controller.close(); + // AsyncGenerator.return() cannot interrupt a pending next(). Cleanup is + // cooperative via request.signal, so stream cancellation must not await it. + void generator.return(undefined).catch(() => {}); + } - if (!settled) { - // SSE comments (lines starting with `:`) are ignored by the client - yield ': keep-alive\n\n'; - continue; - } + return new Response( + new ReadableStream({ + start(controller) { + stream_controller = controller; + schedule_keep_alive(); + }, + async pull() { + if (!open || pulling) return; + pulling = true; + + try { + while (open) { + const { value, done } = await generator.next(); + + if (!open) return; + + if (done) { + teardown(false); + return; + } - const winner = await pending; - pending = null; + if (result !== (result = stringify(value))) { + send({ type: 'result', result }); + return; + } + } + } catch (error) { + if (!open) return; - if (winner.done) return; + if (error instanceof Redirect) { + send({ type: 'redirect', location: error.location }); + } else { + const transformed = await handle_error_and_jsonify(event, state, options, error); - // only send changed data - if (result !== (result = stringify(winner.value))) { - yield frame({ type: 'result', result }); - } - } - } catch (error) { - if (!live_event.request.signal.aborted) { - if (error instanceof Redirect) { - yield frame({ type: 'redirect', location: error.location }); - } else { - yield frame({ - type: 'error', - error: await handle_error_and_jsonify(event, state, options, error) - }); + send({ type: 'error', error: transformed }); + } + + teardown(false); + } finally { + pulling = false; } + }, + cancel() { + teardown(true); } - } finally { - cancellation.abort(); - await generator.return(undefined); - } - } - - return new Response( - stream_from_iterator(frames(), () => cancellation.abort()), + }), { headers: { 'cache-control': 'private, no-store', @@ -189,6 +134,24 @@ function handle_live_query(event, state, options, internals) { } ); } + +/** @type {typeof handle_remote_call_internal} */ +export async function handle_remote_call(event, state, options, manifest, id) { + return record_span({ + name: 'sveltekit.remote.call', + attributes: { + 'sveltekit.remote.call.id': id + }, + fn: async (current) => { + const traced_event = merge_tracing(event, current); + const response = await with_request_store({ event: traced_event, state }, () => + handle_remote_call_internal(traced_event, state, options, manifest, id) + ); + return with_version_header(response); + } + }); +} + /** * @param {RequestEvent} event * @param {RequestState} state @@ -197,7 +160,18 @@ function handle_live_query(event, state, options, internals) { * @param {string} id */ async function handle_remote_call_internal(event, state, options, manifest, id) { - const { fn, internals, additional_args } = await resolve_remote_function(manifest, id); + const [hash, name, additional_args] = id.split('/'); + const remotes = manifest._.remotes; + + if (!Object.hasOwn(remotes, hash)) error(404); + + const module = await remotes[hash](); + const fn = Object.hasOwn(module.default, name) ? module.default[name] : undefined; + + if (!fn) error(404); + + /** @type {RemoteInternals} */ + const internals = fn.__; event.tracing.current.setAttributes({ 'sveltekit.remote.call.type': internals.type, @@ -212,8 +186,27 @@ async function handle_remote_call_internal(event, state, options, manifest, id) const data = {}; switch (internals.type) { - case 'query_live': - return handle_live_query(event, state, options, internals); + case 'query_live': { + if (event.request.method !== 'GET') { + throw new SvelteKitError( + 405, + 'Method Not Allowed', + `\`query.live\` functions must be invoked via GET request, not ${event.request.method}` + ); + } + + const payload = /** @type {string} */ ( + new URL(event.request.url).searchParams.get('payload') + ); + + return create_live_query_response( + event, + state, + options, + internals, + parse_remote_arg(payload) + ); + } case 'query_batch': { if (event.request.method !== 'POST') { @@ -274,7 +267,13 @@ async function handle_remote_call_internal(event, state, options, manifest, id) if (data._.issues) { // special case — don't serialize refreshes/reconnects - return result_response(data, headers); + return Response.json( + /** @type {RemoteFunctionResponse} */ ({ + type: 'result', + data: stringify(data) + }), + { headers } + ); } break; @@ -316,12 +315,24 @@ async function handle_remote_call_internal(event, state, options, manifest, id) await collect_remote_data(data, event, state, options); - return result_response(data, headers); + return Response.json( + /** @type {RemoteFunctionResponse} */ ({ + type: 'result', + data: stringify(data) + }), + { headers } + ); } catch (error) { if (error instanceof Redirect) { const data = await collect_remote_data({ redirect: error.location }, event, state, options); - return result_response(data, headers); + return Response.json( + /** @type {RemoteFunctionResponse} */ ({ + type: 'result', + data: stringify(data) + }), + { headers } + ); } const transformed = await handle_error_and_jsonify(event, state, options, error); diff --git a/packages/kit/src/runtime/server/remote-functions.spec.js b/packages/kit/src/runtime/server/remote-functions.spec.js new file mode 100644 index 000000000000..eb6b7785a342 --- /dev/null +++ b/packages/kit/src/runtime/server/remote-functions.spec.js @@ -0,0 +1,72 @@ +import { beforeAll, expect, test, vi } from 'vitest'; + +/** @type {typeof import('./remote-functions.js').create_live_query_response} */ +let create_live_query_response; + +beforeAll(async () => { + vi.stubGlobal('__SVELTEKIT_DEV__', false); + ({ create_live_query_response } = await import('./remote-functions.js')); +}); + +/** + * @param {(event: import('@sveltejs/kit').RequestEvent) => AsyncGenerator} run + */ +function create_response(run) { + const event = /** @type {import('@sveltejs/kit').RequestEvent} */ ({ + request: new Request('http://localhost/_app/remote/test?payload=undefined') + }); + + return create_live_query_response( + event, + /** @type {import('types').RequestState} */ ({}), + /** @type {import('types').SSROptions} */ ({}), + /** @type {import('types').RemoteQueryLiveInternals} */ (/** @type {unknown} */ ({ run })), + undefined + ); +} + +// https://github.com/sveltejs/kit/issues/16778 +test('cancellation ignores a value that arrives after generator.next()', async () => { + /** @type {() => void} */ + let resume = () => {}; + const parked = new Promise((resolve) => (resume = () => resolve(undefined))); + + const response = create_response(async function* () { + yield 'initial'; + await parked; + yield 'late'; + }); + + const reader = /** @type {ReadableStream} */ (response.body).getReader(); + await reader.read(); + const pending = reader.read(); + await Promise.resolve(); + + await expect(reader.cancel()).resolves.toBeUndefined(); + await expect(pending).resolves.toEqual({ value: undefined, done: true }); + + resume(); + await new Promise((resolve) => setTimeout(resolve, 0)); +}); + +test('cancellation aborts the generator request signal and runs cleanup', async () => { + const cleaned_up = vi.fn(); + + const response = create_response(async function* (event) { + try { + yield 'initial'; + await new Promise((resolve) => event.request.signal.addEventListener('abort', resolve)); + } finally { + cleaned_up(); + } + }); + + const reader = /** @type {ReadableStream} */ (response.body).getReader(); + await reader.read(); + const pending = reader.read(); + await Promise.resolve(); + + await reader.cancel(); + await expect(pending).resolves.toEqual({ value: undefined, done: true }); + await vi.waitFor(() => expect(cleaned_up).toHaveBeenCalledOnce()); +}); diff --git a/packages/kit/src/runtime/server/utils.js b/packages/kit/src/runtime/server/utils.js index 67105a351655..990decb06821 100644 --- a/packages/kit/src/runtime/server/utils.js +++ b/packages/kit/src/runtime/server/utils.js @@ -1,39 +1,5 @@ import { text } from '@sveltejs/kit'; import { ENDPOINT_METHODS } from '../../constants.js'; -import { text_encoder } from '../utils.js'; - -/** - * Builds a text stream from an iterator of string chunks. The controller is - * confined here so that a chunk arriving after the stream was torn down is - * dropped instead of hitting a closed controller — either side can tear the - * stream down while `pull` is suspended on `iterator.next()`. - * @param {AsyncIterator} iterator - * @param {() => void} [oncancel] called when the consumer cancels the stream - * @returns {ReadableStream} - */ -export function stream_from_iterator(iterator, oncancel) { - let open = true; - - return new ReadableStream({ - async pull(controller) { - const { value, done } = await iterator.next(); - - if (!open) return; - - if (done) { - open = false; - controller.close(); - } else { - controller.enqueue(text_encoder.encode(value)); - } - }, - async cancel() { - open = false; - oncancel?.(); - await iterator.return?.(undefined); - } - }); -} /** * @param {Partial>} mod diff --git a/packages/kit/src/runtime/server/utils.spec.js b/packages/kit/src/runtime/server/utils.spec.js deleted file mode 100644 index 7f12a67fa81a..000000000000 --- a/packages/kit/src/runtime/server/utils.spec.js +++ /dev/null @@ -1,76 +0,0 @@ -import { expect, test, vi } from 'vitest'; -import { stream_from_iterator } from './utils.js'; - -const decoder = new TextDecoder(); - -/** @param {string[]} chunks */ -async function* from(chunks) { - for (const chunk of chunks) { - await Promise.resolve(); - yield chunk; - } -} - -test('streams encoded chunks and closes when the iterator is done', async () => { - const reader = stream_from_iterator(from(['one', 'two'])).getReader(); - - const first = await reader.read(); - expect(decoder.decode(first.value)).toBe('one'); - - const second = await reader.read(); - expect(decoder.decode(second.value)).toBe('two'); - - await expect(reader.read()).resolves.toEqual({ value: undefined, done: true }); -}); - -test('cancellation notifies the producer and returns the iterator', async () => { - let finished = false; - const oncancel = vi.fn(); - - async function* source() { - try { - await Promise.resolve(); - yield 'one'; - yield 'two'; - } finally { - finished = true; - } - } - - const reader = stream_from_iterator(source(), oncancel).getReader(); - await reader.read(); - - await reader.cancel(); - - expect(oncancel).toHaveBeenCalledOnce(); - expect(finished).toBe(true); -}); - -// https://github.com/sveltejs/kit/issues/16778 -test('cancellation settles while the producer is parked', async () => { - /** @type {(result: IteratorResult) => void} */ - let resolve_next = () => {}; - let returned = false; - - /** @type {AsyncIterator} */ - const iterator = { - next: () => new Promise((resolve) => (resolve_next = resolve)), - return: () => { - returned = true; - return Promise.resolve({ value: undefined, done: true }); - } - }; - - const reader = stream_from_iterator(iterator).getReader(); - const read = reader.read(); // pull is now suspended on `iterator.next()` - await Promise.resolve(); - - // neither cancel() nor the pending read may wait for the parked next() - await reader.cancel(); - expect(returned).toBe(true); - await expect(read).resolves.toEqual({ value: undefined, done: true }); - - // the parked value landing afterwards is a no-op - resolve_next({ value: 'late', done: false }); - await new Promise((resolve) => setTimeout(resolve, 0)); -}); From 404b12a3e9a8ce780fd63c2dff60d2a511107db8 Mon Sep 17 00:00:00 2001 From: Elliott Johnson Date: Wed, 19 Aug 2026 14:31:25 -0600 Subject: [PATCH 2/5] whoops --- .../runtime/server/remote-functions.spec.js | 26 ++++++++++++++++--- 1 file changed, 22 insertions(+), 4 deletions(-) diff --git a/packages/kit/src/runtime/server/remote-functions.spec.js b/packages/kit/src/runtime/server/remote-functions.spec.js index eb6b7785a342..a486697c1bbe 100644 --- a/packages/kit/src/runtime/server/remote-functions.spec.js +++ b/packages/kit/src/runtime/server/remote-functions.spec.js @@ -1,10 +1,14 @@ import { beforeAll, expect, test, vi } from 'vitest'; +import { init_transport } from '#app/internal/transport'; + +const decoder = new TextDecoder(); /** @type {typeof import('./remote-functions.js').create_live_query_response} */ let create_live_query_response; beforeAll(async () => { vi.stubGlobal('__SVELTEKIT_DEV__', false); + init_transport({}); ({ create_live_query_response } = await import('./remote-functions.js')); }); @@ -30,17 +34,24 @@ test('cancellation ignores a value that arrives after generator.next()', async ( /** @type {() => void} */ let resume = () => {}; const parked = new Promise((resolve) => (resume = () => resolve(undefined))); + /** @type {() => void} */ + let did_park = () => {}; + const parked_on_next = new Promise((resolve) => (did_park = () => resolve(undefined))); const response = create_response(async function* () { yield 'initial'; + did_park(); await parked; yield 'late'; }); const reader = /** @type {ReadableStream} */ (response.body).getReader(); - await reader.read(); + const first = await reader.read(); + expect(decoder.decode(first.value)).toBe( + 'data: {"type":"result","result":"[\\"initial\\"]"}\n\n' + ); const pending = reader.read(); - await Promise.resolve(); + await parked_on_next; await expect(reader.cancel()).resolves.toBeUndefined(); await expect(pending).resolves.toEqual({ value: undefined, done: true }); @@ -51,10 +62,14 @@ test('cancellation ignores a value that arrives after generator.next()', async ( test('cancellation aborts the generator request signal and runs cleanup', async () => { const cleaned_up = vi.fn(); + /** @type {() => void} */ + let did_park = () => {}; + const parked_on_next = new Promise((resolve) => (did_park = () => resolve(undefined))); const response = create_response(async function* (event) { try { yield 'initial'; + did_park(); await new Promise((resolve) => event.request.signal.addEventListener('abort', resolve)); } finally { cleaned_up(); @@ -62,9 +77,12 @@ test('cancellation aborts the generator request signal and runs cleanup', async }); const reader = /** @type {ReadableStream} */ (response.body).getReader(); - await reader.read(); + const first = await reader.read(); + expect(decoder.decode(first.value)).toBe( + 'data: {"type":"result","result":"[\\"initial\\"]"}\n\n' + ); const pending = reader.read(); - await Promise.resolve(); + await parked_on_next; await reader.cancel(); await expect(pending).resolves.toEqual({ value: undefined, done: true }); From 5620eef25c85aca88c694b0140260fd4a3d24012 Mon Sep 17 00:00:00 2001 From: Nic Polumeyv <162764842+Nic-Polumeyv@users.noreply.github.com> Date: Wed, 19 Aug 2026 16:46:17 -0400 Subject: [PATCH 3/5] detect the late value actually being dropped --- .../runtime/server/remote-functions.spec.js | 24 ++++++++++++------- 1 file changed, 16 insertions(+), 8 deletions(-) diff --git a/packages/kit/src/runtime/server/remote-functions.spec.js b/packages/kit/src/runtime/server/remote-functions.spec.js index a486697c1bbe..06e0b24e827a 100644 --- a/packages/kit/src/runtime/server/remote-functions.spec.js +++ b/packages/kit/src/runtime/server/remote-functions.spec.js @@ -14,8 +14,9 @@ beforeAll(async () => { /** * @param {(event: import('@sveltejs/kit').RequestEvent) => AsyncGenerator} run + * @param {Partial} [options] */ -function create_response(run) { +function create_response(run, options = {}) { const event = /** @type {import('@sveltejs/kit').RequestEvent} */ ({ request: new Request('http://localhost/_app/remote/test?payload=undefined') }); @@ -23,7 +24,7 @@ function create_response(run) { return create_live_query_response( event, /** @type {import('types').RequestState} */ ({}), - /** @type {import('types').SSROptions} */ ({}), + /** @type {import('types').SSROptions} */ (options), /** @type {import('types').RemoteQueryLiveInternals} */ (/** @type {unknown} */ ({ run })), undefined ); @@ -31,6 +32,7 @@ function create_response(run) { // https://github.com/sveltejs/kit/issues/16778 test('cancellation ignores a value that arrives after generator.next()', async () => { + const handle_error = vi.fn(() => ({ message: 'oops' })); /** @type {() => void} */ let resume = () => {}; const parked = new Promise((resolve) => (resume = () => resolve(undefined))); @@ -38,12 +40,15 @@ test('cancellation ignores a value that arrives after generator.next()', async ( let did_park = () => {}; const parked_on_next = new Promise((resolve) => (did_park = () => resolve(undefined))); - const response = create_response(async function* () { - yield 'initial'; - did_park(); - await parked; - yield 'late'; - }); + const response = create_response( + async function* () { + yield 'initial'; + did_park(); + await parked; + yield 'late'; + }, + { hooks: /** @type {any} */ ({ handleError: handle_error }) } + ); const reader = /** @type {ReadableStream} */ (response.body).getReader(); const first = await reader.read(); @@ -58,6 +63,9 @@ test('cancellation ignores a value that arrives after generator.next()', async ( resume(); await new Promise((resolve) => setTimeout(resolve, 0)); + + // enqueueing the late value would throw and route through handleError + expect(handle_error).not.toHaveBeenCalled(); }); test('cancellation aborts the generator request signal and runs cleanup', async () => { From 20bf6d1d7e82d3f9b071507fa986b264df2f10fd Mon Sep 17 00:00:00 2001 From: Elliott Johnson Date: Wed, 19 Aug 2026 15:10:20 -0600 Subject: [PATCH 4/5] Apply suggestions from code review Co-authored-by: Nic Polumeyv --- packages/kit/src/runtime/server/remote-functions.js | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/packages/kit/src/runtime/server/remote-functions.js b/packages/kit/src/runtime/server/remote-functions.js index f6cfabf37e29..385545a78dd3 100644 --- a/packages/kit/src/runtime/server/remote-functions.js +++ b/packages/kit/src/runtime/server/remote-functions.js @@ -56,7 +56,9 @@ export function create_live_query_response(event, state, options, internals, arg clearTimeout(keep_alive); keep_alive = setTimeout(() => { if (!open) return; - stream_controller.enqueue(text_encoder.encode(': keep-alive\n\n')); + if (stream_controller.desiredSize > 0) { + stream_controller.enqueue(text_encoder.encode(': keep-alive\n\n')); + } schedule_keep_alive(); }, KEEP_ALIVE_INTERVAL); } From dcb5ae3fe8efe3b02dd92c4397a2478927648bca Mon Sep 17 00:00:00 2001 From: Elliott Johnson Date: Thu, 20 Aug 2026 17:17:39 -0600 Subject: [PATCH 5/5] Update packages/kit/src/runtime/server/remote-functions.js Co-authored-by: Nic Polumeyv --- packages/kit/src/runtime/server/remote-functions.js | 2 ++ 1 file changed, 2 insertions(+) diff --git a/packages/kit/src/runtime/server/remote-functions.js b/packages/kit/src/runtime/server/remote-functions.js index 385545a78dd3..c14fca67488f 100644 --- a/packages/kit/src/runtime/server/remote-functions.js +++ b/packages/kit/src/runtime/server/remote-functions.js @@ -82,6 +82,8 @@ export function create_live_query_response(event, state, options, internals, arg void generator.return(undefined).catch(() => {}); } + event.request.signal.addEventListener('abort', () => teardown(true), { once: true }); + return new Response( new ReadableStream({ start(controller) {