diff --git a/app/api/inngest/client.ts b/app/api/inngest/client.ts index caf5854f..4f082923 100644 --- a/app/api/inngest/client.ts +++ b/app/api/inngest/client.ts @@ -29,6 +29,13 @@ export type Events = { "user/cleanup-expired-oauth-sessions": { data: {}; }; + "user/check-oauth-session": { + data: { + identityId: string; + did: string; + tokenCount: number; + }; + }; }; // Create a client to send and receive events diff --git a/app/api/inngest/functions/cleanup_expired_oauth_sessions.ts b/app/api/inngest/functions/cleanup_expired_oauth_sessions.ts index 3da705a4..486cee15 100644 --- a/app/api/inngest/functions/cleanup_expired_oauth_sessions.ts +++ b/app/api/inngest/functions/cleanup_expired_oauth_sessions.ts @@ -2,19 +2,12 @@ import { supabaseServerClient } from "supabase/serverClient"; import { inngest } from "../client"; import { restoreOAuthSession } from "src/atproto-oauth"; +// Main function that fetches identities and publishes events for each one export const cleanup_expired_oauth_sessions = inngest.createFunction( { id: "cleanup_expired_oauth_sessions" }, { event: "user/cleanup-expired-oauth-sessions" }, async ({ step }) => { - const stats = { - totalIdentities: 0, - validSessions: 0, - expiredSessions: 0, - tokensDeleted: 0, - errors: [] as string[], - }; - - // Step 1: Get all identities with an atp_did (OAuth users) that have at least one auth token + // Get all identities with an atp_did (OAuth users) that have at least one auth token const identities = await step.run("fetch-oauth-identities", async () => { const { data, error } = await supabaseServerClient .from("identities") @@ -33,102 +26,98 @@ export const cleanup_expired_oauth_sessions = inngest.createFunction( }) .map((identity) => ({ id: identity.id, - atp_did: identity.atp_did, + atp_did: identity.atp_did!, tokenCount: identity.email_auth_tokens?.[0]?.count ?? 0, })); }); - stats.totalIdentities = identities.length; - console.log(`Found ${identities.length} OAuth identities with active sessions to check`); + console.log( + `Found ${identities.length} OAuth identities with active sessions to check`, + ); - // Step 2: Check identities' OAuth sessions in batched parallel and cleanup if expired - const BATCH_SIZE = 150; - const allResults: { - identityId: string; - valid: boolean; - tokensDeleted: number; - error?: string; - }[] = []; + // Publish events for each identity in batches + const BATCH_SIZE = 100; + let totalSent = 0; for (let i = 0; i < identities.length; i += BATCH_SIZE) { const batch = identities.slice(i, i + BATCH_SIZE); - const batchNum = Math.floor(i / BATCH_SIZE) + 1; - const totalBatches = Math.ceil(identities.length / BATCH_SIZE); - console.log( - `Processing batch ${batchNum}/${totalBatches} (${batch.length} identities)`, - ); + await step.run(`send-events-batch-${i}`, async () => { + const events = batch.map((identity) => ({ + name: "user/check-oauth-session" as const, + data: { + identityId: identity.id, + did: identity.atp_did, + tokenCount: identity.tokenCount, + }, + })); - const batchResults = await Promise.all( - batch.map((identity) => - step.run(`check-session-${identity.id}`, async () => { - console.log( - `Checking OAuth session for DID: ${identity.atp_did} (${identity.tokenCount} tokens)`, - ); - - const sessionResult = await restoreOAuthSession(identity.atp_did!); - - if (sessionResult.ok) { - console.log(` Session valid for ${identity.atp_did}`); - return { identityId: identity.id, valid: true, tokensDeleted: 0 }; - } - - // Session is expired/invalid - delete associated auth tokens - console.log( - ` Session expired for ${identity.atp_did}: ${sessionResult.error.message}`, - ); - - const { error: deleteError } = await supabaseServerClient - .from("email_auth_tokens") - .delete() - .eq("identity", identity.id); - - if (deleteError) { - console.error( - ` Error deleting tokens for identity ${identity.id}: ${deleteError.message}`, - ); - return { - identityId: identity.id, - valid: false, - tokensDeleted: 0, - error: deleteError.message, - }; - } - - console.log( - ` Deleted ${identity.tokenCount} auth tokens for identity ${identity.id}`, - ); - - return { - identityId: identity.id, - valid: false, - tokensDeleted: identity.tokenCount, - }; - }), - ), - ); + await inngest.send(events); + return events.length; + }); - allResults.push(...batchResults); + totalSent += batch.length; } - // Aggregate results - for (const result of allResults) { - if (result.valid) { - stats.validSessions++; - } else { - stats.expiredSessions++; - stats.tokensDeleted += result.tokensDeleted; - if ("error" in result && result.error) { - stats.errors.push(`Identity ${result.identityId}: ${result.error}`); - } + console.log(`Published ${totalSent} check-oauth-session events`); + + return { + success: true, + identitiesQueued: totalSent, + }; + }, +); + +// Function that checks a single identity's OAuth session and cleans up if expired +export const check_oauth_session = inngest.createFunction( + { id: "check_oauth_session" }, + { event: "user/check-oauth-session" }, + async ({ event, step }) => { + const { identityId, did, tokenCount } = event.data; + + const result = await step.run("check-and-cleanup", async () => { + console.log(`Checking OAuth session for DID: ${did} (${tokenCount} tokens)`); + + const sessionResult = await restoreOAuthSession(did); + + if (sessionResult.ok) { + console.log(` Session valid for ${did}`); + return { valid: true, tokensDeleted: 0 }; } - } - console.log("Cleanup completed:", stats); + // Session is expired/invalid - delete associated auth tokens + console.log( + ` Session expired for ${did}: ${sessionResult.error.message}`, + ); + + const { error: deleteError } = await supabaseServerClient + .from("email_auth_tokens") + .delete() + .eq("identity", identityId); + + if (deleteError) { + console.error( + ` Error deleting tokens for identity ${identityId}: ${deleteError.message}`, + ); + return { + valid: false, + tokensDeleted: 0, + error: deleteError.message, + }; + } + + console.log(` Deleted ${tokenCount} auth tokens for identity ${identityId}`); + + return { + valid: false, + tokensDeleted: tokenCount, + }; + }); return { - success: stats.errors.length === 0, - stats, + identityId, + did, + ...result, }; }, ); diff --git a/app/api/inngest/route.tsx b/app/api/inngest/route.tsx index 37fc24c4..e086dc2c 100644 --- a/app/api/inngest/route.tsx +++ b/app/api/inngest/route.tsx @@ -5,7 +5,10 @@ import { come_online } from "./functions/come_online"; import { batched_update_profiles } from "./functions/batched_update_profiles"; import { index_follows } from "./functions/index_follows"; import { migrate_user_to_standard } from "./functions/migrate_user_to_standard"; -import { cleanup_expired_oauth_sessions } from "./functions/cleanup_expired_oauth_sessions"; +import { + cleanup_expired_oauth_sessions, + check_oauth_session, +} from "./functions/cleanup_expired_oauth_sessions"; export const { GET, POST, PUT } = serve({ client: inngest, @@ -16,5 +19,6 @@ export const { GET, POST, PUT } = serve({ index_follows, migrate_user_to_standard, cleanup_expired_oauth_sessions, + check_oauth_session, ], });