""" 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") # Personal token for the "Inference Co-op Bot" account, which is an admin of # the collective. Authenticated as an admin, the GraphQL API exposes member # emails (which are hidden from anonymous access). We use this to match a # member by email — the stable identifier that works even for guest # contributors (who have no Open Collective account and thus no usable slug). OC_PERSONAL_TOKEN = os.environ.get("OC_PERSONAL_TOKEN", "") # 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 fetch_member_emails() -> dict[str, str]: """Return a map of email -> role for active financial contributors. Uses the bot's personal token (admin) so the GraphQL API exposes member emails, which are hidden from anonymous access. This is the authoritative source of truth for "who is a member" — and it works for guest contributors (who have no OC account) because their email is still in the list. """ if not OC_PERSONAL_TOKEN: logger.warning("OC_PERSONAL_TOKEN not set; cannot fetch member emails") return {} query = ( '{ collective(slug: "%s") { members(limit: 100) { nodes { role account { email } } } } }' % OC_COLLECTIVE_SLUG ) async with httpx.AsyncClient() as client: r = await client.post( OC_GRAPHQL_URL, headers={"Personal-Token": OC_PERSONAL_TOKEN}, json={"query": query}, ) if r.status_code != 200: logger.warning("OC member fetch failed: %s", r.status_code) return {} nodes = r.json().get("data", {}).get("collective", {}).get("members", {}).get("nodes", []) result: dict[str, str] = {} for n in nodes: role = n.get("role", "") acct = n.get("account", {}) or {} email = (acct.get("email") or "").strip().lower() if email and role in ("BACKER", "ADMIN"): result[email] = role return result async def is_active_member(email: str) -> bool: """Check whether this email is an active financial contributor. Matches on email (not slug) so guest contributors — who have no OC account and thus no usable slug — are still recognized. The email comes from the OAuth `me` query (with the user's consent); we match it against the admin-visible member list. """ if not email: return False members = await fetch_member_emails() return email.strip().lower() in members async def get_active_member_emails() -> list[str]: """Return the list of active member emails from the portal's own database. We use the local `members` table (populated during OAuth, which captures the member's email) rather than Open Collective's GraphQL, because OC returns `email: null` for privacy — the email field is only exposed via OAuth with the user's consent. The local table is the authoritative source of member emails. """ conn = get_db() rows = conn.execute( "SELECT email FROM members WHERE active = 1" ).fetchall() conn.close() return [r[0] for r in rows] 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", "glm-5-3-flash"], }, ) 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). Annotates each model with OpenWebUI metadata so web search is enabled by default (defaultFeatureIds includes 'web_search'). Without this, OpenWebUI leaves web search off for these models and members must toggle it manually. """ async with httpx.AsyncClient() as client: r = await client.get( f"{LITELLM_BASE}/v1/models", headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"}, ) if r.status_code != 200: return Response( content=r.content, status_code=r.status_code, headers={"content-type": r.headers.get("content-type", "application/json")}, ) payload = r.json() for model in payload.get("data", []): model.setdefault("info", {}).setdefault("meta", {}) model["info"]["meta"].setdefault("capabilities", {})["web_search"] = True model["info"]["meta"].setdefault("defaultFeatureIds", []) if "web_search" not in model["info"]["meta"]["defaultFeatureIds"]: model["info"]["meta"]["defaultFeatureIds"].append("web_search") return JSONResponse(content=payload) @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.
Contributed as a guest? Create an Open Collective account using the same email you used to contribute, then return here and connect it.
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.
Contributed as a guest? Make sure you've created an Open Collective account with the same email you used to contribute, then try again.
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.