refactor(spf): drive live reload via a RecurringRunner instead of an epoch signal (WIP)

Replace the signal-as-event live-reload scheduler with a runner-driven model.
The `resolveTrack` loader schedules its resolve work on a new `RecurringRunner`
that re-runs the task on an injected `reschedule` policy; the separate
`scheduleTrackReload` behavior and its per-type reload-epoch signals are deleted.

Core (`core/tasks`):
- `Task.run()` is now memoized (runs once, shares the result across calls) and
  gains `clone()` (fresh, pending, structurally identical) — added to `TaskLike`.
- `RecurringRunner`: single-slot, id-keyed (dedup same id / abort-and-replace on
  new id), time-free. Each cycle runs `Promise.all([task.run(), reschedule(task,
  previous, signal)])` and re-runs a `clone()` while reschedule resolves `true`.
- `Reschedule<T> = (task, previous, signal) => PromiseLike<boolean>` — invoked
  concurrently with the run, observes it via the memoized `run()`, owns its delay.
- `delayedReschedule(cadence)` builds a Reschedule from a pure ms-cadence fn,
  start-anchored (subtracts the run's elapsed) so reloads are measured from
  load-start per RFC 8216 §6.3.4, preserving half-on-unchanged.

Supporting:
- `@videojs/utils/time`: add cancellable `sleep(ms, signal)`.
- `media/hls/reload-policy`: `mediaPlaylistReloadDelay` (pure cadence; relocated
  scheduler logic — target-duration, half-on-unchanged, stop-on-ENDLIST, retry).
- `resolve-track`: baked universal completeness gate + injected `reschedule`;
  engine composes `delayedReschedule(mediaPlaylistReloadDelay)`.

WIP: not fully validated end-to-end against a live stream through this path.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
Christian Pillsbury
2026-06-25 09:59:24 -07:00
co-authored by Claude Opus 4.8
parent 1d51582d8c
commit 557d8bd72f
16 changed files with 926 additions and 624 deletions
@@ -0,0 +1,35 @@
import { sleep } from '@videojs/utils/time';
import type { Reschedule } from './task';
/**
* Build a {@link Reschedule} from a pure cadence function — the common
* timer-based, *start-anchored* implementation.
*
* Invoked concurrently with the run, it observes the result, then waits
* `cadence(current, previous)` milliseconds **measured from when it was invoked**
* (≈ the run's start): it subtracts the run's own elapsed time, so consecutive
* runs begin one cadence apart regardless of how long each run takes (per
* RFC 8216 §6.3.4's "measured from the last time the client began loading"). If
* the run takes longer than the cadence, the next run starts immediately.
*
* A `null` cadence stops the recurrence. An errored run passes `current` as
* `undefined`, so the cadence function can choose to retry (return a delay) or
* stop (return `null`).
*/
export function delayedReschedule<TValue>(
cadence: (current: TValue | undefined, previous: TValue | undefined) => number | null
): Reschedule<TValue> {
return async (task, previous, signal) => {
const startedAt = Date.now();
let current: TValue | undefined;
try {
current = await task.run();
} catch {
current = undefined;
}
const ms = cadence(current, previous);
if (ms === null) return false;
await sleep(Math.max(0, ms - (Date.now() - startedAt)), signal);
return true;
};
}
+175 -12
View File
@@ -46,8 +46,11 @@ export interface TaskLike<TValue = void, TError = unknown> {
readonly status: TaskStatus;
readonly value: DeepReadonly<TValue> | undefined;
readonly error: DeepReadonly<TError> | undefined;
/** Run the work, memoized: repeated calls share one execution + result. */
run(): Promise<TValue>;
abort(): void;
/** A fresh, structurally identical task (same work + id) in a pending state — for re-running. */
clone(): TaskLike<TValue, TError>;
}
/**
@@ -58,6 +61,11 @@ export interface TaskLike<TValue = void, TError = unknown> {
* propagates into the task's work without requiring the caller to track the
* task separately.
*
* `run()` is memoized: the work runs at most once per instance, and every call
* returns the same promise (so observers can `await run()` to read the result
* without re-triggering the work). To re-run the *same* work, take a `clone()` —
* a fresh instance with its own AbortController and a pending state.
*
* Ordering guarantee: `value` is written before `status` transitions to `'done'`;
* `error` is written before `status` transitions to `'error'`. Any reader
* observing `status === 'done'` is guaranteed `value` is already present.
@@ -65,17 +73,20 @@ export interface TaskLike<TValue = void, TError = unknown> {
export class Task<TValue = void, TError = unknown> implements TaskLike<TValue, TError> {
readonly id: string;
readonly #runFn: (signal: AbortSignal) => Promise<TValue>;
readonly #externalSignal: AbortSignal | undefined;
readonly #abortController = new AbortController();
readonly #signal: AbortSignal;
#status: TaskStatus = 'pending';
#value: TValue | undefined = undefined;
#error: TError | undefined = undefined;
#promise: Promise<TValue> | undefined = undefined;
constructor(runFn: (signal: AbortSignal) => Promise<TValue>, config?: TaskConfig) {
this.#runFn = runFn;
const rawId = config?.id;
this.id = typeof rawId === 'function' ? rawId() : (rawId ?? generateId());
this.#externalSignal = config?.signal;
this.#signal = config?.signal
? anyAbortSignal([this.#abortController.signal, config.signal])
: this.#abortController.signal;
@@ -93,23 +104,37 @@ export class Task<TValue = void, TError = unknown> implements TaskLike<TValue, T
return this.#error as DeepReadonly<TError> | undefined;
}
async run(): Promise<TValue> {
this.#status = 'running';
try {
const result = await this.#runFn(this.#signal);
this.#value = result; // value before status — ordering guarantee
this.#status = 'done';
return result;
} catch (e) {
this.#error = e as TError; // error before status — ordering guarantee
this.#status = 'error';
throw e;
}
run(): Promise<TValue> {
// Memoized: run the work once; repeated calls share the same promise (which
// resolves/rejects immediately once settled). Re-running needs a `clone()`.
this.#promise ??= (async () => {
this.#status = 'running';
try {
const result = await this.#runFn(this.#signal);
this.#value = result; // value before status — ordering guarantee
this.#status = 'done';
return result;
} catch (e) {
this.#error = e as TError; // error before status — ordering guarantee
this.#status = 'error';
throw e;
}
})();
return this.#promise;
}
abort(): void {
this.#abortController.abort();
}
/**
* A fresh task with the same work, id, and external signal, in a pending state
* (its own AbortController, no memoized result) — so it can be run again. Used
* to re-run structurally identical work (e.g. `RecurringRunner` reloads).
*/
clone(): Task<TValue, TError> {
return new Task<TValue, TError>(this.#runFn, { id: this.id, signal: this.#externalSignal });
}
}
// =============================================================================
@@ -285,3 +310,141 @@ export class SerialRunner {
this.abortAll();
}
}
// =============================================================================
// RecurringRunner
// =============================================================================
/**
* Decides whether — and *when* — a {@link RecurringRunner} re-runs its task.
* Invoked **concurrently with the run** (so the inter-run interval can be
* measured from when the run *started*, not when it finished), with:
* - `task` — the in-flight run, observable via the memoized `task.run()` (does
* not re-trigger work), e.g. to read its result for a cadence/stop decision.
* - `previous` — the prior successful run's value (`undefined` on the first),
* for decisions that compare consecutive results.
* - `signal` — aborts the wait (and the recurrence).
*
* Resolves `true` to re-run (after whatever delay it owns) or `false` to stop.
* The runner deals only in this awaitable verdict — *how* the delay is produced
* (a timer, a frame, an event) and *when* it's measured from live entirely in
* the reschedule function, so the runner itself knows nothing about time. See
* `delayedReschedule` for the common timer-based, start-anchored implementation.
*/
export type Reschedule<TValue> = (
task: TaskLike<TValue>,
previous: TValue | undefined,
signal: AbortSignal
) => PromiseLike<boolean>;
/**
* Runs a task, then re-runs it whenever a {@link Reschedule} function says to,
* until it says stop (or it's aborted) — the recurring sibling of
* {@link ConcurrentRunner} / {@link SerialRunner}, and like them it's handed a
* {@link TaskLike} to run.
*
* The runner has no notion of time: it just awaits whatever `reschedule`
* returns (a promise → re-run when it resolves; `null` → stop). With no
* reschedule it runs the task exactly once (the non-recurring default).
*
* Single-slot, keyed by task **id**: there is always at most one identified
* active task for re-running. Scheduling a task whose id matches the active one
* is a no-op — the existing recurrence keeps running (dedup by id). Scheduling a
* task with a *different* id aborts the prior task's in-flight run and pending
* reschedule, then takes over the slot (abort-and-replace) — the right shape
* when there's one logical unit of recurring work (e.g. reloading the *selected*
* track's media playlist).
*
* Each re-run is a fresh `clone()` of the task (since `Task.run()` is memoized —
* the same instance won't re-execute), carrying the same id so the slot's
* identity is stable across cycles. The run function should read any inputs that
* change between cycles at call time rather than capturing them once.
* `abortAll()` aborts the in-flight task and the pending reschedule; an aborted
* (or stopped) recurrence frees the slot, so a later schedule of the same id
* starts fresh.
*/
export class RecurringRunner<TValue = unknown> {
readonly #reschedule: Reschedule<TValue> | undefined;
// The identified active task — the one being (re)run. Held across cycles
// (including the inter-cycle wait); null once the recurrence stops/aborts.
#active: TaskLike<TValue, unknown> | null = null;
// Aborts the active recurrence: its in-flight run (via the task) and the
// pending reschedule await (via the signal passed to `reschedule`).
#abort: AbortController | null = null;
#destroyed = false;
constructor(reschedule?: Reschedule<TValue>) {
this.#reschedule = reschedule;
}
schedule(task: TaskLike<TValue, unknown>): void {
if (this.#destroyed) return;
// Dedup by id: this id is already the active re-run target, so the existing
// recurrence continues uninterrupted (don't restart it).
if (this.#active?.id === task.id) return;
// Different id supersedes: abort the prior recurrence, then take over as the
// one identified active task.
this.#cancel();
this.#active = task;
const ac = new AbortController();
this.#abort = ac;
void this.#loop(task, ac);
}
async #loop(task: TaskLike<TValue, unknown>, ac: AbortController): Promise<void> {
const signal = ac.signal;
let current = task;
let previous: TValue | undefined;
while (!signal.aborted) {
let result: TValue | undefined;
let again = false;
try {
// Start the run and the reschedule together: reschedule observes the run
// (via the memoized `task.run()`) and owns the inter-run delay, so the
// interval can be measured from the run's *start*. `Promise.all` waits
// for both — the result (the next cycle's `previous`) and the verdict +
// delay. An errored run resolves to `undefined`, leaving the retry/stop
// choice to reschedule.
[result, again] = await Promise.all([
current.run().then(
(value) => value,
() => undefined
),
this.#reschedule ? this.#reschedule(current, previous, signal) : Promise.resolve(false),
]);
} catch {
return; // reschedule rejected (aborted during its delay) — stop without freeing
}
if (signal.aborted || !again) break;
// Only a successful run advances the comparison baseline.
if (result !== undefined) previous = result;
// Re-run the same work as a fresh task — `run()` is memoized, so re-running
// `current` wouldn't re-execute. The clone keeps the id (stable slot
// identity); track it as active so an abort hits the in-flight instance.
current = current.clone();
this.#active = current;
}
// Loop ended on its own terms (stop / aborted-but-not-cancelled): free the
// slot iff this loop still owns it (a supersede/abortAll swapped #abort).
if (this.#abort === ac) {
this.#active = null;
this.#abort = null;
}
}
#cancel(): void {
this.#abort?.abort();
this.#abort = null;
this.#active?.abort();
this.#active = null;
}
abortAll(): void {
this.#cancel();
}
destroy(): void {
this.#destroyed = true;
this.abortAll();
}
}
@@ -0,0 +1,95 @@
import { afterEach, describe, expect, it, vi } from 'vitest';
import { delayedReschedule } from '../delayed-reschedule';
import { Task } from '../task';
afterEach(() => {
vi.useRealTimers();
});
describe('delayedReschedule', () => {
it('observes the run result + previous, then waits the cadence before resolving true', async () => {
vi.useFakeTimers();
const task = new Task<number>(async () => 5);
const cadence = vi.fn(() => 100);
const reschedule = delayedReschedule<number>(cadence);
let resolved: boolean | undefined;
const done = reschedule(task, 4, new AbortController().signal).then((v) => {
resolved = v;
});
await vi.advanceTimersByTimeAsync(0); // run settles → cadence consulted
expect(cadence).toHaveBeenCalledWith(5, 4);
expect(resolved).toBeUndefined(); // still waiting out the cadence
await vi.advanceTimersByTimeAsync(100);
await done;
expect(resolved).toBe(true);
});
it('start-anchors: subtracts the run elapsed from the cadence', async () => {
vi.useFakeTimers();
// A run that takes 40ms; cadence 100 → next run ~60ms after the run settles
// (so the interval is 100ms measured from the run's start).
const task = new Task<number>(async () => {
await new Promise<void>((resolve) => setTimeout(resolve, 40));
return 1;
});
const reschedule = delayedReschedule<number>(() => 100);
let resolved = false;
const done = reschedule(task, undefined, new AbortController().signal).then(() => {
resolved = true;
});
await vi.advanceTimersByTimeAsync(40); // run completes; 40ms elapsed
await vi.advanceTimersByTimeAsync(59);
expect(resolved).toBe(false); // 99ms from start — not yet
await vi.advanceTimersByTimeAsync(1);
await done;
expect(resolved).toBe(true); // 100ms from start
});
it('resolves false (stop) without waiting when the cadence returns null', async () => {
vi.useFakeTimers();
const task = new Task<number>(async () => 1);
const reschedule = delayedReschedule<number>(() => null);
await expect(reschedule(task, undefined, new AbortController().signal)).resolves.toBe(false);
});
it('passes undefined to the cadence when the run errors (so it can retry)', async () => {
vi.useFakeTimers();
const task = new Task<number>(async () => {
throw new Error('boom');
});
const cadence = vi.fn(() => 50); // retry on error
const reschedule = delayedReschedule<number>(cadence);
let resolved: boolean | undefined;
const done = reschedule(task, undefined, new AbortController().signal).then((v) => {
resolved = v;
});
await vi.advanceTimersByTimeAsync(0);
expect(cadence).toHaveBeenCalledWith(undefined, undefined);
await vi.advanceTimersByTimeAsync(50);
await done;
expect(resolved).toBe(true);
});
it('rejects when aborted during the wait', async () => {
vi.useFakeTimers();
const task = new Task<number>(async () => 1);
const reschedule = delayedReschedule<number>(() => 100);
const ac = new AbortController();
const done = reschedule(task, undefined, ac.signal);
await vi.advanceTimersByTimeAsync(0); // run settles → into the wait
ac.abort(new DOMException('Aborted', 'AbortError'));
await expect(done).rejects.toBeInstanceOf(DOMException);
});
});
+244 -1
View File
@@ -1,5 +1,5 @@
import { describe, expect, it, vi } from 'vitest';
import { ConcurrentRunner, SerialRunner, Task } from '../task';
import { ConcurrentRunner, RecurringRunner, type Reschedule, SerialRunner, Task } from '../task';
// =============================================================================
// Task
@@ -170,6 +170,66 @@ describe('Task', () => {
expect(t1.id).not.toBe(t2.id);
});
});
describe('memoization', () => {
it('runs the work at most once and shares the result across run() calls', async () => {
const work = vi.fn(async () => 42);
const task = new Task(work);
const [a, b] = await Promise.all([task.run(), task.run()]);
const c = await task.run(); // after settle
expect(work).toHaveBeenCalledTimes(1);
expect([a, b, c]).toEqual([42, 42, 42]);
});
it('shares the rejection across run() calls', async () => {
const err = new Error('boom');
const work = vi.fn(async () => {
throw err;
});
const task = new Task<void, Error>(work);
await expect(task.run()).rejects.toBe(err);
await expect(task.run()).rejects.toBe(err);
expect(work).toHaveBeenCalledTimes(1);
});
});
describe('clone', () => {
it('produces a fresh, pending task with the same id and work', async () => {
let runs = 0;
const original = new Task<number>(async () => ++runs, { id: 'x' });
await original.run();
const cloned = original.clone();
expect(cloned).not.toBe(original);
expect(cloned.id).toBe('x');
expect(cloned.status).toBe('pending');
// The clone re-executes the same work (a fresh memoization).
await expect(cloned.run()).resolves.toBe(2);
expect(runs).toBe(2);
});
it('gives the clone an independent abort scope', async () => {
const signals: AbortSignal[] = [];
const original = new Task(async (signal) => {
signals.push(signal);
await new Promise<void>((resolve) => setTimeout(resolve, 10));
});
const cloned = original.clone();
const run = cloned.run();
cloned.abort();
await run;
// Aborting the clone aborts only the clone's signal, not the original's.
original.abort();
expect(signals).toHaveLength(1);
expect(signals[0]?.aborted).toBe(true);
});
});
});
// =============================================================================
@@ -573,3 +633,186 @@ describe('SerialRunner', () => {
expect(taskSignal?.aborted).toBe(true);
});
});
describe('RecurringRunner', () => {
/** A reschedule that parks forever, rejecting only when its signal aborts — so a
* recurrence stays "live" (awaiting) until superseded or aborted. */
const parkUntilAborted: Reschedule<number> = (_task, _previous, signal) =>
new Promise<boolean>((_resolve, reject) => {
signal.addEventListener('abort', () => reject(signal.reason), { once: true });
});
/** Flush pending macrotasks so "did NOT happen" assertions are meaningful. */
const flush = () => new Promise((resolve) => setTimeout(resolve, 0));
it('runs the task once when no reschedule is supplied', async () => {
let runs = 0;
const task = new Task<number>(async () => ++runs, { id: 'x' });
const runner = new RecurringRunner<number>();
runner.schedule(task);
await vi.waitFor(() => expect(runs).toBe(1));
await flush();
expect(runs).toBe(1);
});
it('re-runs (a clone of) the task while reschedule resolves true, stops on false', async () => {
let runs = 0;
// The runner clones per cycle; the clones share this run fn's `runs` counter.
const task = new Task<number>(async () => ++runs, { id: 'x' });
const runner = new RecurringRunner<number>(async (t) => (await t.run()) < 3); // continue while < 3
runner.schedule(task);
await vi.waitFor(() => expect(runs).toBe(3));
await flush();
expect(runs).toBe(3);
});
it('observes the run and receives the previous successful value', async () => {
let n = 0;
const task = new Task<number>(async () => ++n, { id: 'x' });
const seen: Array<[number, number | undefined]> = [];
const runner = new RecurringRunner<number>(async (t, previous) => {
const current = await t.run(); // observe via the memoized run
seen.push([current, previous]);
return current < 2;
});
runner.schedule(task);
await vi.waitFor(() => expect(seen.length).toBe(2));
expect(seen).toEqual([
[1, undefined],
[2, 1],
]);
});
it('keeps the loop alive across a transient error (observed value is undefined)', async () => {
let n = 0;
const task = new Task<number>(async () => {
n += 1;
if (n === 1) throw new Error('boom');
return n;
});
// Retry on error (observed value undefined); keep going until a value reaches 3.
const runner = new RecurringRunner<number>(async (t) => {
const current = await t.run().catch(() => undefined);
return current === undefined || current < 3;
});
runner.schedule(task);
await vi.waitFor(() => expect(n).toBe(3));
runner.destroy();
});
it('aborts an in-flight run when superseded by a new id', async () => {
let aborted = false;
const slow = new Task<number>(
(signal) =>
new Promise<number>((_resolve, reject) => {
signal.addEventListener('abort', () => {
aborted = true;
reject(new DOMException('Aborted', 'AbortError'));
});
}),
{ id: 'a' }
);
let ranB = false;
const taskB = new Task<number>(
async () => {
ranB = true;
return 1;
},
{ id: 'b' }
);
const runner = new RecurringRunner<number>(parkUntilAborted);
runner.schedule(slow); // parks mid-run, listening for abort
runner.schedule(taskB); // new id → abort slow's in-flight run, run B
expect(aborted).toBe(true);
await vi.waitFor(() => expect(ranB).toBe(true));
runner.destroy();
});
it('ignores a schedule with the same id — the existing recurrence keeps running', async () => {
let runsA = 0;
const taskA = new Task<number>(async () => ++runsA, { id: 'x' });
let runsB = 0;
const taskB = new Task<number>(async () => ++runsB, { id: 'x' }); // same id
const runner = new RecurringRunner<number>(parkUntilAborted);
runner.schedule(taskA);
await vi.waitFor(() => expect(runsA).toBe(1)); // A ran once, parked
runner.schedule(taskB); // same id while A is live → ignored
await flush();
expect(runsB).toBe(0);
expect(runsA).toBe(1);
runner.destroy();
});
it('a new id aborts the prior recurrence and takes over the slot', async () => {
let runsA = 0;
const taskA = new Task<number>(async () => ++runsA, { id: 'a' });
let runsB = 0;
const taskB = new Task<number>(async () => ++runsB, { id: 'b' });
const runner = new RecurringRunner<number>(parkUntilAborted);
runner.schedule(taskA);
await vi.waitFor(() => expect(runsA).toBe(1)); // A parked
runner.schedule(taskB); // new id → abort A, run B
await vi.waitFor(() => expect(runsB).toBe(1));
await flush();
expect(runsA).toBe(1); // A did not re-run
runner.destroy();
});
it('frees the slot when a recurrence stops, so the same id can start fresh', async () => {
let runs = 0;
// Fresh instances per schedule (the real pattern — callers build a new task
// each time); a shared counter observes runs across both.
const make = () => new Task<number>(async () => ++runs, { id: 'x' });
const runner = new RecurringRunner<number>(async () => false); // stop after first run
runner.schedule(make());
await vi.waitFor(() => expect(runs).toBe(1));
await flush(); // let the loop run reschedule → false → free the slot
runner.schedule(make()); // not deduped (recurrence ended) → runs fresh
await vi.waitFor(() => expect(runs).toBe(2));
runner.destroy();
});
it('abortAll stops the recurrence; no re-run', async () => {
let runs = 0;
const task = new Task<number>(async () => ++runs, { id: 'x' });
const runner = new RecurringRunner<number>(parkUntilAborted);
runner.schedule(task);
await vi.waitFor(() => expect(runs).toBe(1));
runner.abortAll();
await flush();
expect(runs).toBe(1);
});
it('does not run after destroy', async () => {
let runs = 0;
const task = new Task<number>(async () => ++runs, { id: 'x' });
const runner = new RecurringRunner<number>(async () => true);
runner.destroy();
runner.schedule(task);
await flush();
expect(runs).toBe(0);
});
});