diff --git a/src/pages/constituent/blur/artwork-controller/_applet.astro b/src/pages/constituent/blur/artwork-controller/_applet.astro index 5d90509..dedff2d 100644 --- a/src/pages/constituent/blur/artwork-controller/_applet.astro +++ b/src/pages/constituent/blur/artwork-controller/_applet.astro @@ -532,6 +532,8 @@ import "@styles/diffuse/fonts.css"; // ORCHESTRATED context.settled().then(() => { + console.log("READY"); + if (isMainGroup() && context.isMainInstance()) { orchestrator.primary .sendAction("insert_tracks_into_queue", undefined, { @@ -768,11 +770,11 @@ import "@styles/diffuse/fonts.css"; } function previous() { - engine.queue.sendAction("unshift", undefined, { worker: true }); + engine.queue.sendAction("unshift", { groupId: context.groupId }, { worker: true }); } function next() { - engine.queue.sendAction("shift", undefined, { worker: true }); + engine.queue.sendAction("shift", { groupId: context.groupId }, { worker: true }); } controller.appendChild(Controls); diff --git a/src/pages/core/types.d.ts b/src/pages/core/types.d.ts index 306eb4f..5fc2d64 100644 --- a/src/pages/core/types.d.ts +++ b/src/pages/core/types.d.ts @@ -11,7 +11,7 @@ export type Consult = | { supported: true; consult: "undetermined" | boolean }; export type ConsultGrouping = - | { available: false; reason: string } + | { available: false; reason: string; tracks: Track[] } | { available: true; tracks: Track[] }; export type GroupConsult = Record; diff --git a/src/pages/engine/queue/_applet.astro b/src/pages/engine/queue/_applet.astro index 5d68db8..739ec64 100644 --- a/src/pages/engine/queue/_applet.astro +++ b/src/pages/engine/queue/_applet.astro @@ -19,12 +19,13 @@ // Register applet const context = register({ mode: "shared-worker", worker }); + const groupId = context.groupId || "main"; // Initial state - context.data = await worker.data(); + context.data = await worker.data(groupId); // Keep applet data with worker data in sync - sync(context, port); + sync(context, port, { groupId }); //////////////////////////////////////////// // ACTIONS @@ -35,18 +36,18 @@ context.setActionHandler("unshift", unshift); async function add(items: Track[]) { - await worker.add(transfer(items)); + await worker.add({ groupId, items }); } - async function pool(items: Track[]) { - await worker.pool(transfer(items)); + async function pool(tracks: Track[]) { + await worker.pool({ groupId, tracks }); } async function shift() { - await worker.shift(); + await worker.shift({ groupId }); } async function unshift() { - await worker.unshift(); + await worker.unshift({ groupId }); } diff --git a/src/pages/orchestrator/primary/_applet.astro b/src/pages/orchestrator/primary/_applet.astro index 16b1044..7174f7b 100644 --- a/src/pages/orchestrator/primary/_applet.astro +++ b/src/pages/orchestrator/primary/_applet.astro @@ -98,7 +98,7 @@ audio, (data) => data.items[queue.data.now?.id ?? Infinity]?.hasEnded ?? false, (hasEnded) => { - if (hasEnded) queue.sendAction("shift", undefined, { worker: true }); + if (hasEnded) queue.sendAction("shift", { groupId: context.groupId }, { worker: true }); }, ); } @@ -195,6 +195,8 @@ async function insertTracksIntoQueue() { await context.settled(); + console.log("SETTLED"); + // Add tracks to the queue once the tracks have been loaded; // and every time the collection changes. @@ -204,6 +206,8 @@ await wait(output, (d) => d?.tracks.state === "loaded"); + console.log("TRACKS LOADED"); + reactive( output, (data) => data.tracks.cacheId, @@ -214,6 +218,8 @@ { timeoutDuration: 60000 * 5, worker: true }, ); + console.log("CONSULTED"); + // Available tracks const tracks = Object.values(groups).reduce((acc: Track[], value) => { if (value.available === false) return acc; @@ -221,10 +227,14 @@ }, []); // Set pool - await queue.sendAction("pool", tracks, { - timeoutDuration: 60000, - worker: true, - }); + await queue.sendAction( + "pool", + { groupId: context.groupId, tracks }, + { + timeoutDuration: 60000, + worker: true, + }, + ); }, ); } diff --git a/src/scripts/common.ts b/src/scripts/common.ts index d83794e..b237ddf 100644 --- a/src/scripts/common.ts +++ b/src/scripts/common.ts @@ -181,9 +181,13 @@ export function provide< export function sync( context: DiffuseApplet, port: MessagePort | Worker, + options: { groupId?: string } = {}, ) { port.onmessage = (event) => { - if (event.data?.type === "data") { + if ( + event.data?.type === "data" && + (options.groupId ? event.data?.groupId === options.groupId : true) + ) { context.data = event.data.data; } }; diff --git a/src/scripts/configurator/input/worker.ts b/src/scripts/configurator/input/worker.ts index 25675ec..4dc43da 100644 --- a/src/scripts/configurator/input/worker.ts +++ b/src/scripts/configurator/input/worker.ts @@ -1,6 +1,12 @@ import * as URI from "uri-js"; -import type { Consult, GroupConsult, InputWorkerTasks, Track } from "@applets/core/types"; +import type { + Consult, + ConsultGrouping, + GroupConsult, + InputWorkerTasks, + Track, +} from "@applets/core/types"; import { groupTracksPerScheme, initialConnections, provide } from "@scripts/common"; //////////////////////////////////////////// @@ -62,7 +68,11 @@ async function groupConsult(tracks: Track[]) { Object.keys(groups).map(async (scheme) => { if (!isSupportedScheme(scheme)) { return { - [scheme]: { available: false, reason: "Unsupported scheme" }, + [scheme]: { + available: false, + reason: "Unsupported scheme", + tracks: groups[scheme] || [], + }, }; } @@ -79,13 +89,14 @@ async function groupConsult(tracks: Track[]) { } async function list(cachedTracks: Track[] = []) { - const groups = groupTracks(cachedTracks); + const groups = await groupConsult(cachedTracks); const promises = Object.entries(groups).map( - async ([scheme, cachedTracksGroup]: [string, Track[]]) => { - if (!isSupportedScheme(scheme)) return cachedTracksGroup; + async ([scheme, { available, tracks }]: [string, ConsultGrouping]) => { + if (!available) return tracks; + if (!isSupportedScheme(scheme)) return tracks; const conn = await connections[scheme].promise; - return conn.list(cachedTracksGroup); + return conn.list(tracks); }, ); diff --git a/src/scripts/engine/queue/worker.ts b/src/scripts/engine/queue/worker.ts index 8a91d29..b0b831c 100644 --- a/src/scripts/engine/queue/worker.ts +++ b/src/scripts/engine/queue/worker.ts @@ -3,7 +3,6 @@ import { getTransferables } from "@okikio/transferables"; import type { Track } from "@applets/core/types.js"; import type { Item, State } from "./types"; import { arrayShuffle, postMessages, provide, transfer } from "@scripts/common.ts"; -import { effect, signal } from "@scripts/signal"; //////////////////////////////////////////// // SETUP @@ -30,87 +29,103 @@ export type Tasks = typeof tasks; const QUEUE_SIZE = 25; -const internal: { pool: Track[] } = { - pool: [], -}; +const _internal: Record = {}; +const _state: Record = {}; + +function data(groupId: string) { + return state(groupId); +} + +function emptyState(groupId: string): State { + return { + future: [], + now: null, + past: [], + }; +} -const [future, setFuture] = signal([]); -const [past, setPast] = signal([]); -const [now, setNow] = signal(null); +function notify(groupId: string) { + const d = data(groupId); -effect(() => { postMessages({ data: { type: "data", - data: state(), + data: d, + groupId, }, ports: ports.applets, - transfer: getTransferables(state), + transfer: getTransferables(d), }); -}); +} -function data() { - return state(); +function internal(groupId: string) { + _internal[groupId] ??= { pool: [] }; + return _internal[groupId]; } -function state(): State { - return { - future: future(), - past: past(), - now: now(), - }; +function state(groupId: string) { + _state[groupId] ??= emptyState(groupId); + return _state[groupId]; } //////////////////////////////////////////// // ACTIONS //////////////////////////////////////////// -function add(items: Item[]) { - setFuture([...future(), ...items]); +function add({ groupId, items }: { groupId: string; items: Item[] }) { + state(groupId).future = [...state(groupId).future, ...items]; + notify(groupId); } -function pool(tracks: Track[]) { - internal.pool = tracks; +function pool({ groupId, tracks }: { groupId: string; tracks: Track[] }) { + internal(groupId).pool = tracks; + const queue = state(groupId); // TODO: If the pool changes, only remove non-existing tracks // instead of resetting the whole future queue. // // What about past queue items? - setFuture([]); - fill(); + queue.future = []; + fill(groupId); // Automatically insert track if there isn't any - if (!now()) return shift(); + if (!queue.now) return shift({ groupId }); + else notify(groupId); } -function shift() { - const now = future()[0] ?? null; - setNow(now); +function shift({ groupId }: { groupId: string }) { + const queue = state(groupId); + const now = queue.future[0] ?? null; + queue.now = now; - setFuture(future().slice(1)); - setPast(now ? [...past(), now] : past()); + queue.future = queue.future.slice(1); + queue.past = now ? [...queue.past, now] : queue.past; - fill(); + fill(groupId); } -function unshift() { - if (past().length === 0) return; +function unshift({ groupId }: { groupId: string }) { + const queue = state(groupId); + if (queue.past.length === 0) return; - const [last] = past().splice(past().length - 1, 1); + const [last] = queue.past.splice(queue.past.length - 1, 1); const now = last ?? null; - setNow(now); - setFuture(now ? [now, ...future()] : future()); + queue.now = now; + queue.future = now ? [now, ...queue.future] : queue.future; + + notify(groupId); } // 🛠️ // TODO: Most likely there's a more performant solution -function fill() { - if (future().length >= QUEUE_SIZE) return; +function fill(groupId: string) { + const queue = state(groupId); + if (queue.future.length >= QUEUE_SIZE) return; - let reducedPool = internal.pool.reduce( + let reducedPool = internal(groupId).pool.reduce( ({ past, pool }: { past: Set; pool: Track[] }, track: Track) => { if (past.has(track.id)) return { @@ -123,13 +138,13 @@ function fill() { pool: [...pool, track], }; }, - { past: new Set(past().map((t) => t.id)), pool: [] }, + { past: new Set(queue.past.map((t) => t.id)), pool: [] }, ).pool; if (reducedPool.length === 0) { - reducedPool = internal.pool; + reducedPool = internal(groupId).pool; } - const poolSelection = arrayShuffle(reducedPool).slice(0, QUEUE_SIZE - future().length); - add(poolSelection); + const poolSelection = arrayShuffle(reducedPool).slice(0, QUEUE_SIZE - queue.future.length); + add({ groupId, items: poolSelection }); } diff --git a/src/scripts/input/native-fs/worker.ts b/src/scripts/input/native-fs/worker.ts index a770747..cd8d845 100644 --- a/src/scripts/input/native-fs/worker.ts +++ b/src/scripts/input/native-fs/worker.ts @@ -55,7 +55,7 @@ async function groupConsult(tracks: Track[]): Promise { const handle = handles[handleId]; const grouping: ConsultGrouping = handle ? { available: true, tracks } - : { available: false, reason: "Handle not available" }; + : { available: false, reason: "Handle not available", tracks }; return { key: URI.serialize({ scheme: SCHEME, host: handleId }), diff --git a/src/scripts/input/opensubsonic/worker.ts b/src/scripts/input/opensubsonic/worker.ts index 04bd461..2e4b905 100644 --- a/src/scripts/input/opensubsonic/worker.ts +++ b/src/scripts/input/opensubsonic/worker.ts @@ -55,7 +55,7 @@ async function groupConsult(tracks: Track[]): Promise { const available = await consultServer(server); const grouping: ConsultGrouping = available ? { available, tracks } - : { available, reason: "Server ping failed" }; + : { available, reason: "Server ping failed", tracks }; return { key: `${SCHEME}:${serverId}`, diff --git a/src/scripts/input/s3/worker.ts b/src/scripts/input/s3/worker.ts index 349a590..f58e281 100644 --- a/src/scripts/input/s3/worker.ts +++ b/src/scripts/input/s3/worker.ts @@ -52,7 +52,7 @@ async function groupConsult(tracks: Track[]): Promise { const available = await consultBucket(bucket); const grouping: ConsultGrouping = available ? { available, tracks } - : { available, reason: "Bucket unavailable" }; + : { available, reason: "Bucket unavailable", tracks }; return { key: `${SCHEME}:${bucketId}`, diff --git a/src/scripts/theme/pilot/index.ts b/src/scripts/theme/pilot/index.ts index 8049f44..bda2bc9 100644 --- a/src/scripts/theme/pilot/index.ts +++ b/src/scripts/theme/pilot/index.ts @@ -56,7 +56,7 @@ reactive( // Automatically start playing something if nothing is playing yet. if (!audioId) { if (isPlaying) { - const now = await engine.queue.sendAction("shift"); + const now = await engine.queue.sendAction("shift", { groupId: "main" }); if (!now) { console.warn("No tracks available yet, try again later."); await ui.audio.sendAction("modifyIsPlaying", false);