478 lines
17 KiB
Python
478 lines
17 KiB
Python
"""
|
|
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
|
|
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 LibreChat sends as a header on every request, so the
|
|
# portal can verify the request genuinely came through LibreChat (which is
|
|
# behind Cloudron SSO) rather than a direct, spoofed request.
|
|
PORTAL_SECRET = os.environ.get("PORTAL_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"
|
|
|
|
# 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", "")
|
|
|
|
# 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, "
|
|
"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'))"
|
|
")"
|
|
)
|
|
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 store_member(email: str, key_token: str, cloudron_user_id: str) -> None:
|
|
conn = get_db()
|
|
conn.execute(
|
|
"INSERT INTO members (email, key_token, cloudron_user_id, active) "
|
|
"VALUES (?, ?, ?, 1) "
|
|
"ON CONFLICT(email) DO UPDATE SET key_token=excluded.key_token, "
|
|
"cloudron_user_id=excluded.cloudron_user_id, active=1",
|
|
(email, key_token, cloudron_user_id),
|
|
)
|
|
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 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)
|
|
r = await client.post(
|
|
f"{CLOUDRON_API}/api/v1/users",
|
|
headers=cloudron_headers(),
|
|
json={
|
|
"username": email.split("@")[0],
|
|
"email": email,
|
|
"displayName": name,
|
|
"role": "user",
|
|
"active": True,
|
|
},
|
|
)
|
|
r.raise_for_status()
|
|
return r.json()["id"]
|
|
|
|
|
|
async def cloudron_set_group(user_id: str) -> None:
|
|
"""Assign a user to the members group (grants chat app access)."""
|
|
if not MEMBERS_GROUP_ID:
|
|
logger.warning("MEMBERS_GROUP_ID not set; 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": [MEMBERS_GROUP_ID]},
|
|
)
|
|
r.raise_for_status()
|
|
|
|
|
|
async def cloudron_set_active(user_id: str, active: bool) -> None:
|
|
async with httpx.AsyncClient() as client:
|
|
r = await client.post(
|
|
f"{CLOUDRON_API}/api/v1/users/{user_id}",
|
|
headers=cloudron_headers(),
|
|
json={"active": active},
|
|
)
|
|
r.raise_for_status()
|
|
|
|
|
|
async def litellm_create_key(email: str, budget: float) -> str:
|
|
"""Create a LiteLLM virtual key for a member with a budget cap."""
|
|
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": f"member:{email}",
|
|
"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", {})
|
|
member = data.get("member", {}) or data.get("fromCollective", {})
|
|
# 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.
|
|
name = member.get("name", "Member")
|
|
slug = member.get("slug", "")
|
|
|
|
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 → deactivate (we may not have their email yet)
|
|
return JSONResponse({"status": "deactivated"})
|
|
|
|
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 LibreChat (shared secret), so a direct
|
|
# caller can't spoof the X-User-Email header and use another member's key.
|
|
_verify_portal_secret(request)
|
|
|
|
email = request.headers.get("x-user-email", "")
|
|
if not email:
|
|
raise HTTPException(401, "No member identity (x-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 = """<!DOCTYPE html>
|
|
<html lang="en">
|
|
<head>
|
|
<meta charset="UTF-8">
|
|
<meta name="viewport" content="width=device-width, initial-scale=1.0">
|
|
<title>Finish your setup — Inference Cooperative</title>
|
|
<style>
|
|
body { font-family: -apple-system, Segoe UI, Roboto, sans-serif; background: #faf8f5; color: #2d3327; display: flex; align-items: center; justify-content: center; min-height: 100vh; margin: 0; padding: 24px; }
|
|
.card { background: #fff; border: 1px solid #e8e2d8; border-radius: 12px; padding: 40px; max-width: 480px; text-align: center; }
|
|
h1 { font-size: 24px; margin: 0 0 12px; }
|
|
p { color: #6b7a62; line-height: 1.5; margin: 0 0 20px; }
|
|
.btn { display: inline-block; background: #5b8c5a; color: #fff; text-decoration: none; padding: 14px 32px; border-radius: 10px; font-weight: 600; }
|
|
.btn:hover { opacity: 0.9; }
|
|
.note { font-size: 13px; color: #94a08c; margin-top: 20px; }
|
|
.note a { color: #6b7a62; }
|
|
</style>
|
|
</head>
|
|
<body>
|
|
<div class="card">
|
|
<h1>Finish your setup</h1>
|
|
<p>Thanks for joining the Inference Cooperative! To set up your account, connect your Open Collective account so we can verify your membership.</p>
|
|
<a class="btn" href="/oauth/start">Connect Open Collective</a>
|
|
<p class="note">You'll receive an account-setup email shortly after connecting.<br>Didn't get it? Contact <a href="mailto:info@inference.coop">info@inference.coop</a>.</p>
|
|
</div>
|
|
</body>
|
|
</html>"""
|
|
return Response(content=html, media_type="text/html")
|
|
|
|
|
|
@app.get("/oauth/start")
|
|
async def oauth_start():
|
|
"""Redirect the user to Open Collective's consent screen."""
|
|
state = secrets.token_urlsafe(16)
|
|
params = {
|
|
"client_id": OC_OAUTH_CLIENT_ID,
|
|
"response_type": "code",
|
|
"redirect_uri": OC_OAUTH_REDIRECT_URI,
|
|
"scope": "email",
|
|
"state": state,
|
|
}
|
|
qs = "&".join(f"{k}={v}" for k, v in params.items())
|
|
return Response(
|
|
status_code=302,
|
|
headers={"Location": f"{OC_AUTHORIZE_URL}?{qs}"},
|
|
)
|
|
|
|
|
|
@app.get("/oauth/callback")
|
|
async def oauth_callback(request: Request):
|
|
"""Exchange the OAuth code for a token, fetch the email, and provision."""
|
|
code = request.query_params.get("code", "")
|
|
if not code:
|
|
raise HTTPException(400, "Missing code")
|
|
|
|
# Exchange code for access token
|
|
async with httpx.AsyncClient() as client:
|
|
r = await client.post(
|
|
OC_TOKEN_URL,
|
|
data={
|
|
"grant_type": "authorization_code",
|
|
"client_id": OC_OAUTH_CLIENT_ID,
|
|
"client_secret": OC_OAUTH_CLIENT_SECRET,
|
|
"code": code,
|
|
"redirect_uri": OC_OAUTH_REDIRECT_URI,
|
|
},
|
|
)
|
|
r.raise_for_status()
|
|
token = r.json().get("access_token", "")
|
|
|
|
if not token:
|
|
raise HTTPException(500, "No access token returned")
|
|
|
|
# Fetch the user's email
|
|
r = await client.post(
|
|
OC_GRAPHQL_URL,
|
|
headers={"Authorization": f"Bearer {token}"},
|
|
json={"query": "{ me { id name email } }"},
|
|
)
|
|
r.raise_for_status()
|
|
me = r.json().get("data", {}).get("me", {})
|
|
email = me.get("email", "")
|
|
name = me.get("name", "Member")
|
|
|
|
if not email:
|
|
raise HTTPException(400, "No email returned — did you grant the email scope?")
|
|
|
|
# Provision the member
|
|
user_id = await cloudron_create_user(email, name)
|
|
await cloudron_set_group(user_id)
|
|
await cloudron_set_active(user_id, True)
|
|
key_token = await litellm_create_key(email, MEMBER_BUDGET)
|
|
store_member(email, key_token, user_id)
|
|
|
|
logger.info("Provisioned member %s (user_id=%s)", email, user_id)
|
|
|
|
html = """<!DOCTYPE html>
|
|
<html lang="en">
|
|
<head><meta charset="UTF-8"><meta name="viewport" content="width=device-width, initial-scale=1.0">
|
|
<title>You're in — Inference Cooperative</title>
|
|
<style>
|
|
body { font-family: -apple-system, Segoe UI, Roboto, sans-serif; background: #faf8f5; color: #2d3327; display: flex; align-items: center; justify-content: center; min-height: 100vh; margin: 0; padding: 24px; }
|
|
.card { background: #fff; border: 1px solid #e8e2d8; border-radius: 12px; padding: 40px; max-width: 480px; text-align: center; }
|
|
h1 { font-size: 24px; margin: 0 0 12px; }
|
|
p { color: #6b7a62; line-height: 1.5; margin: 0 0 20px; }
|
|
.note { font-size: 13px; color: #94a08c; margin-top: 20px; }
|
|
.note a { color: #6b7a62; }
|
|
</style>
|
|
</head>
|
|
<body>
|
|
<div class="card">
|
|
<h1>You're in!</h1>
|
|
<p>Your account is being set up. Check your email for a link to create your account and start chatting.</p>
|
|
<p class="note">Didn't get it? Contact <a href="mailto:info@inference.coop">info@inference.coop</a>.</p>
|
|
</div>
|
|
</body>
|
|
</html>"""
|
|
return Response(content=html, media_type="text/html")
|
|
|
|
|
|
# --- D. Admin / health ---
|
|
|
|
@app.get("/health")
|
|
async def health():
|
|
return {"status": "ok"}
|
|
|
|
|
|
@app.get("/")
|
|
async def index():
|
|
return {"service": "Inference Cooperative Member Portal", "version": "0.1.0"}
|