A music player that connects to your cloud/distributed storage. diffuse.sh
Something went wrong. Try again.
JavaScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291/// <reference lib="webworker" />import { getTransferables } from "@okikio/transferables";import { debounceMicrotask } from "@vicary/debounce-microtask";import { xxh32 } from "xxh32";
import { RpcChannel } from "./worker/rpc-channel.js";
export { getTransferables } from "@okikio/transferables";
/** * @import {Announcement, MessengerRealm, ProxiedActions, Tunnel} from "./worker.d.ts" */
// Early message buffer for regular Workers.//// If a Worker module (or its dependencies) contains a top-level `await`, the// browser can deliver queued incoming messages to `globalThis` while the module// evaluation is paused — before `ostiary`/`rpc()` has had a chance to register// a handler. Those messages would otherwise be silently dropped.//// This buffer captures such messages the moment this module is imported (which// happens before any top-level `await` pause) and replays them once `ostiary`// sets up the real handler.//// Detection: regular Workers are instances of DedicatedWorkerGlobalScope.// Previously we checked `globalThis.onmessage === null`, but Safari initialises// that property as `undefined` rather than `null`, causing the check to fail.
/** @type {MessageEvent[]} */const _earlyMessages = [];
/** @type {null | (() => void)} */let _flushEarlyMessages = null;
if ( typeof DedicatedWorkerGlobalScope !== "undefined" && globalThis instanceof DedicatedWorkerGlobalScope) { const handler = /** @type {EventListener} */ ((event) => { _earlyMessages.push(/** @type {MessageEvent} */ (event)); });
globalThis.addEventListener("message", handler);
_flushEarlyMessages = () => { globalThis.removeEventListener("message", handler); };}
////////////////////////////////////////////// MISC////////////////////////////////////////////
/** * Manage incoming connections for a shared worker. * If a regular worker is used instead, it'll just execute the callback immediately. * * @template {MessagePort | Worker | MessengerRealm} T * @param {(context: MessagePort | T, firstConnection: boolean, connectionId: string) => void} callback * @param {T} [context] Uses `globalThis` by default. */export function ostiary( callback, context = /** @type {T} */ (/** @type {unknown} */ (globalThis)),) { if ( typeof DedicatedWorkerGlobalScope !== "undefined" && context instanceof DedicatedWorkerGlobalScope ) { callback(context, true, crypto.randomUUID());
// Replay any messages that arrived before the handler was registered. if (_flushEarlyMessages) { _flushEarlyMessages(); _flushEarlyMessages = null; const ctx = /** @type {EventTarget} */ (/** @type {unknown} */ (context)); _earlyMessages.splice(0).forEach((e) => { ctx.dispatchEvent( new MessageEvent("message", { data: e.data, ports: [...e.ports] }), ); }); }
return; }
const c = /** @type {any} */ (context); c.__id ??= crypto.randomUUID();
context.addEventListener( "connect", /** * @param {any} event */ (event) => { /** @type {MessagePort} */ const port = event.ports[0]; port.start();
// Initiate setup callback(port, !(c.__initiated ?? false), c.__id); c.__initiated = true; }, );}
/** * @param {Worker | SharedWorker} worker */export function workerLink(worker) { if (worker instanceof SharedWorker) { worker.port.start(); return worker.port; } else { return worker; }}
/** * @template {Record<string, (...args: any[]) => any>} Actions * @param {() => MessagePort | Worker} workerLinkCreator * @returns {ProxiedActions<Actions>} */export function workerProxy(workerLinkCreator) { /** @type {RpcChannel<{}, Actions> | undefined} */ let channel;
const proxy = new Proxy(/** @type {any} */ ({}), { get: (_target, /** @type {string} */ prop) => { /** @param {Parameters<Actions[any]>} args */ return (...args) => { channel ??= new RpcChannel(workerLinkCreator()); return channel.callMethod(prop, args); }; }, });
return /** @type {ProxiedActions<Actions>} */ (proxy);}
/** * @param {() => MessagePort | Worker | SharedWorker} workerCreator * @param {{ fromWorker?: (message: any) => Promise<{ data: any, transfer?: Transferable[] }>; toWorker?: (message: any) => Promise<{ data: any, transfer?: Transferable[] }> }} [hooks] * @returns {Tunnel} */export function workerTunnel(workerCreator, hooks = {}) { /** @type {MessagePort | Worker | undefined} */ let link;
const channel = new MessageChannel();
function ensureLink() { if (link) return link;
const workerOrLink = workerCreator();
link = workerOrLink instanceof SharedWorker ? workerLink(workerOrLink) : workerOrLink;
link.addEventListener("message", workerListener);
return link; }
channel.port1.addEventListener("message", async (event) => { // Send to worker const { data, transfer } = await hooks?.toWorker?.(event.data) ?? { data: event.data }; ensureLink().postMessage(data, { transfer }); });
/** * @param {Event} event */ const workerListener = async (event) => { // Receive from worker const msgEvent = /** @type {MessageEvent} */ (event); const { data, transfer } = await hooks?.fromWorker?.(msgEvent.data) ?? { data: msgEvent.data }; channel.port1.postMessage(data, { transfer }); };
channel.port1.start(); channel.port2.start();
return { disconnect: () => { link?.removeEventListener("message", workerListener); channel.port1.close(); channel.port2.close(); }, port: channel.port2, };}
////////////////////////////////////////////// RAW////////////////////////////////////////////
/** * @template T * @param {string} name * @param {T} args * @param {MessagePort | Worker | MessengerRealm} [context] Uses `globalThis` by default. */export function announce( name, args, context,) { const a = announcement(name, args); const transferables = getTransferables(a); (context ?? globalThis).postMessage(a, { transfer: transferables });}
/** * @template T * @param {string} name * @param {(args: T) => void} fn * @param {MessagePort | Worker | MessengerRealm} [context] */export function listen( name, fn, context = /** @type {MessengerRealm} */ (globalThis),) { const c = /** @type {any} */ (context);
if (!c.__incoming) { context.addEventListener("message", incomingAnnouncementsHandler(context)); c.__incoming = {}; }
c.__incoming[name] = debounceMicrotask(fn, { updateArguments: true });}
////////////////////////////////////////////// RPC////////////////////////////////////////////
/** * @template {Record<string, (...args: any[]) => any>} LocalAPI * @template {Record<string, (...args: any[]) => any>} RemoteAPI * @param {MessagePort | Worker | MessengerRealm} context * @param {RemoteAPI} actions * @returns {RpcChannel<{}, RemoteAPI>} */export function rpc(context, actions) { /** @type {RpcChannel<{}, RemoteAPI>} */ const channel = new RpcChannel(context, { expose: actions }); return channel;}
////////////////////////////////////////////// ⛔️////////////////////////////////////////////
const ANNOUNCEMENT = "announcement";
/** * @template T * @param {string} name * @param {T} args * @returns {Announcement<T>} */function announcement(name, args) { return { ns: ANNOUNCEMENT, name, key: xxh32(crypto.randomUUID()),
type: ANNOUNCEMENT, args, };}
/** * @param {MessagePort | Worker | MessengerRealm} context */function incomingAnnouncementsHandler(context) { /** @param {any} event */ return (event) => { const { ns, type } = event.data; if (ns !== ANNOUNCEMENT || type !== ANNOUNCEMENT) return; const announcement = /** @type {Announcement<any>} */ (event.data); const c = /** @type {any} */ (context); c.__incoming[announcement.name]?.(announcement.args); };}