From 70e6083cd6a8e787527d6f3a1e9feacded1c232b Mon Sep 17 00:00:00 2001 From: Jared Pereira Date: Sat, 24 Jan 2026 23:39:57 -0500 Subject: [PATCH 1/9] add fix incorrect site values function --- app/api/inngest/client.ts | 5 + .../functions/fix_incorrect_site_values.ts | 300 ++++++++++++++++++ app/api/inngest/route.tsx | 2 + 3 files changed, 307 insertions(+) create mode 100644 app/api/inngest/functions/fix_incorrect_site_values.ts diff --git a/app/api/inngest/client.ts b/app/api/inngest/client.ts index 27dfa5ac..f3850470 100644 --- a/app/api/inngest/client.ts +++ b/app/api/inngest/client.ts @@ -41,6 +41,11 @@ export type Events = { documentUris: string[]; }; }; + "documents/fix-incorrect-site-values": { + data: { + did: string; + }; + }; }; // Create a client to send and receive events diff --git a/app/api/inngest/functions/fix_incorrect_site_values.ts b/app/api/inngest/functions/fix_incorrect_site_values.ts new file mode 100644 index 00000000..ac83e091 --- /dev/null +++ b/app/api/inngest/functions/fix_incorrect_site_values.ts @@ -0,0 +1,300 @@ +import { supabaseServerClient } from "supabase/serverClient"; +import { inngest } from "../client"; +import { restoreOAuthSession } from "src/atproto-oauth"; +import { AtpBaseClient, SiteStandardDocument } from "lexicons/api"; +import { AtUri } from "@atproto/syntax"; +import { Json } from "supabase/database.types"; + +async function createAuthenticatedAgent(did: string): Promise { + const result = await restoreOAuthSession(did); + if (!result.ok) { + throw new Error(`Failed to restore OAuth session: ${result.error.message}`); + } + const credentialSession = result.value; + return new AtpBaseClient( + credentialSession.fetchHandler.bind(credentialSession), + ); +} + +/** + * Build set of valid site values for a publication. + * A site value is valid if it matches the publication or its legacy equivalent. + */ +function buildValidSiteValues(pubUri: string): Set { + const validValues = new Set([pubUri]); + + try { + const aturi = new AtUri(pubUri); + + if (pubUri.includes("/site.standard.publication/")) { + // Also accept legacy pub.leaflet.publication + validValues.add( + `at://${aturi.hostname}/pub.leaflet.publication/${aturi.rkey}`, + ); + } else if (pubUri.includes("/pub.leaflet.publication/")) { + // Also accept new site.standard.publication + validValues.add( + `at://${aturi.hostname}/site.standard.publication/${aturi.rkey}`, + ); + } + } catch (e) { + // Invalid URI, just use the original + } + + return validValues; +} + +/** + * This function finds and fixes documents that have incorrect site values. + * A document has an incorrect site value if its `site` field doesn't match + * the publication it belongs to (via documents_in_publications). + * + * Takes a DID as input and processes publications owned by that identity. + */ +export const fix_incorrect_site_values = inngest.createFunction( + { id: "fix_incorrect_site_values" }, + { event: "documents/fix-incorrect-site-values" }, + async ({ event, step }) => { + const { did } = event.data; + + const stats = { + publicationsChecked: 0, + documentsChecked: 0, + documentsWithIncorrectSite: 0, + documentsFixed: 0, + documentsMissingSite: 0, + errors: [] as string[], + }; + + // Step 1: Get all publications owned by this identity + const publications = await step.run("fetch-publications", async () => { + const { data, error } = await supabaseServerClient + .from("publications") + .select("uri") + .eq("identity_did", did); + + if (error) { + throw new Error(`Failed to fetch publications: ${error.message}`); + } + return data || []; + }); + + stats.publicationsChecked = publications.length; + + if (publications.length === 0) { + return { + success: true, + message: "No publications found for this identity", + stats, + }; + } + + // Step 2: Get all documents_in_publications entries for these publications + const publicationUris = publications.map((p) => p.uri); + + const joinEntries = await step.run( + "fetch-documents-in-publications", + async () => { + const { data, error } = await supabaseServerClient + .from("documents_in_publications") + .select("document, publication") + .in("publication", publicationUris); + + if (error) { + throw new Error( + `Failed to fetch documents_in_publications: ${error.message}`, + ); + } + return data || []; + }, + ); + + if (joinEntries.length === 0) { + return { + success: true, + message: "No documents found in publications", + stats, + }; + } + + // Create a map of document URI -> expected publication URI + const documentToPublication = new Map(); + for (const row of joinEntries) { + documentToPublication.set(row.document, row.publication); + } + + // Step 3: Fetch all document records + const documentUris = Array.from(documentToPublication.keys()); + + const allDocuments = await step.run("fetch-documents", async () => { + const { data, error } = await supabaseServerClient + .from("documents") + .select("uri, data") + .in("uri", documentUris); + + if (error) { + throw new Error(`Failed to fetch documents: ${error.message}`); + } + return data || []; + }); + + stats.documentsChecked = allDocuments.length; + + // Step 4: Find documents with incorrect site values + const documentsToFix: Array<{ + uri: string; + currentSite: string | null; + correctSite: string; + docData: SiteStandardDocument.Record; + }> = []; + + for (const doc of allDocuments) { + const expectedPubUri = documentToPublication.get(doc.uri); + if (!expectedPubUri) continue; + + const data = doc.data as unknown as SiteStandardDocument.Record; + const currentSite = data?.site; + + if (!currentSite) { + stats.documentsMissingSite++; + continue; + } + + const validSiteValues = buildValidSiteValues(expectedPubUri); + + if (!validSiteValues.has(currentSite)) { + // Document has incorrect site value - determine the correct one + // Prefer the site.standard.publication format if the doc is site.standard.document + let correctSite = expectedPubUri; + + if (doc.uri.includes("/site.standard.document/")) { + // For site.standard.document, use site.standard.publication format + try { + const pubAturi = new AtUri(expectedPubUri); + if (expectedPubUri.includes("/pub.leaflet.publication/")) { + correctSite = `at://${pubAturi.hostname}/site.standard.publication/${pubAturi.rkey}`; + } + } catch (e) { + // Use as-is + } + } + + documentsToFix.push({ + uri: doc.uri, + currentSite, + correctSite, + docData: data, + }); + } + } + + stats.documentsWithIncorrectSite = documentsToFix.length; + + if (documentsToFix.length === 0) { + return { + success: true, + message: "All documents have correct site values", + stats, + }; + } + + // Step 5: Group documents by author DID for efficient OAuth session handling + const docsByDid = new Map(); + for (const doc of documentsToFix) { + try { + const aturi = new AtUri(doc.uri); + const authorDid = aturi.hostname; + const existing = docsByDid.get(authorDid) || []; + existing.push(doc); + docsByDid.set(authorDid, existing); + } catch (e) { + stats.errors.push(`Invalid URI: ${doc.uri}`); + } + } + + // Step 6: Process each author's documents + for (const [authorDid, docs] of docsByDid) { + // Verify OAuth session for this author + const oauthValid = await step.run( + `verify-oauth-${authorDid.slice(-8)}`, + async () => { + const result = await restoreOAuthSession(authorDid); + return result.ok; + }, + ); + + if (!oauthValid) { + stats.errors.push(`No valid OAuth session for ${authorDid}`); + continue; + } + + // Fix each document for this author + for (const docToFix of docs) { + const result = await step.run( + `fix-doc-${docToFix.uri.slice(-12)}`, + async () => { + try { + const docAturi = new AtUri(docToFix.uri); + + // Build updated record + const updatedRecord: SiteStandardDocument.Record = { + ...docToFix.docData, + site: docToFix.correctSite, + }; + + // Update on PDS + const agent = await createAuthenticatedAgent(authorDid); + await agent.com.atproto.repo.putRecord({ + repo: authorDid, + collection: docAturi.collection, + rkey: docAturi.rkey, + record: updatedRecord, + validate: false, + }); + + // Update in database + const { error: dbError } = await supabaseServerClient + .from("documents") + .update({ data: updatedRecord as Json }) + .eq("uri", docToFix.uri); + + if (dbError) { + return { + success: false as const, + error: `Database update failed: ${dbError.message}`, + }; + } + + return { + success: true as const, + oldSite: docToFix.currentSite, + newSite: docToFix.correctSite, + }; + } catch (e) { + return { + success: false as const, + error: e instanceof Error ? e.message : String(e), + }; + } + }, + ); + + if (result.success) { + stats.documentsFixed++; + } else { + stats.errors.push(`${docToFix.uri}: ${result.error}`); + } + } + } + + return { + success: stats.errors.length === 0, + stats, + documentsToFix: documentsToFix.map((d) => ({ + uri: d.uri, + oldSite: d.currentSite, + newSite: d.correctSite, + })), + }; + }, +); diff --git a/app/api/inngest/route.tsx b/app/api/inngest/route.tsx index b867e731..8699e8b7 100644 --- a/app/api/inngest/route.tsx +++ b/app/api/inngest/route.tsx @@ -6,6 +6,7 @@ import { batched_update_profiles } from "./functions/batched_update_profiles"; import { index_follows } from "./functions/index_follows"; import { migrate_user_to_standard } from "./functions/migrate_user_to_standard"; import { fix_standard_document_publications } from "./functions/fix_standard_document_publications"; +import { fix_incorrect_site_values } from "./functions/fix_incorrect_site_values"; import { cleanup_expired_oauth_sessions, check_oauth_session, @@ -20,6 +21,7 @@ export const { GET, POST, PUT } = serve({ index_follows, migrate_user_to_standard, fix_standard_document_publications, + fix_incorrect_site_values, cleanup_expired_oauth_sessions, check_oauth_session, ], -- 2.51.2 From 1626669b17a97b115facd15c36d347201fc78dc2 Mon Sep 17 00:00:00 2001 From: Jared Pereira Date: Sun, 25 Jan 2026 19:04:51 -0500 Subject: [PATCH 2/9] add sort_data column and use it --- .../dashboard/PublishedPostsLists.tsx | 1 + supabase/database.types.ts | 1 + .../20260125000000_add_sort_date_column.sql | 14 ++++++++++++++ 3 files changed, 16 insertions(+) create mode 100644 supabase/migrations/20260125000000_add_sort_date_column.sql diff --git a/app/lish/[did]/[publication]/dashboard/PublishedPostsLists.tsx b/app/lish/[did]/[publication]/dashboard/PublishedPostsLists.tsx index 69f9735f..b8a24406 100644 --- a/app/lish/[did]/[publication]/dashboard/PublishedPostsLists.tsx +++ b/app/lish/[did]/[publication]/dashboard/PublishedPostsLists.tsx @@ -111,6 +111,7 @@ function PublishedPostItem(props: { documents: { uri: doc.uri, indexed_at: doc.indexed_at, + sort_date: doc.sort_date, data: doc.data, }, }, diff --git a/supabase/database.types.ts b/supabase/database.types.ts index a385b699..ed234f44 100644 --- a/supabase/database.types.ts +++ b/supabase/database.types.ts @@ -337,6 +337,7 @@ export type Database = { Row: { data: Json indexed_at: string + sort_date: string uri: string } Insert: { diff --git a/supabase/migrations/20260125000000_add_sort_date_column.sql b/supabase/migrations/20260125000000_add_sort_date_column.sql new file mode 100644 index 00000000..483fcf2a --- /dev/null +++ b/supabase/migrations/20260125000000_add_sort_date_column.sql @@ -0,0 +1,14 @@ +-- Add sort_date computed column to documents table +-- This column stores the older of publishedAt (from JSON data) or indexed_at +-- Used for sorting feeds chronologically by when content was actually published + +ALTER TABLE documents +ADD COLUMN sort_date timestamptz GENERATED ALWAYS AS ( + LEAST( + COALESCE((data->>'publishedAt')::timestamptz, indexed_at), + indexed_at + ) +) STORED; + +-- Create index on sort_date for efficient ordering +CREATE INDEX documents_sort_date_idx ON documents (sort_date DESC, uri DESC); -- 2.51.2 From 1c680551b603698134ec8a327b7cb243d9775016 Mon Sep 17 00:00:00 2001 From: Jared Pereira Date: Sun, 25 Jan 2026 19:12:35 -0500 Subject: [PATCH 3/9] use sort_date to order --- app/(home-pages)/discover/PubListing.tsx | 2 +- app/(home-pages)/discover/getPublications.ts | 16 ++++++++-------- .../p/[didOrHandle]/getProfilePosts.ts | 10 +++++----- app/(home-pages)/reader/getReaderFeed.ts | 10 +++++----- app/(home-pages)/reader/getSubscriptions.ts | 4 ++-- app/(home-pages)/tag/[tag]/getDocumentsByTag.ts | 4 ++-- app/api/rpc/[command]/get_publication_data.ts | 1 + feeds/index.ts | 6 +++--- 8 files changed, 27 insertions(+), 26 deletions(-) diff --git a/app/(home-pages)/discover/PubListing.tsx b/app/(home-pages)/discover/PubListing.tsx index eff5df20..a899c0f4 100644 --- a/app/(home-pages)/discover/PubListing.tsx +++ b/app/(home-pages)/discover/PubListing.tsx @@ -62,7 +62,7 @@ export const PubListing = (

