diff --git a/actions/publications/subscribeEmail.tsx b/actions/publications/subscribeEmail.tsx index 158abc37..554a84cb 100644 --- a/actions/publications/subscribeEmail.tsx +++ b/actions/publications/subscribeEmail.tsx @@ -18,6 +18,7 @@ import { deleteSuppression, } from "src/utils/postmarkSuppressions"; import { + backfillAtprotoSubscriptionsForIdentity, publishAtprotoSubscriptionForDid, unsubscribeToPublication, } from "app/lish/subscribeToPublication"; @@ -360,6 +361,7 @@ async function linkEmailToCurrentIdentity( console.error("[subscribeEmail] attach email failed:", error); return Err("database_error"); } + await backfillAtprotoSubscriptionsForIdentity(current.id, current.atp_did); return Ok(current.id); } diff --git a/app/api/oauth/[route]/route.ts b/app/api/oauth/[route]/route.ts index d5a8c95e..c71fd7d5 100644 --- a/app/api/oauth/[route]/route.ts +++ b/app/api/oauth/[route]/route.ts @@ -1,4 +1,7 @@ -import { subscribeToPublication } from "app/lish/subscribeToPublication"; +import { + backfillAtprotoSubscriptionsForIdentity, + subscribeToPublication, +} from "app/lish/subscribeToPublication"; import { cookies } from "next/headers"; import { redirect } from "next/navigation"; import { NextRequest, NextResponse } from "next/server"; @@ -114,6 +117,10 @@ export async function GET( .from("identities") .update({ atp_did: session.did }) .eq("id", currentIdentity.id); + await backfillAtprotoSubscriptionsForIdentity( + currentIdentity.id, + session.did, + ); return handleAction(s.action, redirectPath); } if (identity.id !== currentIdentity.id) { @@ -148,6 +155,10 @@ export async function GET( .from("identities") .update({ atp_did: session.did }) .eq("id", currentIdentity.id); + await backfillAtprotoSubscriptionsForIdentity( + currentIdentity.id, + session.did, + ); return handleAction(s.action, redirectPath); } const { data } = await supabaseServerClient diff --git a/app/lish/subscribeToPublication.ts b/app/lish/subscribeToPublication.ts index db64f7af..35d71e9b 100644 --- a/app/lish/subscribeToPublication.ts +++ b/app/lish/subscribeToPublication.ts @@ -192,6 +192,28 @@ export async function publishAtprotoSubscriptionForDid( } } +// Publishes atproto subscription records for every confirmed email +// subscription tied to `identityId`. Call after an identity gains an atp_did +// (account-linking, merge) so existing email-only subscriptions also live as +// records on the user's PDS — otherwise they'd be invisible to atproto-side +// consumers and would silently drop if a publication ever leaves email-only +// mode. Best-effort: errors are swallowed by publishAtprotoSubscriptionForDid. +export async function backfillAtprotoSubscriptionsForIdentity( + identityId: string, + atp_did: string, +): Promise { + const { data: subs } = await supabaseServerClient + .from("publication_email_subscribers") + .select("publication") + .eq("identity_id", identityId) + .eq("state", "confirmed"); + if (!subs || subs.length === 0) return; + + await Promise.all( + subs.map((s) => publishAtprotoSubscriptionForDid(atp_did, s.publication)), + ); +} + type UnsubscribeResult = | { success: true } | { success: false; error: OAuthSessionError }; diff --git a/src/mergeIdentity.ts b/src/mergeIdentity.ts index e941a3a9..4e6ad5cc 100644 --- a/src/mergeIdentity.ts +++ b/src/mergeIdentity.ts @@ -11,6 +11,7 @@ import { } from "drizzle/schema"; import { pool } from "supabase/pool"; import { Err, Ok, type Result } from "src/result"; +import { backfillAtprotoSubscriptionsForIdentity } from "app/lish/subscribeToPublication"; export type MergeError = | "merge_not_pending" @@ -34,6 +35,7 @@ export async function mergeEmailIdentityIntoAtpIdentity(args: { if (sourceId === targetId) return Err("same_identity"); const client = await pool.connect(); + let targetAtpDid: string; try { const db = drizzle(client); await db.transaction(async (tx) => { @@ -55,6 +57,7 @@ export async function mergeEmailIdentityIntoAtpIdentity(args: { if (lockedTarget.atp_did === null) throw new Error("merge: target missing atp_did"); if (!lockedSource.email) throw new Error("merge: source missing email"); + targetAtpDid = lockedTarget.atp_did; const sourceEmail = lockedSource.email; @@ -152,5 +155,7 @@ export async function mergeEmailIdentityIntoAtpIdentity(args: { client.release(); } + await backfillAtprotoSubscriptionsForIdentity(targetId, targetAtpDid!); + return Ok(null); }