Skip to content
Draft
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
5 changes: 5 additions & 0 deletions .changeset/pre/live-query-stream-teardown.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@sveltejs/kit': patch
---

fix: don't touch the `query.live` stream controller after teardown, and make response cancellation observable via the generator's `request.signal`
311 changes: 156 additions & 155 deletions packages/kit/src/runtime/server/remote-functions.js
Original file line number Diff line number Diff line change
@@ -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, RequestState, SSROptions } from 'types' */
/** @import { RemoteFormInternals, RemoteFunctionData, RemoteFunctionResponse, RemoteInternals, RemoteQueryLiveInternals, RequestState, SSROptions } from 'types' */

import { json, error } from '@sveltejs/kit';
import { Redirect, SvelteKitError } from '@sveltejs/kit/internal';
Expand All @@ -15,7 +15,7 @@ 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 { with_version_header } from './utils.js';
import { stream_from_iterator, with_version_header } from './utils.js';

/**
* How long (in milliseconds) to wait after the last message was sent before
Expand All @@ -24,6 +24,8 @@ import { with_version_header } from './utils.js';
*/
const KEEP_ALIVE_INTERVAL = 30_000;

const KEEP_ALIVE = Symbol('keep-alive');

/** @type {typeof handle_remote_call_internal} */
export async function handle_remote_call(event, state, options, manifest, id) {
return record_span({
Expand All @@ -42,13 +44,11 @@ export async function handle_remote_call(event, state, options, manifest, id) {
}

/**
* @param {RequestEvent} event
* @param {RequestState} state
* @param {SSROptions} options
* Looks a remote function up in the manifest by its request id.
* @param {SSRManifest} manifest
* @param {string} id
*/
async function handle_remote_call_internal(event, state, options, manifest, id) {
async function resolve_remote_function(manifest, id) {
const [hash, name, additional_args] = id.split('/');
const remotes = manifest._.remotes;

Expand All @@ -59,149 +59,168 @@ async function handle_remote_call_internal(event, state, options, manifest, id)

if (!fn) error(404);

/** @type {RemoteInternals} */
const internals = fn.__;
return { fn, internals: /** @type {RemoteInternals} */ (fn.__), additional_args };
}

event.tracing.current.setAttributes({
'sveltekit.remote.call.type': internals.type,
'sveltekit.remote.call.name': internals.name
});
/**
* @param {RemoteFunctionData} data
* @param {HeadersInit | undefined} headers
*/
function result_response(data, headers) {
return json(
/** @type {RemoteFunctionResponse} */ ({
type: 'result',
data: stringify(data)
}),
{ headers }
);
}

/** @type {HeadersInit | undefined} */
const headers = state.prerendering ? undefined : { 'cache-control': 'private, no-store' };
/**
* 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
*/
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}`
);
}

try {
/** @type {RemoteFunctionData} */
const data = {};
const payload = /** @type {string} */ (new URL(event.request.url).searchParams.get('payload'));

switch (internals.type) {
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}`
);
}
// 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
const cancellation = new AbortController();

const payload = /** @type {string} */ (
new URL(event.request.url).searchParams.get('payload')
);
if (event.request.signal.aborted) {
cancellation.abort();
} else {
event.request.signal.addEventListener('abort', () => cancellation.abort(), {
once: true
});
}

const generator = internals.run(event, state, parse_remote_arg(payload));

const encoder = new TextEncoder();

let closed = false;

/** @type {ReturnType<typeof setTimeout> | undefined} */
let keep_alive;

/**
* (Re)schedule the keep-alive comment. Called whenever a message is sent, so
* that a keep-alive is only emitted once `KEEP_ALIVE_INTERVAL` has elapsed
* without any other activity.
* @param {ReadableStreamDefaultController} controller
*/
function schedule_keep_alive(controller) {
clearTimeout(keep_alive);
keep_alive = setTimeout(() => {
if (closed || event.request.signal.aborted) return;
// SSE comments (lines starting with `:`) are ignored by the client
controller.enqueue(encoder.encode(': keep-alive\n\n'));
schedule_keep_alive(controller);
}, KEEP_ALIVE_INTERVAL);
}
const live_event = {
...event,
request: new Request(event.request, { signal: cancellation.signal })
};

/**
* @param {ReadableStreamDefaultController} controller
* @param {any} payload
*/
function send(controller, payload) {
controller.enqueue(encoder.encode('data: ' + JSON.stringify(payload) + '\n\n'));
schedule_keep_alive(controller);
const generator = internals.run(live_event, state, parse_remote_arg(payload));

/** @param {any} payload */
const frame = (payload) => 'data: ' + JSON.stringify(payload) + '\n\n';

/**
* Resolves with the next iterator result, or with `KEEP_ALIVE` once
* `KEEP_ALIVE_INTERVAL` has elapsed without one.
* @param {Promise<IteratorResult<any>>} pending
* @returns {Promise<IteratorResult<any> | typeof KEEP_ALIVE>}
*/
function next_or_keep_alive(pending) {
return new Promise((resolve, reject) => {
const timer = setTimeout(() => resolve(KEEP_ALIVE), KEEP_ALIVE_INTERVAL);
pending.then(
(result) => {
clearTimeout(timer);
resolve(result);
},
(error) => {
clearTimeout(timer);
reject(error);
}
);
});
}

/** @type {string | undefined} */
let result = undefined;

/** @type {string | undefined} */
let result = undefined;
// 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<IteratorResult<any>> | null} */
let pending = null;

async function cancel() {
if (closed) return;
closed = true;
clearTimeout(keep_alive);
await generator.return(undefined);
try {
while (true) {
pending ??= generator.next();
const winner = await next_or_keep_alive(pending);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Re-attaching a new .then to the same unsettled pending promise every keep-alive interval accumulates unresolved reactions, leaking memory on idle long-lived query.live SSE streams.

Fix on Vercel


if (winner === KEEP_ALIVE) {
// SSE comments (lines starting with `:`) are ignored by the client
yield ': keep-alive\n\n';
continue;
}

event.request.signal.addEventListener('abort', cancel, { once: true });
pending = null;

return new Response(
new ReadableStream({
start(controller) {
schedule_keep_alive(controller);
},
async pull(controller) {
if (event.request.signal.aborted) {
await cancel();
controller.close();
return;
}
if (winner.done) return;

try {
while (true) {
const { value, done } = await generator.next();

if (done) {
await cancel();
controller.close();
return;
}

// only send changed data
if (result !== (result = stringify(value))) {
send(controller, {
type: 'result',
result
});

return;
}
}
} catch (error) {
if (!event.request.signal.aborted) {
if (error instanceof Redirect) {
send(controller, {
type: 'redirect',
location: error.location
});
} else {
const transformed = await handle_error_and_jsonify(
event,
state,
options,
error
);

send(controller, {
type: 'error',
error: transformed
});
}
}

await cancel();
controller.close();
}
},
cancel
}),
{
headers: {
'cache-control': 'private, no-store',
'content-type': 'text/event-stream'
}
}
);
// only send changed data
if (result !== (result = stringify(winner.value))) {
yield frame({ type: 'result', result });
}
}
} catch (error) {
if (!cancellation.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)
});
}
}
} finally {
cancellation.abort();
await generator.return(undefined);
}
}

