From c55c92646622fe84db5a24a515949dd6071dfa8c Mon Sep 17 00:00:00 2001 From: Jared Pereira Date: Fri, 24 Apr 2026 16:08:22 -0400 Subject: [PATCH] fix login thing and add merge email/atp account flow --- actions/mergeIdentity.ts | 191 ++++++++++++++++++ .../p/[didOrHandle]/getProfilePosts.ts | 4 +- app/api/oauth/[route]/route.ts | 118 +++++++++-- .../dashboard/PublishedPostsLists.tsx | 1 + app/merge-accounts/page.tsx | 159 +++++++++++++++ components/ActionBar/DesktopNavigation.tsx | 2 +- components/ActionBar/MobileNavigation.tsx | 2 +- components/LoginButton.tsx | 72 ++++++- drizzle/relations.ts | 62 ++++-- drizzle/schema.ts | 96 ++++++++- src/auth.ts | 44 +++- supabase/database.types.ts | 20 +- 12 files changed, 724 insertions(+), 47 deletions(-) create mode 100644 actions/mergeIdentity.ts create mode 100644 app/merge-accounts/page.tsx diff --git a/actions/mergeIdentity.ts b/actions/mergeIdentity.ts new file mode 100644 index 00000000..2bb76aa4 --- /dev/null +++ b/actions/mergeIdentity.ts @@ -0,0 +1,191 @@ +"use server"; + +import { cookies } from "next/headers"; +import { drizzle } from "drizzle-orm/node-postgres"; +import { eq, inArray, sql } from "drizzle-orm"; +import { + custom_domains, + email_auth_tokens, + identities, + permission_token_on_homepage, + publication_email_subscribers, + user_entitlements, + user_subscriptions, +} from "drizzle/schema"; +import { pool } from "supabase/pool"; +import { supabaseServerClient } from "supabase/serverClient"; +import { + AUTH_TOKEN_COOKIE, + PENDING_MERGE_TOKEN_COOKIE, + removePendingMergeToken, + resolveAuthToken, + setAuthToken, +} from "src/auth"; +import { Err, Ok, type Result } from "src/result"; + +type MergeError = + | "merge_not_pending" + | "invalid_source" + | "invalid_target" + | "same_identity" + | "database_error"; + +export async function confirmIdentityMerge(): Promise> { + const jar = await cookies(); + const sourceTokenId = jar.get(AUTH_TOKEN_COOKIE)?.value; + const pendingTokenId = jar.get(PENDING_MERGE_TOKEN_COOKIE)?.value; + + const [source, target] = await Promise.all([ + resolveAuthToken(sourceTokenId), + resolveAuthToken(pendingTokenId), + ]); + if (!source || !target) return Err("merge_not_pending"); + + if (source.identity.id === target.identity.id) return Err("same_identity"); + // Source must be an unlinked email account. If it already has atp_did, a + // merge would silently drop that Bluesky link — refuse. + if (source.identity.atp_did || !source.identity.email) + return Err("invalid_source"); + if (!target.identity.atp_did) return Err("invalid_target"); + + const sourceId = source.identity.id; + const targetId = target.identity.id; + const sourceEmail = source.identity.email; + + const client = await pool.connect(); + try { + const db = drizzle(client); + await db.transaction(async (tx) => { + // Re-verify invariants under a row lock. Protects against racing + // concurrent merges or an identity being mutated between cookie + // resolution and this transaction. + const locked = await tx + .select({ + id: identities.id, + email: identities.email, + atp_did: identities.atp_did, + }) + .from(identities) + .where(inArray(identities.id, [sourceId, targetId])) + .for("update"); + const lockedSource = locked.find((r) => r.id === sourceId); + const lockedTarget = locked.find((r) => r.id === targetId); + if (!lockedSource || !lockedTarget) + throw new Error("merge: identity disappeared under lock"); + if (lockedSource.atp_did !== null) + throw new Error("merge: source has atp_did"); + if (lockedTarget.atp_did === null) + throw new Error("merge: target missing atp_did"); + if (!lockedSource.email) throw new Error("merge: source missing email"); + + // email_auth_tokens: caller is about to swap auth_token to the pending + // token (which already points at target). Source's tokens are now stale. + await tx + .delete(email_auth_tokens) + .where(eq(email_auth_tokens.identity, sourceId)); + + // Target wins on (identity_id) PK collision. + await tx.execute(sql` + delete from user_subscriptions + where identity_id = ${sourceId} + and exists (select 1 from user_subscriptions where identity_id = ${targetId}) + `); + await tx + .update(user_subscriptions) + .set({ identity_id: targetId }) + .where(eq(user_subscriptions.identity_id, sourceId)); + + // Target wins on (identity_id, entitlement_key) PK collision. + await tx.execute(sql` + delete from user_entitlements + where identity_id = ${sourceId} + and entitlement_key in ( + select entitlement_key from user_entitlements where identity_id = ${targetId} + ) + `); + await tx + .update(user_entitlements) + .set({ identity_id: targetId }) + .where(eq(user_entitlements.identity_id, sourceId)); + + // Target wins on (token, identity) PK collision. + await tx.execute(sql` + delete from permission_token_on_homepage + where identity = ${sourceId} + and token in ( + select token from permission_token_on_homepage where identity = ${targetId} + ) + `); + await tx + .update(permission_token_on_homepage) + .set({ identity: targetId }) + .where(eq(permission_token_on_homepage.identity, sourceId)); + + // Target wins on unique (publication, email) collision. + await tx.execute(sql` + delete from ${publication_email_subscribers} + where ${publication_email_subscribers.identity_id} = ${sourceId} + and (${publication_email_subscribers.publication}, ${publication_email_subscribers.email}) in ( + select ${publication_email_subscribers.publication}, ${publication_email_subscribers.email} + from ${publication_email_subscribers} + where ${publication_email_subscribers.identity_id} = ${targetId} + ) + `); + await tx + .update(publication_email_subscribers) + .set({ identity_id: targetId }) + .where(eq(publication_email_subscribers.identity_id, sourceId)); + + await tx + .update(custom_domains) + .set({ identity_id: targetId }) + .where(eq(custom_domains.identity_id, sourceId)); + + // identities.email is unique — step via NULL so we don't collide + // mid-swap. custom_domains.identity is nullable and cascades on update; + // we re-set it explicitly after the final value lands because NULL→value + // cascades don't fire. + await tx + .update(identities) + .set({ email: null }) + .where(eq(identities.id, targetId)); + await tx + .update(identities) + .set({ email: null }) + .where(eq(identities.id, sourceId)); + await tx + .update(identities) + .set({ email: sourceEmail }) + .where(eq(identities.id, targetId)); + await tx + .update(custom_domains) + .set({ identity: sourceEmail }) + .where(eq(custom_domains.identity_id, targetId)); + + await tx.delete(identities).where(eq(identities.id, sourceId)); + }); + } catch (e) { + console.error("[mergeIdentity] transaction failed:", e); + return Err("database_error"); + } finally { + client.release(); + } + + await setAuthToken(pendingTokenId!); + await removePendingMergeToken(); + return Ok(null); +} + +export async function cancelIdentityMerge(): Promise> { + const jar = await cookies(); + const pendingTokenId = jar.get(PENDING_MERGE_TOKEN_COOKIE)?.value; + if (pendingTokenId) { + // Drop the token row so the orphan credential can't be used later. + await supabaseServerClient + .from("email_auth_tokens") + .delete() + .eq("id", pendingTokenId); + } + await removePendingMergeToken(); + return Ok(null); +} diff --git a/app/(home-pages)/p/[didOrHandle]/getProfilePosts.ts b/app/(home-pages)/p/[didOrHandle]/getProfilePosts.ts index 07ab25b7..bb469e87 100644 --- a/app/(home-pages)/p/[didOrHandle]/getProfilePosts.ts +++ b/app/(home-pages)/p/[didOrHandle]/getProfilePosts.ts @@ -23,8 +23,8 @@ export async function getProfilePosts( let [{ data: rawFeed, error }, { data: profile }] = await Promise.all([ supabaseServerClient.rpc("get_profile_posts", { p_did: did, - p_cursor_sort_date: cursor?.sort_date ?? null, - p_cursor_uri: cursor?.uri ?? null, + p_cursor_sort_date: cursor?.sort_date ?? undefined, + p_cursor_uri: cursor?.uri ?? undefined, p_limit: limit, }), supabaseServerClient diff --git a/app/api/oauth/[route]/route.ts b/app/api/oauth/[route]/route.ts index 05f16579..aa72fa0c 100644 --- a/app/api/oauth/[route]/route.ts +++ b/app/api/oauth/[route]/route.ts @@ -3,7 +3,11 @@ import { cookies } from "next/headers"; import { redirect } from "next/navigation"; import { NextRequest, NextResponse } from "next/server"; import { createOauthClient } from "src/atproto-oauth"; -import { setAuthToken } from "src/auth"; +import { + AUTH_TOKEN_COOKIE, + setAuthToken, + setPendingMergeToken, +} from "src/auth"; import { supabaseServerClient } from "supabase/serverClient"; import { URLSearchParams } from "url"; @@ -16,6 +20,7 @@ import { inngest } from "app/api/inngest/client"; type OauthRequestClientState = { redirect: string | null; action: ActionAfterSignIn | null; + link?: boolean; }; export async function GET( @@ -33,11 +38,12 @@ export async function GET( const searchParams = req.nextUrl.searchParams; const handle = searchParams.get("handle") as string; const signup = searchParams.get("signup") === "true"; + const link = searchParams.get("link") === "true"; // Put originating page here! let redirect = searchParams.get("redirect_url"); if (redirect) redirect = decodeURIComponent(redirect); let action = parseActionFromSearchParam(searchParams.get("action")); - let state: OauthRequestClientState = { redirect, action }; + let state: OauthRequestClientState = { redirect, action, link }; // Revoke any pending authentication requests if the connection is closed (optional) const ac = new AbortController(); @@ -65,20 +71,53 @@ export async function GET( .select() .eq("atp_did", session.did) .single(); - if (!identity) { - let existingIdentity = (await cookies()).get("auth_token"); - if (existingIdentity) { - let data = await supabaseServerClient - .from("email_auth_tokens") - .select("*, identities(*)") - .eq("id", existingIdentity.value) - .single(); - if (data.data?.identity && data.data.confirmed) - await supabaseServerClient - .from("identities") - .update({ atp_did: session.did }) - .eq("id", data.data.identity); + const currentAuthToken = (await cookies()).get(AUTH_TOKEN_COOKIE)?.value; + let currentIdentity: { + id: string; + email: string | null; + atp_did: string | null; + } | null = null; + if (currentAuthToken) { + const { data: currentTokenRow } = await supabaseServerClient + .from("email_auth_tokens") + .select("confirmed, identities(id, email, atp_did)") + .eq("id", currentAuthToken) + .single(); + if (currentTokenRow?.confirmed) + currentIdentity = currentTokenRow.identities; + } + + // Explicit link flow from the LoginModal for an email-only user. Never + // fall through to a normal DID login — we must either attach the atp_did + // to the existing email identity or route to /merge-accounts. + if ( + s.link && + currentIdentity && + currentIdentity.email && + !currentIdentity.atp_did + ) { + if (!identity) { + await supabaseServerClient + .from("identities") + .update({ atp_did: session.did }) + .eq("id", currentIdentity.id); + return handleAction(s.action, redirectPath); + } + if (identity.id !== currentIdentity.id) { + await stagePendingMerge(identity.id, redirectPath); + // Only reached if the token insert failed. Fall through to the + // normal sign-in flow rather than blocking the user. + } + // Same identity already linked — fall through to refresh session. + } + + if (!identity) { + if (currentIdentity && !currentIdentity.atp_did) { + await supabaseServerClient + .from("identities") + .update({ atp_did: session.did }) + .eq("id", currentIdentity.id); return handleAction(s.action, redirectPath); } const { data } = await supabaseServerClient @@ -87,6 +126,17 @@ export async function GET( .select() .single(); identity = data; + } else if ( + currentIdentity && + currentIdentity.id !== identity.id && + !currentIdentity.atp_did && + currentIdentity.email + ) { + // DID already has an identity row. Caller is currently signed in as a + // *different* email-only identity. Stage a pending merge and let the + // user confirm on /merge-accounts before we touch either account. + await stagePendingMerge(identity.id, redirectPath); + // Only reached if the token insert failed — fall through. } // Trigger migration if identity needs it @@ -117,6 +167,16 @@ export async function GET( console.log("User authenticated as:", session.did); return handleAction(s.action, redirectPath); } catch (e) { + // `redirect()` throws a NEXT_REDIRECT error that Next.js needs to see + // at the framework boundary — don't swallow it. + if ( + e && + typeof e === "object" && + "digest" in e && + typeof (e as { digest?: unknown }).digest === "string" && + (e as { digest: string }).digest.startsWith("NEXT_REDIRECT") + ) + throw e; console.log(e); redirect(redirectPath); } @@ -126,6 +186,34 @@ export async function GET( } } +// Mints a confirmed email_auth_token for `identityId`, stores it as the +// pending_merge_token cookie, and redirects to /merge-accounts. Throws on +// success (via `redirect()`), so callers only reach the line after the call +// when the token insert failed — at which point they should fall through to +// the normal sign-in flow rather than blocking the user. +const stagePendingMerge = async ( + identityId: string, + redirectPath: string, +) => { + const { data: targetToken, error } = await supabaseServerClient + .from("email_auth_tokens") + .insert({ + identity: identityId, + confirmed: true, + confirmation_code: "", + }) + .select("id") + .single(); + if (error) + console.error( + "[oauth/callback] pending merge token insert failed:", + error, + ); + if (!targetToken) return; + await setPendingMergeToken(targetToken.id); + redirect(`/merge-accounts?redirect=${encodeURIComponent(redirectPath)}`); +}; + const handleAction = async ( action: ActionAfterSignIn | null, redirectPath: string, diff --git a/app/lish/[did]/[publication]/dashboard/PublishedPostsLists.tsx b/app/lish/[did]/[publication]/dashboard/PublishedPostsLists.tsx index e6deccc6..e5c4f979 100644 --- a/app/lish/[did]/[publication]/dashboard/PublishedPostsLists.tsx +++ b/app/lish/[did]/[publication]/dashboard/PublishedPostsLists.tsx @@ -113,6 +113,7 @@ function PublishedPostItem(props: { bsky_like_count: doc.bsky_like_count ?? 0, indexed: true, recommend_count: doc.recommendsCount ?? 0, + identity_did: null, }, }, ], diff --git a/app/merge-accounts/page.tsx b/app/merge-accounts/page.tsx new file mode 100644 index 00000000..d45e7269 --- /dev/null +++ b/app/merge-accounts/page.tsx @@ -0,0 +1,159 @@ +import { cookies } from "next/headers"; +import { redirect } from "next/navigation"; +import { supabaseServerClient } from "supabase/serverClient"; +import { ButtonPrimary, ButtonSecondary } from "components/Buttons"; +import { + confirmIdentityMerge, + cancelIdentityMerge, +} from "actions/mergeIdentity"; +import { + AUTH_TOKEN_COOKIE, + PENDING_MERGE_TOKEN_COOKIE, + resolveAuthToken, +} from "src/auth"; + +type SearchParams = { [key: string]: string | string[] | undefined }; + +export default async function MergeAccountsPage(props: { + searchParams: Promise; +}) { + const search = await props.searchParams; + const redirectParam = typeof search.redirect === "string" ? search.redirect : "/"; + + const jar = await cookies(); + const [source, target] = await Promise.all([ + resolveAuthToken(jar.get(AUTH_TOKEN_COOKIE)?.value), + resolveAuthToken(jar.get(PENDING_MERGE_TOKEN_COOKIE)?.value), + ]); + + const mergeValid = + source && + target && + source.identity.id !== target.identity.id && + !source.identity.atp_did && + source.identity.email && + target.identity.atp_did; + + if (!mergeValid) { + return ( +
+
+

Merge no longer pending

+

+ This merge link has expired or is no longer valid. Please try + signing in again. +

+
+
{ + "use server"; + await cancelIdentityMerge(); + redirect(redirectParam); + }} + > + Continue +
+
+
+
+ ); + } + + const [{ data: bsky }, { count: docCount }] = await Promise.all([ + supabaseServerClient + .from("bsky_profiles") + .select("handle") + .eq("did", target.identity.atp_did!) + .maybeSingle(), + supabaseServerClient + .from("permission_token_on_homepage") + .select("token", { count: "exact", head: true }) + .eq("identity", source.identity.id) + .not("archived", "is", true), + ]); + const targetHandle = bsky?.handle ?? target.identity.atp_did!; + const sourceEmail = source.identity.email!; + const targetEmail = target.identity.email; + const documents = docCount ?? 0; + + async function onConfirm() { + "use server"; + const res = await confirmIdentityMerge(); + if (!res.ok) redirect("/merge-accounts?error=" + res.error); + redirect(redirectParam); + } + + async function onCancel() { + "use server"; + await cancelIdentityMerge(); + redirect(redirectParam); + } + + return ( +
+
+

Merge accounts?

+
+

+ You're signed in as{" "} + {sourceEmail}. +

+

+ You just signed in with Bluesky as{" "} + @{targetHandle} + {targetEmail ? ( + <> + , which is already a Leaflet account with email{" "} + {targetEmail} + + ) : ( + <>, which is already a Leaflet account + )} + . +

+

+ Merging moves everything from{" "} + {sourceEmail} into{" "} + @{targetHandle}. Nothing + will be lost. +

+
+ +
+
+ After merging +
+ {documents > 0 && ( +
+ Your {documents} {documents === 1 ? "document" : "documents"}{" "} + will belong to{" "} + @{targetHandle} +
+ )} +
+ You'll sign in to{" "} + @{targetHandle} with{" "} + {sourceEmail} + {targetEmail && ( + <> + {" "} + ({targetEmail} will no + longer work) + + )}{" "} + or Bluesky +
+
+ +
+
+ Cancel +
+
+ Merge accounts +
+
+
+
+ ); +} diff --git a/components/ActionBar/DesktopNavigation.tsx b/components/ActionBar/DesktopNavigation.tsx index 82ea1734..a63cc222 100644 --- a/components/ActionBar/DesktopNavigation.tsx +++ b/components/ActionBar/DesktopNavigation.tsx @@ -29,7 +29,7 @@ export const DesktopNavigation = (props: { return (
- {identity?.atp_did ? ( + {identity ? ( <> )}
- {identity?.atp_did ? ( + {identity ? (
diff --git a/components/LoginButton.tsx b/components/LoginButton.tsx index a2245089..a02c46be 100644 --- a/components/LoginButton.tsx +++ b/components/LoginButton.tsx @@ -62,7 +62,16 @@ export const LoginContent = (props: { let [loading, setLoading] = useState(false); let toaster = useToaster(); - if (identityData.identity) return null; + if (identityData.identity?.atp_did) return null; + if (identityData.identity?.email && !identityData.identity.atp_did) { + return ( + + ); + } const handleEmailSubmit = async () => { setLoading(true); @@ -240,3 +249,64 @@ export const LoginContent = (props: {
); }; + +const LinkAtmosphereContent = (props: { + pageView?: boolean; + redirectRoute?: string; + open?: boolean; +}) => { + let [loading, setLoading] = useState(false); + return ( +
+
+

+ Link your
+ Atmosphere account +

+ + What's the Atmosphere? +
+ } + /> +
+ } + loading={loading} + onSubmit={(handle) => { + setLoading(true); + let redirectUrl: string; + if (props.redirectRoute) { + redirectUrl = props.redirectRoute; + } else { + let url = new URL(window.location.href); + url.searchParams.set("refreshAuth", ""); + redirectUrl = url.toString(); + } + window.location.href = `/api/oauth/login?handle=${encodeURIComponent(handle)}&redirect_url=${encodeURIComponent(redirectUrl)}&link=true`; + }} + /> +
+
+ + + + +
+ + ); +}; diff --git a/drizzle/relations.ts b/drizzle/relations.ts index 87c4981e..38c37f5d 100644 --- a/drizzle/relations.ts +++ b/drizzle/relations.ts @@ -1,5 +1,5 @@ import { relations } from "drizzle-orm/relations"; -import { identities, notifications, publications, documents, comments_on_documents, bsky_profiles, entity_sets, entities, facts, email_auth_tokens, recommends_on_documents, poll_votes_on_entity, permission_tokens, user_subscriptions, phone_rsvps_to_entity, site_standard_publications, custom_domains, custom_domain_routes, site_standard_documents, email_subscriptions_to_entity, atp_poll_records, atp_poll_votes, bsky_follows, subscribers_to_publications, site_standard_documents_in_publications, documents_in_publications, document_mentions_in_bsky, bsky_posts, permission_token_on_homepage, publication_domains, publication_subscriptions, site_standard_subscriptions, user_entitlements, permission_token_rights, leaflets_to_documents, leaflets_in_publications } from "./schema"; +import { identities, notifications, publications, documents, comments_on_documents, bsky_profiles, entity_sets, entities, facts, email_auth_tokens, recommends_on_documents, poll_votes_on_entity, permission_tokens, user_subscriptions, phone_rsvps_to_entity, site_standard_publications, custom_domains, custom_domain_routes, site_standard_documents, email_subscriptions_to_entity, atp_poll_records, atp_poll_votes, publication_newsletter_settings, publication_email_subscribers, publication_email_subscriber_events, bsky_follows, site_standard_documents_in_publications, documents_in_publications, document_mentions_in_bsky, bsky_posts, permission_token_on_homepage, publication_domains, publication_subscriptions, site_standard_subscriptions, user_entitlements, permission_token_rights, publication_post_sends, leaflets_to_documents, leaflets_in_publications } from "./schema"; export const notificationsRelations = relations(notifications, ({one}) => ({ identity: one(identities, { @@ -27,13 +27,13 @@ export const identitiesRelations = relations(identities, ({one, many}) => ({ custom_domains_identity_id: many(custom_domains, { relationName: "custom_domains_identity_id_identities_id" }), + publication_email_subscribers: many(publication_email_subscribers), bsky_follows_follows: many(bsky_follows, { relationName: "bsky_follows_follows_identities_atp_did" }), bsky_follows_identity: many(bsky_follows, { relationName: "bsky_follows_identity_identities_atp_did" }), - subscribers_to_publications: many(subscribers_to_publications), permission_token_on_homepages: many(permission_token_on_homepage), publication_domains: many(publication_domains), publication_subscriptions: many(publication_subscriptions), @@ -46,10 +46,13 @@ export const publicationsRelations = relations(publications, ({one, many}) => ({ fields: [publications.identity_did], references: [identities.atp_did] }), - subscribers_to_publications: many(subscribers_to_publications), + publication_newsletter_settings: many(publication_newsletter_settings), + publication_email_subscribers: many(publication_email_subscribers), + publication_email_subscriber_events: many(publication_email_subscriber_events), documents_in_publications: many(documents_in_publications), publication_domains: many(publication_domains), publication_subscriptions: many(publication_subscriptions), + publication_post_sends: many(publication_post_sends), leaflets_in_publications: many(leaflets_in_publications), })); @@ -69,6 +72,7 @@ export const documentsRelations = relations(documents, ({many}) => ({ recommends_on_documents: many(recommends_on_documents), documents_in_publications: many(documents_in_publications), document_mentions_in_bskies: many(document_mentions_in_bsky), + publication_post_sends: many(publication_post_sends), leaflets_to_documents: many(leaflets_to_documents), leaflets_in_publications: many(leaflets_in_publications), })); @@ -245,6 +249,36 @@ export const atp_poll_recordsRelations = relations(atp_poll_records, ({many}) => atp_poll_votes: many(atp_poll_votes), })); +export const publication_newsletter_settingsRelations = relations(publication_newsletter_settings, ({one}) => ({ + publication: one(publications, { + fields: [publication_newsletter_settings.publication], + references: [publications.uri] + }), +})); + +export const publication_email_subscribersRelations = relations(publication_email_subscribers, ({one, many}) => ({ + identity: one(identities, { + fields: [publication_email_subscribers.identity_id], + references: [identities.id] + }), + publication: one(publications, { + fields: [publication_email_subscribers.publication], + references: [publications.uri] + }), + publication_email_subscriber_events: many(publication_email_subscriber_events), +})); + +export const publication_email_subscriber_eventsRelations = relations(publication_email_subscriber_events, ({one}) => ({ + publication: one(publications, { + fields: [publication_email_subscriber_events.publication], + references: [publications.uri] + }), + publication_email_subscriber: one(publication_email_subscribers, { + fields: [publication_email_subscriber_events.subscriber], + references: [publication_email_subscribers.id] + }), +})); + export const bsky_followsRelations = relations(bsky_follows, ({one}) => ({ identity_follows: one(identities, { fields: [bsky_follows.follows], @@ -258,17 +292,6 @@ export const bsky_followsRelations = relations(bsky_follows, ({one}) => ({ }), })); -export const subscribers_to_publicationsRelations = relations(subscribers_to_publications, ({one}) => ({ - identity: one(identities, { - fields: [subscribers_to_publications.identity], - references: [identities.email] - }), - publication: one(publications, { - fields: [subscribers_to_publications.publication], - references: [publications.uri] - }), -})); - export const site_standard_documents_in_publicationsRelations = relations(site_standard_documents_in_publications, ({one}) => ({ site_standard_document: one(site_standard_documents, { fields: [site_standard_documents_in_publications.document], @@ -372,6 +395,17 @@ export const permission_token_rightsRelations = relations(permission_token_right }), })); +export const publication_post_sendsRelations = relations(publication_post_sends, ({one}) => ({ + document: one(documents, { + fields: [publication_post_sends.document], + references: [documents.uri] + }), + publication: one(publications, { + fields: [publication_post_sends.publication], + references: [publications.uri] + }), +})); + export const leaflets_to_documentsRelations = relations(leaflets_to_documents, ({one}) => ({ document: one(documents, { fields: [leaflets_to_documents.document], diff --git a/drizzle/schema.ts b/drizzle/schema.ts index 237c22bc..38cc380b 100644 --- a/drizzle/schema.ts +++ b/drizzle/schema.ts @@ -1,4 +1,4 @@ -import { pgTable, pgEnum, text, jsonb, foreignKey, timestamp, boolean, uuid, index, bigint, uniqueIndex, unique, smallint, integer, primaryKey } from "drizzle-orm/pg-core" +import { pgTable, pgEnum, text, jsonb, foreignKey, timestamp, boolean, uuid, index, bigint, unique, uniqueIndex, smallint, integer, primaryKey } from "drizzle-orm/pg-core" import { sql } from "drizzle-orm" export const aal_level = pgEnum("aal_level", ['aal1', 'aal2', 'aal3']) @@ -84,6 +84,29 @@ export const facts = pgTable("facts", { } }); +export const mention_services = pgTable("mention_services", { + uri: text("uri").primaryKey().notNull(), + identity_did: text("identity_did").notNull(), + record: jsonb("record").notNull(), +}, +(table) => { + return { + idx_mention_services_did: index("idx_mention_services_did").on(table.identity_did), + } +}); + +export const mention_service_configs = pgTable("mention_service_configs", { + uri: text("uri").primaryKey().notNull(), + identity_did: text("identity_did").notNull(), + record: jsonb("record").notNull(), +}, +(table) => { + return { + idx_mention_service_configs_did: index("idx_mention_service_configs_did").on(table.identity_did), + mention_service_configs_identity_did_key: unique("mention_service_configs_identity_did_key").on(table.identity_did), + } +}); + export const replicache_clients = pgTable("replicache_clients", { client_id: text("client_id").primaryKey().notNull(), client_group: text("client_group").notNull(), @@ -304,24 +327,63 @@ export const oauth_session_store = pgTable("oauth_session_store", { session: jsonb("session").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" } ), +export const publication_newsletter_settings = pgTable("publication_newsletter_settings", { + publication: text("publication").primaryKey().notNull().references(() => publications.uri, { onDelete: "cascade" } ), + enabled: boolean("enabled").default(false).notNull(), + reply_to_email: text("reply_to_email"), + reply_to_verified_at: timestamp("reply_to_verified_at", { withTimezone: true, mode: 'string' }), + created_at: timestamp("created_at", { withTimezone: true, mode: 'string' }).defaultNow().notNull(), + updated_at: timestamp("updated_at", { withTimezone: true, mode: 'string' }).defaultNow().notNull(), + confirmation_code: text("confirmation_code"), }, (table) => { return { - bsky_follows_pkey: primaryKey({ columns: [table.identity, table.follows], name: "bsky_follows_pkey"}), + enabled_idx: index("publication_newsletter_settings_enabled_idx").on(table.publication), } }); -export const subscribers_to_publications = pgTable("subscribers_to_publications", { - identity: text("identity").notNull().references(() => identities.email, { onUpdate: "cascade" } ), - publication: text("publication").notNull().references(() => publications.uri), +export const publication_email_subscribers = pgTable("publication_email_subscribers", { + id: uuid("id").defaultRandom().primaryKey().notNull(), + publication: text("publication").notNull().references(() => publications.uri, { onDelete: "cascade" } ), + email: text("email").notNull(), + identity_id: uuid("identity_id").references(() => identities.id, { onDelete: "set null" } ), + state: text("state").default('pending').notNull(), + confirmation_code: text("confirmation_code"), + unsubscribe_token: uuid("unsubscribe_token").defaultRandom().notNull(), created_at: timestamp("created_at", { withTimezone: true, mode: 'string' }).defaultNow().notNull(), + confirmed_at: timestamp("confirmed_at", { withTimezone: true, mode: 'string' }), + unsubscribed_at: timestamp("unsubscribed_at", { withTimezone: true, mode: 'string' }), +}, +(table) => { + return { + confirmed_idx: index("publication_email_subscribers_confirmed_idx").on(table.publication), + publication_email_subscribers_publication_email_key: unique("publication_email_subscribers_publication_email_key").on(table.publication, table.email), + publication_email_subscribers_unsubscribe_token_key: unique("publication_email_subscribers_unsubscribe_token_key").on(table.unsubscribe_token), + } +}); + +export const publication_email_subscriber_events = pgTable("publication_email_subscriber_events", { + id: uuid("id").defaultRandom().primaryKey().notNull(), + subscriber: uuid("subscriber").notNull().references(() => publication_email_subscribers.id, { onDelete: "cascade" } ), + publication: text("publication").notNull().references(() => publications.uri, { onDelete: "cascade" } ), + event_type: text("event_type").notNull(), + occurred_at: timestamp("occurred_at", { withTimezone: true, mode: 'string' }).defaultNow().notNull(), + metadata: jsonb("metadata"), +}, +(table) => { + return { + subscriber_idx: index("publication_email_subscriber_events_subscriber_idx").on(table.subscriber, table.occurred_at), + publication_type_idx: index("publication_email_subscriber_events_publication_type_idx").on(table.publication, table.event_type, table.occurred_at), + } +}); + +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" } ), }, (table) => { return { - subscribers_to_publications_pkey: primaryKey({ columns: [table.identity, table.publication], name: "subscribers_to_publications_pkey"}), + bsky_follows_pkey: primaryKey({ columns: [table.identity, table.follows], name: "bsky_follows_pkey"}), } }); @@ -449,6 +511,22 @@ export const permission_token_rights = pgTable("permission_token_rights", { } }); +export const publication_post_sends = pgTable("publication_post_sends", { + publication: text("publication").notNull().references(() => publications.uri, { onDelete: "cascade" } ), + document: text("document").notNull().references(() => documents.uri, { onDelete: "cascade" } ), + status: text("status").default('pending').notNull(), + subscriber_count: integer("subscriber_count"), + started_at: timestamp("started_at", { withTimezone: true, mode: 'string' }).defaultNow().notNull(), + completed_at: timestamp("completed_at", { withTimezone: true, mode: 'string' }), + error: text("error"), +}, +(table) => { + return { + publication_started_idx: index("publication_post_sends_publication_started_idx").on(table.publication, table.started_at), + publication_post_sends_pkey: primaryKey({ columns: [table.publication, table.document], name: "publication_post_sends_pkey"}), + } +}); + export const leaflets_to_documents = pgTable("leaflets_to_documents", { leaflet: uuid("leaflet").notNull().references(() => permission_tokens.id, { onDelete: "cascade", onUpdate: "cascade" } ), document: text("document").notNull().references(() => documents.uri, { onDelete: "cascade", onUpdate: "cascade" } ), diff --git a/src/auth.ts b/src/auth.ts index 37a57006..4fea47c9 100644 --- a/src/auth.ts +++ b/src/auth.ts @@ -1,10 +1,17 @@ import { cookies, headers } from "next/headers"; +import { supabaseServerClient } from "supabase/serverClient"; import { isProductionDomain } from "./utils/isProductionDeployment"; +export const AUTH_TOKEN_COOKIE = "auth_token"; +// pending_merge_token carries a freshly-minted email_auth_token for the target +// identity while the user decides whether to merge. Short-lived; confirm turns +// it into the real auth_token, cancel just clears it. +export const PENDING_MERGE_TOKEN_COOKIE = "pending_merge_token"; + export async function setAuthToken(tokenID: string) { let c = await cookies(); let host = (await headers()).get("host"); - c.set("auth_token", tokenID, { + c.set(AUTH_TOKEN_COOKIE, tokenID, { maxAge: 60 * 60 * 24 * 365, secure: process.env.NODE_ENV === "production", domain: isProductionDomain() ? host! : undefined, @@ -13,10 +20,43 @@ export async function setAuthToken(tokenID: string) { }); } +export async function setPendingMergeToken(tokenID: string) { + let c = await cookies(); + let host = (await headers()).get("host"); + c.set(PENDING_MERGE_TOKEN_COOKIE, tokenID, { + maxAge: 60 * 15, + secure: process.env.NODE_ENV === "production", + domain: isProductionDomain() ? host! : undefined, + httpOnly: true, + sameSite: "lax", + }); +} + export async function removeAuthToken() { let c = await cookies(); c.delete({ - name: "auth_token", + name: AUTH_TOKEN_COOKIE, domain: isProductionDomain() ? ".leaflet.pub" : undefined, }); } + +export async function removePendingMergeToken() { + let c = await cookies(); + c.delete({ + name: PENDING_MERGE_TOKEN_COOKIE, + domain: isProductionDomain() ? ".leaflet.pub" : undefined, + }); +} + +// Resolves a cookie-held email_auth_token id to the identity it authenticates. +// Returns null unless the token is confirmed and the identity row exists. +export async function resolveAuthToken(tokenId: string | undefined) { + if (!tokenId) return null; + const { data } = await supabaseServerClient + .from("email_auth_tokens") + .select("id, confirmed, identities(id, email, atp_did)") + .eq("id", tokenId) + .maybeSingle(); + if (!data?.confirmed || !data.identities) return null; + return { tokenId: data.id, identity: data.identities }; +} diff --git a/supabase/database.types.ts b/supabase/database.types.ts index 7ce471d3..60d5bb81 100644 --- a/supabase/database.types.ts +++ b/supabase/database.types.ts @@ -337,6 +337,7 @@ export type Database = { Row: { bsky_like_count: number data: Json + identity_did: string | null indexed: boolean indexed_at: string recommend_count: number @@ -346,6 +347,7 @@ export type Database = { Insert: { bsky_like_count?: number data: Json + identity_did?: string | null indexed?: boolean indexed_at?: string recommend_count?: number @@ -355,6 +357,7 @@ export type Database = { Update: { bsky_like_count?: number data?: Json + identity_did?: string | null indexed?: boolean indexed_at?: string recommend_count?: number @@ -1618,11 +1621,24 @@ export type Database = { like: unknown }[] } + get_leaflet_page_data: { + Args: { + p_token_id: string + } + Returns: { + permission_token: Json + permission_token_rights: Json + leaflets_in_publications: Json + leaflets_to_documents: Json + custom_domain_routes: Json + facts: Json + }[] + } get_profile_posts: { Args: { p_did: string - p_cursor_sort_date?: string | null - p_cursor_uri?: string | null + p_cursor_sort_date?: string + p_cursor_uri?: string p_limit?: number } Returns: { -- 2.51.2