Updated{" "} {timeAgo( - props.documents_in_publications?.[0]?.documents?.indexed_at || + props.documents_in_publications?.[0]?.documents?.sort_date || "", )}

diff --git a/app/(home-pages)/discover/getPublications.ts b/app/(home-pages)/discover/getPublications.ts index ce649276..aa633c1a 100644 --- a/app/(home-pages)/discover/getPublications.ts +++ b/app/(home-pages)/discover/getPublications.ts @@ -8,7 +8,7 @@ import { import { deduplicateByUri } from "src/utils/deduplicateRecords"; export type Cursor = { - indexed_at?: string; + sort_date?: string; count?: number; uri: string; }; @@ -32,7 +32,7 @@ export async function getPublications( .or( "record->preferences->showInDiscover.is.null,record->preferences->>showInDiscover.eq.true", ) - .order("indexed_at", { + .order("documents(sort_date)", { referencedTable: "documents_in_publications", ascending: false, }) @@ -64,10 +64,10 @@ export async function getPublications( } else { // recentlyUpdated const aDate = new Date( - a.documents_in_publications[0]?.indexed_at || 0, + a.documents_in_publications[0]?.documents?.sort_date || 0, ).getTime(); const bDate = new Date( - b.documents_in_publications[0]?.indexed_at || 0, + b.documents_in_publications[0]?.documents?.sort_date || 0, ).getTime(); if (bDate !== aDate) { return bDate - aDate; @@ -89,11 +89,11 @@ export async function getPublications( (pubCount === cursor.count && pub.uri < cursor.uri) ); } else { - const pubDate = pub.documents_in_publications[0]?.indexed_at || ""; + const pubDate = pub.documents_in_publications[0]?.documents?.sort_date || ""; // Find first pub after cursor return ( - pubDate < (cursor.indexed_at || "") || - (pubDate === cursor.indexed_at && pub.uri < cursor.uri) + pubDate < (cursor.sort_date || "") || + (pubDate === cursor.sort_date && pub.uri < cursor.uri) ); } }); @@ -117,7 +117,7 @@ export async function getPublications( normalizedPage.length > 0 && startIndex + limit < allPubs.length ? order === "recentlyUpdated" ? { - indexed_at: lastItem.documents_in_publications[0]?.indexed_at, + sort_date: lastItem.documents_in_publications[0]?.documents?.sort_date, uri: lastItem.uri, } : { diff --git a/app/(home-pages)/p/[didOrHandle]/getProfilePosts.ts b/app/(home-pages)/p/[didOrHandle]/getProfilePosts.ts index 29f98afb..131b821f 100644 --- a/app/(home-pages)/p/[didOrHandle]/getProfilePosts.ts +++ b/app/(home-pages)/p/[didOrHandle]/getProfilePosts.ts @@ -10,7 +10,7 @@ import { import { deduplicateByUriOrdered } from "src/utils/deduplicateRecords"; export type Cursor = { - indexed_at: string; + sort_date: string; uri: string; }; @@ -29,13 +29,13 @@ export async function getProfilePosts( documents_in_publications(publications(*))`, ) .like("uri", `at://${did}/%`) - .order("indexed_at", { ascending: false }) + .order("sort_date", { ascending: false }) .order("uri", { ascending: false }) .limit(limit); if (cursor) { query = query.or( - `indexed_at.lt.${cursor.indexed_at},and(indexed_at.eq.${cursor.indexed_at},uri.lt.${cursor.uri})`, + `sort_date.lt.${cursor.sort_date},and(sort_date.eq.${cursor.sort_date},uri.lt.${cursor.uri})`, ); } @@ -79,7 +79,7 @@ export async function getProfilePosts( documents: { data: normalizedData, uri: doc.uri, - indexed_at: doc.indexed_at, + sort_date: doc.sort_date, comments_on_documents: doc.comments_on_documents, document_mentions_in_bsky: doc.document_mentions_in_bsky, }, @@ -99,7 +99,7 @@ export async function getProfilePosts( const nextCursor = posts.length === limit ? { - indexed_at: posts[posts.length - 1].documents.indexed_at, + sort_date: posts[posts.length - 1].documents.sort_date, uri: posts[posts.length - 1].documents.uri, } : null; diff --git a/app/(home-pages)/reader/getReaderFeed.ts b/app/(home-pages)/reader/getReaderFeed.ts index 753142c7..d3996eb6 100644 --- a/app/(home-pages)/reader/getReaderFeed.ts +++ b/app/(home-pages)/reader/getReaderFeed.ts @@ -38,12 +38,12 @@ export async function getReaderFeed( "documents_in_publications.publications.publication_subscriptions.identity", auth_res.atp_did, ) - .order("indexed_at", { ascending: false }) + .order("sort_date", { ascending: false }) .order("uri", { ascending: false }) .limit(25); if (cursor) { query = query.or( - `indexed_at.lt.${cursor.timestamp},and(indexed_at.eq.${cursor.timestamp},uri.lt.${cursor.uri})`, + `sort_date.lt.${cursor.timestamp},and(sort_date.eq.${cursor.timestamp},uri.lt.${cursor.uri})`, ); } let { data: rawFeed, error } = await query; @@ -78,7 +78,7 @@ export async function getReaderFeed( document_mentions_in_bsky: post.document_mentions_in_bsky, data: normalizedData, uri: post.uri, - indexed_at: post.indexed_at, + sort_date: post.sort_date, }, }; return p; @@ -88,7 +88,7 @@ export async function getReaderFeed( const nextCursor = posts.length > 0 ? { - timestamp: posts[posts.length - 1].documents.indexed_at, + timestamp: posts[posts.length - 1].documents.sort_date, uri: posts[posts.length - 1].documents.uri, } : null; @@ -109,7 +109,7 @@ export type Post = { documents: { data: NormalizedDocument | null; uri: string; - indexed_at: string; + sort_date: string; comments_on_documents: { count: number }[] | undefined; document_mentions_in_bsky: { count: number }[] | undefined; }; diff --git a/app/(home-pages)/reader/getSubscriptions.ts b/app/(home-pages)/reader/getSubscriptions.ts index 3d2d7765..d1b3e6ff 100644 --- a/app/(home-pages)/reader/getSubscriptions.ts +++ b/app/(home-pages)/reader/getSubscriptions.ts @@ -32,7 +32,7 @@ export async function getSubscriptions( .select(`*, publications(*, documents_in_publications(*, documents(*)))`) .order(`created_at`, { ascending: false }) .order(`uri`, { ascending: false }) - .order("indexed_at", { + .order("documents(sort_date)", { ascending: false, referencedTable: "publications.documents_in_publications", }) @@ -85,6 +85,6 @@ export type PublicationSubscription = { record: NormalizedPublication; uri: string; documents_in_publications: { - documents: { data?: Json; indexed_at: string } | null; + documents: { data?: Json; sort_date: string } | null; }[]; }; diff --git a/app/(home-pages)/tag/[tag]/getDocumentsByTag.ts b/app/(home-pages)/tag/[tag]/getDocumentsByTag.ts index b0e266e8..ebaffbe0 100644 --- a/app/(home-pages)/tag/[tag]/getDocumentsByTag.ts +++ b/app/(home-pages)/tag/[tag]/getDocumentsByTag.ts @@ -24,7 +24,7 @@ export async function getDocumentsByTag( documents_in_publications(publications(*))`, ) .contains("data->tags", `["${tag}"]`) - .order("indexed_at", { ascending: false }) + .order("sort_date", { ascending: false }) .limit(50); if (error) { @@ -69,7 +69,7 @@ export async function getDocumentsByTag( document_mentions_in_bsky: doc.document_mentions_in_bsky, data: normalizedData, uri: doc.uri, - indexed_at: doc.indexed_at, + sort_date: doc.sort_date, }, }; return post; diff --git a/app/api/rpc/[command]/get_publication_data.ts b/app/api/rpc/[command]/get_publication_data.ts index 46a04bad..b33b9203 100644 --- a/app/api/rpc/[command]/get_publication_data.ts +++ b/app/api/rpc/[command]/get_publication_data.ts @@ -83,6 +83,7 @@ export const get_publication_data = makeRoute({ uri: dip.documents.uri, record: normalized, indexed_at: dip.documents.indexed_at, + sort_date: dip.documents.sort_date, data: dip.documents.data, commentsCount: dip.documents.comments_on_documents[0]?.count || 0, mentionsCount: dip.documents.document_mentions_in_bsky[0]?.count || 0, diff --git a/feeds/index.ts b/feeds/index.ts index 564dd576..279f5d7c 100644 --- a/feeds/index.ts +++ b/feeds/index.ts @@ -116,12 +116,12 @@ app.get("/xrpc/app.bsky.feed.getFeedSkeleton", async (c) => { } query = query .or("data->postRef.not.is.null,data->bskyPostRef.not.is.null") - .order("indexed_at", { ascending: false }) + .order("sort_date", { ascending: false }) .order("uri", { ascending: false }) .limit(25); if (parsedCursor) query = query.or( - `indexed_at.lt.${parsedCursor.date},and(indexed_at.eq.${parsedCursor.date},uri.lt.${parsedCursor.uri})`, + `sort_date.lt.${parsedCursor.date},and(sort_date.eq.${parsedCursor.date},uri.lt.${parsedCursor.uri})`, ); let { data, error } = await query; @@ -131,7 +131,7 @@ app.get("/xrpc/app.bsky.feed.getFeedSkeleton", async (c) => { posts = posts || []; let lastPost = posts[posts.length - 1]; - let newCursor = lastPost ? `${lastPost.indexed_at}::${lastPost.uri}` : null; + let newCursor = lastPost ? `${lastPost.sort_date}::${lastPost.uri}` : null; return c.json({ cursor: newCursor || cursor, feed: posts.flatMap((p) => { -- 2.51.2 From e27f0d54501305e0ba5e6b8e3d476dc532cdcbc4 Mon Sep 17 00:00:00 2001 From: Jared Pereira Date: Sun, 25 Jan 2026 19:18:19 -0500 Subject: [PATCH 4/9] make timestamp cast immutable --- supabase/migrations/20260125000000_add_sort_date_column.sql | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/supabase/migrations/20260125000000_add_sort_date_column.sql b/supabase/migrations/20260125000000_add_sort_date_column.sql index 483fcf2a..c79aff4e 100644 --- a/supabase/migrations/20260125000000_add_sort_date_column.sql +++ b/supabase/migrations/20260125000000_add_sort_date_column.sql @@ -2,10 +2,12 @@ -- This column stores the older of publishedAt (from JSON data) or indexed_at -- Used for sorting feeds chronologically by when content was actually published +-- Note: We use ::timestamp AT TIME ZONE 'UTC' to make the expression immutable +-- (direct ::timestamptz cast is not immutable as it depends on session timezone) ALTER TABLE documents ADD COLUMN sort_date timestamptz GENERATED ALWAYS AS ( LEAST( - COALESCE((data->>'publishedAt')::timestamptz, indexed_at), + COALESCE((data->>'publishedAt')::timestamp AT TIME ZONE 'UTC', indexed_at), indexed_at ) ) STORED; -- 2.51.2 From abc790522c4e4a7bbd20d5b07e1ddac3192a5db8 Mon Sep 17 00:00:00 2001 From: Jared Pereira Date: Sun, 25 Jan 2026 20:14:36 -0500 Subject: [PATCH 5/9] add function to get immutable timestamp --- .../20260125000000_add_sort_date_column.sql | 18 +++++++++++++++--- 1 file changed, 15 insertions(+), 3 deletions(-) diff --git a/supabase/migrations/20260125000000_add_sort_date_column.sql b/supabase/migrations/20260125000000_add_sort_date_column.sql index c79aff4e..cb8133d4 100644 --- a/supabase/migrations/20260125000000_add_sort_date_column.sql +++ b/supabase/migrations/20260125000000_add_sort_date_column.sql @@ -2,12 +2,24 @@ -- This column stores the older of publishedAt (from JSON data) or indexed_at -- Used for sorting feeds chronologically by when content was actually published --- Note: We use ::timestamp AT TIME ZONE 'UTC' to make the expression immutable --- (direct ::timestamptz cast is not immutable as it depends on session timezone) +-- Create an immutable function to parse ISO 8601 timestamps from text +-- This is needed because direct ::timestamp cast is not immutable (accepts 'now', 'today', etc.) +-- The regex validates the format before casting to ensure immutability +CREATE OR REPLACE FUNCTION parse_iso_timestamp(text) RETURNS timestamptz +LANGUAGE sql IMMUTABLE STRICT AS $$ + SELECT CASE + -- Match ISO 8601 format: YYYY-MM-DDTHH:MM:SS with optional fractional seconds and Z/timezone + WHEN $1 ~ '^\d{4}-\d{2}-\d{2}[T ]\d{2}:\d{2}:\d{2}(\.\d+)?(Z|[+-]\d{2}:?\d{2})?$' THEN + $1::timestamptz + ELSE + NULL + END +$$; + ALTER TABLE documents ADD COLUMN sort_date timestamptz GENERATED ALWAYS AS ( LEAST( - COALESCE((data->>'publishedAt')::timestamp AT TIME ZONE 'UTC', indexed_at), + COALESCE(parse_iso_timestamp(data->>'publishedAt'), indexed_at), indexed_at ) ) STORED; -- 2.51.2 From 8c91fcf1a4ec892c6e1b43e6bce0d640cc1d0594 Mon Sep 17 00:00:00 2001 From: Jared Pereira Date: Sun, 25 Jan 2026 20:36:03 -0500 Subject: [PATCH 6/9] publish bsky post properly for new records --- app/[leaflet_id]/publish/publishBskyPost.ts | 11 ++++------- 1 file changed, 4 insertions(+), 7 deletions(-) diff --git a/app/[leaflet_id]/publish/publishBskyPost.ts b/app/[leaflet_id]/publish/publishBskyPost.ts index be612396..7edce1f1 100644 --- a/app/[leaflet_id]/publish/publishBskyPost.ts +++ b/app/[leaflet_id]/publish/publishBskyPost.ts @@ -8,11 +8,8 @@ import { import sharp from "sharp"; import { TID } from "@atproto/common"; import { getIdentityData } from "actions/getIdentityData"; -import { AtpBaseClient, PubLeafletDocument } from "lexicons/api"; -import { - restoreOAuthSession, - OAuthSessionError, -} from "src/atproto-oauth"; +import { AtpBaseClient, SiteStandardDocument } from "lexicons/api"; +import { restoreOAuthSession, OAuthSessionError } from "src/atproto-oauth"; import { supabaseServerClient } from "supabase/serverClient"; import { Json } from "supabase/database.types"; import { @@ -30,7 +27,7 @@ export async function publishPostToBsky(args: { url: string; title: string; description: string; - document_record: PubLeafletDocument.Record; + document_record: SiteStandardDocument.Record; rkey: string; facets: AppBskyRichtextFacet.Main[]; }): Promise { @@ -115,7 +112,7 @@ export async function publishPostToBsky(args: { }, ); let record = args.document_record; - record.postRef = post; + record.bskyPostRef = post; let { data: result } = await agent.com.atproto.repo.putRecord({ rkey: args.rkey, -- 2.51.2 From 5626c866b760fd08a73c4be6283d2c2df3df48cf Mon Sep 17 00:00:00 2001 From: Jared Pereira Date: Sun, 25 Jan 2026 20:54:31 -0500 Subject: [PATCH 7/9] add inngest function to fix postref issue --- app/api/inngest/client.ts | 5 + .../fix_standard_document_postref.ts | 196 ++++++++++++++++++ app/api/inngest/route.tsx | 2 + 3 files changed, 203 insertions(+) create mode 100644 app/api/inngest/functions/fix_standard_document_postref.ts diff --git a/app/api/inngest/client.ts b/app/api/inngest/client.ts index f3850470..2c635edb 100644 --- a/app/api/inngest/client.ts +++ b/app/api/inngest/client.ts @@ -46,6 +46,11 @@ export type Events = { did: string; }; }; + "documents/fix-postref": { + data: { + documentUris?: string[]; + }; + }; }; // Create a client to send and receive events diff --git a/app/api/inngest/functions/fix_standard_document_postref.ts b/app/api/inngest/functions/fix_standard_document_postref.ts new file mode 100644 index 00000000..e38cc798 --- /dev/null +++ b/app/api/inngest/functions/fix_standard_document_postref.ts @@ -0,0 +1,196 @@ +import { supabaseServerClient } from "supabase/serverClient"; +import { inngest } from "../client"; +import { restoreOAuthSession } from "src/atproto-oauth"; +import { + AtpBaseClient, + SiteStandardDocument, + ComAtprotoRepoStrongRef, +} from "lexicons/api"; +import { AtUri } from "@atproto/syntax"; +import { Json } from "supabase/database.types"; + +async function createAuthenticatedAgent(did: string): Promise { + const result = await restoreOAuthSession(did); + if (!result.ok) { + throw new Error(`Failed to restore OAuth session: ${result.error.message}`); + } + const credentialSession = result.value; + return new AtpBaseClient( + credentialSession.fetchHandler.bind(credentialSession), + ); +} + +/** + * Fixes site.standard.document records that have the legacy `postRef` field set. + * Migrates the value to `bskyPostRef` (the correct field for site.standard.document) + * and removes the legacy `postRef` field. + * + * Can be triggered with specific document URIs, or will find all affected documents + * if no URIs are provided. + */ +export const fix_standard_document_postref = inngest.createFunction( + { id: "fix_standard_document_postref" }, + { event: "documents/fix-postref" }, + async ({ event, step }) => { + const { documentUris: providedUris } = event.data as { + documentUris?: string[]; + }; + + const stats = { + documentsFound: 0, + documentsFixed: 0, + documentsSkipped: 0, + errors: [] as string[], + }; + + // Step 1: Find documents to fix (either provided or query for them) + const documentUris = await step.run("find-documents", async () => { + if (providedUris && providedUris.length > 0) { + return providedUris; + } + + // Find all site.standard.document records with postRef set + const { data: documents, error } = await supabaseServerClient + .from("documents") + .select("uri") + .like("uri", "at://%/site.standard.document/%") + .not("data->postRef", "is", null); + + if (error) { + throw new Error(`Failed to query documents: ${error.message}`); + } + + return (documents || []).map((d) => d.uri); + }); + + stats.documentsFound = documentUris.length; + + if (documentUris.length === 0) { + return { + success: true, + message: "No documents found with postRef field", + stats, + }; + } + + // Step 2: Group documents by DID for efficient OAuth session handling + const docsByDid = new Map(); + for (const uri of documentUris) { + try { + const aturi = new AtUri(uri); + const did = aturi.hostname; + const existing = docsByDid.get(did) || []; + existing.push(uri); + docsByDid.set(did, existing); + } catch (e) { + stats.errors.push(`Invalid URI: ${uri}`); + } + } + + // Step 3: Process each DID's documents + for (const [did, uris] of docsByDid) { + // Verify OAuth session for this user + const oauthValid = await step.run( + `verify-oauth-${did.slice(-8)}`, + async () => { + const result = await restoreOAuthSession(did); + return result.ok; + }, + ); + + if (!oauthValid) { + stats.errors.push(`No valid OAuth session for ${did}`); + stats.documentsSkipped += uris.length; + continue; + } + + // Fix each document + for (const docUri of uris) { + const result = await step.run( + `fix-doc-${docUri.slice(-12)}`, + async () => { + // Fetch the document + const { data: doc, error: fetchError } = await supabaseServerClient + .from("documents") + .select("uri, data") + .eq("uri", docUri) + .single(); + + if (fetchError || !doc) { + return { + success: false as const, + error: `Document not found: ${fetchError?.message || "no data"}`, + }; + } + + const data = doc.data as Record; + const postRef = data.postRef as + | ComAtprotoRepoStrongRef.Main + | undefined; + + if (!postRef) { + return { + success: false as const, + skipped: true, + error: "Document does not have postRef field", + }; + } + + // Build updated record: move postRef to bskyPostRef + const { postRef: _, ...restData } = data; + let updatedRecord: SiteStandardDocument.Record = { + ...(restData as SiteStandardDocument.Record), + }; + + updatedRecord.bskyPostRef = data.bskyPostRef + ? (data.bskyPostRef as ComAtprotoRepoStrongRef.Main) + : postRef; + + // Write to PDS + const docAturi = new AtUri(docUri); + const agent = await createAuthenticatedAgent(did); + await agent.com.atproto.repo.putRecord({ + repo: did, + collection: "site.standard.document", + rkey: docAturi.rkey, + record: updatedRecord, + validate: false, + }); + + // Update database + const { error: dbError } = await supabaseServerClient + .from("documents") + .update({ data: updatedRecord as Json }) + .eq("uri", docUri); + + if (dbError) { + return { + success: false as const, + error: `Database update failed: ${dbError.message}`, + }; + } + + return { + success: true as const, + postRef, + bskyPostRef: updatedRecord.bskyPostRef, + }; + }, + ); + + if (result.success) { + stats.documentsFixed++; + } else if ("skipped" in result && result.skipped) { + stats.documentsSkipped++; + } else { + stats.errors.push(`${docUri}: ${result.error}`); + } + } + } + + return { + success: stats.errors.length === 0, + stats, + }; + }, +); diff --git a/app/api/inngest/route.tsx b/app/api/inngest/route.tsx index 8699e8b7..7e32e2c1 100644 --- a/app/api/inngest/route.tsx +++ b/app/api/inngest/route.tsx @@ -7,6 +7,7 @@ import { index_follows } from "./functions/index_follows"; import { migrate_user_to_standard } from "./functions/migrate_user_to_standard"; import { fix_standard_document_publications } from "./functions/fix_standard_document_publications"; import { fix_incorrect_site_values } from "./functions/fix_incorrect_site_values"; +import { fix_standard_document_postref } from "./functions/fix_standard_document_postref"; import { cleanup_expired_oauth_sessions, check_oauth_session, @@ -22,6 +23,7 @@ export const { GET, POST, PUT } = serve({ migrate_user_to_standard, fix_standard_document_publications, fix_incorrect_site_values, + fix_standard_document_postref, cleanup_expired_oauth_sessions, check_oauth_session, ], -- 2.51.2 From 5bb63f90313f23191423c4eeda1eb51bed54ee50 Mon Sep 17 00:00:00 2001 From: Jared Pereira Date: Sun, 25 Jan 2026 20:56:45 -0500 Subject: [PATCH 8/9] fix return type of publishtopublication --- actions/publishToPublication.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/actions/publishToPublication.ts b/actions/publishToPublication.ts index cf1b2a82..db3c789e 100644 --- a/actions/publishToPublication.ts +++ b/actions/publishToPublication.ts @@ -65,7 +65,7 @@ import { } from "src/utils/collectionHelpers"; type PublishResult = - | { success: true; rkey: string; record: PubLeafletDocument.Record } + | { success: true; rkey: string; record: SiteStandardDocument.Record } | { success: false; error: OAuthSessionError }; export async function publishToPublication({ -- 2.51.2 From 17f80efe813e141916ca3f02029a32a6b9bc6a59 Mon Sep 17 00:00:00 2001 From: Jared Pereira Date: Mon, 26 Jan 2026 13:38:14 -0500 Subject: [PATCH 9/9] use rkey and fallback to name for pub dashboard urls --- app/lish/createPub/getPublicationURL.ts | 10 +++++++--- 1 file changed, 7 insertions(+), 3 deletions(-) diff --git a/app/lish/createPub/getPublicationURL.ts b/app/lish/createPub/getPublicationURL.ts index da8a77b6..613bef1b 100644 --- a/app/lish/createPub/getPublicationURL.ts +++ b/app/lish/createPub/getPublicationURL.ts @@ -25,7 +25,11 @@ export function getPublicationURL(pub: PublicationInput): string { } // Fall back to checking raw record for legacy base_path - if (isLeafletPublication(pub.record) && pub.record.base_path && isProductionDomain()) { + if ( + isLeafletPublication(pub.record) && + pub.record.base_path && + isProductionDomain() + ) { return `https://${pub.record.base_path}`; } @@ -36,7 +40,7 @@ export function getBasePublicationURL(pub: PublicationInput): string { const normalized = normalizePublicationRecord(pub.record); const aturi = new AtUri(pub.uri); - // Use normalized name if available, fall back to rkey - const name = normalized?.name || aturi.rkey; + //use rkey, fallback to name + const name = aturi.rkey || normalized?.name; return `/lish/${aturi.host}/${encodeURIComponent(name || "")}`; } -- 2.51.2