return new Response(
stream_from_iterator(frames(), () => cancellation.abort()),
{
headers: {
'cache-control': 'private, no-store',
'content-type': 'text/event-stream'
}
}
);
}
/**
* @param {RequestEvent} event
* @param {RequestState} state
* @param {SSROptions} options
* @param {SSRManifest} manifest
* @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);

event.tracing.current.setAttributes({
'sveltekit.remote.call.type': internals.type,
'sveltekit.remote.call.name': internals.name
});

/** @type {HeadersInit | undefined} */
const headers = state.prerendering ? undefined : { 'cache-control': 'private, no-store' };

try {
/** @type {RemoteFunctionData} */
const data = {};

switch (internals.type) {
case 'query_live':
return handle_live_query(event, state, options, internals);

case 'query_batch': {
if (event.request.method !== 'POST') {
Expand Down Expand Up @@ -262,13 +281,7 @@ async function handle_remote_call_internal(event, state, options, manifest, id)

if (data._.issues) {
// special case — don't serialize refreshes/reconnects
return json(
/** @type {RemoteFunctionResponse} */ ({
type: 'result',
data: stringify(data)
}),
{ headers }
);
return result_response(data, headers);
}

break;
Expand Down Expand Up @@ -310,24 +323,12 @@ async function handle_remote_call_internal(event, state, options, manifest, id)

await collect_remote_data(data, event, state, options);

return json(
/** @type {RemoteFunctionResponse} */ ({
type: 'result',
data: stringify(data)
}),
{ headers }
);
return result_response(data, headers);
} catch (error) {
if (error instanceof Redirect) {
const data = await collect_remote_data({ redirect: error.location }, event, state, options);

return json(
/** @type {RemoteFunctionResponse} */ ({
type: 'result',
data: stringify(data)
}),
{ headers }
);
return result_response(data, headers);
}

const transformed = await handle_error_and_jsonify(event, state, options, error);
Expand Down
Loading
Loading