From 69e2ad9394bf161a2088c657539a10d5ad62085c Mon Sep 17 00:00:00 2001 From: Roscoe Rubin-Rottenberg Date: Tue, 12 May 2026 17:32:38 -0400 Subject: [PATCH] feat: plyr tracks as sounds --- data-plane/indexing/index.ts | 3 + data-plane/indexing/plugins/audio.ts | 2 +- data-plane/indexing/plugins/plyr-track.ts | 84 ++++++++++ data-plane/indexing/processor.ts | 22 +-- data-plane/subscription.ts | 1 + hydration/feed.ts | 32 +++- lexicons/fm/plyr/track.json | 119 ++++++++++++++ tests/sounds_test.ts | 179 ++++++++++++++++++++++ views/index.ts | 65 +++++--- 9 files changed, 460 insertions(+), 47 deletions(-) create mode 100644 data-plane/indexing/plugins/plyr-track.ts create mode 100644 lexicons/fm/plyr/track.json diff --git a/data-plane/indexing/index.ts b/data-plane/indexing/index.ts index 9d540c7..9c31cde 100644 --- a/data-plane/indexing/index.ts +++ b/data-plane/indexing/index.ts @@ -28,6 +28,7 @@ import * as Profile from "./plugins/profile.ts"; import * as Repost from "./plugins/repost.ts"; import * as Story from "./plugins/story.ts"; import * as Audio from "./plugins/audio.ts"; +import * as PlyrTrack from "./plugins/plyr-track.ts"; import * as Labeler from "./plugins/labeler.ts"; import { RecordProcessor } from "./processor.ts"; import { ServerConfig } from "../../config.ts"; @@ -45,6 +46,7 @@ export class IndexingService { generator: Generator.PluginType; story: Story.PluginType; audio: Audio.PluginType; + plyrTrack: PlyrTrack.PluginType; labeler: Labeler.PluginType; }; private pushService?: PushService; @@ -68,6 +70,7 @@ export class IndexingService { generator: Generator.makePlugin(this.db, this.background), story: Story.makePlugin(this.db, this.background), audio: Audio.makePlugin(this.db, this.background), + plyrTrack: PlyrTrack.makePlugin(this.db, this.background), labeler: Labeler.makePlugin(this.db, this.background), }; diff --git a/data-plane/indexing/plugins/audio.ts b/data-plane/indexing/plugins/audio.ts index 524a398..843eda8 100644 --- a/data-plane/indexing/plugins/audio.ts +++ b/data-plane/indexing/plugins/audio.ts @@ -33,7 +33,7 @@ const insertFn = async ( // Use findOneAndUpdate with upsert to handle potential duplicate key errors const insertedAudio = await db.models.Audio.findOneAndUpdate( { uri: uri.toString() }, - { $set: audio }, + { $set: audio, $setOnInsert: { useCount: 0 } }, { upsert: true, returnDocument: "after", includeResultMetadata: false }, ); return insertedAudio; diff --git a/data-plane/indexing/plugins/plyr-track.ts b/data-plane/indexing/plugins/plyr-track.ts new file mode 100644 index 0000000..106a3f6 --- /dev/null +++ b/data-plane/indexing/plugins/plyr-track.ts @@ -0,0 +1,84 @@ +import { Cid } from "@atp/lex"; +import { AtUri, normalizeDatetimeAlways } from "@atp/syntax"; +import * as fm from "../../../lex/fm.ts"; +import { BackgroundQueue } from "../../background.ts"; +import { Database } from "../../db/index.ts"; +import { AudioDocument, MediaRef } from "../../db/models.ts"; +import { RecordProcessor } from "../processor.ts"; + +const schema = fm.plyr.track.main; +type PlyrTrackRecord = fm.plyr.track.Main; +type IndexedAudio = AudioDocument; + +const insertFn = async ( + db: Database, + uri: AtUri, + cid: Cid, + obj: PlyrTrackRecord, + timestamp: string, +): Promise => { + if (obj.supportGate || !obj.audioBlob) { + return null; + } + + const audio: Omit = { + uri: uri.toString(), + cid: cid.toString(), + authorDid: uri.host, + sound: obj.audioBlob as unknown as MediaRef, + title: obj.title, + details: { + artist: obj.artist, + title: obj.title, + }, + createdAt: normalizeDatetimeAlways(obj.createdAt), + indexedAt: timestamp, + }; + + const insertedAudio = await db.models.Audio.findOneAndUpdate( + { uri: uri.toString() }, + { $set: audio, $setOnInsert: { useCount: 0 } }, + { upsert: true, returnDocument: "after", includeResultMetadata: false }, + ); + return insertedAudio; +}; + +const findDuplicate = (): AtUri | null => { + return null; +}; + +const notifsForInsert = () => { + return []; +}; + +const deleteFn = async ( + db: Database, + uri: AtUri, +): Promise => { + const deleted = await db.models.Audio.findOneAndDelete({ + uri: uri.toString(), + }); + return deleted; +}; + +const notifsForDelete = () => { + return { notifs: [], toDelete: [] }; +}; + +export type PluginType = RecordProcessor; + +export const makePlugin = ( + db: Database, + background: BackgroundQueue, +): PluginType => { + return new RecordProcessor(db, background, { + schema, + insertFn, + findDuplicate, + deleteFn, + notifsForInsert, + notifsForDelete, + }); +}; + +export default makePlugin; diff --git a/data-plane/indexing/processor.ts b/data-plane/indexing/processor.ts index affbc2d..99014db 100644 --- a/data-plane/indexing/processor.ts +++ b/data-plane/indexing/processor.ts @@ -29,7 +29,6 @@ type RecordProcessorOptions = { replacedBy: TRow | null, ) => { notifs: Notif[]; toDelete: string[] }; updateAggregates?: (db: Database, obj: TRow) => Promise; - deleteRecordIfInsertReturnsNull?: boolean; }; type Notif = { @@ -154,12 +153,6 @@ export class RecordProcessor { timestamp, ); if (!inserted) { - if (this.options.deleteRecordIfInsertReturnsNull) { - await this.db.models.Record.deleteOne({ uri: uri.toString() }); - await this.db.models.DuplicateRecord.deleteOne({ - uri: uri.toString(), - }); - } return; } @@ -239,19 +232,10 @@ export class RecordProcessor { timestamp, ); if (!inserted) { - if (this.options.deleteRecordIfInsertReturnsNull) { - await this.db.models.Record.deleteOne({ uri: uri.toString() }); - await this.db.models.DuplicateRecord.deleteOne({ - uri: uri.toString(), - }); - if (!opts?.disableNotifs) { - await this.handleNotifs({ deleted }); - } - return; + if (!opts?.disableNotifs) { + await this.handleNotifs({ deleted }); } - throw new Error( - "Record update failed: removed from index but could not be replaced", - ); + return; } this.aggregateOnCommit(inserted); if (!opts?.disableNotifs) { diff --git a/data-plane/subscription.ts b/data-plane/subscription.ts index 1445667..fa9e252 100644 --- a/data-plane/subscription.ts +++ b/data-plane/subscription.ts @@ -207,6 +207,7 @@ function createFirehose(opts: { }, filterCollections: [ "so.sprk.*", + "fm.plyr.track", ], }); return { firehose, runner }; diff --git a/hydration/feed.ts b/hydration/feed.ts index 8b4a62d..8734d21 100644 --- a/hydration/feed.ts +++ b/hydration/feed.ts @@ -1,4 +1,6 @@ import * as so from "../lex/so.ts"; +import * as fm from "../lex/fm.ts"; +import { AtUri } from "@atp/syntax"; import { uriToDid as didFromUri } from "../utils/uris.ts"; import { HydrationMap, @@ -9,6 +11,7 @@ import { split, } from "./util.ts"; import { DataPlane } from "../data-plane/index.ts"; +import { Record as DataPlaneRecord } from "../data-plane/routes/records.ts"; export type FeedGenRecord = so.sprk.feed.generator.Main; export type LikeRecord = so.sprk.feed.like.Main; @@ -16,12 +19,14 @@ export type PostRecord = so.sprk.feed.post.Main; export type ReplyRecord = so.sprk.feed.reply.Main; export type RepostRecord = so.sprk.feed.repost.Main; export type AudioRecord = so.sprk.sound.audio.Main; +export type PlyrTrackRecord = fm.plyr.track.Main; +export type SoundRecord = AudioRecord | PlyrTrackRecord; export type Post = RecordInfo; export type Posts = HydrationMap; export type Reply = RecordInfo; export type Replies = HydrationMap; -export type Sound = RecordInfo; +export type Sound = RecordInfo; export type Sounds = HydrationMap; export type SoundAgg = { @@ -168,11 +173,7 @@ export class FeedHydrator { const res = await this.dataplane.records.getRecords(need); return need.reduce((acc, uri, i) => { - const record = parseRecord( - so.sprk.sound.audio.main, - res.records[i], - includeTakedowns, - ); + const record = parseSoundRecord(res.records[i], includeTakedowns); return acc.set( uri, record ? record : null, @@ -381,3 +382,22 @@ export class FeedHydrator { }, new HydrationMap()); } } + +const parseSoundRecord = ( + record: DataPlaneRecord, + includeTakedowns: boolean, +): Sound | undefined => { + const collection = new AtUri(record.uri).collection; + if (collection === fm.plyr.track.$type) { + return parseRecord( + fm.plyr.track.main, + record, + includeTakedowns, + ); + } + return parseRecord( + so.sprk.sound.audio.main, + record, + includeTakedowns, + ); +}; diff --git a/lexicons/fm/plyr/track.json b/lexicons/fm/plyr/track.json new file mode 100644 index 0000000..8e87f35 --- /dev/null +++ b/lexicons/fm/plyr/track.json @@ -0,0 +1,119 @@ +{ + "lexicon": 1, + "id": "fm.plyr.track", + "defs": { + "main": { + "type": "record", + "description": "A music track published to the AT Protocol network.", + "key": "tid", + "record": { + "type": "object", + "required": ["title", "artist", "fileType", "createdAt"], + "properties": { + "title": { + "type": "string", + "description": "The title of the track.", + "minLength": 1, + "maxLength": 256 + }, + "artist": { + "type": "string", + "description": "The artist name (display name of the uploader).", + "minLength": 1, + "maxLength": 256 + }, + "audioUrl": { + "type": "string", + "format": "uri", + "description": "URL to the audio file. Optional when audioBlob is present." + }, + "fileType": { + "type": "string", + "description": "Audio file format extension (e.g., mp3, wav, flac).", + "minLength": 1, + "maxLength": 16 + }, + "album": { + "type": "string", + "description": "Album name this track belongs to.", + "maxLength": 256 + }, + "duration": { + "type": "integer", + "description": "Duration in seconds.", + "minimum": 0 + }, + "features": { + "type": "array", + "description": "Featured artists on this track.", + "items": { + "type": "ref", + "ref": "#featuredArtist" + }, + "maxLength": 10 + }, + "imageUrl": { + "type": "string", + "format": "uri", + "description": "URL to cover artwork image." + }, + "createdAt": { + "type": "string", + "format": "datetime", + "description": "Timestamp when the track was uploaded." + }, + "supportGate": { + "type": "ref", + "ref": "#supportGate", + "description": "If set, this track requires viewer to be a supporter of the artist via atprotofans." + }, + "description": { + "type": "string", + "description": "Track description (liner notes, show notes, etc.).", + "maxLength": 5000 + }, + "audioBlob": { + "type": "blob", + "description": "Audio file stored on the user's PDS. When present, this is the canonical source; audioUrl is the CDN fallback.", + "accept": ["audio/*"], + "maxSize": 104857600 + } + } + } + }, + "supportGate": { + "type": "object", + "description": "Configuration for supporter-gated content.", + "required": ["type"], + "properties": { + "type": { + "type": "string", + "description": "The type of support required to access this content.", + "knownValues": ["any"] + } + } + }, + "featuredArtist": { + "type": "object", + "description": "A featured artist on a track. Identified by DID — the canonical, immutable identifier. handle and displayName are legacy denormalized snapshots; readers should resolve fresh metadata from the DID rather than trusting these fields.", + "required": ["did"], + "properties": { + "did": { + "type": "string", + "format": "did", + "description": "The DID of the featured artist. Canonical, stable identifier." + }, + "handle": { + "type": "string", + "format": "handle", + "description": "DEPRECATED snapshot — mutable, may be stale. Resolve from `did` instead. Older records include this; new records may omit it." + }, + "displayName": { + "type": "string", + "maxLength": 256, + "description": "DEPRECATED snapshot — mutable, may be stale. Resolve from `did` instead. Older records include this; new records may omit it." + } + } + } + } +} diff --git a/tests/sounds_test.ts b/tests/sounds_test.ts index c3149f8..be6ef19 100644 --- a/tests/sounds_test.ts +++ b/tests/sounds_test.ts @@ -1,6 +1,12 @@ import { assertEquals } from "@std/assert"; +import { parseCid } from "@atp/lex/data"; +import { AtUri } from "@atp/syntax"; +import { type RepoRecord, WriteOpAction } from "@atp/repo"; +import { BackgroundQueue } from "../data-plane/background.ts"; +import { IndexingService } from "../data-plane/indexing/index.ts"; import { createTestApp, TEST_USERS } from "./util.ts"; import { $OutputBody as SearchAudiosOutput } from "../lex/so/sprk/sound/searchAudios.ts"; +import { $OutputBody as GetAudiosOutput } from "../lex/so/sprk/sound/getAudios.ts"; const VALID_BLOB_CID = "bafyreihdwdcefgh4dqkjv67uzcmw7ojee6xedzdetojuzjevtenxquvyku"; @@ -189,6 +195,179 @@ Deno.test({ const body = await res.json() as SearchAudiosOutput; assertEquals(body.audios, []); }); + + await t.step("preserves Spark audio use counts on reindex", async () => { + const indexingService = new IndexingService( + ctx.db, + ctx.cfg, + ctx.idResolver, + new BackgroundQueue(ctx.db), + ctx.pushService, + ); + await indexingService.indexRecord( + new AtUri(chillUri), + parseCid(VALID_BLOB_CID), + ({ + ...chillRecord, + sound: { + ...chillRecord.sound, + ref: parseCid(VALID_BLOB_CID), + }, + title: "Chill Beats Remastered", + }) as unknown as RepoRecord, + WriteOpAction.Create, + nowIso, + ); + + const audio = await ctx.db.models.Audio.findOne({ uri: chillUri }) + .lean(); + assertEquals(audio?.title, "Chill Beats Remastered"); + assertEquals(audio?.useCount, 10); + }); + + await t.step("indexes and searches Plyr tracks as sounds", async () => { + const plyrAuthorDid = TEST_USERS[2].did; + const plyrUri = `at://${plyrAuthorDid}/fm.plyr.track/governor`; + const lockedBlobTrackUri = + `at://${plyrAuthorDid}/fm.plyr.track/locked-blob-governor`; + const imageUrl = "https://cdn.plyr.example/images/governor.jpg"; + + await ctx.db.models.Actor.create({ + did: plyrAuthorDid, + handle: "dame.is", + indexedAt: nowIso, + keys: [], + services: "[]", + }); + await ctx.db.models.Profile.create({ + uri: `at://${plyrAuthorDid}/app.bsky.actor.profile/self`, + cid: VALID_BLOB_CID, + authorDid: plyrAuthorDid, + createdAt: nowIso, + indexedAt: nowIso, + displayName: "Dame", + labels: [], + postsCount: 0, + followersCount: 0, + followsCount: 0, + }); + + const indexingService = new IndexingService( + ctx.db, + ctx.cfg, + ctx.idResolver, + new BackgroundQueue(ctx.db), + ctx.pushService, + ); + await indexingService.indexRecord( + new AtUri(plyrUri), + parseCid(VALID_BLOB_CID), + ({ + $type: "fm.plyr.track", + title: "Governor's Ball Symphony", + artist: "Dame", + audioBlob: { + $type: "blob", + ref: parseCid(VALID_BLOB_CID), + mimeType: "audio/mpeg", + size: 12345, + }, + fileType: "mp3", + imageUrl, + createdAt: nowIso, + }) as unknown as RepoRecord, + WriteOpAction.Create, + nowIso, + ); + + await ctx.db.models.Audio.updateOne( + { uri: plyrUri }, + { $set: { useCount: 7 } }, + ); + await indexingService.indexRecord( + new AtUri(plyrUri), + parseCid(VALID_BLOB_CID), + ({ + $type: "fm.plyr.track", + title: "Governor's Ball Symphony Redux", + artist: "Dame", + audioBlob: { + $type: "blob", + ref: parseCid(VALID_BLOB_CID), + mimeType: "audio/mpeg", + size: 12345, + }, + fileType: "mp3", + imageUrl, + createdAt: nowIso, + }) as unknown as RepoRecord, + WriteOpAction.Create, + nowIso, + ); + const reindexedPlyrAudio = await ctx.db.models.Audio.findOne({ + uri: plyrUri, + }).lean(); + assertEquals( + reindexedPlyrAudio?.title, + "Governor's Ball Symphony Redux", + ); + assertEquals(reindexedPlyrAudio?.useCount, 7); + + await indexingService.indexRecord( + new AtUri(lockedBlobTrackUri), + parseCid(VALID_BLOB_CID), + ({ + $type: "fm.plyr.track", + title: "Locked Blob Governor Symphony", + artist: "Dame", + audioUrl: "https://cdn.plyr.example/audio/locked-blob.mp3", + audioBlob: { + $type: "blob", + ref: parseCid(VALID_BLOB_CID), + mimeType: "audio/mpeg", + size: 12345, + }, + fileType: "mp3", + supportGate: { type: "any" }, + createdAt: nowIso, + }) as unknown as RepoRecord, + WriteOpAction.Create, + nowIso, + ); + + const res = await app.request( + "/xrpc/so.sprk.sound.searchAudios?q=governor", + ); + assertEquals(res.status, 200); + + const body = await res.json() as SearchAudiosOutput; + assertEquals(body.audios.length, 1); + assertEquals(body.audios[0].uri, plyrUri); + assertEquals(body.audios[0].title, "Governor's Ball Symphony Redux"); + assertEquals( + body.audios[0].audio, + `https://media.sprk.so/sound/${encodeURIComponent(plyrAuthorDid)}/${ + encodeURIComponent(VALID_BLOB_CID) + }`, + ); + assertEquals(body.audios[0].coverArt, imageUrl); + assertEquals(body.audios[0].details?.artist, "Dame"); + + const params = new URLSearchParams(); + params.append("uris", lockedBlobTrackUri); + const getAudiosRes = await app.request( + `/xrpc/so.sprk.sound.getAudios?${params.toString()}`, + ); + assertEquals(getAudiosRes.status, 200); + + const getAudiosBody = await getAudiosRes.json() as GetAudiosOutput; + assertEquals(getAudiosBody.audios.length, 1); + const audio = getAudiosBody.audios[0]; + const record = audio.record as Record; + assertEquals(audio.audio, undefined); + assertEquals(record.audioUrl, undefined); + assertEquals(record.audioBlob, undefined); + }); } finally { await cleanup(); } diff --git a/views/index.ts b/views/index.ts index 827e6fc..079cd5f 100644 --- a/views/index.ts +++ b/views/index.ts @@ -11,6 +11,7 @@ import type { import { AtUri, INVALID_HANDLE, normalizeDatetimeAlways } from "@atp/syntax"; import { mapDefined } from "@atp/common"; import * as so from "../lex/so.ts"; +import * as fm from "../lex/fm.ts"; import { cidFromBlobJson } from "./util.ts"; import { uriToDid } from "../utils/uris.ts"; import { FeedItem, Like, Post, Reply, Repost } from "../hydration/feed.ts"; @@ -1014,10 +1015,20 @@ export class Views { } const soundAgg = state.soundAggs?.get(uri); + const record = soundInfo.record; + const isPlyrTrack = fm.plyr.track.$matches(record); + const plyrRecord = isPlyrTrack ? record as fm.plyr.track.Main : undefined; + const sparkRecord = !isPlyrTrack + ? record as so.sprk.sound.audio.Main + : undefined; let coverArtUrl: UriString; - const coverArt = (soundInfo.record as { coverArt?: BlobRef }).coverArt; - if (coverArt) { + const coverArt = sparkRecord + ? (sparkRecord as { coverArt?: BlobRef }).coverArt + : undefined; + if (plyrRecord?.imageUrl) { + coverArtUrl = asUri(plyrRecord.imageUrl); + } else if (coverArt) { const coverArtCid = cidFromBlobJson(coverArt); coverArtUrl = asUri( `${this.mediaCdn}/img/medium/${authorDid}/${coverArtCid}/webp`, @@ -1026,27 +1037,39 @@ export class Views { coverArtUrl = author.avatar ?? asUri("https://media.sprk.so"); } - const details = soundInfo.record.details + const details = plyrRecord + ? { artist: plyrRecord.artist, title: plyrRecord.title } + : sparkRecord?.details ? { - artist: soundInfo.record.details.artist, - title: soundInfo.record.details.title, + artist: sparkRecord.details.artist, + title: sparkRecord.details.title, } : undefined; - const record = { - title: soundInfo.record.title, - origin: soundInfo.record.origin ?? undefined, - sound: soundInfo.record.sound ?? undefined, - labels: soundInfo.record.labels ?? undefined, - createdAt: soundInfo.record.createdAt, - } as Record; - - const audioCid = cidFromBlobJson(soundInfo.record.sound); - const audioUrl = asUri( - `https://media.sprk.so/sound/${encodeURIComponent(authorDid)}/${ - encodeURIComponent(audioCid) - }`, - ); + const isGatedPlyrTrack = !!plyrRecord?.supportGate; + const viewRecord = plyrRecord + ? isGatedPlyrTrack + ? (({ audioBlob: _audioBlob, audioUrl: _audioUrl, ...record }) => + record)(plyrRecord) + : plyrRecord as Record + : { + title: sparkRecord?.title, + origin: sparkRecord?.origin ?? undefined, + sound: sparkRecord?.sound ?? undefined, + labels: sparkRecord?.labels ?? undefined, + createdAt: sparkRecord?.createdAt, + } as Record; + + let audioUrl: UriString | undefined; + const audioBlob = plyrRecord ? plyrRecord.audioBlob : sparkRecord?.sound; + if (audioBlob && !isGatedPlyrTrack) { + const audioCid = cidFromBlobJson(audioBlob); + audioUrl = asUri( + `https://media.sprk.so/sound/${encodeURIComponent(authorDid)}/${ + encodeURIComponent(audioCid) + }`, + ); + } const indexedAt = asDatetime( this.indexedAt(soundInfo)?.toISOString() ?? new Date().toISOString(), @@ -1056,9 +1079,9 @@ export class Views { uri: asAtUri(uri), cid: asCid(soundInfo.cid), author, - title: soundInfo.record.title, + title: record.title, coverArt: coverArtUrl, - record, + record: viewRecord, useCount: soundAgg?.uses ?? 0, details, indexedAt, -- 2.51.2