""" Member Portal — membership middleware for the Inference Cooperative. Reconciles three systems: 1. Open Collective — who is a paying member (billing) 2. Cloudron — who can log in (identity/SSO) 3. LiteLLM — who can use the models, and how much (inference) Three responsibilities: A. Webhook handler — listen for Open Collective membership events B. Key injector — read the member's email from a header, inject their LiteLLM key, and forward the request to the gateway C. Admin endpoints — health, status, manual reconciliation """ import os import json import logging import sqlite3 import secrets import httpx import jwt from fastapi import FastAPI, Request, Response, HTTPException from fastapi.responses import JSONResponse logging.basicConfig(level=logging.INFO) logger = logging.getLogger("member-portal") app = FastAPI(title="Inference Cooperative Member Portal") # --- Configuration (from environment) --- CLOUDRON_API = os.environ.get("CLOUDRON_API_ORIGIN", "https://my.inference.coop") CLOUDRON_TOKEN = os.environ.get("CLOUDRON_TOKEN", "") LITELLM_BASE = os.environ.get("LITELLM_BASE", "https://gateway.inference.coop") LITELLM_MASTER_KEY = os.environ.get("LITELLM_MASTER_KEY", "") OPENCOLLECTIVE_SECRET = os.environ.get("OPENCOLLECTIVE_WEBHOOK_SECRET", "") # Shared secret that the chat frontend uses to authenticate requests to the # portal. LibreChat sends it as an X-Portal-Secret header; OpenWebUI signs the # user's email as a JWT (see FORWARD_USER_INFO_HEADER_JWT_SECRET below). PORTAL_SECRET = os.environ.get("PORTAL_SECRET", "") # OpenWebUI can sign the forwarded user identity as a JWT using this shared # secret. When set, the portal verifies the JWT (HS256) to confirm the email # genuinely came from OpenWebUI (behind Cloudron SSO) rather than a spoofed # header. This is stronger than the plain X-Portal-Secret header because the # email itself is tamper-proof. OWUI_JWT_SECRET = os.environ.get("OWUI_JWT_SECRET", "") # Secret token required in the Open Collective webhook URL path. Open # Collective's generic webhooks are not HMAC-signed, so a secret in the URL # is the standard way to authenticate them. WEBHOOK_TOKEN = os.environ.get("WEBHOOK_TOKEN", "") # Open Collective OAuth app credentials (for the "connect your account" flow, # which is how we obtain a member's email — the webhook strips it for privacy). OC_OAUTH_CLIENT_ID = os.environ.get("OC_OAUTH_CLIENT_ID", "") OC_OAUTH_CLIENT_SECRET = os.environ.get("OC_OAUTH_CLIENT_SECRET", "") OC_OAUTH_REDIRECT_URI = os.environ.get( "OC_OAUTH_REDIRECT_URI", "https://portal.inference.coop/oauth/callback" ) OC_AUTHORIZE_URL = "https://opencollective.com/oauth/authorize" OC_TOKEN_URL = "https://opencollective.com/oauth/token" OC_GRAPHQL_URL = "https://opencollective.com/api/graphql/v2" OC_COLLECTIVE_SLUG = os.environ.get("OC_COLLECTIVE_SLUG", "inference-cooperative") # Monthly credit budget (in USD of tokens) for all members. # Single sliding-scale tier: everyone gets the same $15/month in credits, # regardless of their $10/15/20 contribution. Governance decision (Loomio). MEMBER_BUDGET = float(os.environ.get("MEMBER_BUDGET", "15.0")) # The "members" group in Cloudron (group-based access control). # Members are assigned to this group, which grants access to the chat app. MEMBERS_GROUP_ID = os.environ.get("MEMBERS_GROUP_ID", "") # The "inactive" group in Cloudron. Lapsed members are moved here (instead of # being deleted) so they lose access but can reclaim their data on reactivation. INACTIVE_GROUP_ID = os.environ.get("INACTIVE_GROUP_ID", "") # Loomio integration (forum.inference.coop). The portal reconciles the Loomio # group to the current active-member list via the User API (/api/b2). LOOMIO_BASE = os.environ.get("LOOMIO_BASE", "https://forum.inference.coop") LOOMIO_API_KEY = os.environ.get("LOOMIO_API_KEY", "") LOOMIO_GROUP_ID = os.environ.get("LOOMIO_GROUP_ID", "") # Operator accounts that must always remain in the Loomio group (never removed # by the remove_absent reconciliation), comma-separated. LOOMIO_ALWAYS_KEEP = [ e.strip() for e in os.environ.get("LOOMIO_ALWAYS_KEEP", "info@inference.coop,bot@inference.coop").split(",") if e.strip() ] # Persistent store (SQLite) for email → LiteLLM key token mapping. # Lives in /app/data (Cloudron localstorage addon persists this). DB_PATH = os.environ.get("DB_PATH", "/app/data/members.db") def _verify_portal_secret(request: Request) -> None: """Reject requests that didn't come through LibreChat (shared secret).""" if not PORTAL_SECRET: # If no secret is configured, refuse to inject keys (fail closed). raise HTTPException(503, "Portal secret not configured") provided = request.headers.get("x-portal-secret", "") if not secrets.compare_digest(provided, PORTAL_SECRET): raise HTTPException(401, "Invalid portal secret") def get_db() -> sqlite3.Connection: conn = sqlite3.connect(DB_PATH) conn.execute( "CREATE TABLE IF NOT EXISTS members (" "email TEXT PRIMARY KEY, " "key_token TEXT, " "cloudron_user_id TEXT, " "slug TEXT, " "active INTEGER DEFAULT 1" ")" ) conn.execute( "CREATE TABLE IF NOT EXISTS pending_members (" "slug TEXT PRIMARY KEY, " "name TEXT, " "created_at TEXT DEFAULT (datetime('now'))" ")" ) # Migration: add slug column if the members table predates it. cols = [r[1] for r in conn.execute("PRAGMA table_info(members)").fetchall()] if "slug" not in cols: conn.execute("ALTER TABLE members ADD COLUMN slug TEXT") return conn def store_pending_member(slug: str, name: str) -> None: conn = get_db() conn.execute( "INSERT INTO pending_members (slug, name) VALUES (?, ?) " "ON CONFLICT(slug) DO UPDATE SET name=excluded.name", (slug, name), ) conn.commit() conn.close() def is_pending_member(slug: str) -> bool: """True if this slug has a pending membership (i.e. the webhook saw a contribution from them). Used to gate the OAuth flow so only paying members can provision an account.""" if not slug: return False conn = get_db() row = conn.execute( "SELECT 1 FROM pending_members WHERE slug = ?", (slug,) ).fetchone() conn.close() return row is not None async def is_active_member(slug: str) -> bool: """Check the LIVE Open Collective membership list (authoritative source of truth) for whether this slug is a financial contributor (BACKER/ADMIN), not just the fiscal host. More reliable than the webhook, which only fires on new events and can miss existing members.""" if not slug: return False query = ( '{ collective(slug: "%s") { members(limit: 100) { nodes { role account { slug } } } } }' % OC_COLLECTIVE_SLUG ) async with httpx.AsyncClient() as client: r = await client.post(OC_GRAPHQL_URL, json={"query": query}) if r.status_code != 200: logger.warning("OC membership check failed: %s", r.status_code) return False nodes = r.json().get("data", {}).get("collective", {}).get("members", {}).get("nodes", []) for n in nodes: role = n.get("role", "") acct = n.get("account", {}) or {} if acct.get("slug") == slug and role in ("BACKER", "ADMIN"): return True return False async def get_active_member_emails() -> list[str]: """Return the list of active financial contributor emails from Open Collective.""" query = ( '{ collective(slug: "%s") { members(limit: 100) { nodes { role account { email slug } } } } }' % OC_COLLECTIVE_SLUG ) async with httpx.AsyncClient() as client: r = await client.post(OC_GRAPHQL_URL, json={"query": query}) if r.status_code != 200: logger.warning("OC member list fetch failed: %s", r.status_code) return [] nodes = r.json().get("data", {}).get("collective", {}).get("members", {}).get("nodes", []) emails = [] for n in nodes: role = n.get("role", "") acct = n.get("account", {}) or {} email = acct.get("email") if role in ("BACKER", "ADMIN") and email: emails.append(email) return emails async def sync_loomio_memberships() -> None: """Reconcile the Loomio group to the current active-member list. Uses POST /api/b2/memberships with remove_absent=1, which adds new members and removes anyone not in the list. Operator accounts (LOOMIO_ALWAYS_KEEP) are always included so they are never removed. """ if not LOOMIO_API_KEY or not LOOMIO_GROUP_ID: logger.warning("Loomio not configured; skipping membership sync") return emails = await get_active_member_emails() # Always keep operator accounts in the group. for keep in LOOMIO_ALWAYS_KEEP: if keep not in emails: emails.append(keep) async with httpx.AsyncClient() as client: r = await client.post( f"{LOOMIO_BASE}/api/b2/memberships", headers={"Authorization": f"Bearer {LOOMIO_API_KEY}"}, json={"group_id": int(LOOMIO_GROUP_ID), "emails": emails, "remove_absent": 1}, ) if r.status_code != 200: logger.warning("Loomio membership sync failed: %s %s", r.status_code, r.text[:200]) return result = r.json() logger.info( "Loomio sync: added=%s removed=%s", result.get("added_emails", []), result.get("removed_emails", []), ) def store_member(email: str, key_token: str, cloudron_user_id: str, slug: str = "") -> None: conn = get_db() conn.execute( "INSERT INTO members (email, key_token, cloudron_user_id, slug, active) " "VALUES (?, ?, ?, ?, 1) " "ON CONFLICT(email) DO UPDATE SET key_token=excluded.key_token, " "cloudron_user_id=excluded.cloudron_user_id, slug=excluded.slug, active=1", (email, key_token, cloudron_user_id, slug), ) conn.commit() conn.close() def get_member_key(email: str) -> str | None: conn = get_db() row = conn.execute( "SELECT key_token FROM members WHERE email = ? AND active = 1", (email,) ).fetchone() conn.close() return row[0] if row else None def get_member_by_slug(slug: str) -> dict | None: """Look up a member (email, key_token, cloudron_user_id) by their OC slug.""" conn = get_db() row = conn.execute( "SELECT email, key_token, cloudron_user_id FROM members WHERE slug = ?", (slug,) ).fetchone() conn.close() if not row: return None return {"email": row[0], "key_token": row[1], "cloudron_user_id": row[2]} def deactivate_member(email: str) -> str | None: """Mark a member inactive and return their key token (for deletion).""" conn = get_db() row = conn.execute( "SELECT key_token FROM members WHERE email = ?", (email,) ).fetchone() conn.execute("UPDATE members SET active = 0 WHERE email = ?", (email,)) conn.commit() conn.close() return row[0] if row else None # --- Helpers --- def cloudron_headers() -> dict: return {"Authorization": f"Bearer {CLOUDRON_TOKEN}"} async def cloudron_create_user(email: str, name: str) -> str: """Create (or return existing) Cloudron user, assigned to the members group. Verified user shape (from live API): {id, username, email, fallbackEmail, displayName, role, active, groupIds} Roles: "owner", "admin", "user". Group assignment is a SEPARATE call: PUT /api/v1/users/:userId/groups with body {"groupIds": [...]}. """ async with httpx.AsyncClient() as client: # Check if user exists r = await client.get( f"{CLOUDRON_API}/api/v1/users", headers=cloudron_headers(), ) r.raise_for_status() for u in r.json().get("users", []): if u.get("email") == email: return u["id"] # Create user (role "user" = regular member). Only email is required; # Cloudron lets the user choose their own username during account setup. r = await client.post( f"{CLOUDRON_API}/api/v1/users", headers=cloudron_headers(), json={ "email": email, "displayName": name, "role": "user", "active": True, }, ) r.raise_for_status() return r.json()["id"] async def cloudron_set_group(user_id: str, group_id: str | None = None) -> None: """Assign a user to a group (defaults to the members group).""" gid = group_id or MEMBERS_GROUP_ID if not gid: logger.warning("No group id; skipping group assignment") return async with httpx.AsyncClient() as client: r = await client.put( f"{CLOUDRON_API}/api/v1/users/{user_id}/groups", headers=cloudron_headers(), json={"groupIds": [gid]}, ) r.raise_for_status() async def cloudron_set_active(user_id: str, active: bool) -> None: async with httpx.AsyncClient() as client: r = await client.put( f"{CLOUDRON_API}/api/v1/users/{user_id}/active", headers=cloudron_headers(), json={"active": active}, ) r.raise_for_status() async def cloudron_send_invite(user_id: str, email: str) -> None: """Send the account-setup invite email to the member.""" async with httpx.AsyncClient() as client: r = await client.post( f"{CLOUDRON_API}/api/v1/users/{user_id}/send_invite_email", headers=cloudron_headers(), json={"email": email}, ) r.raise_for_status() async def litellm_find_key_by_alias(alias: str) -> str | None: """Find an existing key's token by its alias (key/list returns tokens).""" async with httpx.AsyncClient() as client: r = await client.get( f"{LITELLM_BASE}/key/list", headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"}, ) r.raise_for_status() tokens = r.json().get("keys", []) for token in tokens: info = await client.get( f"{LITELLM_BASE}/key/info", headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"}, params={"key": token}, ) if info.status_code == 200: data = info.json().get("info", {}) if data.get("key_alias") == alias: return token return None async def litellm_create_key(email: str, budget: float) -> str: """Create a LiteLLM virtual key for a member with a budget cap. Idempotent: if a key with this alias already exists, return its token. """ alias = f"member:{email}" # If the key already exists (e.g. from a prior partial attempt), reuse it. existing = await litellm_find_key_by_alias(alias) if existing: logger.info("Reusing existing LiteLLM key for %s", email) return existing async with httpx.AsyncClient() as client: r = await client.post( f"{LITELLM_BASE}/key/generate", headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"}, json={ "key_alias": alias, "max_budget": budget, "budget_duration": "30d", "models": ["deepseek-v4-flash", "gpt-oss-120b"], }, ) r.raise_for_status() return r.json().get("key", "") async def litellm_disable_key(key_token: str) -> None: """Delete a member's LiteLLM key (on payment lapse).""" if not key_token: return async with httpx.AsyncClient() as client: await client.post( f"{LITELLM_BASE}/key/delete", headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"}, json={"keys": [key_token]}, ) # --- A. Open Collective webhook --- @app.post("/webhook/opencollective/{token}") async def opencollective_webhook(request: Request, token: str): """Handle Open Collective membership events. Authenticated by a secret token in the URL path (Open Collective's generic webhooks are not HMAC-signed, so a secret URL is the standard way to authenticate them). """ if not WEBHOOK_TOKEN or not secrets.compare_digest(token, WEBHOOK_TOKEN): raise HTTPException(401, "Invalid webhook token") payload = await request.json() event_type = payload.get("type", "") data = payload.get("data", {}) # The webhook does NOT include email (Open Collective strips it for # privacy). We get name + slug, store a pending member, and obtain the # email later via the OAuth "connect your account" flow. # # Slug/name live in different places depending on the event type: # - order.processed / transaction.created → data.fromCollective.{slug,name} # - collective.member.created → data.member.memberCollective.{slug,name} from_collective = data.get("fromCollective", {}) member_collective = (data.get("member", {}) or {}).get("memberCollective", {}) slug = from_collective.get("slug") or member_collective.get("slug", "") name = from_collective.get("name") or member_collective.get("name", "Member") logger.info("Open Collective event: %s (name=%s, slug=%s)", event_type, name, slug) # Open Collective webhook events (verified): # - "order.processed" → fires on EVERY payment (incl. monthly recurring) # - "new member" → fires on FIRST contribution only # - payload has "firstPayment" boolean to distinguish new vs recurring # - "collective.transaction.created" is DEPRECATED (being removed) if event_type in ("order.processed", "new.member", "collective.member.created"): # New or renewed member → store as pending; they complete via OAuth if slug: store_pending_member(slug, name) return JSONResponse({"status": "pending", "name": name, "slug": slug}) if event_type in ("collective.member.deleted", "collective.transaction.deleted"): # Lapsed member → move to the inactive group (lose access, keep data). # The webhook gives us the slug; we map it to the member's email/user. if slug: member = get_member_by_slug(slug) if member: # Move to inactive group (revokes chat access, keeps account + data) if INACTIVE_GROUP_ID and member.get("cloudron_user_id"): await cloudron_set_group(member["cloudron_user_id"], INACTIVE_GROUP_ID) # Mark inactive in our DB (key injector will refuse requests) deactivate_member(member["email"]) await sync_loomio_memberships() logger.info("Deactivated member %s (slug=%s)", member["email"], slug) return JSONResponse({"status": "deactivated", "email": member["email"]}) return JSONResponse({"status": "deactivated", "note": "no matching member"}) return JSONResponse({"status": "ignored", "type": event_type}) # --- B. Key injector (proxy) --- @app.get("/v1/models") async def list_models(): """List models (same for everyone — no per-user auth needed).""" async with httpx.AsyncClient() as client: r = await client.get( f"{LITELLM_BASE}/v1/models", headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"}, ) return Response( content=r.content, status_code=r.status_code, headers={"content-type": r.headers.get("content-type", "application/json")}, ) @app.api_route("/v1/{path:path}", methods=["GET", "POST", "PUT", "DELETE", "PATCH"]) async def inject_key(request: Request, path: str): """Read the member's email from a header, inject their key, forward to LiteLLM.""" # Verify the request came through the chat frontend, so a direct caller # can't spoof the user-email header and use another member's key. # # Two supported auth paths: # 1. OpenWebUI signs the user identity as a JWT (X-OpenWebUI-User-Jwt) # with a shared secret — the email is tamper-proof, so no separate # secret header is needed. # 2. LibreChat sends a plain X-User-Email header plus an X-Portal-Secret # shared-secret header. email = "" jwt_header = request.headers.get("x-openwebui-user-jwt", "") if jwt_header and OWUI_JWT_SECRET: try: claims = jwt.decode(jwt_header, OWUI_JWT_SECRET, algorithms=["HS256"]) email = claims.get("email", "") except jwt.PyJWTError: raise HTTPException(401, "Invalid user identity token") else: _verify_portal_secret(request) email = ( request.headers.get("x-user-email", "") or request.headers.get("x-openwebui-user-email", "") ) if not email: raise HTTPException(401, "No member identity (user-email header)") # Look up the member's key from the persistent store member_key = get_member_key(email) if not member_key: raise HTTPException(403, "No active membership key") # Forward the request to LiteLLM with the member's key body = await request.body() headers = dict(request.headers) headers["authorization"] = f"Bearer {member_key}" headers.pop("host", None) headers.pop("content-length", None) async with httpx.AsyncClient() as client: upstream = await client.request( method=request.method, url=f"{LITELLM_BASE}/v1/{path}", headers=headers, content=body, ) return Response( content=upstream.content, status_code=upstream.status_code, headers={"content-type": upstream.headers.get("content-type", "application/json")}, ) # --- C. OAuth "connect your account" flow --- @app.get("/join") async def join_page(): """The 'finish your setup' page — where a new member connects their Open Collective account so we can obtain their email (with consent).""" html = """
Thanks for joining the Inference Cooperative! To set up your account, connect your Open Collective account so we can verify your membership.
Connect Open CollectiveYou'll receive an account-setup email shortly after connecting.
Didn't get it? Contact info@inference.coop.
We couldn't find an active membership for your account. To join, contribute on our Open Collective page first, then return here to finish setup.
Questions? Contact info@inference.coop.
Your account is being set up. Check your email for a link to create your account and start chatting.
Didn't get it? Contact info@inference.coop.