diff --git a/actions/getIdentityData.ts b/actions/getIdentityData.ts index b2d16130..6ed8e690 100644 --- a/actions/getIdentityData.ts +++ b/actions/getIdentityData.ts @@ -3,6 +3,7 @@ import { cookies } from "next/headers"; import { supabaseServerClient } from "supabase/serverClient"; import { cache } from "react"; +import { deduplicateByUri } from "src/utils/deduplicateRecords"; export const getIdentityData = cache(uncachedGetIdentityData); export async function uncachedGetIdentityData() { let cookieStore = await cookies(); @@ -44,13 +45,15 @@ export async function uncachedGetIdentityData() { if (!auth_res?.data?.identities) return null; if (auth_res.data.identities.atp_did) { //I should create a relationship table so I can do this in the above query - let { data: publications } = await supabaseServerClient + let { data: rawPublications } = await supabaseServerClient .from("publications") .select("*") .eq("identity_did", auth_res.data.identities.atp_did); + // Deduplicate records that may exist under both pub.leaflet and site.standard namespaces + const publications = deduplicateByUri(rawPublications || []); return { ...auth_res.data.identities, - publications: publications || [], + publications, }; } diff --git a/app/(home-pages)/discover/getPublications.ts b/app/(home-pages)/discover/getPublications.ts index 85fb345f..ce649276 100644 --- a/app/(home-pages)/discover/getPublications.ts +++ b/app/(home-pages)/discover/getPublications.ts @@ -5,6 +5,7 @@ import { normalizePublicationRow, hasValidPublication, } from "src/utils/normalizeRecords"; +import { deduplicateByUri } from "src/utils/deduplicateRecords"; export type Cursor = { indexed_at?: string; @@ -42,8 +43,11 @@ export async function getPublications( return { publications: [], nextCursor: null }; } + // Deduplicate records that may exist under both pub.leaflet and site.standard namespaces + const dedupedPublications = deduplicateByUri(publications || []); + // Filter out publications without documents - const allPubs = (publications || []).filter( + const allPubs = dedupedPublications.filter( (pub) => pub.documents_in_publications.length > 0, ); diff --git a/app/(home-pages)/p/[didOrHandle]/getProfilePosts.ts b/app/(home-pages)/p/[didOrHandle]/getProfilePosts.ts index 067f8e3e..29f98afb 100644 --- a/app/(home-pages)/p/[didOrHandle]/getProfilePosts.ts +++ b/app/(home-pages)/p/[didOrHandle]/getProfilePosts.ts @@ -7,6 +7,7 @@ import { normalizeDocumentRecord, normalizePublicationRecord, } from "src/utils/normalizeRecords"; +import { deduplicateByUriOrdered } from "src/utils/deduplicateRecords"; export type Cursor = { indexed_at: string; @@ -38,7 +39,7 @@ export async function getProfilePosts( ); } - let [{ data: docs }, { data: pubs }, { data: profile }] = await Promise.all([ + let [{ data: rawDocs }, { data: rawPubs }, { data: profile }] = await Promise.all([ query, supabaseServerClient .from("publications") @@ -51,6 +52,10 @@ export async function getProfilePosts( .single(), ]); + // Deduplicate records that may exist under both pub.leaflet and site.standard namespaces + const docs = deduplicateByUriOrdered(rawDocs || []); + const pubs = deduplicateByUriOrdered(rawPubs || []); + // Build a map of publications for quick lookup let pubMap = new Map[number]>(); for (let pub of pubs || []) { diff --git a/app/(home-pages)/reader/getReaderFeed.ts b/app/(home-pages)/reader/getReaderFeed.ts index f3ada41d..753142c7 100644 --- a/app/(home-pages)/reader/getReaderFeed.ts +++ b/app/(home-pages)/reader/getReaderFeed.ts @@ -14,6 +14,7 @@ import { type NormalizedDocument, type NormalizedPublication, } from "src/utils/normalizeRecords"; +import { deduplicateByUriOrdered } from "src/utils/deduplicateRecords"; export type Cursor = { timestamp: string; @@ -45,11 +46,14 @@ export async function getReaderFeed( `indexed_at.lt.${cursor.timestamp},and(indexed_at.eq.${cursor.timestamp},uri.lt.${cursor.uri})`, ); } - let { data: feed, error } = await query; + let { data: rawFeed, error } = await query; + + // Deduplicate records that may exist under both pub.leaflet and site.standard namespaces + const feed = deduplicateByUriOrdered(rawFeed || []); let posts = ( await Promise.all( - feed?.map(async (post) => { + feed.map(async (post) => { let pub = post.documents_in_publications[0].publications!; let uri = new AtUri(post.uri); let handle = await idResolver.did.resolve(uri.host); diff --git a/app/(home-pages)/tag/[tag]/getDocumentsByTag.ts b/app/(home-pages)/tag/[tag]/getDocumentsByTag.ts index b22417e8..b0e266e8 100644 --- a/app/(home-pages)/tag/[tag]/getDocumentsByTag.ts +++ b/app/(home-pages)/tag/[tag]/getDocumentsByTag.ts @@ -9,12 +9,13 @@ import { normalizeDocumentRecord, normalizePublicationRecord, } from "src/utils/normalizeRecords"; +import { deduplicateByUriOrdered } from "src/utils/deduplicateRecords"; export async function getDocumentsByTag( tag: string, ): Promise<{ posts: Post[] }> { // Query documents that have this tag - const { data: documents, error } = await supabaseServerClient + const { data: rawDocuments, error } = await supabaseServerClient .from("documents") .select( `*, @@ -31,6 +32,9 @@ export async function getDocumentsByTag( return { posts: [] }; } + // Deduplicate records that may exist under both pub.leaflet and site.standard namespaces + const documents = deduplicateByUriOrdered(rawDocuments || []); + const posts = await Promise.all( documents.map(async (doc) => { const pub = doc.documents_in_publications[0]?.publications; diff --git a/app/api/inngest/functions/migrate_user_to_standard.ts b/app/api/inngest/functions/migrate_user_to_standard.ts index c27b14ec..7c630e42 100644 --- a/app/api/inngest/functions/migrate_user_to_standard.ts +++ b/app/api/inngest/functions/migrate_user_to_standard.ts @@ -1,10 +1,18 @@ import { supabaseServerClient } from "supabase/serverClient"; import { inngest } from "../client"; import { restoreOAuthSession } from "src/atproto-oauth"; -import { AtpBaseClient, SiteStandardPublication, SiteStandardDocument, SiteStandardGraphSubscription } from "lexicons/api"; +import { + AtpBaseClient, + SiteStandardPublication, + SiteStandardDocument, + SiteStandardGraphSubscription, +} from "lexicons/api"; import { AtUri } from "@atproto/syntax"; import { Json } from "supabase/database.types"; -import { normalizePublicationRecord, normalizeDocumentRecord } from "src/utils/normalizeRecords"; +import { + normalizePublicationRecord, + normalizeDocumentRecord, +} from "src/utils/normalizeRecords"; type MigrationResult = | { success: true; oldUri: string; newUri: string; skipped?: boolean } @@ -17,7 +25,7 @@ async function createAuthenticatedAgent(did: string): Promise { } const credentialSession = result.value; return new AtpBaseClient( - credentialSession.fetchHandler.bind(credentialSession) + credentialSession.fetchHandler.bind(credentialSession), ); } @@ -39,374 +47,439 @@ export const migrate_user_to_standard = inngest.createFunction( await step.run("verify-oauth-session", async () => { const result = await restoreOAuthSession(did); if (!result.ok) { - throw new Error(`Failed to restore OAuth session: ${result.error.message}`); + throw new Error( + `Failed to restore OAuth session: ${result.error.message}`, + ); } return { success: true }; }); // Step 2: Get user's pub.leaflet.publication records - const oldPublications = await step.run("fetch-old-publications", async () => { - const { data, error } = await supabaseServerClient - .from("publications") - .select("*") - .eq("identity_did", did) - .like("uri", `at://${did}/pub.leaflet.publication/%`); - - if (error) throw new Error(`Failed to fetch publications: ${error.message}`); - return data || []; - }); - - // Step 3: Migrate each publication - const publicationUriMap: Record = {}; // old URI -> new URI + const oldPublications = await step.run( + "fetch-old-publications", + async () => { + const { data, error } = await supabaseServerClient + .from("publications") + .select("*") + .eq("identity_did", did) + .like("uri", `at://${did}/pub.leaflet.publication/%`); - for (const pub of oldPublications) { - const aturi = new AtUri(pub.uri); + if (error) + throw new Error(`Failed to fetch publications: ${error.message}`); + return data || []; + }, + ); - // Skip if already a site.standard.publication - if (aturi.collection === "site.standard.publication") { - publicationUriMap[pub.uri] = pub.uri; - continue; - } + // Step 3: Migrate all publications in parallel + const publicationUriMap: Record = {}; // old URI -> new URI - const rkey = aturi.rkey; - const normalized = normalizePublicationRecord(pub.record); + // Prepare publications that need migration + const publicationsToMigrate = oldPublications + .map((pub) => { + const aturi = new AtUri(pub.uri); - if (!normalized) { - stats.errors.push(`Publication ${pub.uri}: Failed to normalize publication record`); - continue; - } + // Skip if already a site.standard.publication + if (aturi.collection === "site.standard.publication") { + publicationUriMap[pub.uri] = pub.uri; + return null; + } - // Build site.standard.publication record - const newRecord: SiteStandardPublication.Record = { - $type: "site.standard.publication", - name: normalized.name, - url: normalized.url, - description: normalized.description, - icon: normalized.icon, - theme: normalized.theme, - basicTheme: normalized.basicTheme, - preferences: normalized.preferences, - }; - - // Step: Write to PDS - const pdsResult = await step.run(`pds-write-publication-${pub.uri}`, async () => { - const agent = await createAuthenticatedAgent(did); - const putResult = await agent.com.atproto.repo.putRecord({ - repo: did, - collection: "site.standard.publication", - rkey, - record: newRecord, - validate: false, - }); - return { newUri: putResult.data.uri }; - }); + const rkey = aturi.rkey; + const normalized = normalizePublicationRecord(pub.record); - const newUri = pdsResult.newUri; + if (!normalized) { + stats.errors.push( + `Publication ${pub.uri}: Failed to normalize publication record`, + ); + return null; + } - // Step: Write to database - const dbResult = await step.run(`db-write-publication-${pub.uri}`, async () => { - const { error: dbError } = await supabaseServerClient - .from("publications") - .upsert({ - uri: newUri, - identity_did: did, - name: normalized.name, - record: newRecord as Json, + const newRecord: SiteStandardPublication.Record = { + $type: "site.standard.publication", + name: normalized.name, + url: normalized.url, + description: normalized.description, + icon: normalized.icon, + theme: normalized.theme, + basicTheme: normalized.basicTheme, + preferences: normalized.preferences, + }; + + return { pub, rkey, normalized, newRecord }; + }) + .filter((x) => x !== null); + + // Run all PDS writes in parallel + const pubPdsResults = await Promise.all( + publicationsToMigrate.map(({ pub, rkey, newRecord }) => + step.run(`pds-write-publication-${pub.uri}`, async () => { + const agent = await createAuthenticatedAgent(did); + const putResult = await agent.com.atproto.repo.putRecord({ + repo: did, + collection: "site.standard.publication", + rkey, + record: newRecord, + validate: false, }); + return { oldUri: pub.uri, newUri: putResult.data.uri }; + }), + ), + ); + + // Run all DB writes in parallel + const pubDbResults = await Promise.all( + publicationsToMigrate.map(({ pub, normalized, newRecord }, index) => { + const newUri = pubPdsResults[index].newUri; + return step.run(`db-write-publication-${pub.uri}`, async () => { + const { error: dbError } = await supabaseServerClient + .from("publications") + .upsert({ + uri: newUri, + identity_did: did, + name: normalized.name, + record: newRecord as Json, + }); - if (dbError) { - return { success: false as const, error: dbError.message }; - } - return { success: true as const }; - }); + if (dbError) { + return { + success: false as const, + oldUri: pub.uri, + newUri, + error: dbError.message, + }; + } + return { success: true as const, oldUri: pub.uri, newUri }; + }); + }), + ); - if (dbResult.success) { - publicationUriMap[pub.uri] = newUri; + // Process results + for (const result of pubDbResults) { + if (result.success) { + publicationUriMap[result.oldUri] = result.newUri; stats.publicationsMigrated++; } else { - stats.errors.push(`Publication ${pub.uri}: Database error: ${dbResult.error}`); + stats.errors.push( + `Publication ${result.oldUri}: Database error: ${result.error}`, + ); } } - // Step 4: Get ALL user's pub.leaflet.document records (both in publications and standalone) - const oldDocuments = await step.run("fetch-old-documents", async () => { - const { data, error } = await supabaseServerClient - .from("documents") - .select("uri, data") - .like("uri", `at://${did}/pub.leaflet.document/%`); - - if (error) throw new Error(`Failed to fetch documents: ${error.message}`); - return data || []; - }); - - // Also fetch publication associations for documents - const documentPublicationMap = await step.run("fetch-document-publications", async () => { - const docUris = oldDocuments.map(d => d.uri); - if (docUris.length === 0) return {}; - - const { data, error } = await supabaseServerClient - .from("documents_in_publications") - .select("document, publication") - .in("document", docUris); - - if (error) throw new Error(`Failed to fetch document publications: ${error.message}`); - - // Create a map of document URI -> publication URI - const map: Record = {}; - for (const row of data || []) { - map[row.document] = row.publication; - } - return map; - }); + // Step 4: Get ALL user's documents and their publication associations in parallel + const [oldDocuments, allDocumentPublications] = await Promise.all([ + step.run("fetch-old-documents", async () => { + const { data, error } = await supabaseServerClient + .from("documents") + .select("uri, data") + .like("uri", `at://${did}/pub.leaflet.document/%`); + + if (error) + throw new Error(`Failed to fetch documents: ${error.message}`); + return data || []; + }), + step.run("fetch-document-publications", async () => { + const { data, error } = await supabaseServerClient + .from("documents_in_publications") + .select("document, publication") + .like("document", `at://${did}/pub.leaflet.document/%`); + + if (error) + throw new Error( + `Failed to fetch document publications: ${error.message}`, + ); + return data || []; + }), + ]); + + // Create a map of document URI -> publication URI + const documentPublicationMap: Record = {}; + for (const row of allDocumentPublications) { + documentPublicationMap[row.document] = row.publication; + } const documentUriMap: Record = {}; // old URI -> new URI - for (const doc of oldDocuments) { - const aturi = new AtUri(doc.uri); - - // Skip if already a site.standard.document - if (aturi.collection === "site.standard.document") { - documentUriMap[doc.uri] = doc.uri; - continue; - } - - const rkey = aturi.rkey; - const normalized = normalizeDocumentRecord(doc.data, doc.uri); - - if (!normalized) { - stats.errors.push(`Document ${doc.uri}: Failed to normalize document record`); - continue; - } + // Prepare documents that need migration + const documentsToMigrate = oldDocuments + .map((doc) => { + const aturi = new AtUri(doc.uri); - // Determine the site field: - // - If document is in a publication, use the new publication URI (if migrated) or old URI - // - If standalone, use the HTTPS URL format - const oldPubUri = documentPublicationMap[doc.uri]; - let siteValue: string; - - if (oldPubUri) { - // Document is in a publication - use new URI if migrated, otherwise keep old - siteValue = publicationUriMap[oldPubUri] || oldPubUri; - } else { - // Standalone document - use HTTPS URL format - siteValue = `https://leaflet.pub/p/${did}`; - } + // Skip if already a site.standard.document + if (aturi.collection === "site.standard.document") { + documentUriMap[doc.uri] = doc.uri; + return null; + } - // Build site.standard.document record - const newRecord: SiteStandardDocument.Record = { - $type: "site.standard.document", - title: normalized.title || "Untitled", - site: siteValue, - path: rkey, - publishedAt: normalized.publishedAt || new Date().toISOString(), - description: normalized.description, - content: normalized.content, - tags: normalized.tags, - coverImage: normalized.coverImage, - bskyPostRef: normalized.bskyPostRef, - }; - - // Step: Write to PDS - const pdsResult = await step.run(`pds-write-document-${doc.uri}`, async () => { - const agent = await createAuthenticatedAgent(did); - const putResult = await agent.com.atproto.repo.putRecord({ - repo: did, - collection: "site.standard.document", - rkey, - record: newRecord, - validate: false, - }); - return { newUri: putResult.data.uri }; - }); + const rkey = aturi.rkey; + const normalized = normalizeDocumentRecord(doc.data, doc.uri); - const newUri = pdsResult.newUri; + if (!normalized) { + stats.errors.push( + `Document ${doc.uri}: Failed to normalize document record`, + ); + return null; + } - // Step: Write to database - const dbResult = await step.run(`db-write-document-${doc.uri}`, async () => { - const { error: dbError } = await supabaseServerClient - .from("documents") - .upsert({ - uri: newUri, - data: newRecord as Json, - }); + // Determine the site field: + // - If document is in a publication, use the new publication URI (if migrated) or old URI + // - If standalone, use the HTTPS URL format + const oldPubUri = documentPublicationMap[doc.uri]; + let siteValue: string; - if (dbError) { - return { success: false as const, error: dbError.message }; + if (oldPubUri) { + // Document is in a publication - use new URI if migrated, otherwise keep old + siteValue = publicationUriMap[oldPubUri] || oldPubUri; + } else { + // Standalone document - use HTTPS URL format + siteValue = `https://leaflet.pub/p/${did}`; } - // If document was in a publication, add to documents_in_publications with new URIs - if (oldPubUri) { - const newPubUri = publicationUriMap[oldPubUri] || oldPubUri; - await supabaseServerClient - .from("documents_in_publications") + // Build site.standard.document record + const newRecord: SiteStandardDocument.Record = { + $type: "site.standard.document", + title: normalized.title || "Untitled", + site: siteValue, + path: rkey, + publishedAt: normalized.publishedAt || new Date().toISOString(), + description: normalized.description, + content: normalized.content, + tags: normalized.tags, + coverImage: normalized.coverImage, + bskyPostRef: normalized.bskyPostRef, + }; + + return { doc, rkey, normalized, newRecord, oldPubUri }; + }) + .filter((x) => x !== null); + + // Run all PDS writes in parallel + const docPdsResults = await Promise.all( + documentsToMigrate.map(({ doc, rkey, newRecord }) => + step.run(`pds-write-document-${doc.uri}`, async () => { + const agent = await createAuthenticatedAgent(did); + const putResult = await agent.com.atproto.repo.putRecord({ + repo: did, + collection: "site.standard.document", + rkey, + record: newRecord, + validate: false, + }); + return { oldUri: doc.uri, newUri: putResult.data.uri }; + }), + ), + ); + + // Run all DB writes in parallel + const docDbResults = await Promise.all( + documentsToMigrate.map(({ doc, newRecord, oldPubUri }, index) => { + const newUri = docPdsResults[index].newUri; + return step.run(`db-write-document-${doc.uri}`, async () => { + const { error: dbError } = await supabaseServerClient + .from("documents") .upsert({ - publication: newPubUri, - document: newUri, + uri: newUri, + data: newRecord as Json, }); - } - return { success: true as const }; - }); + if (dbError) { + return { + success: false as const, + oldUri: doc.uri, + newUri, + error: dbError.message, + }; + } + + // If document was in a publication, add to documents_in_publications with new URIs + if (oldPubUri) { + const newPubUri = publicationUriMap[oldPubUri] || oldPubUri; + await supabaseServerClient + .from("documents_in_publications") + .upsert({ + publication: newPubUri, + document: newUri, + }); + } + + return { success: true as const, oldUri: doc.uri, newUri }; + }); + }), + ); - if (dbResult.success) { - documentUriMap[doc.uri] = newUri; + // Process results + for (const result of docDbResults) { + if (result.success) { + documentUriMap[result.oldUri] = result.newUri; stats.documentsMigrated++; } else { - stats.errors.push(`Document ${doc.uri}: Database error: ${dbResult.error}`); + stats.errors.push( + `Document ${result.oldUri}: Database error: ${result.error}`, + ); } } - // Step 5: Update references in database tables + // Step 5: Update references in database tables (all in parallel) await step.run("update-references", async () => { - // Update leaflets_in_publications - update publication and doc references - for (const [oldUri, newUri] of Object.entries(publicationUriMap)) { - const { error } = await supabaseServerClient - .from("leaflets_in_publications") - .update({ publication: newUri }) - .eq("publication", oldUri); - - if (!error) stats.referencesUpdated++; - } - - for (const [oldUri, newUri] of Object.entries(documentUriMap)) { - const { error } = await supabaseServerClient - .from("leaflets_in_publications") - .update({ doc: newUri }) - .eq("doc", oldUri); - - if (!error) stats.referencesUpdated++; - } - - // Update leaflets_to_documents - update document references - for (const [oldUri, newUri] of Object.entries(documentUriMap)) { - const { error } = await supabaseServerClient - .from("leaflets_to_documents") - .update({ document: newUri }) - .eq("document", oldUri); - - if (!error) stats.referencesUpdated++; - } - - // Update publication_domains - update publication references - for (const [oldUri, newUri] of Object.entries(publicationUriMap)) { - const { error } = await supabaseServerClient - .from("publication_domains") - .update({ publication: newUri }) - .eq("publication", oldUri); - - if (!error) stats.referencesUpdated++; - } - - // Update comments_on_documents - update document references - for (const [oldUri, newUri] of Object.entries(documentUriMap)) { - const { error } = await supabaseServerClient - .from("comments_on_documents") - .update({ document: newUri }) - .eq("document", oldUri); - - if (!error) stats.referencesUpdated++; - } - - // Update document_mentions_in_bsky - update document references - for (const [oldUri, newUri] of Object.entries(documentUriMap)) { - const { error } = await supabaseServerClient - .from("document_mentions_in_bsky") - .update({ document: newUri }) - .eq("document", oldUri); - - if (!error) stats.referencesUpdated++; - } - - // Update subscribers_to_publications - update publication references - for (const [oldUri, newUri] of Object.entries(publicationUriMap)) { - const { error } = await supabaseServerClient - .from("subscribers_to_publications") - .update({ publication: newUri }) - .eq("publication", oldUri); - - if (!error) stats.referencesUpdated++; - } - - // Update publication_subscriptions - update publication references for incoming subscriptions - for (const [oldUri, newUri] of Object.entries(publicationUriMap)) { - const { error } = await supabaseServerClient - .from("publication_subscriptions") - .update({ publication: newUri }) - .eq("publication", oldUri); - - if (!error) stats.referencesUpdated++; - } - + const pubEntries = Object.entries(publicationUriMap); + const docEntries = Object.entries(documentUriMap); + + const updatePromises = [ + // Update leaflets_in_publications - publication references + ...pubEntries.map(([oldUri, newUri]) => + supabaseServerClient + .from("leaflets_in_publications") + .update({ publication: newUri }) + .eq("publication", oldUri), + ), + // Update leaflets_in_publications - doc references + ...docEntries.map(([oldUri, newUri]) => + supabaseServerClient + .from("leaflets_in_publications") + .update({ doc: newUri }) + .eq("doc", oldUri), + ), + // Update leaflets_to_documents - document references + ...docEntries.map(([oldUri, newUri]) => + supabaseServerClient + .from("leaflets_to_documents") + .update({ document: newUri }) + .eq("document", oldUri), + ), + // Update publication_domains - publication references + ...pubEntries.map(([oldUri, newUri]) => + supabaseServerClient + .from("publication_domains") + .update({ publication: newUri }) + .eq("publication", oldUri), + ), + // Update comments_on_documents - document references + ...docEntries.map(([oldUri, newUri]) => + supabaseServerClient + .from("comments_on_documents") + .update({ document: newUri }) + .eq("document", oldUri), + ), + // Update document_mentions_in_bsky - document references + ...docEntries.map(([oldUri, newUri]) => + supabaseServerClient + .from("document_mentions_in_bsky") + .update({ document: newUri }) + .eq("document", oldUri), + ), + // Update subscribers_to_publications - publication references + ...pubEntries.map(([oldUri, newUri]) => + supabaseServerClient + .from("subscribers_to_publications") + .update({ publication: newUri }) + .eq("publication", oldUri), + ), + // Update publication_subscriptions - publication references + ...pubEntries.map(([oldUri, newUri]) => + supabaseServerClient + .from("publication_subscriptions") + .update({ publication: newUri }) + .eq("publication", oldUri), + ), + ]; + + const results = await Promise.all(updatePromises); + stats.referencesUpdated = results.filter((r) => !r.error).length; return stats.referencesUpdated; }); // Step 6: Migrate user's own subscriptions - subscriptions BY this user to other publications - const userSubscriptions = await step.run("fetch-user-subscriptions", async () => { - const { data, error } = await supabaseServerClient - .from("publication_subscriptions") - .select("*") - .eq("identity", did) - .like("uri", `at://${did}/pub.leaflet.graph.subscription/%`); - - if (error) throw new Error(`Failed to fetch user subscriptions: ${error.message}`); - return data || []; - }); + const userSubscriptions = await step.run( + "fetch-user-subscriptions", + async () => { + const { data, error } = await supabaseServerClient + .from("publication_subscriptions") + .select("*") + .eq("identity", did) + .like("uri", `at://${did}/pub.leaflet.graph.subscription/%`); + + if (error) + throw new Error( + `Failed to fetch user subscriptions: ${error.message}`, + ); + return data || []; + }, + ); const userSubscriptionUriMap: Record = {}; // old URI -> new URI - for (const sub of userSubscriptions) { - const aturi = new AtUri(sub.uri); + // Prepare subscriptions that need migration + const subscriptionsToMigrate = userSubscriptions + .map((sub) => { + const aturi = new AtUri(sub.uri); - // Skip if already a site.standard.graph.subscription - if (aturi.collection === "site.standard.graph.subscription") { - userSubscriptionUriMap[sub.uri] = sub.uri; - continue; - } + // Skip if already a site.standard.graph.subscription + if (aturi.collection === "site.standard.graph.subscription") { + userSubscriptionUriMap[sub.uri] = sub.uri; + return null; + } - const rkey = aturi.rkey; - - // Build site.standard.graph.subscription record - const newRecord: SiteStandardGraphSubscription.Record = { - $type: "site.standard.graph.subscription", - publication: sub.publication, - }; - - // Step: Write to PDS - const pdsResult = await step.run(`pds-write-subscription-${sub.uri}`, async () => { - const agent = await createAuthenticatedAgent(did); - const putResult = await agent.com.atproto.repo.putRecord({ - repo: did, - collection: "site.standard.graph.subscription", - rkey, - record: newRecord, - validate: false, + const rkey = aturi.rkey; + const newRecord: SiteStandardGraphSubscription.Record = { + $type: "site.standard.graph.subscription", + publication: sub.publication, + }; + + return { sub, rkey, newRecord }; + }) + .filter((x) => x !== null); + + // Run all PDS writes in parallel + const subPdsResults = await Promise.all( + subscriptionsToMigrate.map(({ sub, rkey, newRecord }) => + step.run(`pds-write-subscription-${sub.uri}`, async () => { + const agent = await createAuthenticatedAgent(did); + const putResult = await agent.com.atproto.repo.putRecord({ + repo: did, + collection: "site.standard.graph.subscription", + rkey, + record: newRecord, + validate: false, + }); + return { oldUri: sub.uri, newUri: putResult.data.uri }; + }), + ), + ); + + // Run all DB writes in parallel + const subDbResults = await Promise.all( + subscriptionsToMigrate.map(({ sub, newRecord }, index) => { + const newUri = subPdsResults[index].newUri; + return step.run(`db-write-subscription-${sub.uri}`, async () => { + const { error: dbError } = await supabaseServerClient + .from("publication_subscriptions") + .update({ + uri: newUri, + record: newRecord as Json, + }) + .eq("uri", sub.uri); + + if (dbError) { + return { + success: false as const, + oldUri: sub.uri, + newUri, + error: dbError.message, + }; + } + return { success: true as const, oldUri: sub.uri, newUri }; }); - return { newUri: putResult.data.uri }; - }); - - const newUri = pdsResult.newUri; + }), + ); - // Step: Write to database - const dbResult = await step.run(`db-write-subscription-${sub.uri}`, async () => { - const { error: dbError } = await supabaseServerClient - .from("publication_subscriptions") - .update({ - uri: newUri, - record: newRecord as Json, - }) - .eq("uri", sub.uri); - - if (dbError) { - return { success: false as const, error: dbError.message }; - } - return { success: true as const }; - }); - - if (dbResult.success) { - userSubscriptionUriMap[sub.uri] = newUri; + // Process results + for (const result of subDbResults) { + if (result.success) { + userSubscriptionUriMap[result.oldUri] = result.newUri; stats.userSubscriptionsMigrated++; } else { - stats.errors.push(`User subscription ${sub.uri}: Database error: ${dbResult.error}`); + stats.errors.push( + `User subscription ${result.oldUri}: Database error: ${result.error}`, + ); } } @@ -424,5 +497,5 @@ export const migrate_user_to_standard = inngest.createFunction( documentUriMap, userSubscriptionUriMap, }; - } + }, ); diff --git a/app/api/rpc/[command]/get_profile_data.ts b/app/api/rpc/[command]/get_profile_data.ts index 16e81fa3..f5a2216e 100644 --- a/app/api/rpc/[command]/get_profile_data.ts +++ b/app/api/rpc/[command]/get_profile_data.ts @@ -10,6 +10,7 @@ import { normalizePublicationRow, hasValidPublication, } from "src/utils/normalizeRecords"; +import { deduplicateByUri } from "src/utils/deduplicateRecords"; export type GetProfileDataReturnType = Awaited< ReturnType<(typeof get_profile_data)["handler"]> @@ -58,13 +59,16 @@ export const get_profile_data = makeRoute({ .select("*") .eq("identity_did", did); - let [{ data: profile }, { data: publications }] = await Promise.all([ + let [{ data: profile }, { data: rawPublications }] = await Promise.all([ profileReq, publicationsReq, ]); + // Deduplicate records that may exist under both pub.leaflet and site.standard namespaces + const publications = deduplicateByUri(rawPublications || []); + // Normalize publication records before returning - const normalizedPublications = (publications || []) + const normalizedPublications = publications .map(normalizePublicationRow) .filter(hasValidPublication); diff --git a/app/api/rpc/[command]/get_publication_data.ts b/app/api/rpc/[command]/get_publication_data.ts index bde8c685..46a04bad 100644 --- a/app/api/rpc/[command]/get_publication_data.ts +++ b/app/api/rpc/[command]/get_publication_data.ts @@ -52,8 +52,12 @@ export const get_publication_data = makeRoute({ ) )`, ) - .or(`name.eq."${publication_name}", uri.eq."${pubLeafletUri}", uri.eq."${siteStandardUri}"`) + .or( + `name.eq."${publication_name}", uri.eq."${pubLeafletUri}", uri.eq."${siteStandardUri}"`, + ) .eq("identity_did", did) + .order("uri", { ascending: false }) + .limit(1) .single(); let leaflet_data = await getFactsFromHomeLeaflets.handler( @@ -70,7 +74,10 @@ export const get_publication_data = makeRoute({ const documents = (publication?.documents_in_publications || []) .map((dip) => { if (!dip.documents) return null; - const normalized = normalizeDocumentRecord(dip.documents.data, dip.documents.uri); + const normalized = normalizeDocumentRecord( + dip.documents.data, + dip.documents.uri, + ); if (!normalized) return null; return { uri: dip.documents.uri, diff --git a/app/api/rpc/[command]/search_publication_names.ts b/app/api/rpc/[command]/search_publication_names.ts index 72636497..e5ab9ed7 100644 --- a/app/api/rpc/[command]/search_publication_names.ts +++ b/app/api/rpc/[command]/search_publication_names.ts @@ -2,6 +2,7 @@ import { z } from "zod"; import { makeRoute } from "../lib"; import type { Env } from "./route"; import { getPublicationURL } from "app/lish/createPub/getPublicationURL"; +import { deduplicateByUri } from "src/utils/deduplicateRecords"; export type SearchPublicationNamesReturnType = Awaited< ReturnType<(typeof search_publication_names)["handler"]> @@ -15,7 +16,7 @@ export const search_publication_names = makeRoute({ }), handler: async ({ query, limit }, { supabase }: Pick) => { // Search publications by name in record (case-insensitive partial match) - const { data: publications, error } = await supabase + const { data: rawPublications, error } = await supabase .from("publications") .select("uri, record") .ilike("record->>name", `%${query}%`) @@ -25,6 +26,9 @@ export const search_publication_names = makeRoute({ throw new Error(`Failed to search publications: ${error.message}`); } + // Deduplicate records that may exist under both pub.leaflet and site.standard namespaces + const publications = deduplicateByUri(rawPublications || []); + const result = publications.map((p) => { const record = p.record as { name?: string }; return { diff --git a/app/lish/[did]/[publication]/generateFeed.ts b/app/lish/[did]/[publication]/generateFeed.ts index 59fbed31..e12aa75f 100644 --- a/app/lish/[did]/[publication]/generateFeed.ts +++ b/app/lish/[did]/[publication]/generateFeed.ts @@ -19,7 +19,7 @@ export async function generateFeed( let renderToReadableStream = await import("react-dom/server").then( (module) => module.renderToReadableStream, ); - let { data: publications } = await supabaseServerClient + let { data: publications, error } = await supabaseServerClient .from("publications") .select( `*, @@ -31,6 +31,7 @@ export async function generateFeed( .or(publicationNameOrUriFilter(did, publication_name)) .order("uri", { ascending: false }) .limit(1); + console.log(error); let publication = publications?.[0]; const pubRecord = normalizePublicationRecord(publication?.record); @@ -54,7 +55,10 @@ export async function generateFeed( await Promise.all( publication.documents_in_publications.map(async (doc) => { if (!doc.documents) return; - const record = normalizeDocumentRecord(doc.documents?.data, doc.documents?.uri); + const record = normalizeDocumentRecord( + doc.documents?.data, + doc.documents?.uri, + ); const uri = new AtUri(doc.documents?.uri); const rkey = uri.rkey; if (!record) return; diff --git a/lexicons/api/lexicons.ts b/lexicons/api/lexicons.ts index 07822609..409eb985 100644 --- a/lexicons/api/lexicons.ts +++ b/lexicons/api/lexicons.ts @@ -1441,8 +1441,8 @@ export const schemaDict = { properties: { title: { type: 'string', - maxLength: 1280, - maxGraphemes: 128, + maxLength: 5000, + maxGraphemes: 500, }, postRef: { type: 'ref', @@ -1450,8 +1450,8 @@ export const schemaDict = { }, description: { type: 'string', - maxLength: 3000, - maxGraphemes: 300, + maxLength: 30000, + maxGraphemes: 3000, }, publishedAt: { type: 'string', @@ -2128,8 +2128,8 @@ export const schemaDict = { type: 'blob', }, description: { - maxGraphemes: 300, - maxLength: 3000, + maxGraphemes: 3000, + maxLength: 30000, type: 'string', }, path: { @@ -2165,8 +2165,8 @@ export const schemaDict = { type: 'ref', }, title: { - maxGraphemes: 128, - maxLength: 1280, + maxGraphemes: 500, + maxLength: 5000, type: 'string', }, updatedAt: { diff --git a/lexicons/pub/leaflet/document.json b/lexicons/pub/leaflet/document.json index 9043cdfd..2e714226 100644 --- a/lexicons/pub/leaflet/document.json +++ b/lexicons/pub/leaflet/document.json @@ -18,8 +18,8 @@ "properties": { "title": { "type": "string", - "maxLength": 1280, - "maxGraphemes": 128 + "maxLength": 5000, + "maxGraphemes": 500 }, "postRef": { "type": "ref", @@ -27,8 +27,8 @@ }, "description": { "type": "string", - "maxLength": 3000, - "maxGraphemes": 300 + "maxLength": 30000, + "maxGraphemes": 3000 }, "publishedAt": { "type": "string", diff --git a/lexicons/site/standard/document.json b/lexicons/site/standard/document.json index d0035ecd..433f19dc 100644 --- a/lexicons/site/standard/document.json +++ b/lexicons/site/standard/document.json @@ -19,8 +19,8 @@ "type": "blob" }, "description": { - "maxGraphemes": 300, - "maxLength": 3000, + "maxGraphemes": 3000, + "maxLength": 30000, "type": "string" }, "path": { @@ -53,8 +53,8 @@ "type": "ref" }, "title": { - "maxGraphemes": 128, - "maxLength": 1280, + "maxGraphemes": 500, + "maxLength": 5000, "type": "string" }, "updatedAt": { diff --git a/lexicons/src/document.ts b/lexicons/src/document.ts index 5fc33aee..f1beabff 100644 --- a/lexicons/src/document.ts +++ b/lexicons/src/document.ts @@ -16,9 +16,9 @@ export const PubLeafletDocument: LexiconDoc = { type: "object", required: ["pages", "author", "title"], properties: { - title: { type: "string", maxLength: 1280, maxGraphemes: 128 }, + title: { type: "string", maxLength: 5000, maxGraphemes: 500 }, postRef: { type: "ref", ref: "com.atproto.repo.strongRef" }, - description: { type: "string", maxLength: 3000, maxGraphemes: 300 }, + description: { type: "string", maxLength: 30000, maxGraphemes: 3000 }, publishedAt: { type: "string", format: "datetime" }, publication: { type: "string", format: "at-uri" }, author: { type: "string", format: "at-identifier" }, diff --git a/package.json b/package.json index 13be275d..e3f13bed 100644 --- a/package.json +++ b/package.json @@ -4,6 +4,7 @@ "description": "", "main": "index.js", "scripts": { + "lint": "next lint", "dev": "TZ=UTC next dev --turbo", "publish-lexicons": "tsx lexicons/publish.ts", "generate-db-types": "supabase gen types --local > supabase/database.types.ts && drizzle-kit introspect && rm -rf ./drizzle/*.sql ./drizzle/meta",