more fork fixes

entangle-batches-2
Simon Holthausen 2 months ago
parent 078e0c54e1
commit 485dda33ad
No known key found for this signature in database

@ -270,7 +270,8 @@ export class Batch {
/** /**
* Async and block effects that ran or were proven clean inside this fork. * Async and block effects that ran or were proven clean inside this fork.
* When the fork commits, only these effects can retain their speculative results * The version is that of the latest real-world execution when the effect was
* validated. Fork executions do not advance it, so forks remain independent
* @type {Map<Effect, number> | null} * @type {Map<Effect, number> | null}
*/ */
fork_effects = null; fork_effects = null;
@ -531,15 +532,17 @@ export class Batch {
return false; return false;
} }
/** @param {Effect} effect */ /**
record_effect(effect) { * @param {Effect} effect
* @param {boolean} [ran]
*/
record_effect(effect, ran = false) {
if ((effect.f & (ASYNC | BLOCK_EFFECT)) === 0) return; if ((effect.f & (ASYNC | BLOCK_EFFECT)) === 0) return;
var version = effect_version++;
effect_versions.set(effect, version);
if (this.is_fork) { if (this.is_fork) {
(this.fork_effects ??= new Map()).set(effect, version); (this.fork_effects ??= new Map()).set(effect, effect_versions.get(effect) ?? 0);
} else if (ran) {
effect_versions.set(effect, effect_version++);
} }
} }
@ -547,7 +550,7 @@ export class Batch {
if (this.fork_effects === null) return; if (this.fork_effects === null) return;
for (const [effect, version] of this.fork_effects) { for (const [effect, version] of this.fork_effects) {
if ((effect.f & DESTROYED) !== 0 || version !== effect_versions.get(effect)) { if ((effect.f & DESTROYED) !== 0 || version !== (effect_versions.get(effect) ?? 0)) {
this.#dirty_effects?.delete(effect); this.#dirty_effects?.delete(effect);
this.#maybe_dirty_effects?.delete(effect); this.#maybe_dirty_effects?.delete(effect);
} }
@ -556,7 +559,7 @@ export class Batch {
for (const [effect, version] of this.fork_effects) { for (const [effect, version] of this.fork_effects) {
if ((effect.f & DESTROYED) !== 0) continue; if ((effect.f & DESTROYED) !== 0) continue;
if (version !== effect_versions.get(effect)) { if (version !== (effect_versions.get(effect) ?? 0)) {
var owner = effect.batch; var owner = effect.batch;
while (owner !== null && owner.merged_into !== null) owner = owner.merged_into; while (owner !== null && owner.merged_into !== null) owner = owner.merged_into;
@ -1079,6 +1082,8 @@ export class Batch {
} }
discard() { discard() {
this.fork_effects = null;
if (this.#discard_callbacks !== null) { if (this.#discard_callbacks !== null) {
for (const fn of this.#discard_callbacks) fn(this); for (const fn of this.#discard_callbacks) fn(this);
this.#discard_callbacks = null; this.#discard_callbacks = null;
@ -1097,6 +1102,44 @@ export class Batch {
} }
#commit() { #commit() {
/** @type {Map<Batch, Effect[]> | null} */
var forks = null;
// A real-world run supersedes the same effect's speculative validation.
// Once that run commits, revalidate each affected fork against the latest
// real state, while keeping the fork's own writes isolated.
for (var fork = first_batch; fork !== null; fork = fork.#next) {
if (!fork.is_fork || fork.fork_effects === null) continue;
for (const [effect, version] of fork.fork_effects) {
if (
version === (effect_versions.get(effect) ?? 0) ||
(effect.f & (DESTROYED | INERT)) !== 0
) {
continue;
}
var owner = effect.batch && effect.batch.#resolved();
if (owner !== this) continue;
set_signal_status(effect, DIRTY);
fork.schedule(effect);
var effects = (forks ??= new Map()).get(fork);
if (effects === undefined) forks.set(fork, [effect]);
else effects.push(effect);
}
}
if (forks !== null) {
for (const [fork, effects] of forks) {
queue_micro_task(() => {
if (!fork.linked) return;
for (const effect of effects) set_signal_status(effect, DIRTY);
fork.flush();
});
}
}
if (this.stale_readers === null) return; if (this.stale_readers === null) return;
var readers = this.stale_readers; var readers = this.stale_readers;
@ -1226,7 +1269,9 @@ export class Batch {
// changes, we should only see values that have been committed. Overlapping // changes, we should only see values that have been committed. Overlapping
// batches are merged unless one is waiting behind a sealed predecessor // batches are merged unless one is waiting behind a sealed predecessor
for (let other = first_batch; other !== null; other = other.#next) { for (let other = first_batch; other !== null; other = other.#next) {
if (other === batch || other.is_fork) continue; // Forks are based on the latest real world, not an independently
// held-back snapshot of it. Other forks remain separate worlds.
if (other === batch || other.is_fork || batch.is_fork) continue;
for (const [source, previous] of other.previous) { for (const [source, previous] of other.previous) {
if (!values.has(source)) { if (!values.has(source)) {
@ -1255,12 +1300,10 @@ export class Batch {
schedule(effect) { schedule(effect) {
last_scheduled_effect = effect; last_scheduled_effect = effect;
if ((effect.f & (ASYNC | BLOCK_EFFECT)) !== 0) { if (this.is_fork && (effect.f & (ASYNC | BLOCK_EFFECT)) !== 0) {
effect_versions.set(effect, effect_version++); this.fork_effects?.delete(effect);
} }
if (this.is_fork) this.fork_effects?.delete(effect);
if (this.claim(effect)) { if (this.claim(effect)) {
this.#dirty_effects ??= new Set(); this.#dirty_effects ??= new Set();
this.#maybe_dirty_effects ??= new Set(); this.#maybe_dirty_effects ??= new Set();
@ -1596,7 +1639,8 @@ function mark_committed_reactions(value, batch, marked, status) {
while (owner !== null && owner.merged_into !== null) owner = owner.merged_into; while (owner !== null && owner.merged_into !== null) owner = owner.merged_into;
var fork_version = batch.fork_effects?.get(effect); var fork_version = batch.fork_effects?.get(effect);
var validated = fork_version !== undefined && fork_version === effect_versions.get(effect); var validated =
fork_version !== undefined && fork_version === (effect_versions.get(effect) ?? 0);
var stale = var stale =
batch.stale_readers?.has(effect) === true || owner?.stale_readers?.has(effect) === true; batch.stale_readers?.has(effect) === true || owner?.stale_readers?.has(effect) === true;
@ -1734,9 +1778,9 @@ export function fork(fn) {
// Marking can entangle the fork with newer work and replace values in // Marking can entangle the fork with newer work and replace values in
// `current`, so only write through once the final world is known. // `current`, so only write through once the final world is known.
for (var [source, value] of batch.current) { for (var [signal, value] of batch.current) {
source.v = value; signal.v = value;
source.wv = increment_write_version(); signal.wv = increment_write_version();
} }
// Effects retained from the speculative world now belong to the real // Effects retained from the speculative world now belong to the real

@ -487,7 +487,7 @@ export function update_effect(effect) {
var teardown = update_reaction(effect); var teardown = update_reaction(effect);
effect.teardown = typeof teardown === 'function' ? teardown : null; effect.teardown = typeof teardown === 'function' ? teardown : null;
effect.wv = write_version; effect.wv = write_version;
batch?.record_effect(effect); batch?.record_effect(effect, true);
// In DEV, increment versions of any sources that were written to during the effect, // In DEV, increment versions of any sources that were written to during the effect,
// so that they are correctly marked as dirty when the effect re-runs // so that they are correctly marked as dirty when the effect re-runs

@ -0,0 +1,35 @@
import { tick } from 'svelte';
import { test } from '../../test';
export default test({
async test({ assert, target, logs }) {
await tick();
const [x, y, shift, pop, commit] = target.querySelectorAll('button');
const [p] = target.querySelectorAll('p');
logs.length = 0;
x.click();
await tick();
assert.deepEqual(logs, ['called with 1,0']);
logs.length = 0;
y.click();
await tick();
assert.deepEqual(logs, ['called with 0,1']);
assert.htmlEqual(p.innerHTML, '0');
logs.length = 0;
pop.click();
await tick();
assert.deepEqual(logs, ['called with 1,1']);
assert.htmlEqual(p.innerHTML, '1');
logs.length = 0;
commit.click();
await tick();
pop.click();
await tick();
assert.deepEqual(logs, []);
assert.htmlEqual(p.innerHTML, '2');
}
});

@ -0,0 +1,23 @@
<script>
import { fork } from 'svelte';
let x = $state(0);
let y = $state(0);
let f;
const deferred = [];
function delay(_, value) {
if (!value) return value;
return new Promise((resolve) => deferred.push(() => resolve(value)));
}
</script>
<button onclick={() => {f = fork(() => x++)}}>x</button>
<button onclick={() => y++}>y</button>
<button onclick={() => deferred.shift()?.()}>shift</button>
<button onclick={() => deferred.pop()?.()}>pop</button>
<button onclick={() => f.commit()}>commit</button>
<p>{await delay(console.log('called with ' + x + ',' + y), x + y)}</p>

@ -0,0 +1,34 @@
import { tick } from 'svelte';
import { test } from '../../test';
export default test({
async test({ assert, target, logs }) {
await tick();
const [x, y, shift, pop, commit] = target.querySelectorAll('button');
const [p] = target.querySelectorAll('p');
logs.length = 0;
y.click();
await tick();
assert.deepEqual(logs, ['called with 0,1']);
logs.length = 0;
x.click();
await tick();
assert.deepEqual(logs, ['called with 1,1']);
assert.htmlEqual(p.innerHTML, '0');
logs.length = 0;
shift.click();
await tick();
assert.deepEqual(logs, []);
assert.htmlEqual(p.innerHTML, '1');
commit.click();
await tick();
pop.click();
await tick();
assert.deepEqual(logs, []);
assert.htmlEqual(p.innerHTML, '2');
}
});

@ -0,0 +1,23 @@
<script>
import { fork } from 'svelte';
let x = $state(0);
let y = $state(0);
let f;
const deferred = [];
function delay(_, value) {
if (!value) return value;
return new Promise((resolve) => deferred.push(() => resolve(value)));
}
</script>
<button onclick={() => {f = fork(() => x++)}}>x</button>
<button onclick={() => y++}>y</button>
<button onclick={() => deferred.shift()?.()}>shift</button>
<button onclick={() => deferred.pop()?.()}>pop</button>
<button onclick={() => f.commit()}>commit</button>
<p>{await delay(console.log('called with ' + x + ',' + y), x + y)}</p>

@ -27,6 +27,11 @@ export default test({
await tick(); await tick();
assert.htmlEqual(p.innerHTML, 'a 1 | b 0 | c 1'); assert.htmlEqual(p.innerHTML, 'a 1 | b 0 | c 1');
// resolve the fork revalidation triggered by the real batch settling
pop.click();
await tick();
assert.htmlEqual(p.innerHTML, 'a 1 | b 0 | c 1');
commit.click(); commit.click();
await tick(); await tick();
assert.htmlEqual(p.innerHTML, 'a 1 | b 1 | c 1'); assert.htmlEqual(p.innerHTML, 'a 1 | b 1 | c 1');

@ -23,6 +23,10 @@ export default test({
await tick(); await tick();
assert.htmlEqual(p.innerHTML, 'a 0 | b 1 | c 0 | d 1'); assert.htmlEqual(p.innerHTML, 'a 0 | b 1 | c 0 | d 1');
pop.click();
await tick();
assert.htmlEqual(p.innerHTML, 'a 0 | b 1 | c 0 | d 1');
pop.click(); pop.click();
await tick(); await tick();
assert.htmlEqual(p.innerHTML, 'a 1 | b 1 | c 1 | d 1'); assert.htmlEqual(p.innerHTML, 'a 1 | b 1 | c 1 | d 1');
@ -35,6 +39,11 @@ export default test({
await tick(); await tick();
assert.htmlEqual(p.innerHTML, 'a 1 | b 1 | c 1 | d 1'); assert.htmlEqual(p.innerHTML, 'a 1 | b 1 | c 1 | d 1');
// resolve the fork revalidation triggered by the real batch settling
pop.click();
await tick();
assert.htmlEqual(p.innerHTML, 'a 1 | b 1 | c 1 | d 1');
commit.click(); commit.click();
await tick(); await tick();
assert.htmlEqual(p.innerHTML, 'a 1 | b 1 | c 1 | d 1'); assert.htmlEqual(p.innerHTML, 'a 1 | b 1 | c 1 | d 1');

@ -35,6 +35,13 @@ export default test({
await tick(); await tick();
assert.htmlEqual(p.innerHTML, 'a 1 | b 1 | c 1 | d 1'); assert.htmlEqual(p.innerHTML, 'a 1 | b 1 | c 1 | d 1');
// resolve the fork revalidations triggered by the real batches settling
shift.click();
await tick();
shift.click();
await tick();
assert.htmlEqual(p.innerHTML, 'a 1 | b 1 | c 1 | d 1');
commit.click(); commit.click();
await tick(); await tick();
assert.htmlEqual(p.innerHTML, 'a 1 | b 1 | c 1 | d 1'); assert.htmlEqual(p.innerHTML, 'a 1 | b 1 | c 1 | d 1');

Loading…
Cancel
Save