From 6d98596e876162ab5c2bfee52614f33384d681a6 Mon Sep 17 00:00:00 2001 From: Steven Vandevelde Date: Tue, 02 Dec 2025 17:20:40 +0000 Subject: [PATCH] refactor: worker connections with dependencies --- src/common/element.d.ts | 6 ------ src/common/element.js | 230 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------------------------------------------------------------------------------------------------- src/common/worker.d.ts | 24 ++++++++++++++++++++++++ src/common/worker.js | 28 +++++++++++++++++++--------- src/common/constituents/default.js | 7 +++++-- src/components/configurator/input/element.js | 72 ++++++++++++++++-------------------------------------------------------- src/components/configurator/input/types.d.ts | 3 --- src/components/configurator/input/worker.js | 95 +++++++++++++++++++++++++++++++++++++++-------------------------------------------------------- src/components/engine/audio/element.js | 2 +- src/components/engine/queue/element.js | 4 ++-- src/components/orchestrator/process-tracks/element.js | 39 +++++++++++++++------------------------ src/components/orchestrator/process-tracks/types.d.ts | 7 +------ src/components/orchestrator/process-tracks/worker.js | 9 ++++----- src/components/orchestrator/queue-tracks/element.js | 62 +++++++++++++++++++++++++------------------------------------- src/components/orchestrator/queue-tracks/types.d.ts | 5 +---- src/components/orchestrator/queue-tracks/worker.js | 9 ++++----- src/components/orchestrator/search-tracks/element.js | 46 ++++++++++++++-------------------------------- src/components/orchestrator/search-tracks/types.d.ts | 5 +---- src/components/orchestrator/search-tracks/worker.js | 9 ++++----- src/themes/blur/artwork-controller/element.js | 7 +++++++ 20 file(s) changed, 311 insertion(s)(+), 358 deletion(s)(-) diff --git a/src/common/element.d.ts b/src/common/element.d.ts --- a/src/common/element.d.ts +++ b/src/common/element.d.ts @@ -9,12 +9,6 @@ ...values: unknown[] ) => string; -export type ProvisionedWorkers = { - [K in T]: ProvisionedWorker; -}; - -export type ProvisionedWorker = Worker | SharedWorker; - export type RenderArg = { html: HtmlTagFunction; state: State; diff --git a/src/common/element.js b/src/common/element.js --- a/src/common/element.js +++ b/src/common/element.js @@ -1,9 +1,15 @@ import QS from "query-string"; -import { RPCChannel } from "@kunkun/kkrpc"; +import { decodeMessage, encodeMessage, RPCChannel } from "@kunkun/kkrpc"; import { html, render } from "lit-html"; import { effect, signal } from "@common/signal.js"; -import { rpc, transfer, workerLink, workerTunnel } from "./worker.js"; +import { + rpc, + transfer, + workerLink, + workerProxy, + workerTunnel, +} from "./worker.js"; import { BrowserPostMessageIo } from "./worker/rpc.js"; // RE-EXPORT @@ -21,11 +27,8 @@ * around rendering and managing signals. */ export class DiffuseElement extends HTMLElement { + #connected = Promise.withResolvers(); #disposables = /** @type {Array<() => void>} */ ([]); - - #teardown() { - this.#disposables.forEach((fn) => fn()); - } constructor() { super(); @@ -55,6 +58,26 @@ this.#disposables.push(effect(fn)); } + /** */ + forceRender() { + return this.#render(); + } + + /** */ + nameWithGroup() { + return `${this.constructor.prototype.constructor.NAME}/${this.group}`; + } + + /** */ + root() { + return (this.shadowRoot ?? this); + } + + /** */ + whenConnected() { + return this.#connected.promise; + } + /** * Avoid replacing the whole subtree, * morph the existing DOM into the new given tree. @@ -70,19 +93,11 @@ render(tmp, this.root()); } - /** */ - forceRender() { - return this.#render(); - } - - /** */ - root() { - return (this.shadowRoot ?? this); - } - // LIFECYCLE connectedCallback() { + this.#connected.resolve(null); + if (!("render" in this && typeof this.render === "function")) return; this.effect(() => { @@ -95,7 +110,11 @@ this.#teardown(); } - // WORKER + #teardown() { + this.#disposables.forEach((fn) => fn()); + } + + // WORKERS /** @type {undefined | Worker | SharedWorker} */ #worker; @@ -115,7 +134,7 @@ ); // Setup worker - const name = `${NAME}/${this.group}`; + const name = this.nameWithGroup(); const url = import.meta.resolve("./" + WORKER_URL) + `?${query}`; let worker; @@ -129,6 +148,20 @@ return worker; } + /** */ + dependencies() { + return Object.fromEntries( + Array.from(this.children).flatMap((element) => { + if ("nameWithGroup" in element === false) { + return []; + } + + const d = /** @type {DiffuseElement} */ (element); + return [[d.localName, d]]; + }), + ); + } + worker() { this.#worker ??= this.createWorker(); return this.#worker; @@ -137,6 +170,74 @@ workerLink() { const worker = this.worker(); return workerLink(worker); + } + + /** + * @template {Record any>} Actions + * @returns {ProxiedActions} + */ + workerProxy() { + return workerProxy( + () => this.workerTunnel().port, + ); + } + + /** + * @param {{ newWorker?: boolean }} [opts] + */ + workerTunnel({ newWorker } = {}) { + // Creates a MessagePort that is connected to the worker. + // All the dependencies are added automatically. + const worker = newWorker ? this.createWorker() : this.worker(); + const deps = this.dependencies(); + + let toWorker; + + if (Object.keys(deps).length) { + toWorker = + /** + * @param {any} msg + */ + async (msg) => { + /** @type {Array<[string, Tunnel]>} */ + const ports = Object.entries(deps).map( + /** @param {[string, DiffuseElement]} _ */ + ([k, v]) => [k, v.workerTunnel()], + ); + + const decoded = await decodeMessage(msg); + const data = { + data: Array.isArray(decoded.args) ? decoded.args[0] : decoded.args, + ports: Object.fromEntries(ports.map(([k, v]) => { + return [k, v.port]; + })), + }; + + const encoded = encodeMessage( + { + ...decoded, + args: Array.isArray(decoded.args) + ? [data, ...decoded.args.slice(1)] + : decoded.args, + }, + {}, + true, + ports.map(([_k, v]) => v.port), + ); + + this.#disposables.push(() => { + ports.forEach(([_k, v]) => v.disconnect()); + }); + + return { + data: encoded, + transfer: ports.map(([_k, v]) => v.port), + }; + }; + } + + const tunnel = workerTunnel(worker, { toWorker }); + return tunnel; } } @@ -171,13 +272,13 @@ /** * @template {Record any }>} ActionsWithStrategy * @template {{ [K in keyof ActionsWithStrategy]: ActionsWithStrategy[K]["fn"] }} Actions - * @param {string} name + * @param {string} channelName * @param {ActionsWithStrategy} actionsWithStrategy */ - broadcast(name, actionsWithStrategy) { + broadcast(channelName, actionsWithStrategy) { if (this.broadcasted) return; - const channel = new BroadcastChannel(name); + const channel = new BroadcastChannel(channelName); const msg = new MessageChannel(); /** @@ -185,7 +286,7 @@ */ this.broadcasted = true; - this.name = name; + this.channelName = channelName; const _rpc = rpc( msg.port2, @@ -291,7 +392,7 @@ // Grab a lock if it isn't acquired yet, // and hold it until `this.lock.promise` resolves. navigator.locks.request( - `${this.name}/lock`, + `${this.channelName}/lock`, { ifAvailable: true }, (lock) => { this.#status.resolve( @@ -305,15 +406,15 @@ // Additionally, wait for lock if needed. this.#status.promise.then((status) => { if (status.leader) { - console.log(`🧙 Elected leader for: ${this.name}`); + console.log(`🧙 Elected leader for: ${this.channelName}`); } else { - console.log(`🔮 Watching leader: ${this.name}`); + console.log(`🔮 Watching leader: ${this.channelName}`); } // Wait for leadership if (status.leader === false) { navigator.locks.request( - `${this.name}/lock`, + `${this.channelName}/lock`, () => { this.#status = Promise.withResolvers(); this.#status.resolve({ leader: true, initialLeader: false }); @@ -334,56 +435,6 @@ super.disconnectedCallback(); this.#lock.resolve(); } -} - -/** - * @template {string} A - * @template {ProvisionedWorkers} B - * @template {Record} C - * @template R - * @param {Promise | undefined} provisions - * @param {(args: C & { ports: { [K in keyof B]: MessagePort } }) => R} fn - * @param {C} fnArgs - * @returns {Promise} - */ -export async function callWorkerWithProvisions(provisions, fn, fnArgs) { - const workers = await provisions; - if (!workers) throw new Error("Workers not defined"); - - /** @type {Array<[keyof B, Tunnel]>} */ - const tunnels = Object.keys(workers).map( - (value) => { - const key = /** @type {keyof B} */ (value); - const worker = workers[key]; - return [key, workerTunnel(worker)]; - }, - ); - - const ports = /** @type {{ [K in keyof B]: MessagePort }} */ ( - Object.fromEntries( - tunnels.map(([key, tunnel]) => { - return [key, tunnel.port]; - }), - ) - ); - - const args = { - ...fnArgs, - ports, - }; - - const result = await fn(transfer( - args, - tunnels.map(([_key, tunnel]) => { - return tunnel.port; - }), - )); - - tunnels.forEach(([_key, tunnel]) => { - tunnel.disconnect(); - }); - - return result; } /** @@ -433,32 +484,9 @@ } /** - * @template {Record} T - * @param {T} elements + * @param {Record} workers */ -export async function provisionWorkers(elements) { - await whenElementsDefined(elements); - - /** @type {Record} */ - const provisions = {}; - - Object.entries(elements).forEach(([key, element]) => { - const worker = element.createWorker(); - provisions[key] = worker; - }); - - const casted = - /** @type {{ [K in keyof T]: ProvisionedWorker}} */ (provisions); - - return casted; -} - -/** - * @param {ProvisionedWorkers | undefined} workers - */ -export function terminateProvisions(workers) { - if (!workers) return; - +export function terminateWorkers(workers) { Object.values(workers).forEach((worker) => { if (worker instanceof Worker) worker.terminate(); }); diff --git a/src/common/worker.d.ts b/src/common/worker.d.ts --- a/src/common/worker.d.ts +++ b/src/common/worker.d.ts @@ -1,6 +1,16 @@ export type Announcement = MRpcBaseMsg & { type: "announcement"; args: T }; export type IncompleteArray = ["Missing required items", T]; +export type ActionsWithTunnel< + Actions extends Record any>, +> = { + [A in keyof Actions]: WithTunnel; +}; + +export type Dependencies = { + [K in T]: Worker | SharedWorker; +}; + /** * Comes from the `@mys/m-rpc` library, * but it is not exported. Used to identify @@ -34,3 +44,17 @@ disconnect: () => void; port: MessagePort; }; + +/** */ +export type WithTunnel< + Fn extends (...args: any[]) => any, + PromisedReturn = (ReturnType extends Promise ? ReturnType + : Promise>), +> = ( + _: { data: Parameters[0]; ports: Record }, + ...args: Rest> +) => PromisedReturn; + +// 🛑 + +type Rest = T extends [any, ...(infer R)[]] ? R : never; diff --git a/src/common/worker.js b/src/common/worker.js --- a/src/common/worker.js +++ b/src/common/worker.js @@ -1,4 +1,4 @@ -import { RPCChannel } from "@kunkun/kkrpc"; +import { RPCChannel, transfer } from "@kunkun/kkrpc"; import { getTransferables } from "@okikio/transferables"; import { debounceMicrotask } from "@vicary/debounce-microtask"; import { xxh32 } from "xxh32"; @@ -9,7 +9,7 @@ export { transfer } from "@kunkun/kkrpc"; /** - * @import {Announcement, MessengerRealm, ProxiedActions, Tunnel} from "./worker.d.ts" + * @import {Announcement, Dependencies, MessengerRealm, ProxiedActions, Tunnel} from "./worker.d.ts" */ //////////////////////////////////////////// @@ -90,8 +90,11 @@ // Create proxy that creates RPC API when needed const proxy = new Proxy(() => {}, { get: (_target, prop) => { - const api = ensureAPI(); - return api[prop.toString()]; + /** @param {Parameters} args */ + return (...args) => { + const api = ensureAPI(); + return api[prop.toString()](...args); + }; }, }); @@ -100,24 +103,31 @@ /** * @param {MessagePort | Worker | SharedWorker} workerOrLink + * @param {{ fromWorker?: (message: any) => Promise<{ data: any, transfer?: Transferable[] }>; toWorker?: (message: any) => Promise<{ data: any, transfer?: Transferable[] }> }} [hooks] * @returns {Tunnel} */ -export function workerTunnel(workerOrLink) { +export function workerTunnel(workerOrLink, hooks = {}) { const link = workerOrLink instanceof SharedWorker ? workerLink(workerOrLink) : workerOrLink; const channel = new MessageChannel(); - channel.port1.addEventListener("message", (event) => { - link.postMessage(event.data); + channel.port1.addEventListener("message", async (event) => { + // Send to worker + const { data, transfer } = await hooks?.toWorker?.(event.data) ?? + { data: event.data }; + link.postMessage(data, { transfer }); }); /** * @param {Event} event */ - const workerListener = (event) => { + const workerListener = async (event) => { + // Receive from worker const msgEvent = /** @type {MessageEvent} */ (event); - channel.port1.postMessage(msgEvent.data); + const { data, transfer } = await hooks?.fromWorker?.(msgEvent.data) ?? + { data: msgEvent.data }; + channel.port1.postMessage(data, { transfer }); }; link.addEventListener("message", workerListener); diff --git a/src/common/constituents/default.js b/src/common/constituents/default.js --- a/src/common/constituents/default.js +++ b/src/common/constituents/default.js @@ -55,7 +55,10 @@ const trigger = queue.now(); const _other_trigger = queue.poolHash(); - queue.fill({ amount: 10, shuffled: true }); - if (!trigger) queue.shift(); + oqt.isLeader().then((isLeader) => { + if (!isLeader) return; + queue.fill({ amount: 10, shuffled: true }); + if (!trigger) queue.shift(); + }); }); } diff --git a/src/components/configurator/input/element.js b/src/components/configurator/input/element.js --- a/src/components/configurator/input/element.js +++ b/src/components/configurator/input/element.js @@ -1,10 +1,8 @@ -import { DiffuseElement, workerProxy } from "@common/element.js"; -import { transfer, workerLink, workerTunnel } from "@common/worker.js"; +import { DiffuseElement, whenElementsDefined } from "@common/element.js"; /** * @import {ProxiedActions, Tunnel} from "@common/worker.d.ts" * @import {InputActions, InputElement} from "@components/input/types.d.ts" - * @import {AdditionalActions} from "./types.d.ts" */ /** @@ -26,7 +24,7 @@ super(); /** @type {ProxiedActions} */ - const proxy = workerProxy(this.workerLink); + const proxy = this.workerProxy(); this.consult = proxy.consult; this.contextualize = proxy.contextualize; @@ -35,70 +33,32 @@ this.resolve = proxy.resolve; } - // WORKER + // LIFECYCLE /** * @override */ - createWorker() { - const worker = super.createWorker(); - - // Wait for child elements to be rendered - setTimeout(() => this.configureWorker(worker), 0); - - return worker; + async connectedCallback() { + super.connectedCallback(); + await whenElementsDefined(this.inputs()); } - // 🛠️ + // WORKERS /** - * @param {Worker | SharedWorker} worker + * @override */ - async configureWorker(worker) { - const inputs = await this.inputTunnels(); - - // Check if any inputs are present - if (inputs.length === 0) return; - - // Configure worker with input ports - const args = transfer({ - ports: Object.fromEntries(inputs.map((input) => { - return [input.element.SCHEME, input.tunnel.port]; - })), - }, inputs.map((i) => i.tunnel.port)); - - /** @type {ProxiedActions} */ - const proxy = workerProxy(() => workerLink(worker)); - proxy.configure(args); + dependencies() { + return this.inputs(); } - async inputTunnels() { - const inputElements = this.children; - const inputs = await Array.from(inputElements).reduce( - /** - * @param {Promise>} acc - * @param {Element} el - */ - async (acc, el) => { - const rec = await acc; - await customElements.whenDefined(el.localName); - - const element = /** @type {InputElement} */ (el); - const worker = element.worker(); - const tunnel = workerTunnel(worker); - - const item = { - element, - tunnel, - worker, - }; - - return [...rec, item]; - }, - Promise.resolve([]), + inputs() { + return Object.fromEntries( + Array.from(this.children).map((element) => { + const input = /** @type {InputElement} */ (element); + return [input.SCHEME, input]; + }), ); - - return inputs; } } diff --git a/src/components/configurator/input/types.d.ts b/src/components/configurator/input/types.d.ts deleted file mode 100644 --- a/src/components/configurator/input/types.d.ts +++ /dev/null @@ -1,3 +0,0 @@ -export type AdditionalActions = { - configure: (args: { ports: { [S in string]: MessagePort } }) => void; -}; diff --git a/src/components/configurator/input/worker.js b/src/components/configurator/input/worker.js --- a/src/components/configurator/input/worker.js +++ b/src/components/configurator/input/worker.js @@ -6,43 +6,23 @@ /** * @import {Track} from "@definitions/types.d.ts"; * @import {GroupConsult, InputActions} from "@components/input/types.d.ts" - * @import {ProxiedActions} from "@common/worker.d.ts" - * @import {AdditionalActions} from "./types.d.ts" + * @import {ActionsWithTunnel, ProxiedActions} from "@common/worker.d.ts" */ - -/** @type {Record>} */ -const inputs = {}; - -//////////////////////////////////////////// -// ACTIONS -//////////////////////////////////////////// - -/** - * @type {AdditionalActions["configure"]} - */ -export function configure({ ports }) { - Object.keys(ports).forEach((key) => { - inputs[key.toLowerCase()] = workerProxy(() => { - const port = ports[key]; - port.start(); - return port; - }); - }); -} //////////////////////////////////////////// // INPUT ACTIONS //////////////////////////////////////////// /** - * @type {InputActions['consult']} + * @type {ActionsWithTunnel['consult']} */ -export async function consult(fileUriOrScheme) { +export async function consult({ data, ports }) { + const fileUriOrScheme = data; const scheme = fileUriOrScheme.includes(":") ? URI.parse(fileUriOrScheme).scheme || fileUriOrScheme : fileUriOrScheme; - const input = grabInput(scheme); + const input = grabInput(scheme, ports); if (!input) { return { supported: false, reason: "Unsupported scheme" }; @@ -52,13 +32,14 @@ } /** - * @type {InputActions['contextualize']} + * @type {ActionsWithTunnel['contextualize']} */ -export async function contextualize(tracks) { - const groups = groupTracks(tracks); +export async function contextualize({ data, ports }) { + const tracks = data; + const groups = groupTracks(tracks, ports); const promises = Object.entries(groups).map( async ([scheme, tracksGroup]) => { - const input = grabInput(scheme); + const input = grabInput(scheme, ports); if (!input || tracksGroup.length === 0) return; return await input.contextualize(tracksGroup); }, @@ -68,15 +49,16 @@ } /** - * @type {InputActions['groupConsult']} + * @type {ActionsWithTunnel['groupConsult']} */ -export async function groupConsult(tracks) { +export async function groupConsult({ data, ports }) { + const tracks = data; const groups = groupTracksPerScheme(tracks); /** @type {GroupConsult[]} */ const consultations = await Promise.all( Object.keys(groups).map(async (scheme) => { - const input = grabInput(scheme); + const input = grabInput(scheme, ports); if (!input) { return { @@ -98,12 +80,12 @@ } /** - * @type {InputActions['list']} + * @type {ActionsWithTunnel['list']} */ -export async function list(cachedTracks = []) { - const groups = await groupConsult(cachedTracks); +export async function list({ data, ports }) { + const groups = await groupConsult({ data, ports }); - Object.keys(inputs).forEach((scheme) => { + Object.keys(ports).forEach((scheme) => { if (!groups[scheme]) groups[scheme] = { available: true, tracks: [] }; }); @@ -111,7 +93,7 @@ async ([scheme, { available, tracks }]) => { if (!available) return tracks; - const input = grabInput(scheme); + const input = grabInput(scheme, ports); if (!input) return tracks; return await input.list(tracks); }, @@ -124,21 +106,16 @@ } /** - * @type {InputActions['resolve']} + * @type {ActionsWithTunnel['resolve']} */ -export async function resolve(args) { - const scheme = args.uri.split(":", 1)[0]; - const input = grabInput(scheme); +export async function resolve({ data, ports }) { + const uri = data.uri; + const scheme = uri.split(":", 1)[0]; + const input = grabInput(scheme, ports); if (!input) return undefined; - try { - return await input.resolve(args); - } catch (err) { - console.error( - `[configurator/input] Resolve error for scheme '${scheme}'.`, - err, - ); - } + const result = await input.resolve(data); + return result; } //////////////////////////////////////////// @@ -152,9 +129,6 @@ groupConsult, list, resolve, - - // Additional - configure, }); }); @@ -164,19 +138,28 @@ /** * @param {string} scheme + * @param {Record} ports + * @returns {ProxiedActions | null} */ -function grabInput(scheme) { - return inputs[scheme.toLowerCase()]; +function grabInput(scheme, ports) { + const port = ports[scheme]; + if (!port) return null; + + return workerProxy(() => { + port.start(); + return port; + }); } /** * @param {Track[]} tracks + * @param {Record} ports */ -function groupTracks(tracks) { +function groupTracks(tracks, ports) { const grouped = groupTracksPerScheme( tracks, Object.fromEntries( - Object.keys(inputs).map((k) => { + Object.keys(ports).map((k) => { return [k, []]; }), ), diff --git a/src/components/engine/audio/element.js b/src/components/engine/audio/element.js --- a/src/components/engine/audio/element.js +++ b/src/components/engine/audio/element.js @@ -46,7 +46,7 @@ // Setup leader election if shared if (this.hasAttribute("group")) { const actions = this.broadcast( - `${this.constructor.prototype.constructor.NAME}/${this.group}`, + this.nameWithGroup(), { adjustVolume: { strategy: "leaderOnly", fn: this.adjustVolume }, pause: { strategy: "leaderOnly", fn: this.pause }, diff --git a/src/components/engine/queue/element.js b/src/components/engine/queue/element.js --- a/src/components/engine/queue/element.js +++ b/src/components/engine/queue/element.js @@ -1,4 +1,4 @@ -import { DiffuseElement, workerProxy } from "@common/element.js"; +import { DiffuseElement } from "@common/element.js"; import { signal } from "@common/signal.js"; import { listen } from "@common/worker.js"; import { hash } from "@common/index.js"; @@ -23,7 +23,7 @@ super(); /** @type {ProxiedActions} */ - this.proxy = workerProxy(this.workerLink); + this.proxy = this.workerProxy(); this.add = this.proxy.add; this.fill = this.proxy.fill; diff --git a/src/components/orchestrator/process-tracks/element.js b/src/components/orchestrator/process-tracks/element.js --- a/src/components/orchestrator/process-tracks/element.js +++ b/src/components/orchestrator/process-tracks/element.js @@ -1,16 +1,8 @@ -import { - callWorkerWithProvisions, - DiffuseElement, - provisionWorkers, - query, - terminateProvisions, - workerProxy, -} from "@common/element.js"; +import { DiffuseElement, query } from "@common/element.js"; import { signal, untracked } from "@common/signal.js"; /** * @import {Track} from "@definitions/types.d.ts" - * @import {ProvisionedWorkers} from "@common/element.d.ts" * @import {ProxiedActions} from "@common/worker.d.ts" * @import {InputElement} from "@components/input/types.d.ts" * @import {OutputElement} from "@components/output/types.d.ts" @@ -34,12 +26,9 @@ /** @type {ProxiedActions} */ #proxy; - /** @type {Promise> | undefined} */ - #workers = undefined; - constructor() { super(); - this.#proxy = workerProxy(this.workerLink); + this.#proxy = this.workerProxy(); } // SIGNALS @@ -72,9 +61,6 @@ this.output = output; this.metadataProcessor = metadataProcessor; - // Create new workers - this.#workers = provisionWorkers({ input, metadataProcessor }); - // Wait until defined await customElements.whenDefined(output.localName); @@ -87,12 +73,21 @@ }); } + // WORKERS + /** * @override */ - async disconnectedCallback() { - super.disconnectedCallback(); - terminateProvisions(await this.#workers); + dependencies() { + if (!this.input) throw new Error("Input element not defined yet"); + if (!this.metadataProcessor) { + throw new Error("Metadata processor element not defined yet"); + } + + return { + input: this.input, + metadataProcessor: this.metadataProcessor, + }; } // ACTIONS @@ -105,11 +100,7 @@ console.log("🪵 Processing initiated"); const cachedTracks = this.output.tracks.collection(); - const result = await callWorkerWithProvisions( - this.#workers, - this.#proxy.process, - { tracks: cachedTracks }, - ); + const result = await this.#proxy.process(cachedTracks); // Save if collection changed if (result) await this.output.tracks.save(result); diff --git a/src/components/orchestrator/process-tracks/types.d.ts b/src/components/orchestrator/process-tracks/types.d.ts --- a/src/components/orchestrator/process-tracks/types.d.ts +++ b/src/components/orchestrator/process-tracks/types.d.ts @@ -1,10 +1,5 @@ import type { Track } from "@definitions/types.d.ts"; export type Actions = { - process: ( - args: { - ports: { input: MessagePort; metadataProcessor: MessagePort }; - tracks: Track[]; - }, - ) => Promise; + process: (tracks: Track[]) => Promise; }; diff --git a/src/components/orchestrator/process-tracks/worker.js b/src/components/orchestrator/process-tracks/worker.js --- a/src/components/orchestrator/process-tracks/worker.js +++ b/src/components/orchestrator/process-tracks/worker.js @@ -4,7 +4,7 @@ /** * @import {Track} from "@definitions/types.d.ts" - * @import {ProxiedActions} from "@common/worker.d.ts" + * @import {ActionsWithTunnel, ProxiedActions} from "@common/worker.d.ts" * @import {InputActions} from "@components/input/types.d.ts" * @import {Actions as MetadataProcessorActions} from "@components/processor/metadata/types.d.ts" * @import {Actions} from "./types.d.ts" @@ -15,11 +15,10 @@ //////////////////////////////////////////// /** - * @type {Actions["process"]} + * @type {ActionsWithTunnel["process"]} */ -export async function process(args) { - const { ports } = args; - const cachedTracks = args.tracks; +export async function process({ data, ports }) { + const cachedTracks = data; /** @type {ProxiedActions} */ const input = workerProxy(() => ports.input); diff --git a/src/components/orchestrator/queue-tracks/element.js b/src/components/orchestrator/queue-tracks/element.js --- a/src/components/orchestrator/queue-tracks/element.js +++ b/src/components/orchestrator/queue-tracks/element.js @@ -1,16 +1,8 @@ -import { - BroadcastableDiffuseElement, - callWorkerWithProvisions, - query, - terminateProvisions, - whenElementsDefined, - workerProxy, -} from "@common/element.js"; +import { BroadcastableDiffuseElement, query } from "@common/element.js"; import { untracked } from "@common/signal.js"; /** * @import {Track} from "@definitions/types.d.ts" - * @import {ProvisionedWorkers} from "@common/element.d.ts" * @import {ProxiedActions} from "@common/worker.d.ts" * @import {InputElement} from "@components/input/types.d.ts" * @import {OutputElement} from "@components/output/types.d.ts" @@ -34,18 +26,23 @@ /** @type {ProxiedActions} */ #proxy; - /** @type {Promise> | undefined} */ - #workers = undefined; - constructor() { super(); - this.#proxy = workerProxy(this.workerLink); + this.#proxy = this.workerProxy(); } + + // LIFECYCLE /** * @override */ async connectedCallback() { + // Broadcast if needed + if (this.hasAttribute("group")) { + this.broadcast(this.nameWithGroup(), {}); + } + + // Super super.connectedCallback(); /** @type {InputElement} */ @@ -62,17 +59,10 @@ this.output = output; this.queue = queue; - // Create new workers - this.#workers = whenElementsDefined({ input, queue }).then(() => { - return { - input: input.createWorker(), - queue: queue.worker(), - }; - }); - // When defined await customElements.whenDefined(this.input.localName); await customElements.whenDefined(this.output.localName); + await customElements.whenDefined(this.queue.localName); // Watch tracks collection this.effect(() => { @@ -82,31 +72,29 @@ if (!isLeader) return; untracked(() => - this.poolAvailable(tracks.filter((t) => t.kind !== "placeholder")) + this.#proxy.poolAvailable( + tracks.filter((t) => t.kind !== "placeholder"), + ) ); }); }); + + // 🌸 } + + // WORKERS /** * @override */ - async disconnectedCallback() { - super.disconnectedCallback(); - terminateProvisions(await this.#workers); - } + dependencies() { + if (!this.input) throw new Error("Input element not defined yet"); + if (!this.queue) throw new Error("Queue element not defined yet"); - // 🌊 - - /** - * @param {Track[]} cachedTracks - */ - async poolAvailable(cachedTracks) { - return await callWorkerWithProvisions( - this.#workers, - this.#proxy.poolAvailable, - { tracks: cachedTracks }, - ); + return { + input: this.input, + queue: this.queue, + }; } } diff --git a/src/components/orchestrator/queue-tracks/types.d.ts b/src/components/orchestrator/queue-tracks/types.d.ts --- a/src/components/orchestrator/queue-tracks/types.d.ts +++ b/src/components/orchestrator/queue-tracks/types.d.ts @@ -1,8 +1,5 @@ import type { Track } from "@definitions/types.d.ts"; export type Actions = { - poolAvailable(args: { - ports: { input: MessagePort; queue: MessagePort }; - tracks: Track[]; - }): Promise; + poolAvailable(tracks: Track[]): Promise; }; diff --git a/src/components/orchestrator/queue-tracks/worker.js b/src/components/orchestrator/queue-tracks/worker.js --- a/src/components/orchestrator/queue-tracks/worker.js +++ b/src/components/orchestrator/queue-tracks/worker.js @@ -2,7 +2,7 @@ /** * @import {Track} from "@definitions/types.d.ts" - * @import {ProxiedActions} from "@common/worker.d.ts" + * @import {ActionsWithTunnel, ProxiedActions} from "@common/worker.d.ts" * @import {InputActions} from "@components/input/types.d.ts" * @import {Actions as QueueEngineActions} from "@components/engine/queue/types.d.ts" * @import {Actions} from "./types.d.ts" @@ -13,11 +13,10 @@ //////////////////////////////////////////// /** - * @type {Actions["poolAvailable"]} + * @type {ActionsWithTunnel["poolAvailable"]} */ -export async function poolAvailable(args) { - const { ports } = args; - const cachedTracks = args.tracks; +export async function poolAvailable({ data, ports }) { + const cachedTracks = data; /** @type {ProxiedActions} */ const input = workerProxy(() => ports.input); diff --git a/src/components/orchestrator/search-tracks/element.js b/src/components/orchestrator/search-tracks/element.js --- a/src/components/orchestrator/search-tracks/element.js +++ b/src/components/orchestrator/search-tracks/element.js @@ -1,15 +1,7 @@ -import { - callWorkerWithProvisions, - DiffuseElement, - provisionWorkers, - query, - terminateProvisions, - workerProxy, -} from "@common/element.js"; +import { DiffuseElement, query } from "@common/element.js"; /** * @import {Track} from "@definitions/types.d.ts" - * @import {ProvisionedWorkers} from "@common/element.d.ts" * @import {ProxiedActions} from "@common/worker.d.ts" * @import {InputElement} from "@components/input/types.d.ts" * @import {OutputElement} from "@components/output/types.d.ts" @@ -33,13 +25,12 @@ /** @type {ProxiedActions} */ #proxy; - /** @type {Promise> | undefined} */ - #workers = undefined; - constructor() { super(); - this.#proxy = workerProxy(this.workerLink); + this.#proxy = this.workerProxy(); } + + // LIFECYCLE /** * @override @@ -61,9 +52,6 @@ this.output = output; this.search = search; - // Create new workers - this.#workers = provisionWorkers({ input, search }); - // When defined await customElements.whenDefined(this.output.localName); @@ -73,29 +61,23 @@ t.kind !== "placeholder" ); - this.supplyAvailable(tracks); + this.#proxy.supplyAvailable(tracks); }); } + + // WORKERS /** * @override */ - async disconnectedCallback() { - super.disconnectedCallback(); - terminateProvisions(await this.#workers); - } + dependencies() { + if (!this.input) throw new Error("Input element not defined yet"); + if (!this.search) throw new Error("Search element not defined yet"); - // 🚛 - - /** - * @param {Track[]} cachedTracks - */ - async supplyAvailable(cachedTracks) { - return await callWorkerWithProvisions( - this.#workers, - this.#proxy.supplyAvailable, - { tracks: cachedTracks }, - ); + return { + input: this.input, + search: this.search, + }; } } diff --git a/src/components/orchestrator/search-tracks/types.d.ts b/src/components/orchestrator/search-tracks/types.d.ts --- a/src/components/orchestrator/search-tracks/types.d.ts +++ b/src/components/orchestrator/search-tracks/types.d.ts @@ -1,8 +1,5 @@ import type { Track } from "@definitions/types.d.ts"; export type Actions = { - supplyAvailable(args: { - ports: { input: MessagePort; search: MessagePort }; - tracks: Track[]; - }): Promise; + supplyAvailable(tracks: Track[]): Promise; }; diff --git a/src/components/orchestrator/search-tracks/worker.js b/src/components/orchestrator/search-tracks/worker.js --- a/src/components/orchestrator/search-tracks/worker.js +++ b/src/components/orchestrator/search-tracks/worker.js @@ -2,7 +2,7 @@ /** * @import {Track} from "@definitions/types.d.ts" - * @import {ProxiedActions} from "@common/worker.d.ts" + * @import {ActionsWithTunnel, ProxiedActions} from "@common/worker.d.ts" * @import {InputActions} from "@components/input/types.d.ts" * @import {Actions as SearchProcessorActions} from "@components/processor/search/types.d.ts" * @import {Actions} from "./types.d.ts" @@ -13,11 +13,10 @@ //////////////////////////////////////////// /** - * @type {Actions["supplyAvailable"]} + * @type {ActionsWithTunnel["supplyAvailable"]} */ -export async function supplyAvailable(args) { - const { ports } = args; - const cachedTracks = args.tracks; +export async function supplyAvailable({ data, ports }) { + const cachedTracks = data; /** @type {ProxiedActions} */ const input = workerProxy(() => ports.input); diff --git a/src/themes/blur/artwork-controller/element.js b/src/themes/blur/artwork-controller/element.js --- a/src/themes/blur/artwork-controller/element.js +++ b/src/themes/blur/artwork-controller/element.js @@ -85,6 +85,8 @@ * @param {Track | null} track */ async #changeArtwork(track) { + console.log("QUEUE NOW", track); + if (!track) { this.#artwork.value = []; return; @@ -97,6 +99,8 @@ method: "HEAD", uri: track.uri, }); + + console.log(resGet, this.input); if (!resGet) return; @@ -116,6 +120,9 @@ }; const art = await this.artwork?.artwork(request) ?? []; + + console.log(art); + const currCacheId = track ? await trackArtworkCacheId(track) : undefined; if (cacheId === currCacheId) this.#artwork.set(art); } -- tangled.sh