From 8ccdabc231c08bef03379a3ddd0f0145c18ad6fc Mon Sep 17 00:00:00 2001 From: Steven Vandevelde Date: Sat, 12 Jul 2025 14:40:52 +0200 Subject: [PATCH] feat: move queue state to worker --- .../blur/artwork-controller/_applet.astro | 4 +- src/pages/engine/queue/_applet.astro | 27 +- .../orchestrator/queue-audio/_applet.astro | 2 +- .../orchestrator/queue-tracks/_applet.astro | 1 + src/scripts/applet/common.ts | 236 +++++++++++------- src/scripts/common.ts | 64 ++++- src/scripts/engine/queue/worker.ts | 108 ++++---- src/scripts/output/common.ts | 4 +- src/scripts/signal.ts | 65 +++++ 9 files changed, 342 insertions(+), 169 deletions(-) create mode 100644 src/scripts/signal.ts diff --git a/src/pages/constituent/blur/artwork-controller/_applet.astro b/src/pages/constituent/blur/artwork-controller/_applet.astro index 98e0b7e..290ca70 100644 --- a/src/pages/constituent/blur/artwork-controller/_applet.astro +++ b/src/pages/constituent/blur/artwork-controller/_applet.astro @@ -742,11 +742,11 @@ import "@styles/diffuse/fonts.css"; } function previous() { - engine.queue.sendAction("unshift"); + engine.queue.sendAction("unshift", undefined, { worker: true }); } function next() { - engine.queue.sendAction("shift"); + engine.queue.sendAction("shift", undefined, { worker: true }); } controller.appendChild(Controls); diff --git a/src/pages/engine/queue/_applet.astro b/src/pages/engine/queue/_applet.astro index 5b3fd08..f2a1ef7 100644 --- a/src/pages/engine/queue/_applet.astro +++ b/src/pages/engine/queue/_applet.astro @@ -4,21 +4,24 @@ import type { State } from "./types.d.ts"; import { register } from "@scripts/applet/common"; - import { endpoint, SharedWorker, transfer } from "@scripts/common"; + import { endpoint, SharedWorker, sync, transfer } from "@scripts/common"; import manifest from "./_manifest.json"; //////////////////////////////////////////// // SETUP //////////////////////////////////////////// - const worker = endpoint( - new SharedWorker(new URL("../../../scripts/engine/queue/worker", import.meta.url), { - type: "module", - name: manifest.name, - }).port, - ); + const port = new SharedWorker(new URL("../../../scripts/engine/queue/worker", import.meta.url), { + type: "module", + name: manifest.name, + }).port; + + const worker = endpoint(port); // Register applet - const context = register({ worker }); + const context = register({ mode: "shared-worker", worker }); + + // Keep applet data with worker data in sync + sync(context, port); // Initial state context.data = { @@ -36,18 +39,18 @@ context.setActionHandler("unshift", unshift); async function add(items: Track[]) { - context.data = await worker.add(transfer(items)); + await worker.add(transfer(items)); } async function pool(items: Track[]) { - context.data = await worker.pool(transfer(items)); + await worker.pool(transfer(items)); } async function shift() { - context.data = await worker.shift(); + await worker.shift(); } async function unshift() { - context.data = await worker.unshift(); + await worker.unshift(); } diff --git a/src/pages/orchestrator/queue-audio/_applet.astro b/src/pages/orchestrator/queue-audio/_applet.astro index 726fe53..d3a60df 100644 --- a/src/pages/orchestrator/queue-audio/_applet.astro +++ b/src/pages/orchestrator/queue-audio/_applet.astro @@ -41,7 +41,7 @@ engine.audio, (data) => data.items[engine.queue.data.now?.id ?? Infinity]?.hasEnded ?? false, (hasEnded) => { - if (hasEnded) engine.queue.sendAction("shift"); + if (hasEnded) engine.queue.sendAction("shift", undefined, { worker: true }); }, ); diff --git a/src/pages/orchestrator/queue-tracks/_applet.astro b/src/pages/orchestrator/queue-tracks/_applet.astro index 1574c15..a9adb37 100644 --- a/src/pages/orchestrator/queue-tracks/_applet.astro +++ b/src/pages/orchestrator/queue-tracks/_applet.astro @@ -49,6 +49,7 @@ // Clear engine.queue.sendAction("pool", tracks, { timeoutDuration: 60000, + worker: true, }); }, ); diff --git a/src/scripts/applet/common.ts b/src/scripts/applet/common.ts index d3cd5e1..0647688 100644 --- a/src/scripts/applet/common.ts +++ b/src/scripts/applet/common.ts @@ -1,5 +1,5 @@ import type { Applet, AppletEvent, AppletScope } from "@web-applets/sdk"; -import type * as Comlink from "comlink"; +import * as Comlink from "comlink"; import { applets } from "@web-applets/sdk"; import { type ElementConfigurator, h } from "spellcaster/hyperscript.js"; @@ -102,7 +102,7 @@ export function tunnel( //////////////////////////////////////////// // 🪟 Applet registration //////////////////////////////////////////// -export type BroadcastedApplet = { +export type DiffuseApplet = { groupId: string | undefined; scope: AppletScope; @@ -111,18 +111,21 @@ export type BroadcastedApplet = { get instanceId(): string; set data(data: T); - codec: { - decode(data: any): T; - encode(data: T): any; - }; + codec: Codec; - isMainInstance(): boolean; + isMainInstance(): boolean | null; setActionHandler(actionId: string, actionHandler: H): void; }; +export type Codec = { + decode(data: any): T; + encode(data: T): any; +}; + export function register( - options: { worker?: Comlink.Remote } = {}, -): BroadcastedApplet { + options: { mode?: "broadcast" | "shared-worker"; worker?: Comlink.Remote } = {}, +): DiffuseApplet { + const mode = options.mode ?? "broadcast"; const url = new URL(location.href); const scope = applets.register(); @@ -130,7 +133,84 @@ export function register( const channelId = `${location.host}${location.pathname}/${groupId}`; const instanceId = crypto.randomUUID(); - let isMainInstance = true; + // Codec + const codec = { + decode: (data: any) => data as DataType, + encode: (data: DataType) => data as any, + }; + + // Channel + const channelContext = + mode === "broadcast" + ? broadcastChannel({ + channelId, + codec, + instanceId, + scope, + }) + : undefined; + + // Context + const context: DiffuseApplet = { + groupId, + scope, + + settled() { + return channelContext?.promise.then(() => {}) ?? Promise.resolve(); + }, + + get instanceId() { + return instanceId; + }, + + get data() { + return scope.data; + }, + + set data(data: DataType) { + scope.data = data; + }, + + codec, + + isMainInstance() { + return channelContext?.mainSignal[0]() ?? null; + }, + + setActionHandler: (actionId: string, actionHandler: H) => { + switch (mode) { + case "broadcast": + return channelContext?.setActionHandler(actionId, actionHandler); + + case "shared-worker": + return scope.setActionHandler(actionId, actionHandler); + } + }, + }; + + if (options.worker) { + context.scope.onworkerport = (event) => { + if (!event.port) return; + options.worker?._listen(transfer(event.port)); + }; + } + + return context; +} + +function broadcastChannel({ + channelId, + codec, + instanceId, + scope, +}: { + channelId: string; + codec: Codec; + instanceId: string; + scope: AppletScope; +}) { + const mainSignal = signal(true); + const [isMain, setIsMain] = mainSignal; // One instance to rule them all // @@ -149,10 +229,10 @@ export function register( instanceId: event.data.instanceId, }); - if (isMainInstance) { + if (isMain()) { channel.postMessage({ type: "data", - data: context.codec.encode(scope.data), + data: codec.encode(scope.data), }); } break; @@ -160,13 +240,13 @@ export function register( case "PONG": { if (event.data.instanceId === instanceId) { - isMainInstance = false; + setIsMain(false); } break; } case "action": { - if (isMainInstance) { + if (isMain()) { const result = await scope.actionHandlers[event.data.actionId]?.(...event.data.arguments); channel.postMessage({ type: "actioncomplete", @@ -178,7 +258,7 @@ export function register( } case "data": { - scope.data = context.codec.decode(event.data.data); + scope.data = codec.decode(event.data.data); break; } } @@ -206,104 +286,74 @@ export function register( const promise = makeMainPromise(); - // Send out ping - channel.postMessage({ - type: "PING", - instanceId, - }); - // If the data on the main instance changes, // pass it on to other instances. scope.addEventListener("data", async (event: AppletEvent) => { await promise; - if (isMainInstance) { + if (isMain()) { channel.postMessage({ type: "data", - data: context.codec.encode(event.data), + data: codec.encode(event.data), }); } }); - // Context - const context: BroadcastedApplet = { - groupId, - scope, - - settled() { - return promise.then(() => {}); - }, - - get instanceId() { - return instanceId; - }, - - get data() { - return scope.data; - }, - - set data(data: DataType) { - scope.data = data; - }, - - codec: { - decode: (data: any) => data as DataType, - encode: (data: DataType) => data as any, - }, + // Send out ping + channel.postMessage({ + type: "PING", + instanceId, + }); - isMainInstance() { - return isMainInstance; - }, + // Action handler + const setActionHandler = (actionId: string, actionHandler: H) => { + const handler = async (...args: any) => { + if (isMain()) { + return actionHandler(...args); + } - setActionHandler: (actionId: string, actionHandler: H) => { - const handler = async (...args: any) => { - if (isMainInstance) { - return actionHandler(...args); - } + // Check if a main instance is still available, + // if not, then this is the new main. + const promised = await makeMainPromise(); + setIsMain(promised.isMain); - // Check if a main instance is still available, - // if not, then this is the new main. - const { isMain } = await makeMainPromise(); - isMainInstance = isMain; + if (isMain()) { + return actionHandler(...args); + } - if (isMainInstance) { - return actionHandler(...args); - } + const actionMessage = { + actionInstanceId: crypto.randomUUID(), + actionId, + type: "action", + arguments: args, + }; - const actionMessage = { - actionInstanceId: crypto.randomUUID(), - actionId, - type: "action", - arguments: args, + return await new Promise((resolve) => { + const actionCallback = (event: MessageEvent) => { + if ( + event.data?.type === "actioncomplete" && + event.data?.actionInstanceId === actionMessage.actionInstanceId + ) { + channel.removeEventListener("message", actionCallback); + resolve(event.data.result); + } }; - return await new Promise((resolve) => { - const actionCallback = (event: MessageEvent) => { - if ( - event.data?.type === "actioncomplete" && - event.data?.actionInstanceId === actionMessage.actionInstanceId - ) { - channel.removeEventListener("message", actionCallback); - resolve(event.data.result); - } - }; - - channel.addEventListener("message", actionCallback); - channel.postMessage(actionMessage); - }); - }; + channel.addEventListener("message", actionCallback); + channel.postMessage(actionMessage); + }); + }; - scope.setActionHandler(actionId, handler); - }, + scope.setActionHandler(actionId, handler); }; - if (options.worker !== undefined) - context.scope.onworkerport = (event) => { - if (!event.port) return; - options.worker?._listen(transfer(event.port)); - }; - - return context; + // Fin + return { + channel, + mainSignal, + promise, + setActionHandler, + }; } //////////////////////////////////////////// @@ -326,7 +376,7 @@ export function reactive( }); } -export function makeConnect(context: BroadcastedApplet) { +export function makeConnect(context: DiffuseApplet) { return (applet: Applet, dataFn: (data: D) => T, effectFn: (t: T) => void) => { return reactive(applet, dataFn, (t: T) => { if (context.isMainInstance()) effectFn(t); diff --git a/src/scripts/common.ts b/src/scripts/common.ts index 4723004..24ede50 100644 --- a/src/scripts/common.ts +++ b/src/scripts/common.ts @@ -4,6 +4,7 @@ import { xxh32 } from "xxh32"; import { getTransferables } from "@okikio/transferables"; import type { Track } from "@applets/core/types"; +import type { DiffuseApplet } from "./applet/common"; // export { SharedWorkerPolyfill as SharedWorker } from "@okikio/sharedworker"; export const SharedWorker = globalThis.SharedWorker; @@ -66,10 +67,19 @@ export function endpoint = WorkerTasks>(ini: Comli return e; } -export function expose>(tasks: A): A { +export function expose>( + tasks: A, + opts?: { + ports?: { + applets: MessagePort[]; + consumers: MessagePort[]; + }; + }, +): A { if (globalThis.SharedWorkerGlobalScope && self instanceof SharedWorkerGlobalScope) { self.onconnect = (event: MessageEvent) => { const port = event.ports[0]; + opts?.ports?.applets?.push(port); Comlink.expose(tasks, port); port.start(); }; @@ -118,6 +128,20 @@ export function jsonEncode(a: T): Uint8Array { return new TextEncoder().encode(JSON.stringify(a)); } +export function postMessages({ + data, + ports, + transfer, +}: { + data: D; + ports: MessagePort[]; + transfer?: Transferable[]; +}) { + ports.forEach((port) => { + port.postMessage(data, transfer ?? []); + }); +} + export function provide< C extends Record, A extends Record, @@ -131,18 +155,37 @@ export function provide< connections?: Record>>; tasks?: T; }) { - const allTasks = expose({ - _listen: _listen(actions || ({} as A)), - _manage: _manage(connections || {}), - ...(tasks || ({} as T)), - }); + const portsHolder = { + applets: [] as MessagePort[], + consumers: [] as MessagePort[], + }; + + const allTasks = expose( + { + _listen: _listen(actions || ({} as A), portsHolder), + _manage: _manage(connections || {}), + ...(tasks || ({} as T)), + }, + { + ports: portsHolder, + }, + ); return { connections: connections || ({} as Record>>), + ports: portsHolder, tasks: allTasks, }; } +export function sync(context: DiffuseApplet, port: MessagePort) { + port.onmessage = (event) => { + if (event.data?.type === "data") { + context.data = event.data.data; + } + }; +} + export async function trackArtworkCacheId(track: Track): Promise { return await crypto.subtle .digest("SHA-256", new TextEncoder().encode(track.uri)) @@ -156,7 +199,13 @@ export function transfer(a: T) { // PRIVATE -function _listen>(actions: A) { +function _listen>( + actions: A, + portsHolder: { + applets: MessagePort[]; + consumers: MessagePort[]; + }, +) { async function handleAction( port: MessagePort, action: { @@ -185,6 +234,7 @@ function _listen>(actions: A) { return (port: MessagePort) => { Comlink.expose(actions, port); + portsHolder.consumers.push(port); port.onmessage = async (message) => { switch (message.data?.type) { diff --git a/src/scripts/engine/queue/worker.ts b/src/scripts/engine/queue/worker.ts index 09c2c54..38653f1 100644 --- a/src/scripts/engine/queue/worker.ts +++ b/src/scripts/engine/queue/worker.ts @@ -1,26 +1,9 @@ +import { getTransferables } from "@okikio/transferables"; + import type { Track } from "@applets/core/types.js"; import type { Item, State } from "./types"; -import { arrayShuffle, provide, transfer } from "@scripts/common.ts"; - -//////////////////////////////////////////// -// STATE -//////////////////////////////////////////// - -const QUEUE_SIZE = 25; - -const internal: { pool: Track[] } = { - pool: [], -}; - -const state: State = { - future: [], - past: [], - now: null, -}; - -function data() { - return transfer({ ...state }); -} +import { arrayShuffle, postMessages, provide, transfer } from "@scripts/common.ts"; +import { effect, signal } from "@scripts/signal"; //////////////////////////////////////////// // SETUP @@ -33,7 +16,7 @@ const actions = { unshift, }; -const { tasks } = provide({ +const { ports, tasks } = provide({ actions, tasks: actions, }); @@ -42,17 +25,45 @@ export type Actions = typeof actions; export type Tasks = typeof tasks; //////////////////////////////////////////// -// ACTIONS +// STATE //////////////////////////////////////////// -function add(items: Item[]): State { - state.future = [...state.future, ...items]; +const QUEUE_SIZE = 25; + +const internal: { pool: Track[] } = { + pool: [], +}; + +const [future, setFuture] = signal([]); +const [past, setPast] = signal([]); +const [now, setNow] = signal(null); + +effect(() => { + const state: State = { + future: future(), + past: past(), + now: now(), + }; + + postMessages({ + data: { + type: "data", + data: state, + }, + ports: ports.applets, + transfer: getTransferables(state), + }); +}); - // Fin - return data(); +//////////////////////////////////////////// +// ACTIONS +//////////////////////////////////////////// + +function add(items: Item[]) { + setFuture([...future(), ...items]); } -function pool(tracks: Track[]): State { +function pool(tracks: Track[]) { internal.pool = tracks; // TODO: If the pool changes, only remove non-existing tracks @@ -60,44 +71,38 @@ function pool(tracks: Track[]): State { // // What about past queue items? - state.future = []; + setFuture([]); fill(); // Automatically insert track if there isn't any - if (!state.now) return shift(); - - // Fin - return data(); + if (!now()) return shift(); } -function shift(): State { - state.now = state.future[0] || null; - state.future = state.future.slice(1); - state.past = state.now ? [...state.past, state.now] : state.past; +function shift() { + const now = future()[0] ?? null; + setNow(now); - fill(); + setFuture(future().slice(1)); + setPast(now ? [...past(), now] : past()); - // Fin - return data(); + fill(); } -function unshift(): State { - if (state.past.length === 0) return state; +function unshift() { + if (past().length === 0) return; - const past = [...state.past]; - const [last] = past.splice(past.length - 1, 1); - state.now = last ?? null; - state.future = state.now ? [state.now, ...state.future] : state.future; + const [last] = past().splice(past().length - 1, 1); + const now = last ?? null; - // Fin - return data(); + setNow(now); + setFuture(now ? [now, ...future()] : future()); } // 🛠️ // TODO: Most likely there's a more performant solution function fill() { - if (state.future.length >= QUEUE_SIZE) return state; + if (future().length >= QUEUE_SIZE) return; let reducedPool = internal.pool.reduce( ({ past, pool }: { past: Set; pool: Track[] }, track: Track) => { @@ -112,14 +117,13 @@ function fill() { pool: [...pool, track], }; }, - { past: new Set(state.past.map((t) => t.id)), pool: [] }, + { past: new Set(past().map((t) => t.id)), pool: [] }, ).pool; if (reducedPool.length === 0) { reducedPool = internal.pool; } - const poolSelection = arrayShuffle(reducedPool).slice(0, QUEUE_SIZE - state.future.length); - + const poolSelection = arrayShuffle(reducedPool).slice(0, QUEUE_SIZE - future().length); add(poolSelection); } diff --git a/src/scripts/output/common.ts b/src/scripts/output/common.ts index 44ad21b..132c0d0 100644 --- a/src/scripts/output/common.ts +++ b/src/scripts/output/common.ts @@ -1,7 +1,7 @@ import { xxh32r } from "xxh32/dist/raw.js"; import type { ManagedOutput, Track } from "@applets/core/types"; -import type { BroadcastedApplet } from "@scripts/applet/common"; +import type { DiffuseApplet } from "@scripts/applet/common"; import { jsonEncode } from "@scripts/common"; export const INITIAL_MANAGED_OUTPUT: ManagedOutput = { @@ -13,7 +13,7 @@ export const INITIAL_MANAGED_OUTPUT: ManagedOutput = { }; export function outputManager(args: { - context: BroadcastedApplet; + context: DiffuseApplet; /* Indicate if the initial data loader may proceed. */ init?: () => Promise; tracks: { diff --git a/src/scripts/signal.ts b/src/scripts/signal.ts new file mode 100644 index 0000000..0aee4f1 --- /dev/null +++ b/src/scripts/signal.ts @@ -0,0 +1,65 @@ +import { Signal } from "signal-polyfill"; + +// SIGNAL + +export type Signal = () => T; + +export const signal = (initial: T): [Signal, (value: T) => void] => { + const state = new Signal.State(initial); + const get = () => state.get(); + const set = (value: T) => state.set(value); + return [get, set]; +}; + +// EFFECT + +export const throttled = ( + job: () => void, + queue: (callback: () => void) => void = queueMicrotask, +): (() => void) => { + let isScheduled = false; + + const perform = () => { + job(); + isScheduled = false; + }; + + const schedule = () => { + if (!isScheduled) { + isScheduled = true; + queue(perform); + } + }; + + return schedule; +}; + +const watcher = new Signal.subtle.Watcher( + throttled(() => { + for (const signal of watcher.getPending()) { + signal.get(); + } + watcher.watch(); + }), +); + +export type Cancel = () => void; + +export const effect = (perform: () => Cancel | void) => { + let cleanup: Cancel | undefined; + + const signal = new Signal.Computed(() => { + cleanup?.(); + cleanup = perform() ?? undefined; + }); + + watcher.watch(signal); + signal.get(); + + const dispose = () => { + cleanup?.(); + watcher.unwatch(signal); + }; + + return dispose; +}; -- 2.51.2