diff --git a/app/lish/[did]/[publication]/[rkey]/Interactions/Comments/commentAction.ts b/app/lish/[did]/[publication]/[rkey]/Interactions/Comments/commentAction.ts index 3b88898d..fbc85963 100644 --- a/app/lish/[did]/[publication]/[rkey]/Interactions/Comments/commentAction.ts +++ b/app/lish/[did]/[publication]/[rkey]/Interactions/Comments/commentAction.ts @@ -9,6 +9,7 @@ import { AtUri, lexToJson, Un$Typed } from "@atproto/api"; import { supabaseServerClient } from "supabase/serverClient"; import { Json } from "supabase/database.types"; import { Notification } from "src/notifications"; +import { v7 } from "uuid"; export async function publishComment(args: { document: string; @@ -68,12 +69,14 @@ export async function publishComment(args: { .select(); let notifications: Notification[] = [ { + id: v7(), recipient: new AtUri(args.document).host, data: { type: "comment", comment_uri: uri.toString() }, }, ]; if (args.comment.replyTo) notifications.push({ + id: v7(), recipient: new AtUri(args.comment.replyTo).host, data: { type: "comment", comment_uri: uri.toString() }, }); diff --git a/drizzle/relations.ts b/drizzle/relations.ts index f717cc4e..ee948385 100644 --- a/drizzle/relations.ts +++ b/drizzle/relations.ts @@ -1,5 +1,5 @@ import { relations } from "drizzle-orm/relations"; -import { identities, publications, documents, comments_on_documents, bsky_profiles, entity_sets, entities, facts, email_auth_tokens, poll_votes_on_entity, permission_tokens, phone_rsvps_to_entity, custom_domains, custom_domain_routes, email_subscriptions_to_entity, atp_poll_records, atp_poll_votes, notifications, bsky_follows, subscribers_to_publications, permission_token_on_homepage, documents_in_publications, document_mentions_in_bsky, bsky_posts, publication_domains, leaflets_in_publications, publication_subscriptions, permission_token_rights } from "./schema"; +import { identities, publications, documents, comments_on_documents, bsky_profiles, entity_sets, entities, facts, email_auth_tokens, poll_votes_on_entity, permission_tokens, phone_rsvps_to_entity, custom_domains, custom_domain_routes, email_subscriptions_to_entity, atp_poll_records, atp_poll_votes, bsky_follows, subscribers_to_publications, permission_token_on_homepage, documents_in_publications, document_mentions_in_bsky, bsky_posts, publication_domains, leaflets_in_publications, publication_subscriptions, permission_token_rights } from "./schema"; export const publicationsRelations = relations(publications, ({one, many}) => ({ identity: one(identities, { @@ -27,7 +27,6 @@ export const identitiesRelations = relations(identities, ({one, many}) => ({ custom_domains_identity_id: many(custom_domains, { relationName: "custom_domains_identity_id_identities_id" }), - notifications: many(notifications), bsky_follows_follows: many(bsky_follows, { relationName: "bsky_follows_follows_identities_atp_did" }), @@ -194,13 +193,6 @@ export const atp_poll_recordsRelations = relations(atp_poll_records, ({many}) => atp_poll_votes: many(atp_poll_votes), })); -export const notificationsRelations = relations(notifications, ({one}) => ({ - identity: one(identities, { - fields: [notifications.recipient], - references: [identities.atp_did] - }), -})); - export const bsky_followsRelations = relations(bsky_follows, ({one}) => ({ identity_follows: one(identities, { fields: [bsky_follows.follows], diff --git a/drizzle/schema.ts b/drizzle/schema.ts index 4213ff2e..4745bf86 100644 --- a/drizzle/schema.ts +++ b/drizzle/schema.ts @@ -165,13 +165,6 @@ export const phone_rsvps_to_entity = pgTable("phone_rsvps_to_entity", { } }); -export const notification = pgTable("notification", { - recipient: text("recipient").primaryKey().notNull(), - created_at: timestamp("created_at", { withTimezone: true, mode: 'string' }).defaultNow().notNull(), - read: boolean("read").default(false).notNull(), - data: jsonb("data").notNull(), -}); - export const custom_domain_routes = pgTable("custom_domain_routes", { id: uuid("id").defaultRandom().primaryKey().notNull(), domain: text("domain").notNull().references(() => custom_domains.domain), @@ -238,13 +231,6 @@ export const oauth_session_store = pgTable("oauth_session_store", { session: jsonb("session").notNull(), }); -export const notifications = pgTable("notifications", { - recipient: text("recipient").primaryKey().notNull().references(() => identities.atp_did, { onDelete: "cascade", onUpdate: "cascade" } ), - created_at: timestamp("created_at", { withTimezone: true, mode: 'string' }).defaultNow().notNull(), - read: boolean("read").default(false).notNull(), - data: jsonb("data").notNull(), -}); - export const bsky_follows = pgTable("bsky_follows", { identity: text("identity").default('').notNull().references(() => identities.atp_did, { onDelete: "cascade" } ), follows: text("follows").notNull().references(() => identities.atp_did, { onDelete: "cascade" } ), diff --git a/src/notifications.ts b/src/notifications.ts new file mode 100644 index 00000000..62958769 --- /dev/null +++ b/src/notifications.ts @@ -0,0 +1,142 @@ +"use server"; + +import { supabaseServerClient } from "supabase/serverClient"; +import { Tables, TablesInsert } from "supabase/database.types"; + +type NotificationRow = Tables<"notifications">; + +export type Notification = Omit, "data"> & { + data: NotificationData; +}; +// Notification data types (for writing to the notifications table) +export type NotificationData = + | { type: "comment"; comment_uri: string } + | { type: "subscribe"; subscription_uri: string }; + +// Hydrated notification types +export type HydratedCommentNotification = { + id: string; + recipient: string; + created_at: string; + type: "comment"; + comment_uri: string; + commentData?: Tables<"comments_on_documents">; +}; + +export type HydratedSubscribeNotification = { + id: string; + recipient: string; + created_at: string; + type: "subscribe"; + subscription_uri: string; + subscriptionData?: Tables<"publication_subscriptions">; +}; + +export type HydratedNotification = + | HydratedCommentNotification + | HydratedSubscribeNotification; + +// Type guard to extract notification type +type ExtractNotificationType = Extract< + NotificationData, + { type: T } +>; + +// Hydrator function type +type NotificationHydrator = ( + notifications: NotificationRow[], +) => Promise>; + +/** + * Hydrates comment notifications + */ +async function hydrateCommentNotifications( + notifications: NotificationRow[], +): Promise { + 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.map((n) => n.data.comment_uri); + const { data: comments } = await supabaseServerClient + .from("comments_on_documents") + .select("*") + .in("uri", commentUris); + + return commentNotifications.map((notification) => ({ + id: notification.id, + recipient: notification.recipient, + created_at: notification.created_at, + type: "comment" as const, + comment_uri: notification.data.comment_uri, + commentData: comments?.find((c) => c.uri === notification.data.comment_uri), + })); +} + +/** + * Hydrates subscribe notifications + */ +async function hydrateSubscribeNotifications( + notifications: NotificationRow[], +): Promise { + 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 + const subscriptionUris = subscribeNotifications.map( + (n) => n.data.subscription_uri, + ); + const { data: subscriptions } = await supabaseServerClient + .from("publication_subscriptions") + .select("*") + .in("uri", subscriptionUris); + + return subscribeNotifications.map((notification) => ({ + id: notification.id, + recipient: notification.recipient, + created_at: notification.created_at, + type: "subscribe" as const, + subscription_uri: notification.data.subscription_uri, + subscriptionData: subscriptions?.find( + (s) => s.uri === notification.data.subscription_uri, + ), + })); +} + +/** + * Main hydration function that processes all notifications + */ +export async function hydrateNotifications( + notifications: NotificationRow[], +): Promise { + // Call all hydrators in parallel + const [commentNotifications, subscribeNotifications] = await Promise.all([ + hydrateCommentNotifications(notifications), + hydrateSubscribeNotifications(notifications), + ]); + + // Combine all hydrated notifications + const allHydrated = [...commentNotifications, ...subscribeNotifications]; + + // 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; +} diff --git a/supabase/database.types.ts b/supabase/database.types.ts index bd919591..ae0c93b0 100644 --- a/supabase/database.types.ts +++ b/supabase/database.types.ts @@ -628,18 +628,21 @@ export type Database = { Row: { created_at: string data: Json + id: string read: boolean recipient: string } Insert: { created_at?: string data: Json + id: string read?: boolean recipient: string } Update: { created_at?: string data?: Json + id?: string read?: boolean recipient?: string } @@ -647,7 +650,7 @@ export type Database = { { foreignKeyName: "notifications_recipient_fkey" columns: ["recipient"] - isOneToOne: true + isOneToOne: false referencedRelation: "identities" referencedColumns: ["atp_did"] }, diff --git a/supabase/migrations/20251030215033_add_notifications_table.sql b/supabase/migrations/20251030215033_add_notifications_table.sql index 5369ad57..3c5c5497 100644 --- a/supabase/migrations/20251030215033_add_notifications_table.sql +++ b/supabase/migrations/20251030215033_add_notifications_table.sql @@ -2,13 +2,14 @@ create table "public"."notifications" ( "recipient" text not null, "created_at" timestamp with time zone not null default now(), "read" boolean not null default false, - "data" jsonb not null + "data" jsonb not null, + "id" uuid not null ); alter table "public"."notifications" enable row level security; -CREATE UNIQUE INDEX notifications_pkey ON public.notifications USING btree (recipient); +CREATE UNIQUE INDEX notifications_pkey ON public.notifications USING btree (id); alter table "public"."notifications" add constraint "notifications_pkey" PRIMARY KEY using index "notifications_pkey";