Files
member-portal/app/main.py
T

1447 lines
58 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 re
import sqlite3
import secrets
import smtplib
from email.mime.text import MIMEText
from email.mime.multipart import MIMEMultipart
import httpx
import jwt
from fastapi import FastAPI, Request, Response, HTTPException
from fastapi.responses import JSONResponse, StreamingResponse
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", "")
# Shared secret that the member dashboard uses to authenticate broker requests.
# The dashboard is Cloudron-SSO-gated and passes the authenticated member's
# email; the portal trusts that email only because it carries this secret. Kept
# separate from PORTAL_SECRET so the two can be rotated independently.
BROKER_SECRET = os.environ.get("BROKER_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", "")
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"))
# Models every member key/team may access.
MEMBER_MODELS = ["deepseek-v4-1-flash", "gpt-oss-120b", "glm-5-3-flash"]
# 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", "")
# Manually-provisioned members (comma-separated emails) who are NOT in the
# Open Collective member list. Used for the free "founder" tier, which OC
# cannot process as a $0 recurring subscription. These members are exempt from
# the sweep's deactivation (they have no OC membership to lapse).
MANUAL_MEMBERS = {
e.strip().lower()
for e in os.environ.get("MANUAL_MEMBERS", "").split(",")
if e.strip()
}
# 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")
# --- Email (Cloudron sendmail addon) ---
# The sendmail addon exports SMTP credentials as CLOUDRON_MAIL_* env vars.
SMTP_SERVER = os.environ.get("CLOUDRON_MAIL_SMTP_SERVER", "")
SMTP_PORT = int(os.environ.get("CLOUDRON_MAIL_SMTP_PORT", "2525"))
SMTP_USERNAME = os.environ.get("CLOUDRON_MAIL_SMTP_USERNAME", "")
SMTP_PASSWORD = os.environ.get("CLOUDRON_MAIL_SMTP_PASSWORD", "")
MAIL_FROM = os.environ.get("CLOUDRON_MAIL_FROM", "info@inference.coop")
MAIL_FROM_NAME = os.environ.get("CLOUDRON_MAIL_FROM_DISPLAY_NAME", "Inference Cooperative")
# Public base URL for the portal (used in email links).
PORTAL_BASE = os.environ.get("PORTAL_BASE", "https://portal.inference.coop")
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 _verify_admin_token(request: Request) -> None:
"""Authenticate admin/internal endpoints via the X-Admin-Token header.
Replaces the old token-in-URL pattern (tokens in URLs leak into access
logs and Referer headers). Callers pass the admin token in a header instead.
"""
if not WEBHOOK_TOKEN:
raise HTTPException(503, "Admin token not configured")
provided = request.headers.get("x-admin-token", "")
if not secrets.compare_digest(provided, WEBHOOK_TOKEN):
raise HTTPException(401, "Invalid admin token")
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"
")"
)
# Member-managed API keys (created via the broker). Each row is one API key
# under the member's team, tracked so we can list/revoke by name without
# exposing the chat key. `sk_token` is the full sk- key (revealed once).
conn.execute(
"CREATE TABLE IF NOT EXISTS member_api_keys ("
"email TEXT NOT NULL, "
"name TEXT NOT NULL, "
"sk_token TEXT NOT NULL, "
"hash TEXT NOT NULL, "
"created_at TEXT DEFAULT (datetime('now')), "
"PRIMARY KEY (email, name)"
")"
)
# 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")
# Migration: add balance column (USD) — the member's current credit
# allowance. Defaults to MEMBER_BUDGET; can be topped up by payments.
# This is the flexible counter to the fixed $15/month allowance.
if "balance" not in cols:
conn.execute("ALTER TABLE members ADD COLUMN balance REAL")
return conn
def get_member_balance(email: str) -> float:
"""Return the member's current balance (USD), defaulting to MEMBER_BUDGET.
The balance is the flexible allowance — it starts at MEMBER_BUDGET but can
be adjusted (topped up by payments) without changing code.
"""
conn = get_db()
row = conn.execute(
"SELECT balance FROM members WHERE email = ?", (email,)
).fetchone()
conn.close()
if row and row[0] is not None:
return float(row[0])
return MEMBER_BUDGET
def set_member_balance(email: str, balance: float) -> None:
"""Set the member's balance (USD). Used by top-up / billing adjustments."""
conn = get_db()
conn.execute(
"UPDATE members SET balance = ? WHERE email = ?", (balance, email)
)
conn.commit()
conn.close()
async def fetch_members() -> list[dict]:
"""Return the full active-member list: email, name, slug, role.
Uses the bot's personal token (admin) so emails are exposed. This is the
authoritative source of truth for reconciliation (provisioning missing
members, deactivating lapsed ones).
"""
if not OC_PERSONAL_TOKEN:
logger.warning("OC_PERSONAL_TOKEN not set; cannot fetch members")
return []
query = (
'{ collective(slug: "%s") { members(limit: 100) { nodes { role account { name slug ... on Individual { 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: list[dict] = []
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.append({
"email": email,
"name": acct.get("name") or email.split("@")[0],
"slug": acct.get("slug") or "",
"role": role,
})
return result
def send_email(to: str, subject: str, text_body: str, html_body: str | None = None) -> bool:
"""Send an email via Cloudron's sendmail addon (SMTP).
Returns True on success, False on failure (logged, not raised — email is
best-effort and must never break the provisioning flow).
"""
if not SMTP_SERVER:
logger.warning("SMTP not configured; skipping email to %s", to)
return False
try:
msg = MIMEMultipart("alternative")
msg["From"] = f"{MAIL_FROM_NAME} <{MAIL_FROM}>"
msg["To"] = to
msg["Subject"] = subject
msg.attach(MIMEText(text_body, "plain", "utf-8"))
if html_body:
msg.attach(MIMEText(html_body, "html", "utf-8"))
with smtplib.SMTP(SMTP_SERVER, SMTP_PORT, timeout=15) as s:
if SMTP_USERNAME:
s.login(SMTP_USERNAME, SMTP_PASSWORD)
s.sendmail(MAIL_FROM, [to], msg.as_string())
logger.info("Sent email to %s: %s", to, subject)
return True
except Exception as e:
logger.warning("Failed to send email to %s: %s", to, e)
return False
def welcome_email(name: str, setup_url: str) -> tuple[str, str, str]:
"""Build the welcome email (subject, text, html) for a new member.
`setup_url` is the Cloudron account-setup link (from invite_link), which
lets the member create their account directly — no OAuth step needed.
The HTML uses a centered, card-style layout matching the landing page
(cream background, forest-green button, Figtree-like system fonts).
Email HTML is table-based for client compatibility; no external assets.
"""
subject = "Welcome to the Inference Cooperative — finish your setup"
text = (
f"Hi {name},\n\n"
"Thanks for joining the Inference Cooperative!\n\n"
"To finish setting up your account and start using private AI chat, "
"click the link below to create your account:\n\n"
f"{setup_url}\n\n"
"This takes about a minute.\n\n"
"Questions? Reply to this email or write to info@inference.coop.\n\n"
"— The Inference Cooperative\n"
)
html = f"""<!DOCTYPE html>
<html lang="en">
<head><meta charset="UTF-8"><meta name="viewport" content="width=device-width, initial-scale=1.0"></head>
<body style="margin:0;padding:0;background:#faf8f5;font-family:-apple-system,'Segoe UI',Roboto,Helvetica,Arial,sans-serif;">
<table role="presentation" width="100%" cellpadding="0" cellspacing="0" style="background:#faf8f5;padding:40px 16px;">
<tr>
<td align="center">
<table role="presentation" width="100%" cellpadding="0" cellspacing="0" style="max-width:520px;">
<tr>
<td align="center" style="padding-bottom:24px;">
<!-- Pine-tree circle logo (simplified) -->
<svg width="56" height="56" viewBox="0 0 130 130" fill="none" xmlns="http://www.w3.org/2000/svg" aria-hidden="true">
<circle cx="65" cy="14" r="5" fill="#5b8c5a"/>
<circle cx="58" cy="30" r="6" fill="#7ab87a"/>
<circle cx="74" cy="30" r="4" fill="#6ba86b"/>
<circle cx="51.5" cy="46" r="6.5" fill="#5b8c5a"/>
<circle cx="65" cy="46" r="5" fill="#7ab87a"/>
<circle cx="81" cy="46" r="4" fill="#6ba86b"/>
<circle cx="44" cy="62" r="7" fill="#6ba86b"/>
<circle cx="61" cy="62" r="5.5" fill="#5b8c5a"/>
<circle cx="76" cy="62" r="4.5" fill="#7ab87a"/>
<circle cx="89.5" cy="62" r="3.5" fill="#6ba86b"/>
<circle cx="36.5" cy="78" r="7.5" fill="#5b8c5a"/>
<circle cx="54" cy="78" r="6" fill="#7ab87a"/>
<circle cx="70" cy="78" r="5" fill="#6ba86b"/>
<circle cx="84" cy="78" r="4" fill="#5b8c5a"/>
<circle cx="97.5" cy="78" r="3.5" fill="#7ab87a"/>
<circle cx="44" cy="94" r="6.5" fill="#6ba86b"/>
<circle cx="61" cy="94" r="5.5" fill="#5b8c5a"/>
<circle cx="76" cy="94" r="4.5" fill="#7ab87a"/>
<circle cx="89.5" cy="94" r="3.5" fill="#6ba86b"/>
<circle cx="51.5" cy="110" r="5.5" fill="#5b8c5a"/>
<circle cx="65" cy="110" r="4.5" fill="#7ab87a"/>
<circle cx="81" cy="110" r="3.5" fill="#6ba86b"/>
</svg>
</td>
</tr>
<tr>
<td align="center" style="background:#ffffff;border:1px solid #e8e2d8;border-radius:16px;padding:40px 32px;">
<h1 style="margin:0 0 12px;font-size:24px;font-weight:700;color:#2d3327;">Welcome, {name}</h1>
<p style="margin:0 0 8px;font-size:16px;color:#6b7a62;line-height:1.5;">Thanks for joining the <strong style="color:#2d3327;">Inference Cooperative</strong>.</p>
<p style="margin:0 0 28px;font-size:15px;color:#6b7a62;line-height:1.5;">Finish setting up your account to start using private AI chat.</p>
<a href="{setup_url}" style="display:inline-block;background:#5b8c5a;color:#ffffff;text-decoration:none;padding:14px 36px;border-radius:10px;font-size:16px;font-weight:600;">Finish setup</a>
<p style="margin:28px 0 0;font-size:13px;color:#94a08c;line-height:1.5;">This takes about a minute.</p>
</td>
</tr>
<tr>
<td align="center" style="padding-top:24px;font-size:13px;color:#94a08c;line-height:1.6;">
Questions? Reply to this email or write to <a href="mailto:info@inference.coop" style="color:#6b7a62;">info@inference.coop</a>.<br>
<span style="color:#2d3327;">The Inference Cooperative</span> &middot; inference.coop<br>
<span style="font-size:12px;">By creating your account, you agree to our <a href="https://git.inference.coop/co-op/docs/src/branch/main/terms-of-service.md" style="color:#6b7a62;">Terms of Service</a> and <a href="https://git.inference.coop/co-op/docs/src/branch/main/privacy-policy.md" style="color:#6b7a62;">Privacy Policy</a>.</span>
</td>
</tr>
</table>
</td>
</tr>
</table>
</body>
</html>"""
return subject, text, html
async def get_active_member_emails() -> list[str]:
"""Return the list of active member emails from the portal's own database.
The local `members` table is the authoritative source of member emails
(populated during provisioning). Used for the Loomio sync.
"""
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 rekey_member(email: str, key_token: str) -> None:
"""Update only the key token, preserving cloudron_user_id and slug.
Used by the Teams migration to swap in a fresh sk- key without clobbering
the member's existing Cloudron identity.
"""
conn = get_db()
conn.execute(
"UPDATE members SET key_token = ?, active = 1 WHERE email = ?",
(key_token, email),
)
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:
# If the existing user has no display name (e.g. created with a
# placeholder), fill it in. Otherwise leave it alone — the user
# may have set their own name during account setup.
if not u.get("displayName"):
await cloudron_update_profile(u["id"], email, name)
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_update_profile(user_id: str, email: str, display_name: str, fallback_email: str = "") -> None:
"""Update a user's profile (display name, email) after creation.
Uses POST /api/v1/users/:userId/profile (the Cloudron "Update user" route).
Needed because cloudron_create_user() only sets displayName at creation and
returns early if the user already exists — so corrections require this.
"""
async with httpx.AsyncClient() as client:
r = await client.post(
f"{CLOUDRON_API}/api/v1/users/{user_id}/profile",
headers=cloudron_headers(),
json={"email": email, "displayName": display_name, "fallbackEmail": fallback_email},
)
r.raise_for_status()
async def cloudron_get_invite_link(user_id: str) -> str:
"""Return the account-setup link WITHOUT sending Cloudron's own email.
We use this so we can send our own branded email (with the setup link
embedded) instead of Cloudron's default invite email.
"""
async with httpx.AsyncClient() as client:
r = await client.get(
f"{CLOUDRON_API}/api/v1/users/{user_id}/invite_link",
headers=cloudron_headers(),
)
r.raise_for_status()
return r.json().get("inviteLink", "")
async def provision_member(email: str, name: str, slug: str = "") -> str:
"""Provision a member end-to-end and send our own welcome email.
Creates the Cloudron user (members group), issues a LiteLLM key under the
member's team (budget = the member's current balance, not a hardcoded $15),
stores the mapping, and sends our branded welcome email containing the
Cloudron account-setup link (instead of Cloudron's default invite email).
Idempotent: if the member already exists, reuses the existing user/key.
Returns the Cloudron user id.
"""
user_id = await cloudron_create_user(email, name)
await cloudron_set_group(user_id)
await cloudron_set_active(user_id, True)
balance = get_member_balance(email)
team_id = await litellm_get_or_create_team(email, balance)
key_token = await litellm_create_key(email, balance, team_id)
store_member(email, key_token, user_id, slug)
# Send our own welcome email with the account-setup link (no OAuth step).
setup_url = await cloudron_get_invite_link(user_id)
subject, text, html = welcome_email(name, setup_url)
send_email(email, subject, text, html)
await sync_loomio_memberships()
logger.info("Provisioned member %s (user_id=%s)", email, user_id)
return user_id
async def litellm_get_or_create_team(email: str, budget: float) -> str:
"""Return the team_id for a member's budget team, creating it if needed.
Each member gets a LiteLLM "team" (a budget container) with a $15/30d
budget. The team holds the chat key + any member-managed API keys, so all
spend draws from the one pool.
Idempotent: looks up the existing team by alias before creating. Note that
/team/new is NOT idempotent (it creates a new team each call), so we must
check /team/list first.
"""
alias = f"member:{email}"
async with httpx.AsyncClient() as client:
r = await client.get(
f"{LITELLM_BASE}/team/list",
headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"},
)
r.raise_for_status()
for t in r.json():
if t.get("team_alias") == alias:
return t["team_id"]
r = await client.post(
f"{LITELLM_BASE}/team/new",
headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"},
json={
"team_alias": alias,
"max_budget": budget,
"budget_duration": "30d",
"models": MEMBER_MODELS,
},
)
r.raise_for_status()
return r.json()["team_id"]
async def litellm_update_team_budget(team_id: str, budget: float) -> None:
"""Set a team's max_budget (the enforced spend cap).
This is the primitive behind "payment → add to balance": a top-up adjusts
the team's max_budget, which every key under the team (chat + API) draws
from. Keys inherit the team budget (their own max_budget is ignored when a
team_id is set), so the team is the single source of truth for allowance.
"""
async with httpx.AsyncClient() as client:
r = await client.post(
f"{LITELLM_BASE}/team/update",
headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"},
json={"team_id": team_id, "max_budget": budget},
)
r.raise_for_status()
async def sync_team_budget_from_balance(email: str) -> float:
"""Apply the member's stored balance to their team's max_budget.
Reads the balance from the members table and pushes it to LiteLLM. Returns
the applied budget. This keeps the portal DB and LiteLLM in sync whenever
a balance changes (top-up, reset, pricing-model change).
"""
balance = get_member_balance(email)
team_id = await litellm_get_or_create_team(email, balance)
await litellm_update_team_budget(team_id, balance)
return balance
async def litellm_create_key(email: str, budget: float, team_id: str) -> str:
"""Create a fresh LiteLLM virtual key under the member's team.
Returns the ACTUAL sk- key (only revealed at generation time). We always
create a new key rather than reusing — the sk- value is unrecoverable from
a stored hash, so "reuse" would store a useless hash (a bug that previously
broke chat for three members).
NOTE: the key does NOT set its own max_budget — it inherits the team's
budget (LiteLLM ignores a key's max_budget when team_id is set anyway).
The team is the single source of truth for allowance, so a balance change
(top-up) updates the team once and every key follows. The `budget` arg is
retained only for the team-creation fallback in the migration path.
"""
alias = f"member:{email}"
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,
"team_id": team_id,
"models": MEMBER_MODELS,
},
)
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).
/key/delete accepts either the sk- key or its SHA256 hash, so this works
for both current keys and legacy hash-keyed members.
"""
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]},
)
# --- Broker: member-managed API keys (scoped to the member's team) ---
def _api_key_alias(email: str, name: str) -> str:
return f"api:{email}:{name}"
async def broker_create_api_key(email: str, name: str) -> dict:
"""Create a member-managed API key under the member's team.
Returns {name, sk_token, created_at}. The key draws from the member's team
budget and is tracked in member_api_keys so it can be listed/revoked by name.
"""
if not name or not re.fullmatch(r"[A-Za-z0-9._-]{1,64}", name):
raise HTTPException(400, "Invalid key name (1-64 chars, alphanumeric/._-)")
balance = get_member_balance(email)
team_id = await litellm_get_or_create_team(email, balance)
alias = _api_key_alias(email, name)
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,
"team_id": team_id,
"models": MEMBER_MODELS,
},
)
r.raise_for_status()
data = r.json()
sk_token = data.get("key", "")
h = data.get("token") or data.get("token_id") or ""
conn = get_db()
conn.execute(
"INSERT INTO member_api_keys (email, name, sk_token, hash) VALUES (?, ?, ?, ?) "
"ON CONFLICT(email, name) DO UPDATE SET sk_token=excluded.sk_token, hash=excluded.hash, "
"created_at=datetime('now')",
(email, name, sk_token, h),
)
conn.commit()
row = conn.execute(
"SELECT created_at FROM member_api_keys WHERE email = ? AND name = ?",
(email, name),
).fetchone()
conn.close()
return {"name": name, "sk_token": sk_token, "created_at": row[0] if row else None}
async def broker_list_api_keys(email: str) -> list[dict]:
conn = get_db()
rows = conn.execute(
"SELECT name, sk_token, created_at FROM member_api_keys WHERE email = ? ORDER BY created_at DESC",
(email,),
).fetchall()
conn.close()
return [
{"name": r[0], "sk_token": r[1], "created_at": r[2]}
for r in rows
]
async def broker_revoke_api_key(email: str, name: str) -> None:
conn = get_db()
row = conn.execute(
"SELECT hash FROM member_api_keys WHERE email = ? AND name = ?", (email, name)
).fetchone()
if not row:
conn.close()
raise HTTPException(404, "Key not found")
conn.execute("DELETE FROM member_api_keys WHERE email = ? AND name = ?", (email, name))
conn.commit()
conn.close()
# Delete the underlying LiteLLM key by its hash.
await litellm_disable_key(row[0])
# --- 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, then look up the email via the admin token.
#
# 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 member → provision immediately and send our own welcome email.
# The webhook strips email, so look it up by slug via the admin token.
# Only provision on firstPayment (order.processed also fires on monthly
# renewals, which should NOT re-send the welcome email).
first_payment = data.get("firstPayment", True)
if slug and first_payment:
members = await fetch_members()
for m in members:
if m["slug"] == slug:
try:
await provision_member(m["email"], m["name"], m["slug"])
except Exception as e:
logger.warning("Webhook: failed to provision %s: %s", m["email"], e)
break
return JSONResponse({"status": "provisioned", "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, STREAMING the
# response back token-by-token (or byte-by-byte) rather than buffering it.
# Buffering collapsed the model's stream and handed OpenWebUI the full
# answer in one blob, so members saw nothing until generation finished —
# which made chat feel slow even though time-to-first-byte was fine.
body = await request.body()
headers = dict(request.headers)
headers["authorization"] = f"Bearer {member_key}"
headers.pop("host", None)
headers.pop("content-length", None)
headers.pop("accept-encoding", None) # let httpx handle decompression
timeout = httpx.Timeout(connect=15.0, read=600.0, write=30.0, pool=15.0)
# Send the request upstream, but don't buffer the body — use a stream.
client = httpx.AsyncClient(timeout=timeout)
upstream_req = client.build_request(
method=request.method,
url=f"{LITELLM_BASE}/v1/{path}",
headers=headers,
content=body,
)
upstream = await client.send(upstream_req, stream=True)
async def gen():
try:
async for chunk in upstream.aiter_raw():
yield chunk
finally:
await upstream.aclose()
await client.aclose()
# Pass through the content-type (critical for SSE streaming) and status.
return StreamingResponse(
gen(),
status_code=upstream.status_code,
headers={
"content-type": upstream.headers.get("content-type", "application/json"),
"cache-control": "no-cache",
},
)
# --- D. Admin / health ---
@app.get("/health")
async def health():
return {"status": "ok"}
@app.post("/admin/sync-loomio")
async def admin_sync_loomio(request: Request):
"""Manually trigger a Loomio membership sync.
Authenticated by the X-Admin-Token header. Useful for reconciling
pre-existing members (accounts created before the sync existed) or
recovering from a missed webhook. Idempotent — safe to call repeatedly.
"""
_verify_admin_token(request)
await sync_loomio_memberships()
return {"status": "synced"}
@app.post("/admin/provision")
async def admin_provision(request: Request):
"""Manually provision a member by email (full pipeline).
Authenticated by the X-Admin-Token header. Runs the full provisioning
pipeline: Cloudron user + members group + LiteLLM key + our welcome email +
Loomio sync. Body: {"email": "...", "name": "..."}. Idempotent.
"""
_verify_admin_token(request)
body = await request.json()
email = (body.get("email") or "").strip().lower()
name = body.get("name") or email.split("@")[0]
if not email:
raise HTTPException(400, "Missing email")
user_id = await provision_member(email, name, "")
return {"status": "provisioned", "email": email, "user_id": user_id}
@app.post("/admin/set-balance")
async def admin_set_balance(request: Request):
"""Set (or add to) a member's credit balance and sync it to LiteLLM.
The flexible alternative to a fixed $15/month allowance. A payment (from OC
or another system) is recorded by adjusting the balance here, which pushes
the new cap to the member's team (every chat/API key under it follows).
Body: {"email": "...", "balance": 25.0} → set absolute balance
{"email": "...", "add": 10.0} → add to current balance
"""
_verify_admin_token(request)
body = await request.json()
email = (body.get("email") or "").strip().lower()
if not email:
raise HTTPException(400, "Missing email")
current = get_member_balance(email)
if "add" in body:
new_balance = current + float(body["add"])
elif "balance" in body:
new_balance = float(body["balance"])
else:
raise HTTPException(400, "Provide 'balance' or 'add'")
if new_balance < 0:
raise HTTPException(400, "Balance cannot be negative")
set_member_balance(email, new_balance)
applied = await sync_team_budget_from_balance(email)
return {"status": "updated", "email": email, "balance": new_balance, "team_budget": applied}
@app.post("/admin/migrate-teams")
async def admin_migrate_teams(request: Request):
"""One-off migration: put every active member under a LiteLLM team + fresh
sk- key. Fixes a bug where some members had SHA256 hashes (from /key/list)
stored as their key, which cannot authenticate.
For each active member: delete any existing key(s) by alias, create/reuse a
team, generate a fresh sk- key under it, and store the sk- token. Idempotent
but re-issues keys, so call once.
"""
_verify_admin_token(request)
conn = get_db()
rows = conn.execute("SELECT email, key_token FROM members WHERE active = 1").fetchall()
conn.close()
results = []
for email, old_token in rows:
try:
# Delete any existing key(s) for this alias (by hash or sk-).
await litellm_delete_keys_by_alias(email)
# (Re)create the team and a fresh sk- key, at the member's balance.
balance = get_member_balance(email)
team_id = await litellm_get_or_create_team(email, balance)
new_key = await litellm_create_key(email, balance, team_id)
rekey_member(email, new_key) # preserve cloudron_user_id + slug
results.append({"email": email, "status": "rekeyed"})
logger.info("Migrated %s to team %s", email, team_id)
except Exception as e:
results.append({"email": email, "status": "error", "error": str(e)})
logger.warning("Migrate failed for %s: %s", email, e)
return {"status": "migrated", "results": results}
@app.get("/admin/overview")
async def admin_overview(request: Request):
"""Return a co-op-wide overview for the admin dashboard.
Authenticated by the X-Admin-Token header. Returns every member (active +
inactive) with their balance, spend, remaining, and team reset date, plus
co-op aggregates (total members, total spend).
"""
_verify_admin_token(request)
conn = get_db()
rows = conn.execute(
"SELECT email, active, balance, cloudron_user_id FROM members ORDER BY email"
).fetchall()
conn.close()
# Fetch Cloudron users once to map user_id → username + inviteAccepted.
# username is set only after the member completes account setup (chooses a
# username + password via the setup link); inviteAccepted tracks that same
# step, so "activated" = inviteAccepted (username present).
usernames = {}
invite_accepted = {}
try:
async with httpx.AsyncClient(timeout=15.0) as client:
r = await client.get(
f"{CLOUDRON_API}/api/v1/users",
headers={"Authorization": f"Bearer {CLOUDRON_TOKEN}"},
)
r.raise_for_status()
for u in r.json().get("users", []):
uid = u.get("id")
if u.get("username"):
usernames[uid] = u["username"]
invite_accepted[uid] = bool(u.get("inviteAccepted"))
except Exception as e:
logger.warning("admin overview: cloudron users fetch failed: %s", e)
# Fetch team spend/budget from LiteLLM once.
teams = {}
try:
async with httpx.AsyncClient(timeout=15.0) as client:
r = await client.get(
f"{LITELLM_BASE}/team/list",
headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"},
)
r.raise_for_status()
for t in r.json():
teams[t.get("team_alias")] = t
except Exception as e:
logger.warning("admin overview: team/list failed: %s", e)
members = []
total_spend = 0.0
for email, active, balance, user_id in rows:
team = teams.get(f"member:{email}", {})
spend = float(team.get("spend") or 0.0)
total_spend += spend
members.append(
{
"email": email,
"username": usernames.get(user_id),
"activated": invite_accepted.get(user_id, False),
"active": bool(active),
"balance": float(balance) if balance is not None else MEMBER_BUDGET,
"spend": spend,
"remaining": max((balance or MEMBER_BUDGET) - spend, 0.0),
"reset_at": team.get("budget_reset_at"),
"cloudron_user_id": user_id,
}
)
return {
"members": members,
"totals": {
"members": len(members),
"active": sum(1 for m in members if m["active"]),
"total_spend": total_spend,
},
}
async def litellm_delete_keys_by_alias(email: str) -> None:
"""Delete all LiteLLM keys whose alias matches member:<email>.
/key/list returns SHA256 hashes; /key/delete accepts hashes. This clears
both legacy hash-keyed and current sk- keyed entries for a member.
"""
alias = f"member:{email}"
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()
hashes = r.json().get("keys", [])
to_delete = []
for h in hashes:
info = await client.get(
f"{LITELLM_BASE}/key/info",
headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"},
params={"key": h},
)
if info.status_code == 200:
data = info.json().get("info", {})
if data.get("key_alias") == alias:
to_delete.append(h)
if to_delete:
await client.post(
f"{LITELLM_BASE}/key/delete",
headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"},
json={"keys": to_delete},
)
# --- Broker HTTP endpoints (called by the member dashboard) ---
# The dashboard is Cloudron-SSO-gated; it authenticates the member and passes
# their email with a shared secret. The portal trusts the email only because it
# carries BROKER_SECRET. These endpoints expose ONLY the authenticated member's
# own keys — never another member's, and never the chat key.
def _verify_broker_secret(request: Request) -> str:
"""Return the authenticated member's email, or 401.
The member dashboard sends X-Broker-Secret (shared secret) + X-Member-Email
(or X-Member-User, a Cloudron username). Both must be present and the secret
must match, or we refuse (fail closed).
Cloudron's proxyAuth injects the USERNAME (e.g. "ntnsndr"), not the email.
So when the dashboard sends a username (no "@"), we resolve it to an email
via the Cloudron user list (we already hold a Cloudron admin token).
"""
if not BROKER_SECRET:
raise HTTPException(503, "Broker secret not configured")
provided = request.headers.get("x-broker-secret", "")
if not secrets.compare_digest(provided, BROKER_SECRET):
raise HTTPException(401, "Invalid broker secret")
identity = (
request.headers.get("x-member-email")
or request.headers.get("x-member-user")
or ""
).strip()
if not identity:
raise HTTPException(401, "Missing member identity")
email = resolve_member_email(identity)
# The member must be active in our DB before they can manage keys.
if not get_member_key(email):
raise HTTPException(403, "No active membership")
return email
def resolve_member_email(identity: str) -> str:
"""Resolve a member identity (email OR Cloudron username) to their email.
If it looks like an email (contains "@"), return it lowercased. Otherwise
treat it as a Cloudron username and look up the corresponding email from
the Cloudron user list (synchronous, using urllib since this runs outside
the async request path of httpx).
"""
identity = identity.strip()
if "@" in identity:
return identity.lower()
# Username → email via Cloudron user list.
import urllib.request
req = urllib.request.Request(
f"{CLOUDRON_API}/api/v1/users",
headers={"Authorization": f"Bearer {CLOUDRON_TOKEN}"},
)
with urllib.request.urlopen(req, timeout=15) as r:
data = json.loads(r.read().decode())
users = data.get("users", [])
for u in users:
if u.get("username") == identity:
return (u.get("email") or "").lower()
raise HTTPException(403, f"No Cloudron user for identity {identity!r}")
@app.get("/broker/keys")
async def broker_list(request: Request):
email = _verify_broker_secret(request)
return {"keys": await broker_list_api_keys(email)}
@app.post("/broker/keys")
async def broker_create(request: Request):
email = _verify_broker_secret(request)
body = await request.json()
name = (body.get("name") or "").strip()
key = await broker_create_api_key(email, name)
return {"status": "created", **key}
@app.delete("/broker/keys/{name}")
async def broker_revoke(name: str, request: Request):
email = _verify_broker_secret(request)
await broker_revoke_api_key(email, name)
return {"status": "revoked", "name": name}
async def broker_get_usage(email: str) -> dict:
"""Return the member's spend vs. balance, plus per-model and token detail.
Reads the member's team (spend + max_budget) and aggregates their recent
spend logs by model for token counts. This is the read-only primitive the
member dashboard uses (no master key on the client).
"""
team_id = await litellm_get_or_create_team(email, get_member_balance(email))
balance = get_member_balance(email)
spend = 0.0
reset_at = None
async with httpx.AsyncClient() as client:
# 1. Team spend + budget
r = await client.get(
f"{LITELLM_BASE}/team/list",
headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"},
)
r.raise_for_status()
for t in r.json():
if t.get("team_id") == team_id:
balance = float(t.get("max_budget") or balance)
spend = float(t.get("spend") or 0.0)
reset_at = t.get("budget_reset_at")
break
# 2. Per-model token/spend breakdown from spend logs (filtered by team).
# /spend/logs does not reliably filter by team_id server-side, so we
# filter client-side on the returned rows.
models = {}
total_tokens = 0
try:
r2 = await client.get(
f"{LITELLM_BASE}/spend/logs",
headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"},
)
r2.raise_for_status()
rows = r2.json()
if isinstance(rows, dict):
rows = rows.get("data", rows.get("logs", []))
for row in rows if isinstance(rows, list) else []:
if row.get("team_id") != team_id:
continue
model = row.get("model") or "unknown"
m = models.setdefault(model, {"spend": 0.0, "tokens": 0, "calls": 0})
m["spend"] += float(row.get("spend") or 0.0)
m["tokens"] += int(row.get("total_tokens") or 0)
m["calls"] += 1
total_tokens += int(row.get("total_tokens") or 0)
except Exception as e: # spend logs are best-effort; don't fail the whole view
logger.warning("spend/logs fetch failed: %s", e)
# Sort models by spend descending (which models are draining the most).
model_list = [
{"model": name, **stats}
for name, stats in sorted(models.items(), key=lambda kv: kv[1]["spend"], reverse=True)
]
return {
"email": email,
"balance": balance,
"spend": spend,
"remaining": max(balance - spend, 0.0),
"total_tokens": total_tokens,
"reset_at": reset_at,
"models": model_list,
}
@app.get("/broker/usage")
async def broker_usage(request: Request):
email = _verify_broker_secret(request)
return await broker_get_usage(email)
async def reconcile_memberships() -> dict:
"""Reconcile the portal against the live Open Collective member list.
- Provisions members who are BACKER/ADMIN on OC but not yet in our DB
(missed webhook, or contributed before the email step existed).
- Deactivates members who are in our DB but no longer BACKER/ADMIN
(lapsed or cancelled membership).
- Syncs Loomio to the resulting active-member set.
Idempotent and safe to run repeatedly. Returns a summary of actions taken.
"""
active = await fetch_members()
active_emails = {m["email"] for m in active}
# FAIL-SAFE: if the OC fetch returned nothing (token expired, OC outage,
# or a query error), we must NOT deactivate everyone — that would be a
# catastrophic false-positive. Skip reconciliation entirely in that case.
if not active:
logger.warning("Sweep: fetch_members returned empty; skipping reconciliation (fail-safe)")
return {"status": "skipped", "reason": "empty_member_list", "provisioned": 0, "deactivated": 0}
# Currently provisioned members (from our DB).
conn = get_db()
rows = conn.execute("SELECT email, cloudron_user_id FROM members WHERE active = 1").fetchall()
conn.close()
provisioned = {r[0]: r[1] for r in rows}
provisioned_count = 0
deactivated_count = 0
# 1. Provision missing members.
for m in active:
email = m["email"]
if email not in provisioned:
try:
await provision_member(email, m["name"], m["slug"])
provisioned_count += 1
logger.info("Sweep: provisioned %s", email)
except Exception as e:
logger.warning("Sweep: failed to provision %s: %s", email, e)
# 2. Deactivate lapsed members (but never manual/founder members, who have
# no OC membership to lapse).
for email, user_id in provisioned.items():
if email not in active_emails and email not in MANUAL_MEMBERS:
try:
if INACTIVE_GROUP_ID and user_id:
await cloudron_set_group(user_id, INACTIVE_GROUP_ID)
deactivate_member(email)
deactivated_count += 1
logger.info("Sweep: deactivated %s", email)
except Exception as e:
logger.warning("Sweep: failed to deactivate %s: %s", email, e)
# 3. Sync Loomio.
await sync_loomio_memberships()
return {
"status": "reconciled",
"provisioned": provisioned_count,
"deactivated": deactivated_count,
}
@app.post("/admin/sweep")
async def admin_sweep(request: Request):
"""Manually trigger a full membership reconciliation.
Authenticated by the X-Admin-Token header. Idempotent — safe to call
repeatedly. This is also what the Cloudron scheduler invokes on a cron
schedule (see the sweep script).
"""
_verify_admin_token(request)
result = await reconcile_memberships()
return result
@app.get("/")
async def index():
return {"service": "Inference Cooperative Member Portal", "version": "0.1.0"}