diff --git a/app/api/inngest/functions/migrate_user_to_standard.ts b/app/api/inngest/functions/migrate_user_to_standard.ts index cf809b64..fc73f4c6 100644 --- a/app/api/inngest/functions/migrate_user_to_standard.ts +++ b/app/api/inngest/functions/migrate_user_to_standard.ts @@ -44,16 +44,33 @@ export const migrate_user_to_standard = inngest.createFunction( }; // Step 1: Verify OAuth session is valid - await step.run("verify-oauth-session", async () => { + const oauthValid = await step.run("verify-oauth-session", async () => { const result = await restoreOAuthSession(did); if (!result.ok) { - throw new Error( - `Failed to restore OAuth session: ${result.error.message}`, - ); + // Mark identity as needing migration so we can retry later + await supabaseServerClient + .from("identities") + .update({ + metadata: { needsStandardSiteMigration: true }, + }) + .eq("atp_did", did); + + return { success: false, error: result.error.message }; } return { success: true }; }); + if (!oauthValid.success) { + return { + success: false, + error: `Failed to restore OAuth session`, + stats, + publicationUriMap: {}, + documentUriMap: {}, + userSubscriptionUriMap: {}, + }; + } + // Step 2: Get user's pub.leaflet.publication records const oldPublications = await step.run( "fetch-old-publications", @@ -472,6 +489,16 @@ export const migrate_user_to_standard = inngest.createFunction( // 3. The normalization layer handles both schemas transparently for reads // Old records are also kept on the user's PDS so existing AT-URI references remain valid. + // Clear the migration flag on success + if (stats.errors.length === 0) { + await step.run("clear-migration-flag", async () => { + await supabaseServerClient + .from("identities") + .update({ metadata: null }) + .eq("atp_did", did); + }); + } + return { success: stats.errors.length === 0, stats, diff --git a/app/api/oauth/[route]/route.ts b/app/api/oauth/[route]/route.ts index 1cd296e4..903ce627 100644 --- a/app/api/oauth/[route]/route.ts +++ b/app/api/oauth/[route]/route.ts @@ -11,6 +11,7 @@ import { ActionAfterSignIn, parseActionFromSearchParam, } from "./afterSignInActions"; +import { inngest } from "app/api/inngest/client"; type OauthRequestClientState = { redirect: string | null; @@ -84,6 +85,16 @@ export async function GET( .single(); identity = data; } + + // Trigger migration if identity needs it + const metadata = identity?.metadata as Record | null; + if (metadata?.needsStandardSiteMigration) { + await inngest.send({ + name: "user/migrate-to-standard", + data: { did: session.did }, + }); + } + let { data: token } = await supabaseServerClient .from("email_auth_tokens") .insert({ diff --git a/drizzle/schema.ts b/drizzle/schema.ts index 6e1f7f71..65cd9add 100644 --- a/drizzle/schema.ts +++ b/drizzle/schema.ts @@ -140,6 +140,7 @@ export const identities = pgTable("identities", { email: text("email"), atp_did: text("atp_did"), interface_state: jsonb("interface_state"), + metadata: jsonb("metadata"), }, (table) => { return { diff --git a/supabase/database.types.ts b/supabase/database.types.ts index 757213b4..a385b699 100644 --- a/supabase/database.types.ts +++ b/supabase/database.types.ts @@ -551,6 +551,7 @@ export type Database = { home_page: string id: string interface_state: Json | null + metadata: Json | null } Insert: { atp_did?: string | null @@ -559,6 +560,7 @@ export type Database = { home_page?: string id?: string interface_state?: Json | null + metadata?: Json | null } Update: { atp_did?: string | null @@ -567,6 +569,7 @@ export type Database = { home_page?: string id?: string interface_state?: Json | null + metadata?: Json | null } Relationships: [ {