diff --git a/data-plane/indexing/index.ts b/data-plane/indexing/index.ts index c26357e..ed58688 100644 --- a/data-plane/indexing/index.ts +++ b/data-plane/indexing/index.ts @@ -26,7 +26,7 @@ import * as Repost from "./plugins/repost.ts"; import * as Story from "./plugins/story.ts"; import * as Audio from "./plugins/audio.ts"; import { RecordProcessor } from "./processor.ts"; -import { Logger } from "@logtape/logtape"; +import { getLogger, Logger } from "@logtape/logtape"; import { ServerConfig } from "../../config.ts"; export class IndexingService { @@ -42,14 +42,15 @@ export class IndexingService { story: Story.PluginType; audio: Audio.PluginType; }; + logger: Logger; constructor( public db: Database, public cfg: ServerConfig, public idResolver: IdResolver, public background: BackgroundQueue, - public logger: Logger, ) { + this.logger = getLogger(["appview", "indexer"]); this.records = { post: Post.makePlugin(this.db, this.background), reply: Reply.makePlugin(this.db, this.background), @@ -70,7 +71,6 @@ export class IndexingService { this.cfg, this.idResolver, this.background, - this.logger, ); } @@ -157,7 +157,7 @@ export class IndexingService { const actorExists = await this.db.models.Actor.findOne({ did }).lean(); if (!actorExists) { - console.log( + this.logger.info( `indexRepo: No actor record found for ${did}, indexing handle first`, ); await this.indexHandle(did, now); @@ -179,7 +179,7 @@ export class IndexingService { const repoRecords = formatCheckout(did, verifiedRepo); const diff = findDiffFromCheckout(currRecords, repoRecords); - console.log(`Indexing ${diff.length} records for ${did}:`); + this.logger.info(`Indexing ${diff.length} records for ${did}:`); await Promise.all( diff.map(async (op) => { diff --git a/data-plane/indexing/plugins/audio.ts b/data-plane/indexing/plugins/audio.ts index 04afa2d..2e5fd0e 100644 --- a/data-plane/indexing/plugins/audio.ts +++ b/data-plane/indexing/plugins/audio.ts @@ -31,21 +31,12 @@ const insertFn = async ( }; // Use findOneAndUpdate with upsert to handle potential duplicate key errors - try { - const insertedAudio = await db.models.Audio.findOneAndUpdate( - { uri: audio.uri }, - audio, - { upsert: true, new: true }, - ); - return insertedAudio; - } catch (err) { - // Handle duplicate key errors gracefully - const mongoError = err as { code?: number }; - if (mongoError.code === 11000) { - return null; // Silently skip duplicates - } - throw err; - } + const insertedAudio = await db.models.Audio.findOneAndUpdate( + { uri: audio.uri }, + { $set: audio }, + { upsert: true, new: true }, + ); + return insertedAudio; }; const findDuplicate = (): AtUri | null => { diff --git a/data-plane/indexing/plugins/block.ts b/data-plane/indexing/plugins/block.ts index d6538b1..a5a9a98 100644 --- a/data-plane/indexing/plugins/block.ts +++ b/data-plane/indexing/plugins/block.ts @@ -26,35 +26,13 @@ const insertFn = async ( indexedAt: timestamp, }; - // Check if block already exists by URI - const existingBlockByUri = await db.models.Block.findOne({ uri: block.uri }) - .lean(); - if (existingBlockByUri) { - return null; // Block already indexed - } - - // Check if a block with same authorDid+subject already exists - const existingBlockByComposite = await db.models.Block.findOne({ - authorDid: block.authorDid, - subject: block.subject, - }).lean(); - if (existingBlockByComposite) { - return null; // Block with same author+subject already exists - } - - // Insert the new block - try { - const insertedBlock = new db.models.Block(block); - await insertedBlock.save(); - return insertedBlock; - } catch (err) { - // Handle duplicate key errors gracefully - const mongoError = err as { code?: number }; - if (mongoError.code === 11000) { - return null; // Silently skip duplicates - } - throw err; - } + // Use findOneAndUpdate with upsert to handle duplicates gracefully + const insertedBlock = await db.models.Block.findOneAndUpdate( + { uri: block.uri }, + { $set: block }, + { upsert: true, new: true }, + ); + return insertedBlock; }; const findDuplicate = async ( diff --git a/data-plane/indexing/plugins/follow.ts b/data-plane/indexing/plugins/follow.ts index c576ada..8a15f75 100644 --- a/data-plane/indexing/plugins/follow.ts +++ b/data-plane/indexing/plugins/follow.ts @@ -26,24 +26,12 @@ const insertFn = async ( indexedAt: timestamp, }; - // Use findOneAndUpdate with upsert on the compound key to handle potential duplicate key errors - try { - const insertedFollow = await db.models.Follow.findOneAndUpdate( - { - authorDid: follow.authorDid, - subject: follow.subject, - }, - follow, - { upsert: true, new: true }, - ); - return insertedFollow; - } catch (err) { - const mongoError = err as { code?: number }; - if (mongoError.code === 11000) { - return null; // Silently skip duplicates - } - throw err; - } + const insertedFollow = await db.models.Follow.findOneAndUpdate( + { uri: follow.uri }, + { $set: follow }, + { upsert: true, new: true }, + ); + return insertedFollow; }; const findDuplicate = async ( @@ -91,47 +79,42 @@ const notifsForDelete = ( }; const updateAggregates = async (db: Database, follow: IndexedFollow) => { - try { - // Update followers count for the subject (count both types) - const followersCount = await db.models.Follow.countDocuments({ - subject: follow.subject, - }); - - // First check if profile exists to avoid creating one with null URI - const existingSubjectProfile = await db.models.Profile.findOne({ - authorDid: follow.subject, - }); - - if (existingSubjectProfile) { - // Only update existing profiles - await db.models.Profile.findOneAndUpdate( - { authorDid: follow.subject }, - { followersCount }, - { new: true }, - ); - } - - // Update follows count for the author (count both types) - const followsCount = await db.models.Follow.countDocuments({ - authorDid: follow.authorDid, - }); - - // First check if profile exists to avoid creating one with null URI - const existingAuthorProfile = await db.models.Profile.findOne({ - authorDid: follow.authorDid, - }); - - if (existingAuthorProfile) { - // Only update existing profiles - await db.models.Profile.findOneAndUpdate( - { authorDid: follow.authorDid }, - { followsCount }, - { new: true }, - ); - } - } catch (error) { - console.error("Error updating follow aggregates:", error); - // Don't throw - allow processing to continue even if aggregates update fails + // Update followers count for the subject (count both types) + const followersCount = await db.models.Follow.countDocuments({ + subject: follow.subject, + }); + + // First check if profile exists to avoid creating one with null URI + const existingSubjectProfile = await db.models.Profile.findOne({ + authorDid: follow.subject, + }); + + if (existingSubjectProfile) { + // Only update existing profiles + await db.models.Profile.findOneAndUpdate( + { authorDid: follow.subject }, + { $set: { followersCount } }, + { new: true }, + ); + } + + // Update follows count for the author (count both types) + const followsCount = await db.models.Follow.countDocuments({ + authorDid: follow.authorDid, + }); + + // First check if profile exists to avoid creating one with null URI + const existingAuthorProfile = await db.models.Profile.findOne({ + authorDid: follow.authorDid, + }); + + if (existingAuthorProfile) { + // Only update existing profiles + await db.models.Profile.findOneAndUpdate( + { authorDid: follow.authorDid }, + { $set: { followsCount } }, + { new: true }, + ); } }; diff --git a/data-plane/indexing/plugins/generator.ts b/data-plane/indexing/plugins/generator.ts index 3b7f0ae..11d95d4 100644 --- a/data-plane/indexing/plugins/generator.ts +++ b/data-plane/indexing/plugins/generator.ts @@ -40,22 +40,12 @@ const insertFn = async ( indexedAt: timestamp, }; - // Use findOneAndUpdate with upsert to handle potential duplicate key errors - try { - const insertedGenerator = await db.models.Generator.findOneAndUpdate( - { uri: generator.uri }, - generator, - { upsert: true, new: true }, - ); - return insertedGenerator; - } catch (err) { - // Handle duplicate key errors gracefully - const mongoError = err as { code?: number }; - if (mongoError.code === 11000) { - return null; // Silently skip duplicates - } - throw err; - } + const insertedGenerator = await db.models.Generator.findOneAndUpdate( + { uri: generator.uri }, + { $set: generator }, + { upsert: true, new: true }, + ); + return insertedGenerator; }; const findDuplicate = (): AtUri | null => { diff --git a/data-plane/indexing/plugins/like.ts b/data-plane/indexing/plugins/like.ts index ab8c7c5..b8e7e81 100644 --- a/data-plane/indexing/plugins/like.ts +++ b/data-plane/indexing/plugins/like.ts @@ -35,24 +35,12 @@ const insertFn = async ( }; // Use findOneAndUpdate with upsert on the compound key to handle potential duplicate key errors - try { - const insertedLike = await db.models.Like.findOneAndUpdate( - { - authorDid: like.authorDid, - subject: like.subject, - }, - like, - { upsert: true, new: true }, - ); - return insertedLike; - } catch (err) { - // Handle duplicate key errors gracefully - const mongoError = err as { code?: number }; - if (mongoError.code === 11000) { - return null; // Silently skip duplicates - } - throw err; - } + const insertedLike = await db.models.Like.findOneAndUpdate( + { uri: like.uri }, + { $set: like }, + { upsert: true, new: true }, + ); + return insertedLike; }; const findDuplicate = async ( @@ -138,52 +126,48 @@ const notifsForDelete = ( }; const updateAggregates = async (db: Database, like: IndexedLike) => { - try { - const likeCount = await db.models.Like.countDocuments({ - subject: like.subject, + const likeCount = await db.models.Like.countDocuments({ + subject: like.subject, + }); + + const subjectUri = new AtUri(like.subject); + + if (subjectUri.collection === "so.sprk.feed.generator") { + const existingGenerator = await db.models.Generator.findOne({ + uri: like.subject, }); - const subjectUri = new AtUri(like.subject); - - if (subjectUri.collection === "so.sprk.feed.generator") { - const existingGenerator = await db.models.Generator.findOne({ - uri: like.subject, - }); - - if (existingGenerator) { - await db.models.Generator.findOneAndUpdate( - { uri: like.subject }, - { $set: { likeCount } }, - { new: true }, - ); - } - } else { - const existingPost = await db.models.Post.findOne({ - uri: like.subject, - }); - - if (existingPost) { - await db.models.Post.findOneAndUpdate( - { uri: like.subject }, - { $set: { likeCount } }, - { new: true }, - ); - } - - const existingReply = await db.models.Reply.findOne({ - uri: like.subject, - }); - - if (existingReply) { - await db.models.Reply.findOneAndUpdate( - { uri: like.subject }, - { $set: { likeCount } }, - { new: true }, - ); - } + if (existingGenerator) { + await db.models.Generator.findOneAndUpdate( + { uri: like.subject }, + { $set: { likeCount } }, + { new: true }, + ); + } + } else { + const existingPost = await db.models.Post.findOne({ + uri: like.subject, + }); + + if (existingPost) { + await db.models.Post.findOneAndUpdate( + { uri: like.subject }, + { $set: { likeCount } }, + { new: true }, + ); + } + + const existingReply = await db.models.Reply.findOne({ + uri: like.subject, + }); + + if (existingReply) { + await db.models.Reply.findOneAndUpdate( + { uri: like.subject }, + { $set: { likeCount } }, + { new: true }, + ); } - } catch (error) { - console.error("Error updating like aggregates:", error); } }; diff --git a/data-plane/indexing/plugins/post.ts b/data-plane/indexing/plugins/post.ts index 13a2017..4550210 100644 --- a/data-plane/indexing/plugins/post.ts +++ b/data-plane/indexing/plugins/post.ts @@ -59,18 +59,6 @@ const insertFn = async ( obj: PostRecord, timestamp: string, ): Promise => { - console.log("DEBUG: Post indexing started"); - // Ensure actor record exists before creating post - const actorExists = await db.models.Actor.findOne({ did: uri.host }).lean(); - if (!actorExists) { - // This should trigger actor indexing, but for now we'll just log - console.log( - `Post indexing: No actor record found for ${uri.host}, post may have missing handle`, - ); - } - - console.log("DEBUG: Post media:", JSON.stringify(obj.media, null, 2)); - const post = { uri: uri.toString(), cid: cid.toString(), @@ -86,83 +74,73 @@ const insertFn = async ( indexedAt: timestamp, }; - // Use findOneAndUpdate with upsert to handle potential duplicate key errors - try { - const insertedPost = await db.models.Post.findOneAndUpdate( - { uri: post.uri }, - post, - { upsert: true, new: true }, - ); - - const facets = (obj.caption?.facets || []) - .flatMap((facet) => facet.features) - .flatMap((feature) => { - if (isMention(feature)) { - return { - type: "mention" as const, - value: feature.did, - }; - } - if (isLink(feature)) { - return { - type: "link" as const, - value: feature.uri, - }; - } - return []; - }); + const insertedPost = await db.models.Post.findOneAndUpdate( + { uri: post.uri }, + { $set: post }, + { upsert: true, new: true }, + ); - // Media processing - medias are stored inline in the Post model - const medias: Array<{ - position?: number; - imageCid?: string; - alt?: string | null; - thumbCid?: string | null; - videoCid?: string; - }> = []; - const postMedias = separateMedia(obj.media); - for (const postMedia of postMedias) { - if (isMediaImages(postMedia)) { - const { images } = postMedia as MediaImages; - const imagesMedia = images.map(( - img: MediaImage, - i: number, - ) => ({ - position: i, - imageCid: img.image.ref.toString(), - alt: img.alt, - })); - medias.push(...imagesMedia); - } else if (isMediaVideo(postMedia)) { - const media = postMedia as MediaVideo; - const videoMedia = { - postUri: uri.toString(), - videoCid: media.video.ref.toString(), - alt: media.alt ?? null, + const facets = (obj.caption?.facets || []) + .flatMap((facet) => facet.features) + .flatMap((feature) => { + if (isMention(feature)) { + return { + type: "mention" as const, + value: feature.did, }; - medias.push(videoMedia); } - } - - const descendents = await getDescendents(db, { - uri: post.uri, - depth: REPLY_NOTIF_DEPTH, + if (isLink(feature)) { + return { + type: "link" as const, + value: feature.uri, + }; + } + return []; }); - return { - post: insertedPost, - facets, - medias, - descendents, - }; - } catch (err) { - // Handle duplicate key errors gracefully - const mongoError = err as { code?: number }; - if (mongoError.code === 11000) { - return null; // Silently skip duplicates + // Media processing - medias are stored inline in the Post model + const medias: Array<{ + position?: number; + imageCid?: string; + alt?: string | null; + thumbCid?: string | null; + videoCid?: string; + }> = []; + const postMedias = separateMedia(obj.media); + for (const postMedia of postMedias) { + if (isMediaImages(postMedia)) { + const { images } = postMedia as MediaImages; + const imagesMedia = images.map(( + img: MediaImage, + i: number, + ) => ({ + position: i, + imageCid: img.image.ref.toString(), + alt: img.alt, + })); + medias.push(...imagesMedia); + } else if (isMediaVideo(postMedia)) { + const media = postMedia as MediaVideo; + const videoMedia = { + postUri: uri.toString(), + videoCid: media.video.ref.toString(), + alt: media.alt ?? null, + }; + medias.push(videoMedia); } - throw err; } + + const descendents = await getDescendents(db, { + uri: post.uri, + depth: REPLY_NOTIF_DEPTH, + }); + + return { + post: insertedPost, + facets, + medias, + descendents, + }; }; const findDuplicate = (): AtUri | null => { @@ -280,28 +258,23 @@ const notifsForDelete = ( }; const updateAggregates = async (db: Database, postIdx: IndexedPost) => { - try { - // Update posts count for author - const postsCount = await db.models.Post.countDocuments({ - authorDid: postIdx.post.authorDid, - }); + // Update posts count for author + const postsCount = await db.models.Post.countDocuments({ + authorDid: postIdx.post.authorDid, + }); - // First check if profile exists to avoid creating one with null URI - const existingProfile = await db.models.Profile.findOne({ - authorDid: postIdx.post.authorDid, - }); + // First check if profile exists to avoid creating one with null URI + const existingProfile = await db.models.Profile.findOne({ + authorDid: postIdx.post.authorDid, + }); - if (existingProfile) { - // Only update existing profiles - await db.models.Profile.findOneAndUpdate( - { authorDid: postIdx.post.authorDid }, - { postsCount }, - { new: true }, - ); - } - } catch (error) { - console.error("Error updating post aggregates:", error); - // Don't throw - allow processing to continue even if aggregates update fails + if (existingProfile) { + // Only update existing profiles + await db.models.Profile.findOneAndUpdate( + { authorDid: postIdx.post.authorDid }, + { $set: { postsCount } }, + { new: true }, + ); } }; diff --git a/data-plane/indexing/plugins/profile.ts b/data-plane/indexing/plugins/profile.ts index fbefabb..3a6d21d 100644 --- a/data-plane/indexing/plugins/profile.ts +++ b/data-plane/indexing/plugins/profile.ts @@ -32,22 +32,12 @@ const insertFn = async ( indexedAt: timestamp, }; - // Use findOneAndUpdate with upsert to handle potential duplicate key errors - try { - const insertedProfile = await db.models.Profile.findOneAndUpdate( - { uri: profile.uri }, - profile, - { upsert: true, new: true }, - ); - return insertedProfile; - } catch (err) { - // Handle duplicate key errors gracefully - const mongoError = err as { code?: number }; - if (mongoError.code === 11000) { - return null; // Silently skip duplicates - } - throw err; - } + const insertedProfile = await db.models.Profile.findOneAndUpdate( + { uri: profile.uri }, + { $set: profile }, + { upsert: true, new: true }, + ); + return insertedProfile; }; const findDuplicate = (): AtUri | null => { diff --git a/data-plane/indexing/plugins/reply.ts b/data-plane/indexing/plugins/reply.ts index bf12591..be143d5 100644 --- a/data-plane/indexing/plugins/reply.ts +++ b/data-plane/indexing/plugins/reply.ts @@ -56,16 +56,6 @@ const insertFn = async ( obj: ReplyRecord, timestamp: string, ): Promise => { - console.log("DEBUG: Post indexing started"); - // Ensure actor record exists before creating post - const actorExists = await db.models.Actor.findOne({ did: uri.host }).lean(); - if (!actorExists) { - // This should trigger actor indexing, but for now we'll just log - console.log( - `Post indexing: No actor record found for ${uri.host}, post may have missing handle`, - ); - } - const reply = { uri: uri.toString(), cid: cid.toString(), @@ -93,84 +83,75 @@ const insertFn = async ( }; // Use findOneAndUpdate with upsert to handle potential duplicate key errors - try { - const insertedReply = await db.models.Reply.findOneAndUpdate( - { uri: reply.uri }, - reply, - { upsert: true, new: true }, - ); + const insertedReply = await db.models.Reply.findOneAndUpdate( + { uri: reply.uri }, + { $set: reply }, + { upsert: true, new: true }, + ); - if (obj.reply) { - const { invalidReplyRoot } = await validateReply( - db, - obj.reply, + if (obj.reply) { + const { invalidReplyRoot } = await validateReply( + db, + obj.reply, + ); + if (invalidReplyRoot) { + Object.assign(insertedReply, { invalidReplyRoot }); + await db.models.Reply.updateOne( + { uri: reply.uri }, + { $set: { invalidReplyRoot } }, ); - if (invalidReplyRoot) { - Object.assign(insertedReply, { invalidReplyRoot }); - await db.models.Reply.updateOne( - { uri: reply.uri }, - { invalidReplyRoot }, - ); - } - } - - const facets = (obj.facets || []) - .flatMap((facet) => facet.features) - .flatMap((feature) => { - if (isMention(feature)) { - return { - type: "mention" as const, - value: feature.did, - }; - } - if (isLink(feature)) { - return { - type: "link" as const, - value: feature.uri, - }; - } - return []; - }); - - // Embed processing - embeds are stored inline in the Post model - let media: { - postUri?: string; - cid?: string; - alt?: string; - } = {}; - if (isMediaImage(obj.media)) { - const imageMedia = { - postUri: uri.toString(), - cid: obj.media.image.ref.toString(), - alt: obj.media.alt as string, - }; - media = imageMedia; } + } - const ancestors = await getAncestorsAndSelf(db, { - uri: reply.uri, - parentHeight: REPLY_NOTIF_DEPTH, - }); - const descendents = await getDescendents(db, { - uri: reply.uri, - depth: REPLY_NOTIF_DEPTH, + const facets = (obj.facets || []) + .flatMap((facet) => facet.features) + .flatMap((feature) => { + if (isMention(feature)) { + return { + type: "mention" as const, + value: feature.did, + }; + } + if (isLink(feature)) { + return { + type: "link" as const, + value: feature.uri, + }; + } + return []; }); - return { - reply: insertedReply, - facets, - media, - ancestors, - descendents, + // Embed processing - embeds are stored inline in the Post model + let media: { + postUri?: string; + cid?: string; + alt?: string; + } = {}; + if (isMediaImage(obj.media)) { + const imageMedia = { + postUri: uri.toString(), + cid: obj.media.image.ref.toString(), + alt: obj.media.alt as string, }; - } catch (err) { - // Handle duplicate key errors gracefully - const mongoError = err as { code?: number }; - if (mongoError.code === 11000) { - return null; // Silently skip duplicates - } - throw err; + media = imageMedia; } + + const ancestors = await getAncestorsAndSelf(db, { + uri: reply.uri, + parentHeight: REPLY_NOTIF_DEPTH, + }); + const descendents = await getDescendents(db, { + uri: reply.uri, + depth: REPLY_NOTIF_DEPTH, + }); + + return { + reply: insertedReply, + facets, + media, + ancestors, + descendents, + }; }; const findDuplicate = (): AtUri | null => { diff --git a/data-plane/indexing/plugins/repost.ts b/data-plane/indexing/plugins/repost.ts index 71ec212..92ff34f 100644 --- a/data-plane/indexing/plugins/repost.ts +++ b/data-plane/indexing/plugins/repost.ts @@ -17,50 +17,31 @@ const insertFn = async ( obj: Repost.Record, timestamp: string, ): Promise => { - try { - // Handle via property safely with type assertion - const viaObj = obj.via as { uri: string; cid: string } | undefined; - const via = viaObj?.uri || null; - const viaCid = viaObj?.cid || null; - - const repost = { - uri: uri.toString(), - cid: cid.toString(), - authorDid: uri.host, - subject: { - uri: obj.subject.uri, - cid: obj.subject.cid, - }, - via, - viaCid, - createdAt: normalizeDatetimeAlways(obj.createdAt), - indexedAt: timestamp, - }; - - // Use findOneAndUpdate with compound key to handle potential duplicate key errors - try { - const insertedRepost = await db.models.Repost.findOneAndUpdate( - { - authorDid: repost.authorDid, - "subject.uri": repost.subject.uri, - }, - repost, - { upsert: true, new: true }, - ); - return insertedRepost; - } catch (err) { - // Handle duplicate key errors gracefully - const mongoError = err as { code?: number }; - if (mongoError.code === 11000) { - return null; // Silently skip duplicates - } - throw err; - } - } catch (error) { - // Log the error but prevent it from crashing the process - console.error("Error processing repost:", error); - return null; - } + const viaObj = obj.via as { uri: string; cid: string } | undefined; + const via = viaObj?.uri || null; + const viaCid = viaObj?.cid || null; + + const repost = { + uri: uri.toString(), + cid: cid.toString(), + authorDid: uri.host, + subject: { + uri: obj.subject.uri, + cid: obj.subject.cid, + }, + via, + viaCid, + createdAt: normalizeDatetimeAlways(obj.createdAt), + indexedAt: timestamp, + }; + + // Use findOneAndUpdate with compound key to handle potential duplicate key errors + const insertedRepost = await db.models.Repost.findOneAndUpdate( + { uri: repost.uri }, + { $set: repost }, + { upsert: true, new: true }, + ); + return insertedRepost; }; const findDuplicate = async ( @@ -76,130 +57,106 @@ const findDuplicate = async ( }; const notifsForInsert = (obj: IndexedRepost) => { - try { - const subjectUri = new AtUri(obj.subject.uri); - // prevent self-notifications - const isRepostFromSubjectUser = subjectUri.host === obj.authorDid; - if (isRepostFromSubjectUser) { - return []; - } + const subjectUri = new AtUri(obj.subject.uri); + // prevent self-notifications + const isRepostFromSubjectUser = subjectUri.host === obj.authorDid; + if (isRepostFromSubjectUser) { + return []; + } - const notifs: Array<{ - did: string; - reason: string; - author: string; - recordUri: string; - recordCid: string; - sortAt: string; - reasonSubject?: string; - }> = [ - // Notification to the author of the reposted record. - { - did: subjectUri.host, - author: obj.authorDid, - recordUri: obj.uri, - recordCid: obj.cid, - reason: "repost" as const, - reasonSubject: subjectUri.toString(), - sortAt: obj.createdAt, - }, - ]; - - if (obj.via) { - try { - const viaUri = new AtUri(obj.via); - const isRepostFromViaSubjectUser = viaUri.host === obj.authorDid; - // prevent self-notifications - if (!isRepostFromViaSubjectUser) { - notifs.push( - // Notification to the reposter via whose repost the repost was made. - { - did: viaUri.host, - author: obj.authorDid, - recordUri: obj.uri, - recordCid: obj.cid, - reason: "repost-via-repost" as const, - reasonSubject: viaUri.toString(), - sortAt: obj.createdAt, - }, - ); - } - } catch (viaError) { - console.error("Error processing via uri in notification:", viaError); - // Continue with just the main notification - } + const notifs: Array<{ + did: string; + reason: string; + author: string; + recordUri: string; + recordCid: string; + sortAt: string; + reasonSubject?: string; + }> = [ + // Notification to the author of the reposted record. + { + did: subjectUri.host, + author: obj.authorDid, + recordUri: obj.uri, + recordCid: obj.cid, + reason: "repost" as const, + reasonSubject: subjectUri.toString(), + sortAt: obj.createdAt, + }, + ]; + + if (obj.via) { + const viaUri = new AtUri(obj.via); + const isRepostFromViaSubjectUser = viaUri.host === obj.authorDid; + // prevent self-notifications + if (!isRepostFromViaSubjectUser) { + notifs.push( + // Notification to the reposter via whose repost the repost was made. + { + did: viaUri.host, + author: obj.authorDid, + recordUri: obj.uri, + recordCid: obj.cid, + reason: "repost-via-repost" as const, + reasonSubject: viaUri.toString(), + sortAt: obj.createdAt, + }, + ); } - - return notifs; - } catch (error) { - console.error("Error generating notifications for insert:", error); - return []; } + + return notifs; }; const deleteFn = async ( db: Database, uri: AtUri, ): Promise => { - try { - const deleted = await db.models.Repost.findOneAndDelete({ - uri: uri.toString(), - }); - return deleted; - } catch (error) { - console.error("Error deleting repost:", error); - return null; - } + const deleted = await db.models.Repost.findOneAndDelete({ + uri: uri.toString(), + }); + return deleted; }; const notifsForDelete = ( deleted: IndexedRepost, replacedBy: IndexedRepost | null, ) => { - try { - const toDelete = replacedBy ? [] : [deleted.uri]; - return { notifs: [], toDelete }; - } catch (error) { - console.error("Error processing notifications for delete:", error); - return { notifs: [], toDelete: [] }; - } + const toDelete = replacedBy ? [] : [deleted.uri]; + return { notifs: [], toDelete }; }; const updateAggregates = async (db: Database, repost: IndexedRepost) => { - try { - const repostCount = await db.models.Repost.countDocuments({ - "subject.uri": repost.subject.uri, - }); - - const existingPost = await db.models.Post.findOne({ - uri: repost.subject.uri, - }); - - if (existingPost) { - await db.models.Post.findOneAndUpdate( - { uri: repost.subject.uri }, - { $set: { repostCount } }, - { new: true }, - ); - } + const repostCount = await db.models.Repost.countDocuments({ + "subject.uri": repost.subject.uri, + }); - const authorRepostCount = await db.models.Repost.countDocuments({ - authorDid: repost.authorDid, - }); + const existingPost = await db.models.Post.findOne({ + uri: repost.subject.uri, + }); - const existingProfile = await db.models.Profile.findOne({ - authorDid: repost.authorDid, - }); + if (existingPost) { + await db.models.Post.findOneAndUpdate( + { uri: repost.subject.uri }, + { $set: { repostCount } }, + { new: true }, + ); + } - if (existingProfile) { - await db.models.Profile.findOneAndUpdate( - { authorDid: repost.authorDid }, - { $set: { repostCount: authorRepostCount } }, - { new: true }, - ); - } - } catch (error) { - console.error("Error updating repost aggregates:", error); + const authorRepostCount = await db.models.Repost.countDocuments({ + authorDid: repost.authorDid, + }); + + const existingProfile = await db.models.Profile.findOne({ + authorDid: repost.authorDid, + }); + + if (existingProfile) { + await db.models.Profile.findOneAndUpdate( + { authorDid: repost.authorDid }, + { $set: { repostCount: authorRepostCount } }, + { new: true }, + ); } }; diff --git a/data-plane/indexing/plugins/story.ts b/data-plane/indexing/plugins/story.ts index 31f23b8..09f30ba 100644 --- a/data-plane/indexing/plugins/story.ts +++ b/data-plane/indexing/plugins/story.ts @@ -6,10 +6,6 @@ import { BackgroundQueue } from "../../background.ts"; import { Database } from "../../db/index.ts"; import { StoryDocument } from "../../db/models.ts"; import { RecordProcessor } from "../processor.ts"; -import { - normalizeEmbed, - normalizeObject, -} from "../../../utils/embed-normalizer.ts"; const lexId = lex.ids.SoSprkStoryPost; type IndexedStory = StoryDocument; @@ -25,8 +21,8 @@ const insertFn = async ( uri: uri.toString(), cid: cid.toString(), authorDid: uri.host, - media: normalizeEmbed(obj.media) || null, - sound: normalizeObject(obj.sound) || null, + media: obj.media, + sound: obj.sound, labels: obj.labels || null, tags: obj.tags || [], createdAt: normalizeDatetimeAlways(obj.createdAt), @@ -34,21 +30,12 @@ const insertFn = async ( }; // Use findOneAndUpdate with upsert to handle potential duplicate key errors - try { - const insertedStory = await db.models.Story.findOneAndUpdate( - { uri: story.uri }, - story, - { upsert: true, new: true }, - ); - return insertedStory; - } catch (err) { - // Handle duplicate key errors gracefully - const mongoError = err as { code?: number }; - if (mongoError.code === 11000) { - return null; // Silently skip duplicates - } - throw err; - } + const insertedStory = await db.models.Story.findOneAndUpdate( + { uri: story.uri }, + story, + { upsert: true, new: true }, + ); + return insertedStory; }; const findDuplicate = (): AtUri | null => { diff --git a/data-plane/indexing/processor.ts b/data-plane/indexing/processor.ts index 2d14b65..5910ff7 100644 --- a/data-plane/indexing/processor.ts +++ b/data-plane/indexing/processor.ts @@ -85,19 +85,8 @@ export class RecordProcessor { } } - assertValidRecord(obj: unknown, uri: AtUri): asserts obj is T { - if (!this.matchesCollection(uri)) { - throw new Error( - `Record collection mismatch: expected ${this.collection}, got ${uri.collection}`, - ); - } - try { - lexicons.assertValidRecord(this.collection, obj); - } catch (err) { - throw new Error( - `Record validation failed for collection: ${this.collection}. Error: ${err}`, - ); - } + assertValidRecord(obj: unknown): asserts obj is T { + lexicons.assertValidRecord(this.collection, obj); } // Helper method to get the lexId this processor handles @@ -112,7 +101,7 @@ export class RecordProcessor { timestamp: string, opts?: { disableNotifs?: boolean }, ) { - this.assertValidRecord(obj, uri); + this.assertValidRecord(obj); // Insert or update record await this.db.models.Record.findOneAndUpdate( @@ -170,7 +159,7 @@ export class RecordProcessor { timestamp: string, opts?: { disableNotifs?: boolean }, ) { - this.assertValidRecord(obj, uri); + this.assertValidRecord(obj); // Update record await this.db.models.Record.findOneAndUpdate( diff --git a/data-plane/routes/identity.ts b/data-plane/routes/identity.ts index 7373b5d..c838e9c 100644 --- a/data-plane/routes/identity.ts +++ b/data-plane/routes/identity.ts @@ -45,18 +45,13 @@ export class Identity { throw new DataPlaneError(Code.InternalError); } - try { - const doc = await this.idResolver.did.resolve(did); - if (!doc) { - throw new DataPlaneError(Code.NotFound); - } - - const result = getResultFromDoc(doc); - return result; - } catch (error) { - console.error("Error resolving DID:", error); - throw new DataPlaneError(Code.InternalError); + const doc = await this.idResolver.did.resolve(did); + if (!doc) { + throw new DataPlaneError(Code.NotFound); } + + const result = getResultFromDoc(doc); + return result; } async getByHandle(handle: string) { @@ -64,23 +59,18 @@ export class Identity { throw new DataPlaneError(Code.InternalError); } - try { - const did = await this.idResolver.handle.resolve(handle); - if (!did) { - throw new DataPlaneError(Code.NotFound); - } - - const doc = await this.idResolver.did.resolve(did); - if (!doc || did !== getDid(doc)) { - throw new DataPlaneError(Code.NotFound); - } + const did = await this.idResolver.handle.resolve(handle); + if (!did) { + throw new DataPlaneError(Code.NotFound); + } - const result = getResultFromDoc(doc); - return result; - } catch (error) { - console.error("Error resolving handle:", error); - throw new DataPlaneError(Code.InternalError); + const doc = await this.idResolver.did.resolve(did); + if (!doc || did !== getDid(doc)) { + throw new DataPlaneError(Code.NotFound); } + + const result = getResultFromDoc(doc); + return result; } async resolve(identifier: string, type?: "did" | "handle") { @@ -88,40 +78,35 @@ export class Identity { throw new DataPlaneError(Code.InternalError); } - try { - let doc: DidDocument | null = null; - let resolvedDid: string | null = null; - - // Auto-detect type if not specified - const identifierType = type || - (identifier.startsWith("did:") ? "did" : "handle"); - - if (identifierType === "did") { - doc = await this.idResolver.did.resolve(identifier); - resolvedDid = identifier; - } else { - resolvedDid = await this.idResolver.handle.resolve(identifier) || null; - if (resolvedDid) { - doc = await this.idResolver.did.resolve(resolvedDid); - } - } + let doc: DidDocument | null = null; + let resolvedDid: string | null = null; + + // Auto-detect type if not specified + const identifierType = type || + (identifier.startsWith("did:") ? "did" : "handle"); - if (!doc || (resolvedDid && resolvedDid !== getDid(doc))) { - throw new DataPlaneError(Code.NotFound); + if (identifierType === "did") { + doc = await this.idResolver.did.resolve(identifier); + resolvedDid = identifier; + } else { + resolvedDid = await this.idResolver.handle.resolve(identifier) || null; + if (resolvedDid) { + doc = await this.idResolver.did.resolve(resolvedDid); } + } - const result = getResultFromDoc(doc); - return { - ...result, - resolvedFrom: { - identifier, - type: identifierType, - }, - }; - } catch (error) { - console.error("Error resolving identity:", error); - throw new DataPlaneError(Code.InternalError); + if (!doc || (resolvedDid && resolvedDid !== getDid(doc))) { + throw new DataPlaneError(Code.NotFound); } + + const result = getResultFromDoc(doc); + return { + ...result, + resolvedFrom: { + identifier, + type: identifierType, + }, + }; } async resolveBatch( @@ -133,43 +118,32 @@ export class Identity { const results = await Promise.allSettled( identifiers.map(async ({ value, type }) => { - try { - let doc: DidDocument | null = null; - let resolvedDid: string | null = null; - - const identifierType = type || - (value.startsWith("did:") ? "did" : "handle"); - - if (identifierType === "did") { - doc = await this.idResolver!.did.resolve(value); - resolvedDid = value; - } else { - resolvedDid = await this.idResolver!.handle.resolve(value) || null; - if (resolvedDid) { - doc = await this.idResolver!.did.resolve(resolvedDid); - } - } - - if (!doc || (resolvedDid && resolvedDid !== getDid(doc))) { - return { - identifier: value, - type: identifierType, - error: "Identity not found", - }; - } + let doc: DidDocument | null = null; + let resolvedDid: string | undefined; + + const identifierType = type || + (value.startsWith("did:") ? "did" : "handle"); + if (identifierType === "did") { + doc = await this.idResolver!.did.resolve(value); + resolvedDid = value; + } else { + resolvedDid = await this.idResolver!.handle.resolve(value); + if (!resolvedDid) throw new DataPlaneError(Code.NotFound); + doc = await this.idResolver!.did.resolve(resolvedDid); + } + if (!doc || (resolvedDid && resolvedDid !== getDid(doc))) { return { identifier: value, type: identifierType, - ...getResultFromDoc(doc), - }; - } catch (_error) { - return { - identifier: value, - type: type || "unknown", - error: "Failed to resolve identity", + error: "Identity not found", }; } + return { + identifier: value, + type: identifierType, + ...getResultFromDoc(doc), + }; }), ); diff --git a/data-plane/routes/records.ts b/data-plane/routes/records.ts index 177a76d..03f9b91 100644 --- a/data-plane/routes/records.ts +++ b/data-plane/routes/records.ts @@ -136,62 +136,32 @@ export class Records { } async getLikeRecords(uris: string[]) { - try { - const result = await getRecords(this.db, uris, ids.SoSprkFeedLike); - return result; - } catch (error) { - console.error("Error fetching like records:", error); - throw new DataPlaneError(Code.InternalError); - } + const result = await getRecords(this.db, uris, ids.SoSprkFeedLike); + return result; } async getPostRecords(uris: string[]) { - try { - const result = await getPostRecords(this.db, uris); - return result; - } catch (error) { - console.error("Error fetching post records:", error); - throw new DataPlaneError(Code.InternalError); - } + const result = await getPostRecords(this.db, uris); + return result; } async getReplyRecords(uris: string[]) { - try { - const result = await getReplyRecords(this.db, uris); - return result; - } catch (error) { - console.error("Error fetching reply records:", error); - throw new DataPlaneError(Code.InternalError); - } + const result = await getReplyRecords(this.db, uris); + return result; } async getProfileRecords(uris: string[]) { - try { - const result = await getRecords(this.db, uris, ids.SoSprkActorProfile); - return result; - } catch (error) { - console.error("Error fetching profile records:", error); - throw new DataPlaneError(Code.InternalError); - } + const result = await getRecords(this.db, uris, ids.SoSprkActorProfile); + return result; } async getRepostRecords(uris: string[]) { - try { - const result = await getRecords(this.db, uris, ids.AppBskyFeedRepost); - return result; - } catch (error) { - console.error("Error fetching repost records:", error); - throw new DataPlaneError(Code.InternalError); - } + const result = await getRecords(this.db, uris, ids.AppBskyFeedRepost); + return result; } async getRecords(uris: string[]) { - try { - const result = await getRecords(this.db, uris); - return result; - } catch (error) { - console.error("Error fetching records:", error); - throw new DataPlaneError(Code.InternalError); - } + const result = await getRecords(this.db, uris); + return result; } } diff --git a/data-plane/subscription.ts b/data-plane/subscription.ts index 9be90e2..59b0824 100644 --- a/data-plane/subscription.ts +++ b/data-plane/subscription.ts @@ -31,7 +31,6 @@ export class RepoSubscription { cfg, idResolver, this.background, - this.logger, ); const { runner, firehose } = createFirehose({ diff --git a/lex/lexicons.ts b/lex/lexicons.ts index d698308..5aa6496 100644 --- a/lex/lexicons.ts +++ b/lex/lexicons.ts @@ -12868,13 +12868,8 @@ export const schemaDict = { "description": "Combinations of post/repost types to include in response.", "knownValues": [ - "posts_with_replies", - "posts_no_replies", - "posts_with_media", - "posts_and_author_threads", "posts_with_video", ], - "default": "posts_with_replies", }, "includePins": { "type": "boolean", diff --git a/lex/types/so/sprk/feed/getAuthorFeed.ts b/lex/types/so/sprk/feed/getAuthorFeed.ts index 613b405..d6700d7 100644 --- a/lex/types/so/sprk/feed/getAuthorFeed.ts +++ b/lex/types/so/sprk/feed/getAuthorFeed.ts @@ -8,11 +8,7 @@ export type QueryParams = { limit: number; cursor?: string; /** Combinations of post/repost types to include in response. */ - filter: - | "posts_with_replies" - | "posts_no_replies" - | "posts_with_media" - | "posts_and_author_threads" + filter?: | "posts_with_video" | (string & globalThis.Record); includePins: boolean; diff --git a/lexicons/so/sprk/feed/getAuthorFeed.json b/lexicons/so/sprk/feed/getAuthorFeed.json index 8ff1963..58463bf 100644 --- a/lexicons/so/sprk/feed/getAuthorFeed.json +++ b/lexicons/so/sprk/feed/getAuthorFeed.json @@ -21,13 +21,8 @@ "type": "string", "description": "Combinations of post/repost types to include in response.", "knownValues": [ - "posts_with_replies", - "posts_no_replies", - "posts_with_media", - "posts_and_author_threads", "posts_with_video" - ], - "default": "posts_with_replies" + ] }, "includePins": { "type": "boolean", diff --git a/utils/embed-normalizer.ts b/utils/embed-normalizer.ts deleted file mode 100644 index 088007d..0000000 --- a/utils/embed-normalizer.ts +++ /dev/null @@ -1,308 +0,0 @@ -/** - * Utility functions for normalizing embeds to ensure CID objects are converted to $link format - * This ensures consistent storage and retrieval of embed data across the application. - */ - -interface CidRef { - $link?: string; - code?: number; - version?: number; - multihash?: Uint8Array; - bytes?: string; - toString?: () => string; -} - -interface NormalizedCidRef { - $link: string; -} - -interface VideoEmbed { - $type: "so.sprk.embed.video"; - video?: { - $type: "blob"; - ref: CidRef; - mimeType?: string; - size?: number; - }; -} - -interface ImageEmbed { - $type: "so.sprk.embed.images"; - images?: Array<{ - image: { - $type: "blob"; - ref: CidRef; - mimeType?: string; - size?: number; - }; - alt?: string; - aspectRatio?: number; - }>; -} - -interface Profile { - avatar?: { - $type?: "blob"; - ref?: NormalizedCidRef | null; - } | CidRef; - banner?: { - $type?: "blob"; - ref?: NormalizedCidRef | null; - } | CidRef; - [key: string]: unknown; -} - -// Normalize embed to ensure CID objects are converted to $link format -export function normalizeEmbed(embed: unknown): unknown { - if (!embed || typeof embed !== "object") return embed; - - const embedObj = embed as Record; - - if ( - embedObj.$type === "so.sprk.media.video" && embedObj.video && - typeof embedObj.video === "object" - ) { - const video = embedObj.video as Record; - if (video.ref) { - const ref = video.ref; - // If ref is a CID object (has code/version/multihash), convert to $link - if ( - typeof ref === "object" && ref && !(ref as CidRef).$link && - ((ref as CidRef).code || (ref as CidRef).version || - (ref as CidRef).multihash) - ) { - const toStringFn = (ref as CidRef).toString; - - if (toStringFn && typeof toStringFn === "function") { - const cidString = toStringFn.call(ref); - // Return cleaned up structure without 'original' field - return { - $type: "so.sprk.media.video", - video: { - $type: "blob", - ref: { $link: cidString }, - mimeType: video.mimeType, - size: video.size, - }, - }; - } else { - console.error("DEBUG: Could not convert CID object to string:", ref); - return embed; // Return original if we can't convert - } - } else if ((ref as CidRef).$link) { - // Already normalized, return cleaned up structure - return { - $type: "so.sprk.media.video", - video: { - $type: "blob", - ref: { $link: (ref as CidRef).$link }, - mimeType: video.mimeType, - size: video.size, - }, - }; - } - } - } - - if ( - embedObj.$type === "so.sprk.media.images" && Array.isArray(embedObj.images) - ) { - const normalizedImages = embedObj.images.map((img: unknown) => { - if ( - typeof img === "object" && img && (img as Record).image - ) { - const image = (img as Record).image as Record< - string, - unknown - >; - if (image.ref) { - const ref = image.ref; - if ( - typeof ref === "object" && ref && !(ref as CidRef).$link && - ((ref as CidRef).code || (ref as CidRef).version || - (ref as CidRef).multihash) - ) { - const toStringFn = (ref as CidRef).toString; - - if (toStringFn && typeof toStringFn === "function") { - const cidString = toStringFn.call(ref); - return { - image: { - $type: "blob", - ref: { $link: cidString }, - mimeType: image.mimeType, - size: image.size, - }, - alt: (img as Record).alt, - aspectRatio: (img as Record).aspectRatio, - }; - } else { - console.error( - "DEBUG: Could not convert CID object to string:", - ref, - ); - return img; - } - } else if ((ref as CidRef).$link) { - // Already normalized - return { - image: { - $type: "blob", - ref: { $link: (ref as CidRef).$link }, - mimeType: image.mimeType, - size: image.size, - }, - alt: (img as Record).alt, - aspectRatio: (img as Record).aspectRatio, - }; - } - } - } - return img; - }); - - return { - $type: "so.sprk.media.images", - images: normalizedImages, - }; - } - - return embed; -} - -// Normalize a single CID reference to $link format -export function normalizeCidRef( - ref: CidRef | null | undefined, -): NormalizedCidRef | null | undefined { - if (!ref) return ref; - - // If it's already in $link format, return as-is - if (typeof ref === "object" && ref.$link) { - return ref as NormalizedCidRef; - } - - // If it's a CID object (has code/version/multihash), convert to $link - if ( - typeof ref === "object" && !ref.$link && - (ref.code || ref.version || ref.multihash) - ) { - let cidString: string; - const toStringFn = ref.toString; - - if (toStringFn && typeof toStringFn === "function") { - cidString = toStringFn.call(ref); - } else { - console.error("DEBUG: Could not convert CID object to string:", ref); - return ref as NormalizedCidRef; // Return original if we can't convert - } - - return { $link: cidString }; - } - - // If it's a string, wrap it in $link format - if (typeof ref === "string") { - return { $link: ref }; - } - - return ref as NormalizedCidRef; -} - -// Normalize profile data to ensure any CID references are converted to $link format -export function normalizeProfile(profile: unknown): unknown { - if (!profile || typeof profile !== "object") return profile; - - const normalized: Record = { - ...profile as Record, - }; - - // Normalize avatar if present - if (normalized.avatar) { - // If avatar is a BlobRef (has ref property), normalize the ref - if ( - typeof normalized.avatar === "object" && normalized.avatar && - "ref" in normalized.avatar - ) { - const normalizedRef = normalizeCidRef( - (normalized.avatar as Record).ref as CidRef, - ); - if (normalizedRef) { - // Return only the MediaRef format: { $type: "blob", ref: { $link: string } } - normalized.avatar = { - $type: "blob", - ref: normalizedRef, - }; - } - } else { - // If avatar is just a CID object, wrap it in the proper structure - const normalizedRef = normalizeCidRef(normalized.avatar as CidRef); - if (normalizedRef) { - normalized.avatar = { - $type: "blob", - ref: normalizedRef, - }; - } - } - } - - // Normalize banner if present - if (normalized.banner) { - // If banner is a BlobRef (has ref property), normalize the ref - if ( - typeof normalized.banner === "object" && normalized.banner && - "ref" in normalized.banner - ) { - const normalizedRef = normalizeCidRef( - (normalized.banner as Record).ref as CidRef, - ); - if (normalizedRef) { - // Return only the MediaRef format: { $type: "blob", ref: { $link: string } } - normalized.banner = { - $type: "blob", - ref: normalizedRef, - }; - } - } else { - // If banner is just a CID object, wrap it in the proper structure - const normalizedRef = normalizeCidRef(normalized.banner as CidRef); - if (normalizedRef) { - normalized.banner = { - $type: "blob", - ref: normalizedRef, - }; - } - } - } - - return normalized; -} - -// Normalize any object that might contain CID references -export function normalizeObject(obj: unknown): unknown { - if (!obj || typeof obj !== "object") return obj; - - if (Array.isArray(obj)) { - return obj.map((item) => normalizeObject(item)); - } - - const normalized: Record = {}; - - for (const [key, value] of Object.entries(obj as Record)) { - if ( - key === "ref" && typeof value === "object" && value && - !(value as CidRef).$link && - ((value as CidRef).code || (value as CidRef).version || - (value as CidRef).multihash) - ) { - // This looks like a CID object, normalize it - normalized[key] = normalizeCidRef(value as CidRef); - } else if (typeof value === "object" && value !== null) { - // Recursively normalize nested objects - normalized[key] = normalizeObject(value); - } else { - // Keep primitive values as-is - normalized[key] = value; - } - } - - return normalized; -} diff --git a/utils/media-transformer.ts b/utils/media-transformer.ts index c74f875..4b1f1b9 100644 --- a/utils/media-transformer.ts +++ b/utils/media-transformer.ts @@ -2,6 +2,7 @@ import type * as SoSprkMediaImage from "../lex/types/so/sprk/media/image.ts"; import { ImageMedia, PostMedia, + StoryMedia, VideoMappingDocument, } from "../data-plane/db/models.ts"; import { ServerConfig } from "../config.ts"; @@ -80,7 +81,7 @@ export function transformVideoMedia( } export function transformMedia( - media: PostMedia | undefined, + media: PostMedia | StoryMedia | undefined, authorDid: string, cfg: ServerConfig, videoMapping?: VideoMappingDocument | null, @@ -92,12 +93,30 @@ export function transformMedia( } if (media.$type === "so.sprk.media.images") { - return transformImagesMedia(media, authorDid, options); + return transformImagesMedia(media as PostMedia, authorDid, options); + } + + if (media.$type === "so.sprk.media.image") { + // Handle single image (used in stories and replies) + const singleImageMedia = media as StoryMedia; + if (!singleImageMedia.image) { + return undefined; + } + + return { + $type: "so.sprk.media.image#view", + thumb: + `https://media.sprk.so/img/medium/${authorDid}/${singleImageMedia.image.ref.$link}/webp`, + fullsize: + `https://media.sprk.so/img/full/${authorDid}/${singleImageMedia.image.ref.$link}/webp`, + alt: singleImageMedia.image.alt ?? "", + aspectRatio: singleImageMedia.image.aspectRatio || undefined, + } as const; } if (media.$type === "so.sprk.media.video") { return transformVideoMedia( - media, + media as PostMedia, authorDid, cfg, videoMapping, diff --git a/utils/uris.ts b/utils/uris.ts index db8bc8b..6559c66 100644 --- a/utils/uris.ts +++ b/utils/uris.ts @@ -32,20 +32,7 @@ export function postUriToPostgateUri(postUri: string) { } export function uriToDid(uri: string) { - try { - return new AtUri(uri).hostname; - } catch (error) { - console.log(`AtUri parser failed for URI: ${uri}, error:`, error); - // Handle custom collection namespaces that AtUri might not recognize - // Extract DID from URI manually as fallback - const match = uri.match(/^at:\/\/(did:[^\/]+)/); - if (match) { - console.log(`Successfully extracted DID using fallback: ${match[1]}`); - return match[1]; - } - console.error(`Failed to extract DID from URI: ${uri}`); - throw new Error(`Invalid AT URI format: ${uri}`); - } + return new AtUri(uri).hostname; } // @TODO temp fix for proliferation of invalid pinned post values