feat: keep client changes to a `Warp` separate from the values being streamed in

elliott/warp-streaming
Elliott Johnson 13 hours ago
parent 8da518d09c
commit b915fa0c3b
No known key found for this signature in database

@ -159,6 +159,8 @@ response.end('</body></html>');
If you stop iterating over `tail` early, the background work is stopped. Since the streamed `<script>` tags can't be known in advance, streaming can't be used with `csp: { hash: true }` — use a `nonce` instead.
You can still change a `Warp` on the client while values are streaming in. Your changes take precedence over the values from the server, including ones that arrive later — for example, after `warp.clear()`, values streamed in afterwards won't be visible. This means that a component which is still waiting for data from the server will load it again on the client instead.
## CSP
`Warp` adds an inline `<script>` block to the `head` returned from `render`. If you're using [Content Security Policy](https://developer.mozilla.org/en-US/docs/Web/HTTP/Guides/CSP) (CSP), this script will likely fail to run. You can provide a `nonce` to `render`:

@ -1,6 +1,7 @@
import type { Store, WarpKey } from '#shared';
import { STATE_SYMBOL } from './constants.js';
import type { Batch } from './reactivity/batch.js';
import type { Overlay } from './warp.js';
import type { Effect, Source, Value } from './reactivity/types.js';
declare global {
@ -8,6 +9,10 @@ declare global {
__svelte?: {
/** `Warp` values, keyed by the `Warp`'s id */
w?: Map<string, Map<WarpKey, unknown>>;
/** The number of `Warp` streams that haven't finished yet */
s?: number;
/** Changes made to `Warp`s while values were streaming in, keyed by the `Warp`'s id */
m?: Map<string, Overlay>;
};
}
}

