A music player that connects to your cloud/distributed storage. diffuse.sh
Something went wrong. Try again.
16 kB · 530 lines
JavaScript
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531import { decode, encode } from "@atcute/cbor";import { ifDefined } from "lit-html/directives/if-defined.js";import deepDiff from "@fry69/deep-diff";
import "~/components/output/polymorphic/indexed-db/element.js";
import * as CID from "~/common/cid.js";import { diff, strictEquality } from "~/common/compare.js";import { collectionSchema } from "~/common/self-describing.js";import { computed, signal } from "~/common/signal.js";import { compareTimestamps } from "~/common/temporal.js";import { OutputTransformer } from "../../base.js";import { defineElement } from "~/common/element.js";
/** * @import { SignalReader } from "~/common/signal.d.ts"; * @import { RenderArg } from "~/common/element.d.ts" * @import { OutputElement } from "@specs/components/output/types.d.ts" * * @import { Container } from "@specs/components/transformer/output/bytes/dasl-sync/types.d.ts" */
/** @type {Container<any>} */const EMPTY = { cid: undefined, data: [], inventory: { current: {}, removed: [] },};
/** * @extends {OutputTransformer<Uint8Array>} */class DaslBytesSyncOutputTransformer extends OutputTransformer { static NAME = "diffuse/transformer/output/bytes/dasl-sync";
constructor() { super();
const remote = this.base(); const local = this.#localOutput.get;
/** * @template {{ id: string; updatedAt: string }} T * @param {string} kind * @param {SignalReader<{ state: "loading" } | { state: "loaded"; data: Uint8Array | undefined } | { state: "error" }>} localCollection * @param {SignalReader<{ state: "loading" } | { state: "loaded"; data: Uint8Array | undefined } | { state: "error" }>} remoteCollection * @param {{ saveLocal: (bytes: Uint8Array) => Promise<void>; saveRemote: (bytes: Uint8Array) => Promise<void> }} sync */ const state = ( kind, localCollection, remoteCollection, { saveLocal, saveRemote }, ) => { const container = signal( /** @type {Container<T>} */ (EMPTY), { compare: strictEquality }, );
const isReady = signal(false); const merging = signal( { isBusy: false, lastCID: "", lastRemoteCID: "" }, { compare: diff }, );
this.effect(() => { if (!isReady.value) return; if (merging.value.isBusy) return;
const lc = localCollection(); const rc = remote.ready() ? remoteCollection() : undefined;
const lb = lc?.state === "loaded" ? lc.data : undefined; const rb = rc?.state === "loaded" ? rc.data : undefined; const rs = rc?.state;
/** @type {Container<T> | undefined} */ const l = lb ? decode(lb) : undefined;
/** @type {Container<T> | undefined} */ const r = rb && rs === "loaded" ? decode(rb) : undefined;
if (!r) { if (l) { container.value = l;
if (remote.ready() && rs === "loaded") { this.isLeader().then((isLeader) => { if (!isLeader) return; const bytes = this.save(l); saveRemote(bytes); }); } } } else if (!l) { container.value = r;
this.isLeader().then((isLeader) => { if (!isLeader) return; const bytes = this.save(r); saveLocal(bytes); }); } else if ( rs === "loaded" && this.hasDiverged({ local: l, remote: r }) ) { // If we've already reconciled this exact (local, remote) pair and it // made no progress (e.g. the remote write failed, or the merge // resolved to the container we already hold), don't re-enter the // merge. Re-merging here only toggles `merging`'s busy flag, which // re-runs this effect and spins the reactive sync loop, freezing the // page. Wait until local or remote actually changes before trying // again. const prev = merging.value; if (prev.lastCID === (l.cid ?? "") && prev.lastRemoteCID === (r.cid ?? "")) { return; }
// Async merge this.isLeader().then((isLeader) => { if (!isLeader) return;
merging.value = { isBusy: true, lastCID: prev.lastCID, lastRemoteCID: prev.lastRemoteCID, };
this.merge(l, r).then(async (c) => { /** * Remote cid after the (possible) push. Stays at the pre-merge * remote cid when there's nothing to push or the push fails, so * the effect-level guard above knows this reconcile was already * attempted and won't retry it in a loop. */ let resolvedRemoteCID = r.cid ?? ""; let reconciledCID = c.cid ?? "";
try { container.value = c;
const bytes = this.save(c);
if (c.cid !== l.cid) { await saveLocal(bytes); }
if (remote.ready() && rs === "loaded" && c.cid !== r.cid) { await saveRemote(bytes); resolvedRemoteCID = c.cid ?? ""; } } catch (err) { console.error("Merge failed:", err); } finally { merging.value = { isBusy: false, lastCID: reconciledCID, lastRemoteCID: resolvedRemoteCID, }; } }).catch((err) => { // Merge() itself rejected (before producing a reconciled // container). Reset the busy flag so the effect isn't stuck, and // record the exact (local, remote) pair we attempted so a // persistent failure is skipped rather than spinning forever. console.error("Merge failed:", err); merging.value = { isBusy: false, lastCID: l.cid ?? "", lastRemoteCID: r.cid ?? "", }; }); }); } else { container.value = l; } });
return computed(() => { if (!isReady.get()) isReady.value = true; return container.get(); }); };
// Container signals const facets = state( "facets", computed(() => local()?.facets.collection() ?? { state: "loading" }), remote.facets.collection, { saveLocal: async (v) => local()?.facets.save(v), saveRemote: remote.facets.save, }, );
const playlistItems = state( "playlistItems", computed(() => local()?.playlistItems.collection() ?? { state: "loading" } ), remote.playlistItems.collection, { saveLocal: async (v) => local()?.playlistItems.save(v), saveRemote: remote.playlistItems.save, }, );
const settings = state( "settings", computed(() => local()?.settings.collection() ?? { state: "loading" }), remote.settings.collection, { saveLocal: async (v) => local()?.settings.save(v), saveRemote: remote.settings.save, }, );
const tracks = state( "tracks", computed(() => local()?.tracks.collection() ?? { state: "loading" }), remote.tracks.collection, { saveLocal: async (v) => local()?.tracks.save(v), saveRemote: remote.tracks.save, }, );
// Output manager this.facets = this.managerProp( { save: async (v) => local()?.facets.save(v) }, remote.facets, remote.ready, facets, "facets", );
this.playlistItems = this.managerProp( { save: async (v) => local()?.playlistItems.save(v) }, remote.playlistItems, remote.ready, playlistItems, "playlistItems", );
this.settings = this.managerProp( { save: async (v) => local()?.settings.save(v) }, remote.settings, remote.ready, settings, "settings", );
this.tracks = this.managerProp( { save: async (v) => local()?.tracks.save(v) }, remote.tracks, remote.ready, tracks, "tracks", );
this.ready = () => true; }
// SIGNALS
#localOutput = signal( /** @type {OutputElement<any> | undefined} */ (undefined), );
// LIFECYCLE
/** * @override */ async connectedCallback() { // Broadcast if needed if (this.hasAttribute("group")) { this.broadcast(this.identifier, {}); }
super.connectedCallback();
/** @type {OutputElement<any> | null} */ const local = this.root().querySelector("dop-indexed-db"); if (!local) throw new Error("Can't find local output");
customElements.whenDefined(local.localName).then(() => { this.#localOutput.value = local; }); }
// DATA FUNCTIONS
/** * @template {{ id: string; updatedAt: string }} T * @param {{ previous: Container<T>, collection: T[], name?: import("~/common/self-describing.js").CollectionName }} _ * @returns {Promise<Container<T>>} */ async updateContainer({ previous, collection, name }) { const inventory = previous.inventory; const schema = name ? collectionSchema(name) : null;
const collIds = collection.map(({ id }) => id);
const currSet = new Set(Object.keys(inventory.current)); const collSet = new Set(collIds);
const newSet = collSet.difference(currSet); const remSet = currSet.difference(collSet);
const alreadyRemoved = new Set(inventory.removed); const allRemoved = alreadyRemoved.union(remSet);
/** @type {Record<string, string>} */ const current = { ...inventory.current };
remSet.forEach((id) => { delete current[id]; });
/** @type Promise<void>[] */ const promises = [];
collection.forEach((a) => { const encoded = encode(a);
promises.push((async () => { const cid = await CID.create(0x71, encoded); current[a.id] = cid; })()); });
await Promise.all(promises);
const newInventory = { current, removed: Array.from(allRemoved), };
const container = { cid: await CID.create(0x71, encode(newInventory)), data: collection, inventory: newInventory, }; return schema ? { ...container, $schema: schema } : container; }
/** * @template {{ id: string; updatedAt: string }} T * @param {{ local: Container<T>, remote: Container<T> }} _ */ hasDiverged({ local, remote }) { return local.cid !== remote.cid; }
/** * @template {{ id: string; updatedAt: string }} T * @param {Container<T>} a * @param {Container<T>} b * @returns {Promise<Container<T>>} */ async merge(a, b) { const removedA = new Set(a.inventory.removed); const removedB = new Set(b.inventory.removed); const allRemoved = removedA.union(removedB);
const currentA = a.inventory.current; const currentB = b.inventory.current;
const mapA = new Map(a.data.map((item) => [item.id, item])); const mapB = new Map(b.data.map((item) => [item.id, item]));
// Combine all known ids from both sides const allIds = new Set([ ...Object.keys(currentA), ...Object.keys(currentB), ]);
/** @type {Record<string, string>} */ const current = {};
/** @type {T[]} */ const data = [];
// Construct `current` and `data` /** @type {Promise<void>[]} */ const cidPromises = [];
for (const id of allIds) { if (allRemoved.has(id)) continue;
if (id in currentA && id in currentB) { const itemA = mapA.get(id); const itemB = mapB.get(id);
if (!itemA || !itemB) { console.warn("Should have found both items but didn't!"); continue; }
// Items are identical, no merge or CID recomputation needed if (currentA[id] === currentB[id]) { data.push(itemA); current[id] = currentA[id]; continue; }
const hasAUpdatedAt = Boolean(itemA.updatedAt); const hasBUpdatedAt = Boolean(itemB.updatedAt);
// The side that carries an `updatedAt` is the one that was edited, so // it should win over a pristine copy that never recorded a change (e.g. // a legacy seeded/starting-set facet persisted before it gained a // timestamp). Only when both carry one do we compare clocks. const isANewerThanB = (hasAUpdatedAt && !hasBUpdatedAt) ? true : (!hasAUpdatedAt && hasBUpdatedAt) ? false : hasAUpdatedAt && hasBUpdatedAt ? compareTimestamps(itemA.updatedAt, itemB.updatedAt) > 0 : false;
const newestItem = isANewerThanB ? itemA : itemB; const oldItem = isANewerThanB ? itemB : itemA;
/** @type {T} */ const mergedItem = { ...oldItem };
deepDiff.applyDiff(mergedItem, newestItem);
data.push(mergedItem);
cidPromises.push( CID.create(0x71, encode(mergedItem)).then((cid) => { current[id] = cid; }), ); } else { const item = mapA.get(id) ?? mapB.get(id);
if (item) { data.push(item); current[id] = currentA[id] ?? currentB[id]; } } }
await Promise.all(cidPromises);
// New inventory const updatedInventory = { current, removed: Array.from(allRemoved) };
return { cid: await CID.create(0x71, encode(updatedInventory)), data, inventory: updatedInventory, $schema: a.$schema, }; }
/** * @template {{ id: string; updatedAt: string }} T * @param {Container<T>} container * @returns {Uint8Array} */ save(container) { return encode(container); }
// OUTPUT MANAGER FUNCTIONS
/** * @template {{ id: string; updatedAt: string }} T * @param {{ save: (bytes: Uint8Array) => Promise<void> | void }} local * @param {{ collection: SignalReader<{ state: "loading" } | { state: "loaded"; data: Uint8Array | undefined } | { state: "error" }>, reload: () => Promise<void>, save: (bytes: Uint8Array) => Promise<void> }} remote * @param {SignalReader<boolean>} remoteReady * @param {SignalReader<Container<T>>} container * @param {import("~/common/self-describing.js").CollectionName} name * @returns {{ collection: SignalReader<{ state: "loading" } | { state: "loaded"; data: T[] } | { state: "error" }>, reload: () => Promise<void>, save: (items: T[]) => Promise<void> }} */ managerProp(local, remote, remoteReady, container, name) { return { collection: computed(() => { const c = container();
if (c.cid === undefined && remoteReady() && remote.collection().state === "loading") { return { state: "loading" }; }
return { state: "loaded", data: c.data }; }), reload: remote.reload, save: async (/** @type {T[]} */ newItems) => { const adjustedContainer = await this.updateContainer({ collection: newItems, previous: container(), name, });
const bytes = this.save(adjustedContainer); await local.save(bytes); }, }; }
// RENDER
/** * @param {RenderArg} _ */ render({ html }) { return html` <dop-indexed-db group="${ifDefined(this.getAttribute(`group`))}" namespace="${ifDefined(this.getAttribute(`namespace`))}" ></dop-indexed-db> `; }}
export default DaslBytesSyncOutputTransformer;
////////////////////////////////////////////// REGISTER////////////////////////////////////////////
export const CLASS = DaslBytesSyncOutputTransformer;export const NAME = "dtob-dasl-sync";
defineElement(NAME, CLASS);