|
3 | 3 |
|
4 | 4 | from fastapi import APIRouter, Depends, HTTPException, Request, status |
5 | 5 | import stripe |
| 6 | +from sqlalchemy.exc import IntegrityError |
6 | 7 | from sqlalchemy import update, delete, select |
7 | 8 | from sqlalchemy.ext.asyncio import AsyncSession |
8 | 9 | from svix import Webhook, WebhookVerificationError |
|
14 | 15 | from app.models.user import User |
15 | 16 | from app.models.admin_users import AdminUsers |
16 | 17 | from app.services.wallet_service import WalletService |
| 18 | +from app.services.provider_service import create_default_tensorblock_provider_for_user |
17 | 19 |
|
18 | 20 | logger = get_logger(name="webhooks") |
19 | 21 |
|
|
27 | 29 | @router.post("/clerk") |
28 | 30 | async def clerk_webhook_handler(request: Request, db: AsyncSession = Depends(get_async_db)): |
29 | 31 | """ |
30 | | - Handle Clerk webhooks for user events. |
| 32 | + Handle Clerk webhooks for user/organization membership events. |
31 | 33 |
|
32 | 34 | Key events to handle: |
| 35 | + # Organization membership events |
33 | 36 | - organizationMembership.created: Add user to admin users table |
34 | 37 | - organizationMembership.updated: Update user in admin users table |
35 | 38 | - organizationMembership.deleted: Remove user from admin users table |
| 39 | +
|
| 40 | + # User events |
| 41 | + - user.created: Upsert user record |
| 42 | + - user.updated: Upsert user record |
36 | 43 | """ |
37 | 44 | # Get the request body |
38 | 45 | payload = await request.body() |
@@ -71,10 +78,12 @@ async def clerk_webhook_handler(request: Request, db: AsyncSession = Depends(get |
71 | 78 | event_type = event_data.get("type") |
72 | 79 | logger.info(f"Received Clerk webhook: {event_type}") |
73 | 80 |
|
74 | | - if event_type == "organizationMembership.created" or event_type == "organizationMembership.updated": |
| 81 | + if event_type in ["organizationMembership.created", "organizationMembership.updated"]: |
75 | 82 | await handle_organization_membership_created(event_data, db) |
76 | 83 | elif event_type == "organizationMembership.deleted": |
77 | 84 | await handle_organization_membership_deleted(event_data, db) |
| 85 | + elif event_type in ["user.created", "user.updated"]: |
| 86 | + await handle_clerk_user_created(event_data, db) |
78 | 87 | else: |
79 | 88 | logger.warning(f"Unhandled Clerk event type: {event_type}") |
80 | 89 | except json.JSONDecodeError: |
@@ -122,6 +131,55 @@ async def handle_organization_membership_deleted(event_data: dict, db: AsyncSess |
122 | 131 | await db.commit() |
123 | 132 |
|
124 | 133 |
|
| 134 | +async def handle_clerk_user_created(event_data: dict, db: AsyncSession): |
| 135 | + data = event_data['data'] |
| 136 | + clerk_user_id = data['id'] |
| 137 | + |
| 138 | + # extract the primary email address |
| 139 | + if not data.get('primary_email_address_id') or not data.get('email_addresses'): |
| 140 | + logger.error(f"No primary email address or email addresses found for user {clerk_user_id}") |
| 141 | + raise HTTPException(status_code=400, detail="No primary email address or email addresses found for user") |
| 142 | + |
| 143 | + email = None |
| 144 | + primary_email_address_id = data['primary_email_address_id'] |
| 145 | + for email_address in data['email_addresses']: |
| 146 | + if email_address['id'] == primary_email_address_id: |
| 147 | + email = email_address['email_address'] |
| 148 | + break |
| 149 | + |
| 150 | + if not email: |
| 151 | + logger.error(f"No email address found for user {clerk_user_id}") |
| 152 | + raise HTTPException(status_code=400, detail="No email address found for user") |
| 153 | + |
| 154 | + # upsert user record |
| 155 | + try: |
| 156 | + result = await db.execute( |
| 157 | + insert(User).values( |
| 158 | + email=email, |
| 159 | + username=email, # Use email as username |
| 160 | + clerk_user_id=clerk_user_id, |
| 161 | + is_active=True, |
| 162 | + hashed_password="", # Clerk handles authentication |
| 163 | + ).on_conflict_do_update( |
| 164 | + index_elements=[User.clerk_user_id], |
| 165 | + set_=dict( |
| 166 | + email=email, |
| 167 | + username=email, |
| 168 | + is_active=True, |
| 169 | + hashed_password="", # Clerk handles authentication |
| 170 | + ) |
| 171 | + ).returning(User.id) |
| 172 | + ) |
| 173 | + user_id = result.scalar_one() |
| 174 | + await create_default_tensorblock_provider_for_user(user_id, db) |
| 175 | + await db.commit() |
| 176 | + except IntegrityError: |
| 177 | + logger.exception("Error upserting user record for clerk user") |
| 178 | + raise HTTPException(status_code=400, detail="Error upserting user record for clerk user") |
| 179 | + |
| 180 | + logger.info(f"Upserted user record for clerk user {clerk_user_id}/{email}") |
| 181 | + |
| 182 | + |
125 | 183 | @router.post("/stripe") |
126 | 184 | async def stripe_webhook_handler(request: Request, db: AsyncSession = Depends(get_async_db)): |
127 | 185 | """ |
|
0 commit comments