diff --git a/src/components/output/raw/atproto/element.js b/src/components/output/raw/atproto/element.js index 4f493540..b57940ef 100644 --- a/src/components/output/raw/atproto/element.js +++ b/src/components/output/raw/atproto/element.js @@ -1,6 +1,6 @@ import { Client, ClientResponseError, ok } from "@atcute/client"; import { ComAtprotoSyncSubscribeRepos } from "@atcute/atproto"; -import { encode } from "@atcute/cbor"; +import { decode, encode } from "@atcute/cbor"; import { xxh32r } from "xxh32/dist/raw.js"; import * as Repo from "@atcute/repo"; import * as IDB from "idb-keyval"; @@ -94,9 +94,16 @@ class ATProtoOutput extends BroadcastedOutputElement { "sh.diffuse.output.trackBundle", ); - const tracks = bundles.flatMap((bundle) => bundle.tracks ?? []); - lastPersistedTracks = tracks; + /** @type {Track[]} */ + const tracks = []; + + for (const bundle of bundles) { + if (!bundle.data?.ref?.$link) continue; + const bytes = await this.#fetchBlob(bundle.data.ref.$link); + tracks.push(...decode(bytes)); + } + lastPersistedTracks = tracks; return tracks; }, put: async (data) => { @@ -107,21 +114,18 @@ class ATProtoOutput extends BroadcastedOutputElement { return; } - /** @type {TrackBundle[]} */ - const bundles = []; + const bytes = encode(data); + const blob = await this.#uploadBlob(bytes); + const id = xxh32r(bytes).toString(16); - for (let i = 0; i < data.length; i += 100) { - const chunk = data.slice(i, i + 100); - bundles.push({ - $type: "sh.diffuse.output.trackBundle", - id: xxh32r(encode(chunk)).toString(16), - tracks: chunk, - }); - } + /** @type {TrackBundle} */ + const bundle = { + $type: "sh.diffuse.output.trackBundle", + id, + data: blob, + }; - await this.putRecords("sh.diffuse.output.trackBundle", bundles, { - upsertBatchSize: 1, - }); + await this.putRecords("sh.diffuse.output.trackBundle", [bundle]); lastPersistedTracks = data; }, @@ -428,6 +432,34 @@ class ATProtoOutput extends BroadcastedOutputElement { return bytes; } + /** + * @param {Uint8Array} bytes + * @returns {Promise} + */ + async #uploadBlob(bytes) { + const rpc = this.#rpc; + if (!rpc) return; + const result = await ok(rpc.post("com.atproto.repo.uploadBlob", { + input: bytes, + headers: { "content-type": "application/octet-stream" }, + })); + return result.blob; + } + + /** + * @param {string} cid + * @returns {Promise} + */ + async #fetchBlob(cid) { + const rpc = this.#rpc; + const did = this.#did.value; + if (!rpc || !did) return new Uint8Array(); + return await ok(rpc.get("com.atproto.sync.getBlob", { + params: { did, cid }, + as: "bytes", + })); + } + /** * @template T * @param {string} collection @@ -497,9 +529,10 @@ class ATProtoOutput extends BroadcastedOutputElement { this.#writeDraining = true; while (this.#writeQueue.length > 0) { - const { fn, resolve, reject } = /** @type {{ fn: () => Promise, resolve: () => void, reject: (err: unknown) => void }} */ ( - this.#writeQueue.shift() - ); + const { fn, resolve, reject } = + /** @type {{ fn: () => Promise, resolve: () => void, reject: (err: unknown) => void }} */ ( + this.#writeQueue.shift() + ); try { await fn(); resolve(); @@ -533,7 +566,10 @@ class ATProtoOutput extends BroadcastedOutputElement { return this.#enqueueWrite(async () => { if (token.cancelled) return; try { - await this.#doPutRecords(collection, data, { deleteBatchSize, upsertBatchSize }, token); + await this.#doPutRecords(collection, data, { + deleteBatchSize, + upsertBatchSize, + }, token); } finally { if (this.#writeCancels.get(collection) === token) { this.#writeCancels.delete(collection); @@ -548,7 +584,12 @@ class ATProtoOutput extends BroadcastedOutputElement { * @param {{ deleteBatchSize: number, upsertBatchSize: number }} options * @param {{ cancelled: boolean }} token */ - async #doPutRecords(collection, data, { deleteBatchSize, upsertBatchSize }, token) { + async #doPutRecords( + collection, + data, + { deleteBatchSize, upsertBatchSize }, + token, + ) { const rpc = this.#rpc; const did = this.#did.value; if (!rpc || !did) return; @@ -620,7 +661,8 @@ class ATProtoOutput extends BroadcastedOutputElement { if (window.length + batch.length > WRITE_RATE_LIMIT) { const needed = window.length + batch.length - WRITE_RATE_LIMIT; const sorted = [...window].sort((a, b) => a.ts - b.ts); - const waitMs = WRITE_WINDOW_MS - (Date.now() - sorted[needed - 1].ts) + 1; + const waitMs = WRITE_WINDOW_MS - + (Date.now() - sorted[needed - 1].ts) + 1; await new Promise((resolve) => setTimeout(resolve, waitMs)); } diff --git a/src/components/output/raw/atproto/oauth-client-metadata.json b/src/components/output/raw/atproto/oauth-client-metadata.json index ae1451ce..6a14c8de 100644 --- a/src/components/output/raw/atproto/oauth-client-metadata.json +++ b/src/components/output/raw/atproto/oauth-client-metadata.json @@ -3,7 +3,7 @@ "client_name": "Diffuse", "client_uri": "https://elements.diffuse.sh", "redirect_uris": ["https://elements.diffuse.sh/oauth/callback"], - "scope": "atproto repo?collection=sh.diffuse.output.facet&collection=sh.diffuse.output.playlistItem&collection=sh.diffuse.output.playlistItemBundle&collection=sh.diffuse.output.setting&collection=sh.diffuse.output.track&collection=sh.diffuse.output.trackBundle", + "scope": "atproto blob:application/octet-stream repo?collection=sh.diffuse.output.facet&collection=sh.diffuse.output.playlistItem&collection=sh.diffuse.output.setting&collection=sh.diffuse.output.track&collection=sh.diffuse.output.trackBundle", "grant_types": ["authorization_code", "refresh_token"], "response_types": ["code"], "token_endpoint_auth_method": "none", diff --git a/src/definitions/index.ts b/src/definitions/index.ts index 0fc0a59a..1da1738f 100644 --- a/src/definitions/index.ts +++ b/src/definitions/index.ts @@ -1,7 +1,6 @@ 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/playlistItemBundle.json b/src/definitions/output/playlistItemBundle.json deleted file mode 100644 index 6b8769fd..00000000 --- a/src/definitions/output/playlistItemBundle.json +++ /dev/null @@ -1,23 +0,0 @@ -{ - "lexicon": 1, - "id": "sh.diffuse.output.playlistItemBundle", - "defs": { - "main": { - "type": "record", - "record": { - "type": "object", - "required": ["id", "playlistItems"], - "properties": { - "id": { "type": "string" }, - "createdAt": { "type": "string", "format": "datetime" }, - "playlistItems": { - "type": "array", - "description": "A bundle of playlist items", - "items": { "type": "ref", "ref": "sh.diffuse.output.playlistItem" } - }, - "updatedAt": { "type": "string", "format": "datetime" } - } - } - } - } -} diff --git a/src/definitions/output/trackBundle.json b/src/definitions/output/trackBundle.json index 6881824b..2118659d 100644 --- a/src/definitions/output/trackBundle.json +++ b/src/definitions/output/trackBundle.json @@ -6,14 +6,14 @@ "type": "record", "record": { "type": "object", - "required": ["id", "tracks"], + "required": ["id", "data"], "properties": { "id": { "type": "string" }, "createdAt": { "type": "string", "format": "datetime" }, - "tracks": { - "type": "array", - "description": "A bundle of tracks", - "items": { "type": "ref", "ref": "sh.diffuse.output.track" } + "data": { + "type": "blob", + "description": "CBOR-encoded tracks", + "accept": ["application/octet-stream"] }, "updatedAt": { "type": "string", "format": "datetime" } } diff --git a/src/definitions/types.d.ts b/src/definitions/types.d.ts index b9b739a1..467a2bd7 100644 --- a/src/definitions/types.d.ts +++ b/src/definitions/types.d.ts @@ -10,10 +10,6 @@ export type { Transformation, } from "./types/sh/diffuse/output/playlistItem.ts"; -export type { - Main as PlaylistItemBundle, -} from "./types/sh/diffuse/output/playlistItemBundle.ts"; - export type { Main as Theme } from "./types/sh/diffuse/output/theme.ts"; export type { Main as Setting } from "./types/sh/diffuse/output/setting.ts"; diff --git a/src/elements.vto b/src/elements.vto index 9df02f92..01b44430 100644 --- a/src/elements.vto +++ b/src/elements.vto @@ -200,10 +200,6 @@ definitions: desc: > Represents a single item in a playlist. Tracks are matched based on the given criteria. A playlist is formed by grouping items by their playlist property. url: "definitions/output/playlistItem.json" - - title: "Output / Playlist Item Bundle" - desc: > - A bundle of playlist items. - url: "definitions/output/playlistItemBundle.json" - title: "Output / Progress" desc: > Used to track progress of (long) audio playback. diff --git a/tasks/replace-gen-import-extensions.ts b/tasks/replace-gen-import-extensions.ts index be442792..9b771631 100644 --- a/tasks/replace-gen-import-extensions.ts +++ b/tasks/replace-gen-import-extensions.ts @@ -9,5 +9,4 @@ function replace(path: string) { } replace("./src/definitions/index.ts"); -replace("./src/definitions/types/sh/diffuse/output/playlistItemBundle.ts"); replace("./src/definitions/types/sh/diffuse/output/trackBundle.ts");