Something went wrong. Try again.
a tool for shared writing and social publishing
Something went wrong. Try again.
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651"use server";
import { supabaseServerClient } from "supabase/serverClient";import { Tables, TablesInsert } from "supabase/database.types";import { AtUri } from "@atproto/syntax";import { idResolver, getProfiles, type Profile } from "src/identity";import { normalizeDocumentRecord, normalizePublicationRecord,} from "src/utils/normalizeRecords";
type NotificationRow = Tables<"notifications">;
export type Notification = Omit<TablesInsert<"notifications">, "data"> & { data: NotificationData;};
type NotificationProfile = Pick< Profile, "did" | "handle" | "displayName" | "avatar">;
function toNotificationProfile( p: Profile | null | undefined,): NotificationProfile | null { if (!p) return null; return { did: p.did, handle: p.handle, displayName: p.displayName, avatar: p.avatar, };}
export type NotificationData = | { type: "comment"; comment_uri: string; parent_uri?: string } | { type: "subscribe"; subscription_uri: string } | { type: "quote"; bsky_post_uri: string; document_uri: string } | { type: "bsky_post_embed"; document_uri: string; bsky_post_uri: string } | { type: "mention"; document_uri: string; mention_type: "did" } | { type: "mention"; document_uri: string; mention_type: "publication"; mentioned_uri: string } | { type: "mention"; document_uri: string; mention_type: "document"; mentioned_uri: string } | { type: "comment_mention"; comment_uri: string; mention_type: "did" } | { type: "comment_mention"; comment_uri: string; mention_type: "publication"; mentioned_uri: string } | { type: "comment_mention"; comment_uri: string; mention_type: "document"; mentioned_uri: string } | { type: "recommend"; document_uri: string; recommend_uri: string };
export type HydratedNotification = | HydratedCommentNotification | HydratedSubscribeNotification | HydratedQuoteNotification | HydratedBskyPostEmbedNotification | HydratedMentionNotification | HydratedCommentMentionNotification | HydratedRecommendNotification;export async function hydrateNotifications( notifications: NotificationRow[],): Promise<Array<HydratedNotification>> { // Call all hydrators in parallel const [commentNotifications, subscribeNotifications, quoteNotifications, bskyPostEmbedNotifications, mentionNotifications, commentMentionNotifications, recommendNotifications] = await Promise.all([ hydrateCommentNotifications(notifications), hydrateSubscribeNotifications(notifications), hydrateQuoteNotifications(notifications), hydrateBskyPostEmbedNotifications(notifications), hydrateMentionNotifications(notifications), hydrateCommentMentionNotifications(notifications), hydrateRecommendNotifications(notifications), ]);
// Combine all hydrated notifications const allHydrated = [...commentNotifications, ...subscribeNotifications, ...quoteNotifications, ...bskyPostEmbedNotifications, ...mentionNotifications, ...commentMentionNotifications, ...recommendNotifications];
// Sort by created_at to maintain order allHydrated.sort( (a, b) => new Date(b.created_at).getTime() - new Date(a.created_at).getTime(), );
return allHydrated;}
// Type guard to extract notification typetype ExtractNotificationType<T extends NotificationData["type"]> = Extract< NotificationData, { type: T }>;
export type HydratedCommentNotification = Awaited< ReturnType<typeof hydrateCommentNotifications>>[0];
async function hydrateCommentNotifications(notifications: NotificationRow[]) { const commentNotifications = notifications.filter( (n): n is NotificationRow & { data: ExtractNotificationType<"comment"> } => (n.data as NotificationData)?.type === "comment", );
if (commentNotifications.length === 0) { return []; }
// Fetch comment data from the database const commentUris = commentNotifications.flatMap((n) => n.data.parent_uri ? [n.data.comment_uri, n.data.parent_uri] : [n.data.comment_uri], ); const { data: comments } = await supabaseServerClient .from("comments_on_documents") .select( "*, documents(*, documents_in_publications(publications(*)))", ) .in("uri", commentUris);
const commenterDids = Array.from( new Set((comments ?? []).map((c) => new AtUri(c.uri).host)), ); const profiles = await getProfiles(commenterDids);
type CommentRow = NonNullable<typeof comments>[number]; const attachProfile = (c: CommentRow) => ({ ...c, profile: toNotificationProfile(profiles.get(new AtUri(c.uri).host)), });
return commentNotifications .map((notification) => { const commentRow = comments?.find((c) => c.uri === notification.data.comment_uri); if (!commentRow) return null; const commentData = attachProfile(commentRow); const parentRow = notification.data.parent_uri ? comments?.find((c) => c.uri === notification.data.parent_uri) : undefined; return { id: notification.id, recipient: notification.recipient, created_at: notification.created_at, type: "comment" as const, comment_uri: notification.data.comment_uri, parentData: parentRow ? attachProfile(parentRow) : undefined, commentData, normalizedDocument: normalizeDocumentRecord(commentData.documents?.data, commentData.documents?.uri), normalizedPublication: normalizePublicationRecord( commentData.documents?.documents_in_publications[0]?.publications?.record, ), }; }) .filter((n) => n !== null);}
export type HydratedSubscribeNotification = Awaited< ReturnType<typeof hydrateSubscribeNotifications>>[0];
async function hydrateSubscribeNotifications(notifications: NotificationRow[]) { const subscribeNotifications = notifications.filter( ( n, ): n is NotificationRow & { data: ExtractNotificationType<"subscribe"> } => (n.data as NotificationData)?.type === "subscribe", );
if (subscribeNotifications.length === 0) { return []; }
// Fetch subscription data from the database with related data const subscriptionUris = subscribeNotifications.map( (n) => n.data.subscription_uri, ); const { data: subscriptions } = await supabaseServerClient .from("publication_subscriptions") .select("*, identities(atp_did), publications(*)") .in("uri", subscriptionUris);
const subscriberDids = Array.from( new Set( (subscriptions ?? []) .map((s) => s.identities?.atp_did) .filter((d): d is string => !!d), ), ); const profiles = await getProfiles(subscriberDids);
return subscribeNotifications .map((notification) => { const subscriptionData = subscriptions?.find((s) => s.uri === notification.data.subscription_uri); if (!subscriptionData) return null; const subscriberDid = subscriptionData.identities?.atp_did ?? null; return { id: notification.id, recipient: notification.recipient, created_at: notification.created_at, type: "subscribe" as const, subscription_uri: notification.data.subscription_uri, subscriptionData: { ...subscriptionData, profile: subscriberDid ? toNotificationProfile(profiles.get(subscriberDid)) : null, }, normalizedPublication: normalizePublicationRecord(subscriptionData.publications?.record), }; }) .filter((n) => n !== null);}
export type HydratedQuoteNotification = Awaited< ReturnType<typeof hydrateQuoteNotifications>>[0];
async function hydrateQuoteNotifications(notifications: NotificationRow[]) { const quoteNotifications = notifications.filter( (n): n is NotificationRow & { data: ExtractNotificationType<"quote"> } => (n.data as NotificationData)?.type === "quote", );
if (quoteNotifications.length === 0) { return []; }
// Fetch bsky post data and document data const bskyPostUris = quoteNotifications.map((n) => n.data.bsky_post_uri); const documentUris = quoteNotifications.map((n) => n.data.document_uri);
const { data: bskyPosts } = await supabaseServerClient .from("bsky_posts") .select("*") .in("uri", bskyPostUris);
const { data: documents } = await supabaseServerClient .from("documents") .select("*, documents_in_publications(publications(*))") .in("uri", documentUris);
return quoteNotifications .map((notification) => { const bskyPost = bskyPosts?.find((p) => p.uri === notification.data.bsky_post_uri); const document = documents?.find((d) => d.uri === notification.data.document_uri); if (!bskyPost || !document) return null; return { id: notification.id, recipient: notification.recipient, created_at: notification.created_at, type: "quote" as const, bsky_post_uri: notification.data.bsky_post_uri, document_uri: notification.data.document_uri, bskyPost, document, normalizedDocument: normalizeDocumentRecord(document.data, document.uri), normalizedPublication: normalizePublicationRecord( document.documents_in_publications[0]?.publications?.record, ), }; }) .filter((n) => n !== null);}
export type HydratedBskyPostEmbedNotification = Awaited< ReturnType<typeof hydrateBskyPostEmbedNotifications>>[0];
async function hydrateBskyPostEmbedNotifications(notifications: NotificationRow[]) { const bskyPostEmbedNotifications = notifications.filter( (n): n is NotificationRow & { data: ExtractNotificationType<"bsky_post_embed"> } => (n.data as NotificationData)?.type === "bsky_post_embed", );
if (bskyPostEmbedNotifications.length === 0) { return []; }
// Fetch document data (the leaflet that embedded the post) const documentUris = bskyPostEmbedNotifications.map((n) => n.data.document_uri); const bskyPostUris = bskyPostEmbedNotifications.map((n) => n.data.bsky_post_uri);
const [{ data: documents }, { data: cachedBskyPosts }] = await Promise.all([ supabaseServerClient .from("documents") .select("*, documents_in_publications(publications(*))") .in("uri", documentUris), supabaseServerClient .from("bsky_posts") .select("*") .in("uri", bskyPostUris), ]);
// Find which posts we need to fetch from the API const cachedPostUris = new Set(cachedBskyPosts?.map((p) => p.uri) ?? []); const missingPostUris = bskyPostUris.filter((uri) => !cachedPostUris.has(uri));
// Fetch missing posts from Bluesky API const fetchedPosts = new Map<string, { text: string } | null>(); if (missingPostUris.length > 0) { try { const { AtpAgent } = await import("@atproto/api"); const agent = new AtpAgent({ service: "https://public.api.bsky.app" }); const response = await agent.app.bsky.feed.getPosts({ uris: missingPostUris }); for (const post of response.data.posts) { const record = post.record as { text?: string }; fetchedPosts.set(post.uri, { text: record.text ?? "" }); } } catch (error) { console.error("Failed to fetch Bluesky posts:", error); } }
// Extract unique DIDs from document URIs to resolve handles const documentCreatorDids = [...new Set(documentUris.map((uri) => new AtUri(uri).host))];
// Resolve DIDs to handles in parallel const didToHandleMap = new Map<string, string | null>(); await Promise.all( documentCreatorDids.map(async (did) => { try { const resolved = await idResolver.did.resolve(did); const handle = resolved?.alsoKnownAs?.[0] ? resolved.alsoKnownAs[0].slice(5) // Remove "at://" prefix : null; didToHandleMap.set(did, handle); } catch (error) { console.error(`Failed to resolve DID ${did}:`, error); didToHandleMap.set(did, null); } }), );
return bskyPostEmbedNotifications .map((notification) => { const document = documents?.find((d) => d.uri === notification.data.document_uri); if (!document) return null;
const documentCreatorDid = new AtUri(notification.data.document_uri).host; const documentCreatorHandle = didToHandleMap.get(documentCreatorDid) ?? null;
// Get post text from cache or fetched data const cachedPost = cachedBskyPosts?.find((p) => p.uri === notification.data.bsky_post_uri); const postView = cachedPost?.post_view as { record?: { text?: string } } | undefined; const bskyPostText = postView?.record?.text ?? fetchedPosts.get(notification.data.bsky_post_uri)?.text ?? null;
return { id: notification.id, recipient: notification.recipient, created_at: notification.created_at, type: "bsky_post_embed" as const, document_uri: notification.data.document_uri, bsky_post_uri: notification.data.bsky_post_uri, document, documentCreatorHandle, bskyPostText, normalizedDocument: normalizeDocumentRecord(document.data, document.uri), normalizedPublication: normalizePublicationRecord( document.documents_in_publications[0]?.publications?.record, ), }; }) .filter((n) => n !== null);}
export type HydratedMentionNotification = Awaited< ReturnType<typeof hydrateMentionNotifications>>[0];
async function hydrateMentionNotifications(notifications: NotificationRow[]) { const mentionNotifications = notifications.filter( (n): n is NotificationRow & { data: ExtractNotificationType<"mention"> } => (n.data as NotificationData)?.type === "mention", );
if (mentionNotifications.length === 0) { return []; }
// Fetch document data from the database const documentUris = mentionNotifications.map((n) => n.data.document_uri); const { data: documents } = await supabaseServerClient .from("documents") .select("*, documents_in_publications(publications(*))") .in("uri", documentUris);
// Extract unique DIDs from document URIs to resolve handles const documentCreatorDids = [...new Set(documentUris.map((uri) => new AtUri(uri).host))];
// Resolve DIDs to handles in parallel const didToHandleMap = new Map<string, string | null>(); await Promise.all( documentCreatorDids.map(async (did) => { try { const resolved = await idResolver.did.resolve(did); const handle = resolved?.alsoKnownAs?.[0] ? resolved.alsoKnownAs[0].slice(5) // Remove "at://" prefix : null; didToHandleMap.set(did, handle); } catch (error) { console.error(`Failed to resolve DID ${did}:`, error); didToHandleMap.set(did, null); } }), );
// Fetch mentioned publications and documents const mentionedPublicationUris = mentionNotifications .filter((n) => n.data.mention_type === "publication") .map((n) => (n.data as Extract<ExtractNotificationType<"mention">, { mention_type: "publication" }>).mentioned_uri);
const mentionedDocumentUris = mentionNotifications .filter((n) => n.data.mention_type === "document") .map((n) => (n.data as Extract<ExtractNotificationType<"mention">, { mention_type: "document" }>).mentioned_uri);
const [{ data: mentionedPublications }, { data: mentionedDocuments }] = await Promise.all([ mentionedPublicationUris.length > 0 ? supabaseServerClient .from("publications") .select("*") .in("uri", mentionedPublicationUris) : Promise.resolve({ data: [] }), mentionedDocumentUris.length > 0 ? supabaseServerClient .from("documents") .select("*, documents_in_publications(publications(*))") .in("uri", mentionedDocumentUris) : Promise.resolve({ data: [] }), ]);
return mentionNotifications .map((notification) => { const document = documents?.find((d) => d.uri === notification.data.document_uri); if (!document) return null;
const mentionedUri = notification.data.mention_type !== "did" ? (notification.data as Extract<ExtractNotificationType<"mention">, { mentioned_uri: string }>).mentioned_uri : undefined;
const documentCreatorDid = new AtUri(notification.data.document_uri).host; const documentCreatorHandle = didToHandleMap.get(documentCreatorDid) ?? null;
const mentionedPublication = mentionedUri ? mentionedPublications?.find((p) => p.uri === mentionedUri) : undefined; const mentionedDoc = mentionedUri ? mentionedDocuments?.find((d) => d.uri === mentionedUri) : undefined;
return { id: notification.id, recipient: notification.recipient, created_at: notification.created_at, type: "mention" as const, document_uri: notification.data.document_uri, mention_type: notification.data.mention_type, mentioned_uri: mentionedUri, document, documentCreatorHandle, mentionedPublication, mentionedDocument: mentionedDoc, normalizedDocument: normalizeDocumentRecord(document.data, document.uri), normalizedPublication: normalizePublicationRecord( document.documents_in_publications[0]?.publications?.record, ), normalizedMentionedPublication: normalizePublicationRecord(mentionedPublication?.record), normalizedMentionedDocument: normalizeDocumentRecord(mentionedDoc?.data, mentionedDoc?.uri), }; }) .filter((n) => n !== null);}
export type HydratedCommentMentionNotification = Awaited< ReturnType<typeof hydrateCommentMentionNotifications>>[0];
async function hydrateCommentMentionNotifications(notifications: NotificationRow[]) { const commentMentionNotifications = notifications.filter( (n): n is NotificationRow & { data: ExtractNotificationType<"comment_mention"> } => (n.data as NotificationData)?.type === "comment_mention", );
if (commentMentionNotifications.length === 0) { return []; }
// Fetch comment data from the database const commentUris = commentMentionNotifications.map((n) => n.data.comment_uri); const { data: comments } = await supabaseServerClient .from("comments_on_documents") .select( "*, documents(*, documents_in_publications(publications(*)))", ) .in("uri", commentUris);
// Extract unique DIDs from comment URIs to resolve handles const commenterDids = [...new Set(commentUris.map((uri) => new AtUri(uri).host))];
// Resolve DIDs to handles in parallel + batch profile fetch const didToHandleMap = new Map<string, string | null>(); const [profiles] = await Promise.all([ getProfiles(commenterDids), Promise.all( commenterDids.map(async (did) => { try { const resolved = await idResolver.did.resolve(did); const handle = resolved?.alsoKnownAs?.[0] ? resolved.alsoKnownAs[0].slice(5) // Remove "at://" prefix : null; didToHandleMap.set(did, handle); } catch (error) { console.error(`Failed to resolve DID ${did}:`, error); didToHandleMap.set(did, null); } }), ), ]);
// Fetch mentioned publications and documents const mentionedPublicationUris = commentMentionNotifications .filter((n) => n.data.mention_type === "publication") .map((n) => (n.data as Extract<ExtractNotificationType<"comment_mention">, { mention_type: "publication" }>).mentioned_uri);
const mentionedDocumentUris = commentMentionNotifications .filter((n) => n.data.mention_type === "document") .map((n) => (n.data as Extract<ExtractNotificationType<"comment_mention">, { mention_type: "document" }>).mentioned_uri);
const [{ data: mentionedPublications }, { data: mentionedDocuments }] = await Promise.all([ mentionedPublicationUris.length > 0 ? supabaseServerClient .from("publications") .select("*") .in("uri", mentionedPublicationUris) : Promise.resolve({ data: [] }), mentionedDocumentUris.length > 0 ? supabaseServerClient .from("documents") .select("*, documents_in_publications(publications(*))") .in("uri", mentionedDocumentUris) : Promise.resolve({ data: [] }), ]);
return commentMentionNotifications .map((notification) => { const commentRow = comments?.find((c) => c.uri === notification.data.comment_uri); if (!commentRow) return null; const commenterDid = new AtUri(commentRow.uri).host; const commentData = { ...commentRow, profile: toNotificationProfile(profiles.get(commenterDid)), };
const mentionedUri = notification.data.mention_type !== "did" ? (notification.data as Extract<ExtractNotificationType<"comment_mention">, { mentioned_uri: string }>).mentioned_uri : undefined;
const commenterHandle = didToHandleMap.get(commenterDid) ?? null;
const mentionedPublication = mentionedUri ? mentionedPublications?.find((p) => p.uri === mentionedUri) : undefined; const mentionedDoc = mentionedUri ? mentionedDocuments?.find((d) => d.uri === mentionedUri) : undefined;
return { id: notification.id, recipient: notification.recipient, created_at: notification.created_at, type: "comment_mention" as const, comment_uri: notification.data.comment_uri, mention_type: notification.data.mention_type, mentioned_uri: mentionedUri, commentData, commenterHandle, mentionedPublication, mentionedDocument: mentionedDoc, normalizedDocument: normalizeDocumentRecord(commentData.documents?.data, commentData.documents?.uri), normalizedPublication: normalizePublicationRecord( commentData.documents?.documents_in_publications[0]?.publications?.record, ), normalizedMentionedPublication: normalizePublicationRecord(mentionedPublication?.record), normalizedMentionedDocument: normalizeDocumentRecord(mentionedDoc?.data, mentionedDoc?.uri), }; }) .filter((n) => n !== null);}
export type HydratedRecommendNotification = Awaited< ReturnType<typeof hydrateRecommendNotifications>>[0];
async function hydrateRecommendNotifications(notifications: NotificationRow[]) { const recommendNotifications = notifications.filter( (n): n is NotificationRow & { data: ExtractNotificationType<"recommend"> } => (n.data as NotificationData)?.type === "recommend", );
if (recommendNotifications.length === 0) { return []; }
// Fetch recommend data from the database const recommendUris = recommendNotifications.map((n) => n.data.recommend_uri); const documentUris = recommendNotifications.map((n) => n.data.document_uri);
const [{ data: recommends }, { data: documents }] = await Promise.all([ supabaseServerClient .from("recommends_on_documents") .select("*") .in("uri", recommendUris), supabaseServerClient .from("documents") .select("*, documents_in_publications(publications(*))") .in("uri", documentUris), ]);
const recommenderDids = Array.from( new Set( (recommends ?? []) .map((r) => r.recommender_did) .filter((d): d is string => !!d), ), ); const profiles = await getProfiles(recommenderDids);
return recommendNotifications .map((notification) => { const recommendRow = recommends?.find((r) => r.uri === notification.data.recommend_uri); const document = documents?.find((d) => d.uri === notification.data.document_uri); if (!recommendRow || !document) return null; const recommendData = { ...recommendRow, profile: recommendRow.recommender_did ? toNotificationProfile(profiles.get(recommendRow.recommender_did)) : null, }; return { id: notification.id, recipient: notification.recipient, created_at: notification.created_at, type: "recommend" as const, recommend_uri: notification.data.recommend_uri, document_uri: notification.data.document_uri, recommendData, document, normalizedDocument: normalizeDocumentRecord(document.data, document.uri), normalizedPublication: normalizePublicationRecord( document.documents_in_publications[0]?.publications?.record, ), }; }) .filter((n) => n !== null);}
export async function pingIdentityToUpdateNotification(did: string) { let channel = supabaseServerClient.channel(`identity.atp_did:${did}`); await channel.send({ type: "broadcast", event: "notification", payload: { message: "poke" }, }); await supabaseServerClient.removeChannel(channel);}