"""Webhook receiver for gnexus-auth events."""
from datetime import datetime, timezone
import structlog
from fastapi import APIRouter, HTTPException, Request
from gnexus_gauth.exceptions import WebhookPayloadException, WebhookVerificationException
from navi.auth.client import get_gauth_client
from navi.config import settings
log = structlog.get_logger()
router = APIRouter(prefix="/webhooks", tags=["webhooks"])
@router.post("/gnexus-auth")
async def gnexus_auth_webhook(request: Request) -> dict:
"""Receive and handle webhooks from gnexus-auth.
Events handled:
- user.blocked / user.archived / user.deleted → invalidate user sessions
- auth.global_logout → invalidate all sessions
- session.revoked → invalidate matching session
- client.roles_changed / client.permissions_changed → update user role/permissions
- user.profile_updated → update the mirrored navi_users profile
(handbook: profile.update; role/permissions have their own events, and a
local locale override must survive this hook — navi has none yet, the
locale column is a straight mirror, so it is updated together)
"""
if not settings.gnauth_client_id or not settings.gnauth_client_secret:
raise HTTPException(status_code=503, detail="OAuth is not configured")
# Unsigned webhooks let anyone forge session invalidations (global logout,
# user blocked), so in auth mode we refuse to process them at all.
if not settings.gnauth_webhook_secret:
if settings.navi_auth_enabled:
raise HTTPException(status_code=503, detail="Webhook secret is not configured")
log.warning(
"webhook.unverified_mode",
reason="GNAUTH_WEBHOOK_SECRET empty; auth disabled — accepting unsigned webhooks",
)
raw_body = await request.body()
body_text = raw_body.decode("utf-8")
client = get_gauth_client()
if settings.gnauth_webhook_secret:
try:
event = client.verify_and_parse_webhook(
body_text, dict(request.headers), settings.gnauth_webhook_secret
)
except WebhookVerificationException:
raise HTTPException(status_code=403, detail="Invalid webhook signature")
except WebhookPayloadException:
raise HTTPException(status_code=400, detail="Invalid JSON payload")
else:
try:
event = client.parse_webhook(body_text)
except WebhookPayloadException:
raise HTTPException(status_code=400, detail="Invalid JSON payload")
event_type = event.event_type
target = event.target_identifiers
log.info("webhook.received", webhook_event=event_type)
if event_type in ("user.blocked", "user.archived", "user.deleted"):
user_id = target.get("user_id")
if user_id:
await _invalidate_user_sessions(user_id)
log.info("webhook.user_invalidated", user_id=user_id, webhook_event=event_type)
elif event_type == "auth.global_logout":
await _invalidate_all_sessions()
log.info("webhook.all_sessions_invalidated")
elif event_type == "session.revoked":
# gnexus-auth session revoked — we don't have a direct mapping, so invalidate all
# for the user if user_id is present
user_id = target.get("user_id")
if user_id:
await _invalidate_user_sessions(user_id)
log.info("webhook.session_revoked", user_id=user_id)
elif event_type in ("client.roles_changed", "client.permissions_changed"):
# Update affected users' cached role/permissions
user_id = target.get("user_id")
if user_id:
await _update_user_permissions(user_id)
log.info("webhook.permissions_updated", user_id=user_id, webhook_event=event_type)
elif event_type == "user.profile_updated":
profile = (event.metadata or {}).get("profile")
user_id = target.get("user_id") or (event.metadata or {}).get("user", {}).get("id")
if user_id and isinstance(profile, dict):
await _update_user_profile(user_id, profile)
log.info("webhook.profile_updated", user_id=user_id, changed_fields=(event.metadata or {}).get("changed_fields"))
else:
log.warning("webhook.profile_update_unusable", user_id=user_id)
return {"ok": True}
async def _update_user_profile(user_id: str, profile: dict) -> None:
"""Mirror a gnexus-auth profile change into navi_users (role/permissions aside)."""
try:
from navi.api.deps import get_session_store
store = get_session_store()
pool = await store._get_pool()
async with pool.acquire() as conn:
await conn.execute(
"""UPDATE navi_users SET
display_name = COALESCE($2, display_name),
avatar_url = COALESCE($3, avatar_url),
username = COALESCE($4, username),
first_name = COALESCE($5, first_name),
last_name = COALESCE($6, last_name),
phone = COALESCE($7, phone),
birth_date = COALESCE($8, birth_date),
country = COALESCE($9, country),
city = COALESCE($10, city),
locale = COALESCE($11, locale),
updated_at = $1
WHERE id = $12""",
datetime.now(timezone.utc),
profile.get("display_name"),
profile.get("avatar_url"),
profile.get("username"),
profile.get("first_name"),
profile.get("last_name"),
profile.get("phone"),
profile.get("birth_date"),
profile.get("country"),
profile.get("city"),
profile.get("locale"),
user_id,
)
except Exception:
log.warning("webhook.update_profile_failed", user_id=user_id, exc_info=True)
async def _invalidate_user_sessions(user_id: str) -> None:
try:
from navi.api.deps import get_session_store
store = get_session_store()
pool = await store._get_pool()
async with pool.acquire() as conn:
await conn.execute("DELETE FROM user_auth_sessions WHERE user_id = $1", user_id)
except Exception:
log.warning("webhook.invalidate_user_failed", user_id=user_id, exc_info=True)
async def _invalidate_all_sessions() -> None:
try:
from navi.api.deps import get_session_store
store = get_session_store()
pool = await store._get_pool()
async with pool.acquire() as conn:
await conn.execute("DELETE FROM user_auth_sessions")
except Exception:
log.warning("webhook.invalidate_all_failed", exc_info=True)
async def _update_user_permissions(user_id: str) -> None:
"""Clear cached permissions so next request re-fetches from gnexus-auth."""
try:
from navi.api.deps import get_session_store
store = get_session_store()
pool = await store._get_pool()
async with pool.acquire() as conn:
await conn.execute(
"UPDATE navi_users SET permissions = '[]', updated_at = $1 WHERE id = $2",
datetime.now(timezone.utc),
user_id,
)
except Exception:
log.warning("webhook.update_permissions_failed", user_id=user_id, exc_info=True)