@ -27,13 +27,18 @@ export class Warp {
this.#id = id;
}
/** @returns {Map<K, V>} */
/**
* The values for this `Warp`. If changes were made to them while values were still streaming
* in from the server, and the stream has since finished, they're applied first
* @returns {Map<K, V>}
*/
#values() {
if (!async_mode_flag) {
e.experimental_async_required('Warp');
}
const store = ((window.__svelte ??= {}).w ??= new Map());
const svelte = (window.__svelte ??= {});
const store = (svelte.w ??= new Map());
let values = store.get(this.#id);
if (values === undefined) {
@ -41,15 +46,47 @@ export class Warp {
store.set(this.#id, values);
}
const overlay = svelte.m?.get(this.#id);
if (overlay !== undefined && !is_streaming()) {
overlay.apply(values);
svelte.m?.delete(this.#id);
}
return /** @type {Map<K, V>} */ (values);
}
/**
* While values are still streaming in from the server, changes are made to an overlay rather
* than the values themselves, so that they can't interfere with the stream (or vice versa)
* @param {boolean} create
* @returns {Overlay | null}
*/
#overlay(create) {
if (!is_streaming()) return null;
const svelte = /** @type {NonNullable<Window['__svelte']>} */ (window.__svelte);
let overlay = svelte.m?.get(this.#id) ?? null;
if (overlay === null && create) {
overlay = new Overlay();
(svelte.m ??= new Map()).set(this.#id, overlay);
}
return overlay;
}
/**
* @param {K} key
* @returns {V | undefined}
*/
get(key) {
return this.#values().get(key);
const values = this.#values();
const overlay = this.#overlay(false);
return /** @type {V | undefined} */ (
overlay === null ? values.get(key) : overlay.get(values, key)
);
}
/**
@ -57,7 +94,10 @@ export class Warp {
* @returns {boolean}
*/
has(key) {
return this.#values().has(key);
const values = this.#values();
const overlay = this.#overlay(false);
return overlay === null ? values.has(key) : overlay.has(values, key);
}
/**
@ -66,7 +106,15 @@ export class Warp {
* @returns {this}
*/
set(key, value) {
this.#values().set(key, value);
const values = this.#values();
const overlay = this.#overlay(true);
if (overlay === null) {
values.set(key, value);
} else {
overlay.set(key, value);
}
return this;
}
@ -77,9 +125,15 @@ export class Warp {
* @returns {V}
*/
getOrInsert(key, value) {
if (!is_streaming()) {
return get_or_insert(this.#values(), key, value);
}
if (this.has(key)) return /** @type {V} */ (this.get(key));
this.set(key, value);
return value;
}
/**
* Returns the value for `key` if it exists, otherwise calls `callback` and adds the result.
* @param {K} key
@ -87,19 +141,36 @@ export class Warp {
* @returns {V}
*/
getOrInsertComputed(key, callback) {
if (!is_streaming()) {
return get_or_insert_computed(this.#values(), key, callback);
}
if (this.has(key)) return /** @type {V} */ (this.get(key));
const value = callback(key);
this.set(key, value);
return value;
}
/**
* @param {K} key
* @returns {boolean}
*/
delete(key) {
return this.#values().delete(key);
const values = this.#values();
const overlay = this.#overlay(true);
return overlay === null ? values.delete(key) : overlay.delete(values, key);
}
clear() {
this.#values().clear();
const values = this.#values();
const overlay = this.#overlay(true);
if (overlay === null) {
values.clear();
} else {
overlay.clear();
}
}
/**
@ -107,30 +178,170 @@ export class Warp {
* @param {any} [this_arg]
*/
forEach(callback, this_arg) {
this.#values().forEach(callback, this_arg);
for (const [key, value] of this.entries()) {
callback.call(this_arg, value, key, this);
}
}
/** @returns {IterableIterator<[K, V]>} */
entries() {
return this.#values().entries();
const values = this.#values();
const overlay = this.#overlay(false);
return overlay === null
? values.entries()
: /** @type {IterableIterator<[K, V]>} */ (overlay.entries(values));
}
keys() {
return this.#values().keys();
/** @returns {IterableIterator<K>} */
*keys() {
for (const [key] of this.entries()) yield key;
}
values() {
return this.#values().values();
/** @returns {IterableIterator<V>} */
*values() {
for (const [, value] of this.entries()) yield value;
}
get size() {
return this.#values().size;
const values = this.#values();
const overlay = this.#overlay(false);
if (overlay === null) return values.size;
let size = 0;
for (const _ of overlay.entries(values)) size += 1;
return size;
}
[Symbol.iterator]() {
return this.#values()[Symbol.iterator]();
return this.entries();
}
get [Symbol.toStringTag]() {
return 'Warp';
}
}
/**
* Whether values are still streaming in from the server. Streaming can only happen while the document is
* still loading, so if the connection drops before the stream finishes, this still becomes `false` eventually
*/
function is_streaming() {
return (window.__svelte?.s ?? 0) > 0 && document.readyState === 'loading';
}
const TOMBSTONE = Symbol('deleted');
/**
* Changes made to a `Warp` while values are still streaming in from the server. They take precedence
* over the values from the server (including ones that arrive later), and are applied once the stream
* has finished.
*/
export class Overlay {
/** Whether the `Warp` was cleared, which hides all the values from the server */
cleared = false;
/**
* Values that were set (or deleted, as `TOMBSTONE`) since the `Warp` was last cleared
* @type {Map<WarpKey, unknown>}
*/
changes = new Map();
/**
* Keys from the server that were deleted and then set again, which moves them to the end
* @type {Set<WarpKey>}
*/
moved = new Set();
/**
* @param {Map<WarpKey, unknown>} values
* @param {WarpKey} key
*/
get(values, key) {
if (this.changes.has(key)) {
const value = this.changes.get(key);
return value === TOMBSTONE ? undefined : value;
}
return this.cleared ? undefined : values.get(key);
}
/**
* @param {Map<WarpKey, unknown>} values
* @param {WarpKey} key
*/
has(values, key) {
if (this.changes.has(key)) return this.changes.get(key) !== TOMBSTONE;
return !this.cleared && values.has(key);
}
/**
* @param {WarpKey} key
* @param {unknown} value
*/
set(key, value) {
if (this.changes.get(key) === TOMBSTONE) {
// like in a `Map`, a key that was deleted and then set again goes to the end
this.changes.delete(key);
this.moved.add(key);
}
this.changes.set(key, value);
}
/**
* @param {Map<WarpKey, unknown>} values
* @param {WarpKey} key
*/
delete(values, key) {
const existed = this.has(values, key);
// even if it doesn't exist yet, it may still arrive from the server, and must stay deleted
this.changes.set(key, TOMBSTONE);
this.moved.delete(key);
return existed;
}
clear() {
this.cleared = true;
this.changes.clear();
this.moved.clear();
}
/**
* @param {Map<WarpKey, unknown>} values
* @returns {Generator<[WarpKey, unknown]>}
*/
*entries(values) {
if (!this.cleared) {
for (const [key, value] of values) {
if (this.moved.has(key)) continue;
if (this.changes.has(key)) {
const change = this.changes.get(key);
if (change !== TOMBSTONE) yield [key, change];
} else {
yield [key, value];
}
}
}
for (const [key, value] of this.changes) {
if (value === TOMBSTONE) continue;
if (!this.cleared && !this.moved.has(key) && values.has(key)) continue;
yield [key, value];
}
}
/** @param {Map<WarpKey, unknown>} values */
apply(values) {
if (this.cleared) values.clear();
for (const [key, value] of this.changes) {
if (value === TOMBSTONE || this.moved.has(key)) values.delete(key);
if (value !== TOMBSTONE) values.set(key, value);
}
}
}

@ -0,0 +1,158 @@
import { afterEach, beforeEach, describe, expect, test, vi } from 'vitest';
import { Warp } from './warp.js';
import { disable_async_mode_flag, enable_async_mode_flag } from '../flags/index.js';
let window: {
__svelte?: { w?: Map<string, Map<unknown, unknown>>; s?: number; m?: Map<any, any> };
};
let document: { readyState: string };
/** the values the stream writes to */
function base() {
return window.__svelte!.w!.get('test')!;
}
beforeEach(() => {
enable_async_mode_flag();
window = {
__svelte: {
w: new Map([
[
'test',
new Map([
['a', 1],
['b', 2]
])
]
])
}
};
document = { readyState: 'loading' };
vi.stubGlobal('window', window);
vi.stubGlobal('document', document);
});
afterEach(() => {
disable_async_mode_flag();
vi.unstubAllGlobals();
});
describe('when not streaming', () => {
test('changes the values directly', () => {
const warp = new Warp('test');
warp.set('c', 3);
warp.delete('a');
expect([...base()]).toEqual([
['b', 2],
['c', 3]
]);
});
});
describe('while streaming', () => {
beforeEach(() => {
window.__svelte!.s = 1;
});
test('does not change the values the stream writes to', () => {
const warp = new Warp('test');
warp.set('a', 10);
warp.delete('b');
warp.set('c', 3);
expect([...base()]).toEqual([
['a', 1],
['b', 2]
]);
expect(warp.get('a')).toBe(10);
expect(warp.has('b')).toBe(false);
expect(warp.size).toBe(2);
expect([...warp]).toEqual([
['a', 10],
['c', 3]
]);
});
test('shows values that arrive later, unless they were changed', () => {
const warp = new Warp('test');
warp.set('late', 'client');
warp.delete('deleted');
base().set('late', 'server');
base().set('deleted', 'server');
base().set('new', 'server');
expect(warp.get('late')).toBe('client');
expect(warp.has('deleted')).toBe(false);
expect(warp.get('new')).toBe('server');
});
test('handles setting values after clearing', () => {
const warp = new Warp('test');
warp.set('c', 3);
warp.clear();
warp.set('b', 20);
// arrives after the clear, so is hidden by it
base().set('late', 'server');
expect(warp.has('a')).toBe(false);
expect(warp.get('b')).toBe(20);
expect(warp.has('c')).toBe(false);
expect(warp.has('late')).toBe(false);
expect([...warp.keys()]).toEqual(['b']);
});
test('keeps the order a Map would have', () => {
const warp = new Warp('test');
warp.set('c', 3);
warp.delete('a');
warp.set('a', 10);
warp.set('b', 20);
const expected = [
['b', 20],
['c', 3],
['a', 10]
];
expect([...warp]).toEqual(expected);
window.__svelte!.s = 0;
expect([...warp]).toEqual(expected);
expect([...base()]).toEqual(expected);
});
test('getOrInsertComputed uses the overlay', () => {
const warp = new Warp('test');
warp.delete('a');
expect(warp.getOrInsertComputed('a', () => 10)).toBe(10);
expect(warp.getOrInsertComputed('b', () => 20)).toBe(2);
expect(base().get('a')).toBe(1);
});
test('applies the changes once the stream finishes', () => {
const warp = new Warp('test');
warp.clear();
warp.set('c', 3);
window.__svelte!.s = 0;
expect([...warp]).toEqual([['c', 3]]);
expect([...base()]).toEqual([['c', 3]]);
expect(window.__svelte!.m!.size).toBe(0);
});
test('applies the changes once the document has loaded, even if the stream did not finish', () => {
const warp = new Warp('test');
warp.delete('a');
document.readyState = 'interactive';
expect(warp.has('a')).toBe(false);
expect([...base()]).toEqual([['b', 2]]);
});
});

@ -1040,9 +1040,15 @@ export class Renderer {
serialization_failed(serialization_error);
}
const streaming = late !== null && ready.next !== null;
const body = `
{
const w = (window.__svelte ??= {}).w ??= new Map();
const w = (window.__svelte ??= {}).w ??= new Map();${
// while this is above zero, changes made to the client's `Warp`s are kept separate,
// so they can't interfere with the values that are still being streamed in
streaming ? '\n\t\t\t\twindow.__svelte.s = (window.__svelte.s ?? 0) + 1;' : ''
}
for (const [id, values] of ${late === null ? head : `(${head})[0]`}) {
const existing = w.get(id);
@ -1068,10 +1074,14 @@ export class Renderer {
return {
head: `\n\t\t<script${csp_attr}>${body}</script>`,
tail:
late === null || !ready.next
? empty()
: stream_scripts(tail, ready.next, csp_attr, late.close)
tail: streaming
? stream_scripts(
tail,
/** @type {Promise<IteratorResult<string>>} */ (ready.next),
csp_attr,
late.close
)
: empty()
};
}
@ -1226,26 +1236,30 @@ function stream_scripts(tail, next, csp_attr, close) {
* @returns {AsyncGenerator<string>}
*/
async function* generate_scripts(tail, next, csp_attr) {
let done = false;
while (!done) {
let code = '';
let result = await next;
while (true) {
const result = await next;
if (result.done) return;
if (result.done) {
// tell the client that this stream has finished
code += 'window.__svelte.s -= 1;';
done = true;
break;
}
let code = result.value;
code += result.value;
next = tail.next();
while (true) {
const more = await Promise.race([
next,
new Promise((fulfil) => setTimeout(() => fulfil(MACROTASK), 0))
]);
if (more === MACROTASK) break;
const { done, value } = /** @type {IteratorResult<string>} */ (more);
if (done) break;
code += value;
next = tail.next();
result = /** @type {IteratorResult<string>} */ (more);
}
yield `<script${csp_attr}>${code}</script>`;

@ -318,7 +318,7 @@ describe('streaming', () => {
/** Collects the tail, and returns the revived values once everything has been evaluated */
async function revive_streamed(head: string, tail: AsyncIterable<string>) {
const window: { __svelte?: { w?: Map<string, Map<unknown, unknown>> } } = {};
const window: { __svelte?: { w?: Map<string, Map<unknown, unknown>>; s?: number } } = {};
const run = (html: string) => {
for (const [, script] of html.matchAll(/<script(?:\s[^>]*)?>([\s\S]*?)<\/script>/g)) {
new Function('window', script)(window);
@ -329,10 +329,14 @@ describe('streaming', () => {
const chunks: string[] = [];
for await (const chunk of tail) {
// the client is told that a stream is in progress until the last chunk
expect(window.__svelte?.s).toBe(1);
chunks.push(chunk);
run(chunk);
}
expect(window.__svelte?.s).toBe(0);
return { values: window.__svelte?.w, chunks };
}

@ -497,8 +497,11 @@ declare module 'svelte' {
clear(): void;
forEach(callback: (value: V, key: K, map: Map<K, V>) => void, this_arg?: any): void;
entries(): IterableIterator<[K, V]>;
keys(): IterableIterator<K>;
values(): IterableIterator<V>;
get size(): number;
[Symbol.iterator](): IterableIterator<[K, V]>;

Loading…
Cancel
Save