From 059c4205f3b847a17283acfc5a14227ce9f41780 Mon Sep 17 00:00:00 2001 From: Steven Vandevelde Date: Mon, 11 May 2026 11:30:26 +0200 Subject: [PATCH] fix: improve atproto output --- docs/components/output/raw/atproto/SPEC.md | 11 + src/components/output/raw/atproto/element.js | 509 ++++++++---------- src/components/output/raw/atproto/types.d.ts | 2 +- .../output/raw/atproto-sync/element.js | 94 ++-- src/definitions/index.ts | 1 + src/definitions/output/trackBundle.json | 1 + src/facets/connect/atproto/index.inline.js | 2 +- 7 files changed, 297 insertions(+), 323 deletions(-) create mode 100644 docs/components/output/raw/atproto/SPEC.md diff --git a/docs/components/output/raw/atproto/SPEC.md b/docs/components/output/raw/atproto/SPEC.md new file mode 100644 index 00000000..c5787bd8 --- /dev/null +++ b/docs/components/output/raw/atproto/SPEC.md @@ -0,0 +1,11 @@ +# AT Protocol raw output + +This element implements the output element interface using the AT Protocol (PDS). + +## Requirements + +- The definition lexicons are used as the schema for each output type (tracks, playlist items, etc). +- The AT Protocol OAuth flow is preferred as the authentication method. +- The authenticated account must be remembered across browser sessions. +- The atproto pds has strict rate limits, we must opt for data structures that take this into consideration. There can be a large amount of tracks, 25000 for example, so a bundle would probably be preferred. Same for playlist items, maybe 5000 items or more. +- Reading and writing should be done as less as possible, though we don't want to miss out on any updates. diff --git a/src/components/output/raw/atproto/element.js b/src/components/output/raw/atproto/element.js index 3c8d0299..8a7ec53a 100644 --- a/src/components/output/raw/atproto/element.js +++ b/src/components/output/raw/atproto/element.js @@ -67,18 +67,22 @@ class ATProtoOutput extends BroadcastedOutputElement { await this.#restoreSettled.promise; return true; }, - facets: { - empty: () => [], - get: () => this.listRecords("sh.diffuse.output.facet"), - put: (data) => this.putRecords("sh.diffuse.output.facet", data), - }, - playlistItems: this.#blobCollection("sh.diffuse.output.playlistItemBundle", { groupBy: "playlist" }), - settings: { - empty: () => [], - get: () => this.listRecords("sh.diffuse.output.setting"), - put: (data) => this.putRecords("sh.diffuse.output.setting", data), - }, - tracks: this.#blobCollection("sh.diffuse.output.trackBundle"), + facets: this.#recordCollection("sh.diffuse.output.facet"), + playlistItems: this.#blobCollection( + "sh.diffuse.output.playlistItemBundle", + { groupBy: "playlist" }, + ), + settings: this.#recordCollection("sh.diffuse.output.setting"), + tracks: this.#blobCollection("sh.diffuse.output.trackBundle", { + groupBy: "scheme", + keyOf: (item) => { + const uri = String( + /** @type {Record} */ (item)["uri"] ?? "", + ); + const colon = uri.indexOf(":"); + return colon > 0 ? uri.substring(0, colon) : undefined; + }, + }), }); this.facets = this.#manager.facets; @@ -93,7 +97,7 @@ class ATProtoOutput extends BroadcastedOutputElement { #handle = signal(/** @type {string | null} */ (null)); #isOnline = signal(navigator.onLine); #rev = signal(/** @type {string | null} */ (null)); - #firehoseRev = signal(/** @type {string | null} */ (null)); + #firehoseRev = signal(/** @type {{ rev: string, collections: ReadonlySet } | null} */ (null)); /** @type {AsyncIterator | null} */ #firehoseIterator = null; @@ -200,7 +204,9 @@ class ATProtoOutput extends BroadcastedOutputElement { // original 401 response, which ok() wraps as a ClientResponseError. if (err instanceof ClientResponseError && err.status === 401) return true; if (err && typeof err === "object" && "cause" in err) { - return this.#isSessionError((/** @type {{ cause: unknown }} */ (err)).cause); + return this.#isSessionError( + (/** @type {{ cause: unknown }} */ (err)).cause, + ); } return false; } @@ -304,7 +310,8 @@ class ATProtoOutput extends BroadcastedOutputElement { #handleFirehoseCommit(message) { if (message.$type !== "com.atproto.sync.subscribeRepos#commit") return; - const commit = /** @type {{ repo: string, rev: string, ops?: Array<{ path: string }> }} */ (message); + const commit = + /** @type {{ repo: string, rev: string, ops?: Array<{ path: string }> }} */ (message); if (commit.repo !== this.#did.value) return; if (commit.rev === this.#rev.value) return; @@ -318,12 +325,7 @@ class ATProtoOutput extends BroadcastedOutputElement { if (touched.size === 0) return; - this.#firehoseRev.value = commit.rev; - - if (touched.has("sh.diffuse.output.facet")) this.#manager.facets.reload(); - if (touched.has("sh.diffuse.output.playlistItemBundle")) this.#manager.playlistItems.reload(); - if (touched.has("sh.diffuse.output.setting")) this.#manager.settings.reload(); - if (touched.has("sh.diffuse.output.trackBundle")) this.#manager.tracks.reload(); + this.#firehoseRev.value = { rev: commit.rev, collections: touched }; } /** @@ -348,35 +350,144 @@ class ATProtoOutput extends BroadcastedOutputElement { // RECORDS + /** + * Returns `{ empty, get, put }` for a small record collection (facets, settings). + * + * Tracks last-known remote state in closure so `put()` skips the `listRecords` + * round-trip on every write. Uses `putRecord` (upsert) to avoid create/update + * ambiguity; batches deletes via `applyWrites`. + * + * @param {string} nsid + */ + #recordCollection(nsid) { + /** @type {Map> | null} */ + let lastKnown = null; + + return { + empty: () => [], + get: async () => { + const records = await this.listRecords(nsid); + lastKnown = new Map( + /** @type {Array>} */ (records).map((r) => [ + String(r["id"]), + r, + ]), + ); + return records; + }, + put: async (/** @type {unknown[]} */ data) => { + const nsidTyped = /** @type {`${string}.${string}.${string}`} */ (nsid); + + /** @type {Map>} */ + const desired = new Map( + /** @type {Array<{ id: string }>} */ (data).map((r) => [ + r.id, + /** @type {Record} */ ({ $type: nsidTyped, ...r }), + ]), + ); + + const known = lastKnown ?? new Map(); + + /** @type {Array<[string, Record]>} */ + const upserts = []; + for (const [id, record] of desired) { + const existing = known.get(id); + if (existing && JSON.stringify(existing) === JSON.stringify(record)) { + continue; + } + upserts.push([id, record]); + } + + /** @type {WriteOp[]} */ + const deletes = []; + for (const id of known.keys()) { + if (!desired.has(id)) { + deletes.push({ + $type: "com.atproto.repo.applyWrites#delete", + collection: nsidTyped, + rkey: id, + }); + } + } + + if (upserts.length === 0 && deletes.length === 0) return; + + const newKnown = new Map(known); + for (const [id, record] of upserts) newKnown.set(id, record); + for (const { rkey } of deletes) newKnown.delete(rkey); + + const prior = this.#writeCancels.get(nsid); + if (prior) prior.cancelled = true; + const token = { cancelled: false }; + this.#writeCancels.set(nsid, token); + + await this.#enqueueWrite(async () => { + if (token.cancelled) return; + const rpc = this.#rpc; + const did = this.#did.value; + if (!rpc || !did) return; + this.#writing++; + try { + for (const [rkey, record] of upserts) { + const result = await ok(rpc.post("com.atproto.repo.putRecord", { + input: { repo: did, collection: nsidTyped, rkey, record }, + })); + if (result?.commit?.rev) this.#rev.value = result.commit.rev; + } + for (let i = 0; i < deletes.length; i += 100) { + const result = await ok(rpc.post("com.atproto.repo.applyWrites", { + input: { repo: did, writes: deletes.slice(i, i + 100) }, + })); + if (result?.commit?.rev) this.#rev.value = result.commit.rev; + } + lastKnown = newKnown; + } catch (err) { + if (this.#isSessionError(err)) { this.#clearSession(); return; } + throw err; + } finally { + this.#writing--; + if (this.#writeCancels.get(nsid) === token) { + this.#writeCancels.delete(nsid); + } + } + }); + }, + }; + } + /** * Returns `{ empty, get, put }` for a collection stored as CBOR blobs. - * When `groupBy` is provided each distinct value of that field gets its own - * bundle record; otherwise all items are packed into a single bundle. - * Each call gets its own closure-local state. + * Each distinct group key gets its own bundle record. + * + * `groupBy` is the field name stored on the bundle record (used by `get()` + * to reconstruct the key). `keyOf` extracts the group key from each item; + * defaults to `item[groupBy]` when omitted. * * @param {string} nsid - * @param {{ groupBy?: string }} [options] + * @param {{ groupBy: string, keyOf?: (item: unknown) => string | undefined }} options */ - #blobCollection(nsid, { groupBy } = {}) { - if (groupBy) { - /** @type {Map} groupKey → content hash */ - let lastHashes = new Map(); - /** @type {Map} groupKey → blob ref */ - let lastBlobs = new Map(); - - return { + #blobCollection(nsid, { groupBy, keyOf } = /** @type {any} */ ({})) { + /** @type {Map} groupKey → content hash */ + let lastHashes = new Map(); + /** @type {Map} groupKey → blob ref */ + let lastBlobs = new Map(); + + return { empty: () => /** @type {unknown[]} */ ([]), get: async () => { const bundles = await this.listRecords(nsid); /** @type {unknown[]} */ const items = []; - lastHashes = new Map(); - lastBlobs = new Map(); + /** @type {Map} */ + const newHashes = new Map(); + /** @type {Map} */ + const newBlobs = new Map(); for (const bundle of bundles) { if (!bundle.data?.ref?.$link) continue; - const key = /** @type {Record} */ (bundle)[groupBy]; + const key = + /** @type {Record} */ (bundle)[groupBy]; if (typeof key !== "string") continue; const bytes = await this.#fetchBlob(bundle.data.ref.$link); @@ -384,20 +495,33 @@ class ATProtoOutput extends BroadcastedOutputElement { if (!Array.isArray(groupItems)) continue; items.push(...groupItems); - lastHashes.set(key, xxh32r(bytes).toString(16)); - lastBlobs.set(key, bundle.data); + // Hash the re-encoded form so put() compares apples to apples. + // Raw PDS bytes may not equal encode(decode(bytes)) if field order + // or numeric encoding differs, causing spurious re-uploads. + newHashes.set(key, xxh32r(encode(groupItems)).toString(16)); + newBlobs.set(key, bundle.data); } + // Assign atomically after all fetches complete so a concurrent put() + // never sees a partially-populated lastHashes and generates wrong #create ops. + lastHashes = newHashes; + lastBlobs = newBlobs; return items; }, put: async (/** @type {unknown[]} */ data) => { - const nsidTyped = /** @type {`${string}.${string}.${string}`} */ (nsid); + const nsidTyped = + /** @type {`${string}.${string}.${string}`} */ (nsid); + + const extractKey = keyOf ?? + ((/** @type {unknown} */ item) => + /** @type {string | undefined} */ ( + /** @type {Record} */ (item)[groupBy] + )); /** @type {Map} */ const groups = new Map(); for (const item of data) { - const record = /** @type {Record} */ (item); - const key = record[groupBy]; + const key = extractKey(item); if (typeof key !== "string") continue; const group = groups.get(key) ?? []; if (!groups.has(key)) groups.set(key, group); @@ -408,8 +532,9 @@ class ATProtoOutput extends BroadcastedOutputElement { const newHashes = new Map(lastHashes); const newBlobs = new Map(lastBlobs); - /** @type {WriteOp[]} */ - const writes = []; + // Upload blobs for changed groups (outside write queue — not a mutation) + /** @type {Array<{ rkey: string, value: unknown }>} */ + const upserts = []; for (const [key, groupItems] of groups) { const bytes = encode(groupItems); @@ -421,65 +546,74 @@ class ATProtoOutput extends BroadcastedOutputElement { if (!blob) continue; const rkey = xxh32r(encode(key)).toString(16); - const value = { $type: nsidTyped, id: rkey, [groupBy]: key, data: blob }; - - if (lastHashes.has(key)) { - writes.push({ $type: "com.atproto.repo.applyWrites#update", collection: nsidTyped, rkey, value }); - } else { - writes.push({ $type: "com.atproto.repo.applyWrites#create", collection: nsidTyped, rkey, value }); - } - + const value = { + $type: nsidTyped, + id: rkey, + [groupBy]: key, + data: blob, + }; + upserts.push({ rkey, value }); newHashes.set(key, hash); newBlobs.set(key, blob); } + /** @type {WriteOp[]} */ + const deletes = []; for (const key of lastHashes.keys()) { if (!groups.has(key)) { const rkey = xxh32r(encode(key)).toString(16); - writes.push({ $type: "com.atproto.repo.applyWrites#delete", collection: nsidTyped, rkey }); + deletes.push({ + $type: "com.atproto.repo.applyWrites#delete", + collection: nsidTyped, + rkey, + }); newHashes.delete(key); newBlobs.delete(key); } } - if (writes.length === 0) return; - - await this.#applyWriteOps(writes); + if (upserts.length === 0 && deletes.length === 0) return; + + await this.#enqueueWrite(async () => { + const rpc = this.#rpc; + const did = this.#did.value; + if (!rpc || !did) return; + this.#writing++; + try { + // putRecord is a true upsert — no need to track create vs update + for (const { rkey, value } of upserts) { + const result = await ok(rpc.post("com.atproto.repo.putRecord", { + input: { + repo: did, + collection: nsidTyped, + rkey, + record: /** @type {Record} */ (value), + }, + })); + if (result?.commit?.rev) this.#rev.value = result.commit.rev; + } + for (let i = 0; i < deletes.length; i += 100) { + const result = await ok( + rpc.post("com.atproto.repo.applyWrites", { + input: { repo: did, writes: deletes.slice(i, i + 100) }, + }), + ); + if (result?.commit?.rev) this.#rev.value = result.commit.rev; + } + } catch (err) { + if (this.#isSessionError(err)) { + this.#clearSession(); + return; + } + throw err; + } finally { + this.#writing--; + } + }); lastHashes = newHashes; lastBlobs = newBlobs; }, - }; - } - - // Single-blob variant (used for tracks) - /** @type {unknown[] | null} */ - let lastPersisted = null; - - return { - empty: () => /** @type {unknown[]} */ ([]), - get: async () => { - const bundles = await this.listRecords(nsid); - /** @type {unknown[]} */ - const items = []; - for (const bundle of bundles) { - if (!bundle.data?.ref?.$link) continue; - const bytes = await this.#fetchBlob(bundle.data.ref.$link); - items.push(...decode(bytes)); - } - lastPersisted = items; - return items; - }, - put: async (/** @type {unknown[]} */ data) => { - if (xxh32r(encode(lastPersisted ?? [])) === xxh32r(encode(data))) return; - const bytes = encode(data); - const blob = await this.#uploadBlob(bytes); - if (!blob) return; - const id = xxh32r(bytes).toString(16); - const bundle = { id, data: blob }; - await this.putRecords(nsid, [bundle]); - lastPersisted = data; - }, }; } @@ -554,18 +688,19 @@ class ATProtoOutput extends BroadcastedOutputElement { /** @type {string | undefined} */ let cursor; do { - const page = /** @type {{ records: { value: unknown }[], cursor?: string }} */ ( - await ok(this.#rpc.get("com.atproto.repo.listRecords", { - params: { - repo: - /** @type {import("@atcute/lexicons").ActorIdentifier} */ (did), - collection: - /** @type {`${string}.${string}.${string}`} */ (collection), - limit: 100, - cursor, - }, - })) - ); + const page = + /** @type {{ records: { value: unknown }[], cursor?: string }} */ ( + await ok(this.#rpc.get("com.atproto.repo.listRecords", { + params: { + repo: + /** @type {import("@atcute/lexicons").ActorIdentifier} */ (did), + collection: + /** @type {`${string}.${string}.${string}`} */ (collection), + limit: 100, + cursor, + }, + })) + ); records.push(...page.records.map((r) => /** @type {T} */ (r.value))); cursor = page.cursor; } while (cursor); @@ -613,184 +748,6 @@ class ATProtoOutput extends BroadcastedOutputElement { this.#writeDraining = false; } - /** - * @param {string} collection - * @param {Array<{ id: string }>} data - * @param {{ deleteBatchSize?: number, upsertBatchSize?: number }} [options] - */ - async putRecords( - collection, - data, - { deleteBatchSize = 100, upsertBatchSize = deleteBatchSize } = {}, - ) { - if (!this.#rpc || !this.#did.value) return; - - // Supersede any prior write for this collection - const prior = this.#writeCancels.get(collection); - if (prior) prior.cancelled = true; - - const token = { cancelled: false }; - this.#writeCancels.set(collection, token); - - return this.#enqueueWrite(async () => { - if (token.cancelled) return; - try { - await this.#doPutRecords(collection, data, { - deleteBatchSize, - upsertBatchSize, - }, token); - } finally { - if (this.#writeCancels.get(collection) === token) { - this.#writeCancels.delete(collection); - } - } - }); - } - - /** - * @param {string} collection - * @param {Array<{ id: string }>} data - * @param {{ deleteBatchSize: number, upsertBatchSize: number }} options - * @param {{ cancelled: boolean }} token - */ - async #doPutRecords( - collection, - data, - { deleteBatchSize, upsertBatchSize }, - token, - ) { - const rpc = this.#rpc; - const did = this.#did.value; - if (!rpc || !did) return; - - const nsid = /** @type {`${string}.${string}.${string}`} */ (collection); - - this.#writing++; - try { - // 1. Fetch current state - /** @type {Map} */ - const existing = new Map(); - - /** @type {string | undefined} */ - let cursor; - do { - const page = /** @type {{ records: { uri: string, value: unknown }[], cursor?: string }} */ ( - await ok(rpc.get("com.atproto.repo.listRecords", { - params: { repo: did, collection: nsid, limit: 100, cursor }, - })) - ); - for (const { uri, value } of page.records) { - const record = /** @type {{ id: string }} */ (value); - const rkey = /** @type {string} */ (uri.split("/").at(-1)); - existing.set(record.id, { rkey, value: record }); - } - cursor = page.cursor; - } while (cursor); - - // 2. Build desired state - const desired = new Map( - data.map((record) => [record.id, { $type: nsid, ...record }]), - ); - - // 3. Compute diff - /** @type {WriteOp[]} */ - const deletes = []; - - /** @type {WriteOp[]} */ - const upserts = []; - - for (const [id, { rkey }] of existing) { - if (!desired.has(id)) { - deletes.push({ $type: "com.atproto.repo.applyWrites#delete", collection: nsid, rkey }); - } - } - - for (const [id, record] of desired) { - const entry = existing.get(id); - - if (!entry) { - upserts.push({ - $type: "com.atproto.repo.applyWrites#create", - collection: nsid, - rkey: id, - value: record, - }); - } else if (JSON.stringify(entry.value) !== JSON.stringify(record)) { - upserts.push({ - $type: "com.atproto.repo.applyWrites#update", - collection: nsid, - rkey: entry.rkey, - value: record, - }); - } - } - - // 4. Apply batches - /** @param {WriteOp[]} batch */ - const applyBatch = async (batch) => { - const result = await ok(rpc.post("com.atproto.repo.applyWrites", { - input: { repo: did, writes: batch }, - })); - - if (result?.commit?.rev) { - this.#rev.value = result.commit.rev; - } - }; - - for (let i = 0; i < deletes.length; i += deleteBatchSize) { - if (token.cancelled) return; - await applyBatch(deletes.slice(i, i + deleteBatchSize)); - } - - for (let i = 0; i < upserts.length; i += upsertBatchSize) { - if (token.cancelled) return; - await applyBatch(upserts.slice(i, i + upsertBatchSize)); - } - } catch (err) { - if (this.#isSessionError(err)) { - this.#clearSession(); - return; - } - - throw err; - } finally { - this.#writing--; - } - } - - /** - * Apply pre-computed write operations via applyWrites, respecting the write - * queue and #writing guard. - * - * @param {WriteOp[]} writes - * @param {number} [batchSize] - */ - #applyWriteOps(writes, batchSize = 100) { - if (!this.#rpc || !this.#did.value) return Promise.resolve(); - return this.#enqueueWrite(async () => { - const rpc = this.#rpc; - const did = this.#did.value; - if (!rpc || !did) return; - this.#writing++; - try { - for (let i = 0; i < writes.length; i += batchSize) { - const batch = writes.slice(i, i + batchSize); - const result = await ok(rpc.post("com.atproto.repo.applyWrites", { - input: { repo: did, writes: batch }, - })); - if (result?.commit?.rev) this.#rev.value = result.commit.rev; - } - } catch (err) { - if (this.#isSessionError(err)) { - this.#clearSession(); - return; - } - throw err; - } finally { - this.#writing--; - } - }); - } } export default ATProtoOutput; diff --git a/src/components/output/raw/atproto/types.d.ts b/src/components/output/raw/atproto/types.d.ts index b2846ffa..50305595 100644 --- a/src/components/output/raw/atproto/types.d.ts +++ b/src/components/output/raw/atproto/types.d.ts @@ -5,7 +5,7 @@ export type ATProtoOutputElement = & OutputElement & { did: SignalReader; - firehoseRev: SignalReader; + firehoseRev: SignalReader<{ rev: string; collections: ReadonlySet } | null>; handle: SignalReader; rev: SignalReader; diff --git a/src/components/transformer/output/raw/atproto-sync/element.js b/src/components/transformer/output/raw/atproto-sync/element.js index 5fbd8c40..b78f130c 100644 --- a/src/components/transformer/output/raw/atproto-sync/element.js +++ b/src/components/transformer/output/raw/atproto-sync/element.js @@ -20,6 +20,14 @@ const COLLECTIONS = /** @type {const} */ ([ "tracks", ]); +/** @type {Record} */ +const NSID_TO_COLLECTION = { + "sh.diffuse.output.facet": "facets", + "sh.diffuse.output.playlistItemBundle": "playlistItems", + "sh.diffuse.output.setting": "settings", + "sh.diffuse.output.trackBundle": "tracks", +}; + const STORAGE_PREFIX = "diffuse/transformer/output/atproto-sync"; /** @@ -100,9 +108,6 @@ class ATProtoOutputSyncTransformer extends OutputTransformer { ); } - // Update known ids - this.#trackIds(name, newData); - await l[name].save(newData); if (remote.ready()) { @@ -114,6 +119,7 @@ class ATProtoOutputSyncTransformer extends OutputTransformer { ? this.#mergeRecords(name, newData, /** @type {typeof newData} */ (remoteCol.data)) : newData; + this.#markDirty(); remote[name].save(dataForRemote).then(() => { const rev = this.#atproto()?.rev(); if (rev) this.#storeRev(rev); @@ -145,10 +151,16 @@ class ATProtoOutputSyncTransformer extends OutputTransformer { this.effect(async () => { const atproto = this.#atproto(); if (!atproto) return; - if (!atproto.firehoseRev()) return; + const firehose = atproto.firehoseRev(); + if (!firehose) return; if (!remote.ready()) return; if (!(await this.isLeader())) return; - this.#sync(); + const touched = /** @type {string[]} */ ( + [...firehose.collections] + .map((nsid) => NSID_TO_COLLECTION[nsid]) + .filter((name) => name !== undefined) + ); + if (touched.length > 0) this.#sync(touched, firehose.rev); }); }); } @@ -170,7 +182,11 @@ class ATProtoOutputSyncTransformer extends OutputTransformer { // SYNC - async #sync() { + /** + * @param {readonly string[]} collections + * @param {string | null} [knownRev] + */ + async #sync(collections = COLLECTIONS, knownRev = null) { if (this.#syncing) return; this.#syncing = true; @@ -181,63 +197,51 @@ class ATProtoOutputSyncTransformer extends OutputTransformer { if (!l || !atproto || !remote.ready()) return; - const remoteRev = await atproto.getLatestCommit(); - if (!remoteRev) return; + const isFull = collections === COLLECTIONS; + /** @type {Record} */ + const lAny = l; + /** @type {Record} */ + const remoteAny = remote; - const localRev = this.#getStoredRev(); - const dirty = this.#isDirty(); + const remoteRev = knownRev ?? await atproto.getLatestCommit(); + if (!remoteRev) return; - if (localRev === remoteRev && !dirty) { - return; + if (isFull) { + const localRev = this.#getStoredRev(); + const dirty = this.#isDirty(); + if (localRev === remoteRev && !dirty) return; } - // Fetch remote data - for (const name of COLLECTIONS) { - await remote[name].reload(); + // Fetch remote data for the affected collections only + for (const name of collections) { + await remoteAny[name].reload(); } const localCollections = await Promise.all( - COLLECTIONS.map((name) => Output.data(l[name])), + collections.map((name) => Output.data(lAny[name])), ); const localHasData = localCollections.some( (data) => Array.isArray(data) && data.length > 0, ); - // Seed knownIds from local data if empty and not dirty. - // Handles the case where localStorage was cleared but IndexedDB still - // has data — without this, #mergeRecords can't detect remote deletions - // (knownIds.has(id) is always false) and deleted records re-appear locally. - if (!dirty) { - for (let i = 0; i < COLLECTIONS.length; i++) { - const name = COLLECTIONS[i]; - const localData = localCollections[i]; - if ( - this.#getKnownIds(name).size === 0 && - Array.isArray(localData) && localData.length > 0 - ) { - this.#trackIds(name, localData); - } - } - } - - if (!localHasData && !dirty) { + if (!localHasData && !this.#isDirty()) { // Local is empty and clean — just pull remote - for (const name of COLLECTIONS) { - const remoteCol = remote[name].collection(); + for (const name of collections) { + const remoteCol = remoteAny[name].collection(); if ( remoteCol.state === "loaded" && Array.isArray(remoteCol.data) && remoteCol.data.length > 0 ) { this.#trackIds(name, remoteCol.data); - await l[name].save(remoteCol.data); + await lAny[name].save(remoteCol.data); } } } else { // Union merge - for (const name of COLLECTIONS) { - const localCol = l[name].collection(); - const remoteCol = remote[name].collection(); + for (const name of collections) { + const localCol = lAny[name].collection(); + const remoteCol = remoteAny[name].collection(); const localArr = localCol.state === "loaded" && Array.isArray(localCol.data) ? localCol.data @@ -249,17 +253,17 @@ class ATProtoOutputSyncTransformer extends OutputTransformer { const merged = this.#mergeRecords(name, localArr, remoteArr); - this.#trackIds(name, merged); - await l[name].save(merged); + await lAny[name].save(merged); if (this.#differFromRemote(merged, remoteArr)) { - await remote[name].save(merged); + await remoteAny[name].save(merged); } + this.#trackIds(name, merged); } } - this.#storeRev(atproto.rev()); - this.#clearDirty(); + this.#storeRev(atproto.rev() ?? remoteRev); + if (isFull) this.#clearDirty(); } catch (err) { console.warn("Sync failed:", err); } finally { diff --git a/src/definitions/index.ts b/src/definitions/index.ts index 1da1738f..0fc0a59a 100644 --- a/src/definitions/index.ts +++ b/src/definitions/index.ts @@ -1,6 +1,7 @@ export * as ShDiffuseOutputCollaboration from "./types/sh/diffuse/output/collaboration.ts"; export * as ShDiffuseOutputFacet from "./types/sh/diffuse/output/facet.ts"; export * as ShDiffuseOutputPlaylistItem from "./types/sh/diffuse/output/playlistItem.ts"; +export * as ShDiffuseOutputPlaylistItemBundle from "./types/sh/diffuse/output/playlistItemBundle.ts"; export * as ShDiffuseOutputSetting from "./types/sh/diffuse/output/setting.ts"; export * as ShDiffuseOutputTrack from "./types/sh/diffuse/output/track.ts"; export * as ShDiffuseOutputTrackBundle from "./types/sh/diffuse/output/trackBundle.ts"; diff --git a/src/definitions/output/trackBundle.json b/src/definitions/output/trackBundle.json index 2118659d..17db08c7 100644 --- a/src/definitions/output/trackBundle.json +++ b/src/definitions/output/trackBundle.json @@ -9,6 +9,7 @@ "required": ["id", "data"], "properties": { "id": { "type": "string" }, + "scheme": { "type": "string" }, "createdAt": { "type": "string", "format": "datetime" }, "data": { "type": "blob", diff --git a/src/facets/connect/atproto/index.inline.js b/src/facets/connect/atproto/index.inline.js index 0e709670..85663d37 100644 --- a/src/facets/connect/atproto/index.inline.js +++ b/src/facets/connect/atproto/index.inline.js @@ -61,7 +61,7 @@ if (true) { /** @type {HTMLElement} */ (document.querySelector("main")), ); - await atprotoEl.whenRestored(); + await atprotoEl?.whenRestored(); } //////////////////////////////////////////// -- 2.51.2