This example demonstrates an authentication provider webhook process using Upstash Workflow. The workflow handles the user creation, trial management, email reminders and notifications.
Use Case#
Our workflow will:
- Receive a webhook event from an authentication provider (e.g. Firebase, Auth0, Clerk etc.)
- Create a new user in our database
- Create a new user in Stripe
- Start a trial in Stripe
- Send a welcome email
- Send a reminder email if the user hasn’t solved any questions in the last 7 days
- Send a trial warning email if the user hasn’t upgraded 2 days before the trial ends
- Send a trial ended email if the user hasn’t upgraded
Code Example#
import { serve } from "@upstash/workflow/nextjs";import { WorkflowContext } from '@upstash/qstash/workflow'/** * This can be the payload of the user created webhook event coming from your * auth provider (e.g. Firebase, Auth0, Clerk etc.) */type UserCreatedPayload = { name: string; email: string;};export const { POST } = serve<UserCreatedPayload>(async (context) => { const { name, email } = context.requestPayload; const { userid } = await context.run("sync user", async () => { return await createUserInDatabase({ name, email }); }); await context.run("create new user in stripe", async () => { await createNewUserInStripe(email); }); await context.run("start trial in Stripe", async () => { await startTrialInStripe(email); }); await context.run("send welcome email", async () => { await sendEmail( email, "Welcome to our platform!, You have 14 days of free trial." ); }); await context.sleep("wait", 7 * 24 * 60 * 60); // get user stats and send email with them const stats = await context.run("get user stats", async () => { return await getUserStats(userid); }); await sendProblemSolvedEmail({context, email, stats}); // wait until there are two days to the end of trial period // and check upgrade status await context.sleep("wait for trial warning", 5 * 24 * 60 * 60); const isUpgraded = await context.run("check upgraded plan", async () => { return await checkUpgradedPlan(email); }); // end the workflow if upgraded if (isUpgraded) return; await context.run("send trial warning email", async () => { await sendEmail( email, "Your trial is about to end in 2 days. Please upgrade your plan to keep using our platform." ); }); await context.sleep("wait for trial end", 2 * 24 * 60 * 60); await context.run("send trial end email", async () => { await sendEmail( email, "Your trial has ended. Please upgrade your plan to keep using our platform." ); });});async function sendProblemSolvedEmail({ context: WorkflowContext<UserCreatedPayload> email: string, stats: { totalProblemsSolved: number }}) { if (stats.totalProblemsSolved === 0) { await context.run("send no answers email", async () => { await sendEmail( email, "Hey, you haven't solved any questions in the last 7 days..." ); }); } else { await context.run("send stats email", async () => { await sendEmail( email, `You have solved ${stats.totalProblemsSolved} problems in the last 7 days. Keep it up!` ); }); }}async function createUserInDatabase({ name, email,}: { name: string; email: string;}) { console.log("Creating a user in the database:", name, email); return { userid: "12345" };}async function createNewUserInStripe(email: string) { // Implement logic to create a new user in Stripe console.log("Creating a user in Stripe for", email);}async function startTrialInStripe(email: string) { // Implement logic to start a trial in Stripe console.log("Starting a trial of 14 days in Stripe for", email);}async function getUserStats(userid: string) { // Implement logic to get user stats console.log("Getting user stats for", userid); return { totalProblemsSolved: 10_000, mostInterestedTopic: "JavaScript", };}async function checkUpgradedPlan(email: string) { // Implement logic to check if the user has upgraded the plan console.log("Checking if the user has upgraded the plan", email); return false;}async function sendEmail(email: string, content: string) { // Implement logic to send an email console.log("Sending email to", email, content);}from fastapi import FastAPIfrom typing import Dict, TypedDictfrom upstash_workflow.fastapi import Servefrom upstash_workflow import AsyncWorkflowContextapp = FastAPI()serve = Serve(app)class UserCreatedPayload(TypedDict): name: str email: strclass UserStats(TypedDict): total_problems_solved: int most_interested_topic: strasync def create_user_in_database(name: str, email: str) -> Dict[str, str]: print("Creating a user in the database:", name, email) return {"userid": "12345"}async def create_new_user_in_stripe(email: str) -> None: # Implement logic to create a new user in Stripe print("Creating a user in Stripe for", email)async def start_trial_in_stripe(email: str) -> None: # Implement logic to start a trial in Stripe print("Starting a trial of 14 days in Stripe for", email)async def get_user_stats(userid: str) -> UserStats: # Implement logic to get user stats print("Getting user stats for", userid) return {"total_problems_solved": 10000, "most_interested_topic": "Python"}async def check_upgraded_plan(email: str) -> bool: # Implement logic to check if the user has upgraded the plan print("Checking if the user has upgraded the plan", email) return Falseasync def send_email(email: str, content: str) -> None: # Implement logic to send an email print("Sending email to", email, content)async def send_problem_solved_email( context: AsyncWorkflowContext[UserCreatedPayload], email: str, stats: UserStats) -> None: if stats["total_problems_solved"] == 0: async def _send_no_answers_email() -> None: await send_email( email, "Hey, you haven't solved any questions in the last 7 days..." ) await context.run("send no answers email", _send_no_answers_email) else: async def _send_stats_email() -> None: await send_email( email, f"You have solved {stats['total_problems_solved']} problems in the last 7 days. Keep it up!", ) await context.run("send stats email", _send_stats_email)@serve.post("/auth-provider-webhook")async def auth_provider_webhook( context: AsyncWorkflowContext[UserCreatedPayload],) -> None: payload = context.request_payload name = payload["name"] email = payload["email"] async def _sync_user() -> str: return await create_user_in_database(name, email) result = await context.run("sync user", _sync_user) userid = result["userid"] async def _create_new_user_in_stripe() -> None: await create_new_user_in_stripe(email) await context.run("create new user in stripe", _create_new_user_in_stripe) async def _start_trial_in_stripe() -> None: await start_trial_in_stripe(email) await context.run("start trial in Stripe", _start_trial_in_stripe) async def _send_welcome_email() -> None: await send_email( email, "Welcome to our platform!, You have 14 days of free trial." ) await context.run("send welcome email", _send_welcome_email) await context.sleep("wait", 7 * 24 * 60 * 60) # get user stats and send email with them async def _get_user_stats() -> UserStats: return await get_user_stats(userid) stats: UserStats = await context.run("get user stats", _get_user_stats) await send_problem_solved_email(context, email, stats) # wait until there are two days to the end of trial period and check upgrade status await context.sleep("wait for trial warning", 5 * 24 * 60 * 60) async def _check_upgraded_plan() -> bool: return await check_upgraded_plan(email) is_upgraded = await context.run("check upgraded plan", _check_upgraded_plan) # end the workflow if upgraded if is_upgraded: return async def _send_trial_warning_email() -> None: await send_email( email, "Your trial is about to end in 2 days. Please upgrade your plan to keep using our platform.", ) await context.run("send trial warning email", _send_trial_warning_email) await context.sleep("wait for trial end", 2 * 24 * 60 * 60) async def _send_trial_end_email() -> None: await send_email( email, "Your trial has ended. Please upgrade your plan to keep using our platform.", ) await context.run("send trial end email", _send_trial_end_email)Code Breakdown#
1. Sync User#
We start by creating a new user in our database:
const { userid } = await context.run("sync user", async () => { return await createUserInDatabase({ name, email });});async def _sync_user() -> str: return await create_user_in_database(name, email)result = await context.run("sync user", _sync_user)userid = result["userid"]2. Create New User in Stripe#
Next, we create a new user in Stripe:
await context.run("create new user in stripe", async () => { await createNewUserInStripe(email);});async def _create_new_user_in_stripe() -> None: await create_new_user_in_stripe(email)await context.run("create new user in stripe", _create_new_user_in_stripe)3. Start Trial in Stripe#
We start a trial in Stripe:
await context.run("start trial in Stripe", async () => { await startTrialInStripe(email);});async def _start_trial_in_stripe() -> None: await start_trial_in_stripe(email)await context.run("start trial in Stripe", _start_trial_in_stripe)4. Send Welcome Email#
We send a welcome email to the user:
await context.run("send welcome email", async () => { await sendEmail( email, "Welcome to our platform!, You have 14 days of free trial." );});async def _send_welcome_email() -> None: await send_email( email, "Welcome to our platform!, You have 14 days of free trial." )await context.run("send welcome email", _send_welcome_email)5. Send Reminder Email#
After 7 days, we check if the user has solved any questions. If not, we send a reminder email:
await context.sleep("wait", 7 * 24 * 60 * 60);const stats = await context.run("get user stats", async () => { return await getUserStats(userid);});await sendProblemSolvedEmail({context, email, stats});await context.sleep("wait", 7 * 24 * 60 * 60)async def _get_user_stats() -> UserStats: return await get_user_stats(userid)stats: UserStats = await context.run("get user stats", _get_user_stats)await send_problem_solved_email(context, email, stats)The sendProblemSolvedEmail method:
async function sendProblemSolvedEmail({ context: WorkflowContext<UserCreatedPayload> email: string, stats: { totalProblemsSolved: number }}) { if (stats.totalProblemsSolved === 0) { await context.run("send no answers email", async () => { await sendEmail( email, "Hey, you haven't solved any questions in the last 7 days..." ); }); } else { await context.run("send stats email", async () => { await sendEmail( email, `You have solved ${stats.totalProblemsSolved} problems in the last 7 days. Keep it up!` ); }); }}async def send_problem_solved_email( context: AsyncWorkflowContext[UserCreatedPayload], email: str, stats: UserStats) -> None: if stats["total_problems_solved"] == 0: async def _send_no_answers_email() -> None: await send_email( email, "Hey, you haven't solved any questions in the last 7 days..." ) await context.run("send no answers email", _send_no_answers_email) else: async def _send_stats_email() -> None: await send_email( email, f"You have solved {stats['total_problems_solved']} problems in the last 7 days. Keep it up!", ) await context.run("send stats email", _send_stats_email)6. Send Trial Warning Email#
If the user hasn’t upgraded 2 days before the trial ends, we send a trial warning email:
await context.sleep("wait for trial warning", 5 * 24 * 60 * 60);const isUpgraded = await context.run("check upgraded plan", async () => { return await checkUpgradedPlan(email);});if (isUpgraded) return;await context.run("send trial warning email", async () => { await sendEmail( email, "Your trial is about to end in 2 days. Please upgrade your plan to keep using our platform." );});await context.sleep("wait for trial warning", 5 * 24 * 60 * 60)async def _check_upgraded_plan() -> bool: return await check_upgraded_plan(email)is_upgraded = await context.run("check upgraded plan", _check_upgraded_plan)# end the workflow if upgradedif is_upgraded: returnasync def _send_trial_warning_email() -> None: await send_email( email, "Your trial is about to end in 2 days. Please upgrade your plan to keep using our platform.", )await context.run("send trial warning email", _send_trial_warning_email)If they upgraded, we end the workflow by returning.
7. Send Trial Ended Email#
If the user hasn’t upgraded after the trial ends, we send a trial ended email:
await context.sleep("wait for trial end", 2 * 24 * 60 * 60);await context.run("send trial end email", async () => { await sendEmail( email, "Your trial has ended. Please upgrade your plan to keep using our platform." );});await context.sleep("wait for trial end", 2 * 24 * 60 * 60)async def _send_trial_end_email() -> None: await send_email( email, "Your trial has ended. Please upgrade your plan to keep using our platform.", )await context.run("send trial end email", _send_trial_end_email)