Files

2599 lines
108 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 datetime import datetime, timedelta, timezone
from email.mime.text import MIMEText
from email.mime.multipart import MIMEMultipart
import httpx
import jwt
import asyncio
from fastapi import FastAPI, Request, Response, HTTPException
from fastapi.responses import JSONResponse, StreamingResponse
from slowapi import Limiter
from slowapi.util import get_remote_address
from slowapi.errors import RateLimitExceeded
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger("member-portal")
app = FastAPI(title="Inference Cooperative Member Portal")
# --- Rate limiting (control plane only, not the member inference path) ---
# The data plane (chat completions via /v1/*) is capped by each member's
# LiteLLM budget, NOT by rate — a member streaming a long conversation must
# never be throttled. Rate limiting here covers the control-plane surfaces:
# admin endpoints, the broker, the webhook, and the unauthenticated model list.
limiter = Limiter(key_func=get_remote_address)
app.state.limiter = limiter
@app.exception_handler(RateLimitExceeded)
async def _rate_limit_handler(request: Request, exc: RateLimitExceeded):
return JSONResponse(
status_code=429,
content={"detail": "Too many requests. Slow down."},
)
# --- 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", "")
# 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")
# Cache of order idV2 -> tip-free contribution amount in cents, resolved via
# OC's PUBLIC GraphQL API (no token needed; our personal token lacks the
# orders scope but unauthenticated reads of public orders work).
_order_amount_cache: dict[str, int] = {}
async def get_order_amount(order_id_v2: str | None) -> int | None:
"""Look up an order's contribution amount (EXCLUDING the platform tip).
The webhook's `totalAmount` includes the buyer's optional OC tip, which
the co-op never receives — credits must be granted on the tip-free
amount. Returns cents, or None if unresolvable.
"""
if not order_id_v2:
return None
if order_id_v2 in _order_amount_cache:
return _order_amount_cache[order_id_v2]
try:
query = ('{ order(order: {id: "%s"}) { amount { valueInCents } '
'platformTipAmount { valueInCents } } }' % order_id_v2)
async with httpx.AsyncClient(timeout=15.0) as client:
r = await client.post(OC_GRAPHQL_URL,
headers={"Content-Type": "application/json",
"User-Agent": "Mozilla/5.0 (X11; Linux x86_64)"},
json={"query": query})
o = r.json().get("data", {}).get("order") or {}
cents = (o.get("amount") or {}).get("valueInCents")
if cents is not None:
_order_amount_cache[order_id_v2] = int(cents)
return int(cents)
except Exception as e:
logger.warning("Order amount lookup failed for %s: %s", order_id_v2, e)
return None
# 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"))
# Days to wait before re-sending the welcome/invite email to a member who was
# provisioned but never completed account setup (inviteAccepted=False). The
# nightly sweep re-invites such members, throttled by this interval so they
# aren't emailed every single night.
REINVITE_AFTER_DAYS = float(os.environ.get("REINVITE_AFTER_DAYS", "3"))
# Models every member key/team may access.
# Deliberately REMOVED: members get access to every model the gateway serves.
# The gateway's `model_list` is the single source of truth for what's available;
# teams and keys carry no per-model restriction (empty = "all models allowed" in
# LiteLLM). Adding a provider is a gateway-config change only — no team/key
# migration. (Kept as a documented empty value so the intent is explicit.)
MEMBER_MODELS: list[str] = []
# 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", "")
# Server API key (B3) — used to fix auto-generated Loomio usernames after
# invites. Loomio's b2/memberships API only takes emails; new users get a
# mangled auto-handle (email local-part + random suffix). The B3 users API
# lets us rename them to the member's Cloudron username.
LOOMIO_B3_KEY = os.environ.get("LOOMIO_B3_KEY", "")
# Matrix homeserver (self-hosted Synapse at matrix.inference.coop). The bot
# account (admin) invites every active member to the co-op Space on sync.
# Members authenticate via Cloudron SSO; their homeserver MXID is
# @<cloudron-username>:inference.coop. Federated MXIDs (external servers)
# can also be invited directly.
MATRIX_HOMESERVER = os.environ.get("MATRIX_HOMESERVER", "https://matrix.inference.coop")
MATRIX_BOT_TOKEN = os.environ.get("MATRIX_BOT_TOKEN", "")
MATRIX_SPACE_ID = os.environ.get("MATRIX_SPACE_ID", "")
MATRIX_GENERAL_ROOM = os.environ.get("MATRIX_GENERAL_ROOM", "")
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")
# Reply-to address so member replies land at the co-op inbox.
MAIL_REPLY_TO = os.environ.get("MAIL_REPLY_TO", "info@inference.coop")
# Public base URL for the portal (used in email links).
PORTAL_BASE = os.environ.get("PORTAL_BASE", "https://portal.inference.coop")
# Listmonk newsletter integration. The portal keeps the "Members" list in sync:
# subscribe on provision, unsubscribe on deactivation. Auth uses listmonk's
# native API-user token scheme (Authorization: token <user>:<token>).
LISTMONK_BASE = os.environ.get("LISTMONK_BASE", "https://newsletter.inference.coop")
LISTMONK_API_USER = os.environ.get("LISTMONK_API_USER", "portal-sync")
LISTMONK_API_TOKEN = os.environ.get("LISTMONK_API_TOKEN", "")
LISTMONK_LIST_ID = int(os.environ.get("LISTMONK_LIST_ID", "3")) # "Members" list
# Direct read access to LiteLLM's Postgres (for usage aggregations).
# READ-ONLY purpose: never write to LiteLLM's tables — its HTTP API remains the
# only interface for writes/management. Falls back to the HTTP path when unset.
LITELLM_DB_HOST = os.environ.get("LITELLM_DB_HOST", "")
LITELLM_DB_PORT = int(os.environ.get("LITELLM_DB_PORT", "5432"))
LITELLM_DB_USER = os.environ.get("LITELLM_DB_USER", "")
LITELLM_DB_PASSWORD = os.environ.get("LITELLM_DB_PASSWORD", "")
LITELLM_DB_NAME = os.environ.get("LITELLM_DB_NAME", "")
def litellm_db_ready() -> bool:
"""True when LiteLLM DB read credentials are configured."""
return bool(LITELLM_DB_HOST and LITELLM_DB_USER and LITELLM_DB_NAME)
def _pg_connect():
"""Open a short-lived connection to LiteLLM's Postgres (read-only use)."""
import psycopg2
return psycopg2.connect(
host=LITELLM_DB_HOST, port=LITELLM_DB_PORT, user=LITELLM_DB_USER,
password=LITELLM_DB_PASSWORD, dbname=LITELLM_DB_NAME,
connect_timeout=10, sslmode="prefer",
)
def aggregate_model_usage(days: int = 0, team_id: str | None = None) -> dict:
"""Aggregate LiteLLM spend logs per model, straight from Postgres.
Returns {"models": [...], "totals": {...}} in the same shape the HTTP path
produced, so callers don't change. `days` bounds the window (0 = all time);
`team_id` restricts to one member's team (used by the dashboard path).
Prefers the `model_group` column, which carries our provider/model alias
(e.g. greenpt/green-r); falls back to `model` for pre-alias rows — matching
the logic the HTTP path used, expressed in SQL.
"""
if not litellm_db_ready():
raise RuntimeError("LiteLLM DB credentials not configured")
where = []
params: list = []
if days and days > 0:
from datetime import datetime, timedelta, timezone
cutoff = (datetime.now(timezone.utc) - timedelta(days=days)).strftime(
"%Y-%m-%d %H:%M:%S"
)
where.append("\"startTime\" >= %s")
params.append(cutoff)
if team_id:
where.append("team_id = %s")
params.append(team_id)
where_sql = ("WHERE " + " AND ".join(where)) if where else ""
sql = f"""
SELECT
COALESCE(NULLIF(model_group, ''), NULLIF(model, ''), 'unknown') AS model,
COALESCE(SUM(spend), 0) AS spend,
COALESCE(SUM(total_tokens), 0) AS tokens,
COUNT(*) AS requests
FROM "LiteLLM_SpendLogs"
{where_sql}
GROUP BY 1
ORDER BY requests DESC
"""
conn = _pg_connect()
try:
with conn.cursor() as cur:
cur.execute(sql, params)
rows = cur.fetchall()
finally:
conn.close()
models = [
{"model": r[0], "spend": float(r[1] or 0.0), "tokens": int(r[2] or 0),
"requests": int(r[3] or 0)}
for r in rows
]
return {
"models": models,
"totals": {
"total_spend": sum(m["spend"] for m in models),
"total_tokens": sum(m["tokens"] for m in models),
"total_requests": sum(m["requests"] for m in models),
},
}
def get_team_budget_from_db(team_id: str) -> dict | None:
"""Fetch a team's spend/budget/reset directly from LiteLLM's Postgres.
Same data as the /team/list HTTP path (LiteLLM computes `spend`), without
downloading the whole team list. Returns None when the team is unknown.
"""
if not litellm_db_ready():
return None
sql = """
SELECT max_budget, spend, budget_reset_at
FROM "LiteLLM_TeamTable"
WHERE team_id = %s
LIMIT 1
"""
conn = _pg_connect()
try:
with conn.cursor() as cur:
cur.execute(sql, (team_id,))
row = cur.fetchone()
finally:
conn.close()
if not row:
return None
return {
"max_budget": float(row[0]) if row[0] is not None else None,
"spend": float(row[1] or 0.0),
"budget_reset_at": row[2],
}
async def listmonk_sync(email: str, subscribe: bool) -> None:
"""Subscribe or unsubscribe a member to the newsletter list.
Best-effort: never raises, so a newsletter outage can't break provisioning.
"""
if not LISTMONK_API_TOKEN:
logger.warning("Listmonk not configured; skipping sync for %s", email)
return
auth = {"Authorization": f"token {LISTMONK_API_USER}:{LISTMONK_API_TOKEN}"}
try:
async with httpx.AsyncClient(timeout=10.0) as client:
if subscribe:
r = await client.post(
f"{LISTMONK_BASE}/api/subscribers",
headers=auth,
json={
"email": email,
"name": "",
"status": "enabled",
"lists": [LISTMONK_LIST_ID],
"preconfirm_subscriptions": True,
},
)
else:
# Look up the subscriber id by email, then blocklist by id.
# (listmonk's blocklist endpoint takes an id, not an email.)
q = await client.get(
f"{LISTMONK_BASE}/api/subscribers",
headers=auth,
params={"query": f"subscribers.email='{email}'"},
)
if q.status_code != 200:
logger.warning("Listmonk lookup for %s failed: %s", email, q.status_code)
return
results = q.json().get("data", {}).get("results", [])
if not results:
return # not subscribed; nothing to do
sid = results[0]["id"]
r = await client.put(
f"{LISTMONK_BASE}/api/subscribers/{sid}/blocklist",
headers=auth,
)
if r.status_code not in (200, 201):
logger.warning(
"Listmonk sync %s for %s failed: %s",
"subscribe" if subscribe else "unsubscribe",
email,
r.status_code,
)
except Exception as e:
logger.warning("Listmonk sync error for %s: %s", email, e)
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")
# Migration: welcome_sent_at — timestamp of the last welcome/invite email,
# used to throttle re-invites of members who never activated.
if "welcome_sent_at" not in cols:
conn.execute("ALTER TABLE members ADD COLUMN welcome_sent_at TEXT")
# Migration: credit_balance — purchased, NON-EXPIRING credits (USD).
# Distinct from `balance` (the monthly allowance, which resets): credits
# carry over indefinitely and are spent only after the allowance is
# exhausted within a cycle.
if "credit_balance" not in cols:
conn.execute("ALTER TABLE members ADD COLUMN credit_balance REAL NOT NULL DEFAULT 0")
# Credit ledger: audit trail for every credit movement (purchases via the
# OC webhook, admin adjustments, spends drawn against credits). Money needs
# an audit trail.
conn.execute(
"CREATE TABLE IF NOT EXISTS credit_ledger ("
"id INTEGER PRIMARY KEY AUTOINCREMENT, "
"email TEXT NOT NULL, "
"delta REAL NOT NULL, "
"reason TEXT NOT NULL, "
"reference TEXT, "
"created_at TEXT DEFAULT (datetime('now'))"
")"
)
# Held credit packs: purchases by NON-members. Money received but credits
# not granted (a membership is required to use credits). Each hold is
# either refunded manually (OC dashboard, 30-day window) or granted
# automatically when the buyer becomes a member.
conn.execute(
"CREATE TABLE IF NOT EXISTS held_credits ("
"id INTEGER PRIMARY KEY AUTOINCREMENT, "
"slug TEXT NOT NULL, "
"email TEXT, "
"name TEXT, "
"amount REAL NOT NULL, "
"oc_reference TEXT, "
"status TEXT NOT NULL DEFAULT 'held', "
"created_at TEXT DEFAULT (datetime('now')), "
"resolved_at TEXT"
")"
)
# Member-managed API keys: add a fingerprint (last-4 of the key, shown in
# lists) and purge the stored plaintext. Member API keys are revealed ONCE
# (at creation) and never again — the industry standard. Revocation uses the
# LiteLLM hash, not the plaintext, so the portal has no reason to keep it.
# (The `members.key_token` chat key is a different table and IS retained in
# plaintext — the injector needs it on every request.)
ak_cols = [r[1] for r in conn.execute("PRAGMA table_info(member_api_keys)").fetchall()]
if "fingerprint" not in ak_cols:
conn.execute("ALTER TABLE member_api_keys ADD COLUMN fingerprint TEXT")
# Backfill fingerprints for any key still holding plaintext (fingerprint =
# last-4 of the key), then purge the plaintext. Idempotent — runs every
# time, affects only rows where plaintext is still present.
conn.execute(
"UPDATE member_api_keys SET fingerprint = substr(sk_token, -4) "
"WHERE (fingerprint IS NULL OR fingerprint = '') AND sk_token != ''"
)
conn.execute("UPDATE member_api_keys SET sk_token = '' WHERE sk_token != ''")
# Commit the DML above — ALTER TABLE auto-commits, but the UPDATEs open an
# implicit transaction that would otherwise roll back when get_db() closes.
conn.commit()
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()
# --- Credits: purchased, non-expiring usage dollars ---
def get_credit_balance(email: str) -> float:
"""Return the member's non-expiring credit balance (USD, default 0)."""
conn = get_db()
row = conn.execute(
"SELECT credit_balance FROM members WHERE email = ?", (email,)
).fetchone()
conn.close()
if row and row[0] is not None:
return float(row[0])
return 0.0
def add_credits(email: str, delta: float, reason: str, reference: str | None = None) -> float:
"""Adjust a member's credit balance and record the movement in the ledger.
Returns the new balance. Negative deltas are allowed (draws against
credits); the balance is clamped at >= 0.
"""
conn = get_db()
conn.execute(
"UPDATE members SET credit_balance = MAX(0, COALESCE(credit_balance, 0) + ?) "
"WHERE email = ?",
(delta, email),
)
conn.execute(
"INSERT INTO credit_ledger (email, delta, reason, reference) VALUES (?, ?, ?, ?)",
(email, delta, reason, reference),
)
row = conn.execute(
"SELECT credit_balance FROM members WHERE email = ?", (email,)
).fetchone()
conn.commit()
conn.close()
return float(row[0]) if row and row[0] is not None else 0.0
def compute_max_budget(allowance: float, credit_balance: float, team_spend: float | None) -> dict:
"""Compute the effective LiteLLM max_budget from the two buckets.
Model: the member has a monthly ALLOWANCE (resets via LiteLLM's
budget_duration) plus a NON-EXPIRING credit balance. Spend draws the
allowance first; once spend exceeds the allowance within a cycle, credits
cover the overflow.
`team_spend` is the member's spend so far this cycle (from LiteLLM).
Credits already consumed this cycle are excluded from the available pool
so the pushed max_budget reflects exactly what's left:
max_budget = allowance + credits_remaining
Returns {"max_budget", "allowance_used", "credits_used"} for dashboards.
"""
if team_spend is None:
# Unknown spend: be permissive (full allowance + full credits) so we
# never wrongly block a member over a metrics gap.
return {
"max_budget": allowance + credit_balance,
"allowance_used": 0.0,
"credits_used": 0.0,
}
allowance_used = min(max(team_spend, 0.0), allowance)
credits_used = max(0.0, team_spend - allowance)
credits_remaining = max(0.0, credit_balance - credits_used)
return {
"max_budget": allowance + credits_remaining,
"allowance_used": allowance_used,
"credits_used": credits_used,
}
async def fetch_members() -> list[dict]:
"""Return the full active-member list: email, name, slug, role, tier.
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).
Membership is tied to the "Membership" TIER — a one-time donation grants
BACKER role on Open Collective but is NOT a membership. So this returns
only users whose tier name contains "membership" (case-insensitive) or
who have no tier (legacy: contributors who predate the tier system).
"""
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 tier { name } 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()
tier = ((n.get("tier") or {}).get("name") or "").strip()
if not email or role not in ("BACKER", "ADMIN"):
continue
# Admins (the co-op itself) always count. BACKERs must be on the
# Membership tier; no-tier BACKERs are legacy contributors from before
# tiers existed — grandfathered as members.
if role == "ADMIN" or "membership" in tier.lower() or tier == "":
result.append({
"email": email,
"name": acct.get("name") or email.split("@")[0],
"slug": acct.get("slug") or "",
"role": role,
"tier": tier,
})
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["Reply-To"] = MAIL_REPLY_TO
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 loomio_fix_usernames(client: httpx.AsyncClient) -> dict:
"""Rename Loomio users with auto-generated handles to their Cloudron username.
Loomio's b2/memberships invite creates the account with a mangled handle
(email local-part + random suffix, e.g. danishipleydsdk4ibdihka). After
first SSO login the user carries an oauth identity whose uid IS their
Cloudron username — prefer that; fall back to the email local-part.
Uses the server (b3) API. Returns a {fixed, skipped} summary.
"""
H = {"Authorization": f"Bearer {LOOMIO_B3_KEY}", "Content-Type": "application/json"}
r = await client.get(f"{LOOMIO_BASE}/api/b3/users", headers=H)
r.raise_for_status()
users = r.json().get("users", [])
fixed, skipped = 0, 0
for u in users:
target = None
for ident in u.get("identities") or []:
if ident.get("identity_type") == "oauth" and ident.get("uid"):
uid = ident["uid"].lower()
# Loomio usernames allow lowercase letters, numbers, underscores and hyphens.
if uid.replace("_", "").replace("-", "").isalnum() and len(uid) >= 2:
target = uid
break
# uid unusable (dots etc) — use email local-part instead
target = (u.get("email") or "").split("@")[0].lower() or None
break
if not target and u.get("email") and "@" in u["email"]:
target = u["email"].split("@")[0].lower()
if not target or u.get("username") == target:
skipped += 1
continue
rr = await client.patch(
f"{LOOMIO_BASE}/api/b3/users/{u['id']}",
headers=H, json={"username": target},
)
if rr.status_code == 200:
fixed += 1
else:
logger.warning("Loomio username fix for %s failed: %s %s",
u.get("email"), rr.status_code, rr.text[:120])
if fixed:
logger.info("Loomio usernames fixed: %s", fixed)
return {"fixed": fixed, "skipped": skipped}
async def sync_matrix_space() -> dict:
"""Invite every active member to the co-op Matrix Space (and General room).
Members log in via Cloudron SSO, so their local MXID is
@<username>:inference.coop. We resolve the Cloudron username for each
active member email and invite them if not already in the space.
Best-effort: missing config just skips; failures log and continue.
"""
if not (MATRIX_BOT_TOKEN and MATRIX_SPACE_ID):
logger.info("Matrix space sync skipped (not configured)")
return {"skipped": True}
H = {"Authorization": f"Bearer {MATRIX_BOT_TOKEN}", "Content-Type": "application/json"}
async def invite(mxid: str, room: str) -> bool:
import asyncio as _aio
for attempt in range(3):
r = await client.post(
f"{MATRIX_HOMESERVER}/_matrix/client/v3/rooms/{room}/invite",
headers=H, json={"user_id": mxid},
)
if r.status_code == 200:
return True
if r.status_code == 429:
wait = (r.json().get("retry_after_ms", 1000) / 1000) + 0.2
await _aio.sleep(wait)
continue
# 403 with 'already in room' is fine
if "already" in r.text:
return True
logger.warning("Matrix invite %s -> %s failed: %s %s",
mxid, room, r.status_code, r.text[:120])
return False
return False
async with httpx.AsyncClient(timeout=30) as client:
# current members of the space (to avoid duplicate invites)
r = await client.get(
f"{MATRIX_HOMESERVER}/_matrix/client/v3/rooms/{MATRIX_SPACE_ID}/joined_members",
headers=H)
joined = set(r.json().get("joined", {}).keys()) if r.status_code == 200 else set()
# also pending invites
r = await client.get(
f"{MATRIX_HOMESERVER}/_matrix/client/v3/rooms/{MATRIX_SPACE_ID}/members?membership=invite",
headers=H)
invited = {m["state_key"] for m in r.json().get("chunk", [])} if r.status_code == 200 else set()
invited_count = 0
async for email, mxid in matrix_member_mxids(client):
if not mxid or mxid in joined or mxid in invited:
continue
await invite(mxid, MATRIX_SPACE_ID)
if MATRIX_GENERAL_ROOM:
await invite(mxid, MATRIX_GENERAL_ROOM)
invited_count += 1
if invited_count:
logger.info("Matrix space invites sent: %s", invited_count)
return {"invited": invited_count}
async def matrix_member_mxids(client: httpx.AsyncClient):
"""Yield (email, mxid) for each active member.
Local members log in via Cloudron SSO and Synapse names them by Cloudron
username → @<username>:inference.coop. Emails that don't resolve to a
Cloudron user (e.g. federated/external members) are skipped here but can
be added via the MATRIX_EXTRA_MXIDS env (comma list of full MXIDs).
"""
extra = [m.strip() for m in os.environ.get("MATRIX_EXTRA_MXIDS", "").split(",") if m.strip()]
seen: set[str] = set()
# Cloudron paginates /api/v1/users at 25/page — walk every page.
users: dict[str, str] = {}
page = 1
while True:
r = await client.get(
f"{CLOUDRON_API}/api/v1/users",
params={"page": page},
headers=cloudron_headers(),
)
if r.status_code != 200:
logger.warning("Cloudron users lookup page %s failed: %s", page, r.status_code)
break
batch = r.json().get("users", [])
for u in batch:
if u.get("email") and u.get("username"):
users[u["email"].lower()] = u["username"]
if len(batch) < 25:
break
page += 1
for email in await get_active_member_emails():
username = users.get(email.lower())
if username:
mxid = f"@{username}:inference.coop"
seen.add(mxid)
yield email, mxid
for mxid in extra:
if mxid not in seen:
yield mxid, mxid
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", []),
)
# Fix auto-generated usernames (email+randomsuffix) → Cloudron username.
# Runs after every sync so newly invited members get corrected within a day
# (or immediately when the sync follows provisioning). Best-effort: no B3
# key or API failure just logs and moves on.
if LOOMIO_B3_KEY:
try:
await loomio_fix_usernames(client)
except Exception as e:
logger.warning("Loomio username fix pass failed: %s", e)
def store_member(email: str, key_token: str, cloudron_user_id: str, slug: str = "") -> None:
conn = get_db()
# Preserve an existing slug if the caller passed none — re-provisioning an
# already-provisioned member must not wipe their OC slug (the deactivation
# webhook maps slugs back to members).
if not slug:
row = conn.execute("SELECT slug FROM members WHERE email = ?", (email,)).fetchone()
if row and row[0]:
slug = row[0]
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_stored_key(email: str) -> str | None:
"""Return the stored LiteLLM key token regardless of active status.
Unlike get_member_key (which filters active=1), this returns the key even
for a lapsed member being reactivated, so provisioning can reuse it rather
than mint a duplicate alias (LiteLLM rejects a duplicate `member:<email>`
alias with 400).
"""
conn = get_db()
row = conn.execute(
"SELECT key_token FROM members WHERE email = ?", (email,)
).fetchone()
conn.close()
return (row[0] if row and row[0] else None)
def touch_welcome_sent(email: str) -> None:
"""Record that a welcome/invite email was just sent for this member."""
conn = get_db()
conn.execute(
"UPDATE members SET welcome_sent_at = datetime('now') WHERE email = ?",
(email,),
)
conn.commit()
conn.close()
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?per_page=100",
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 cloudron_is_activated(user_id: str) -> bool:
"""Return True if the user has completed account setup (inviteAccepted).
`inviteAccepted` flips true once the member clicks their setup link and
chooses a username/password. Until then they're provisioned but not able to
log in — the state we re-invite for.
"""
async with httpx.AsyncClient() as client:
r = await client.get(
f"{CLOUDRON_API}/api/v1/users?per_page=100",
headers=cloudron_headers(),
)
r.raise_for_status()
for u in r.json().get("users", []):
if u.get("id") == user_id:
return bool(u.get("inviteAccepted"))
return False
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 has a stored LiteLLM key, reuse it rather
than minting a new one (LiteLLM rejects a duplicate `member:<email>` alias
with 400 — the bug that previously made re-provisioning of existing members
fail). Re-sending this way also re-sends the welcome email, so it doubles
as the manual "re-invite" path.
"""
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)
# Reuse an existing key if we already have one (the sk- value is stored in
# our DB at creation time and is the authoritative token for this member's
# chat key). Only mint a fresh key when there is none.
key_token = get_stored_key(email)
if not key_token:
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)
touch_welcome_sent(email)
# Subscribe to the member newsletter (single opt-in; members consented by joining).
await listmonk_sync(email, subscribe=True)
# Auto-grant any HELD credit packs for this member (they pre-purchased
# credits before joining — apply them now as part of provisioning).
granted = 0.0
try:
conn = get_db()
rows = conn.execute(
"SELECT id, amount FROM held_credits "
"WHERE status = 'held' AND (LOWER(email) = ? OR LOWER(slug) = ?)",
(email.lower(), slug.lower() if slug else ""),
).fetchall()
if rows:
hold_ids = [r[0] for r in rows]
granted = sum(float(r[1] or 0.0) for r in rows)
conn.execute(
"UPDATE held_credits SET status = 'granted', resolved_at = datetime('now') "
"WHERE id IN (%s)" % ",".join("?" * len(hold_ids)),
hold_ids,
)
conn.commit()
conn.close()
if granted > 0:
new_balance = add_credits(
email, granted,
reason="held credit pack(s) granted on membership",
reference=f"held:{','.join(str(r[0]) for r in rows)}",
)
await sync_team_budget_from_balance(email)
logger.info("Granted %.2f held credits to new member %s (balance %s)",
granted, email, new_balance)
except Exception as e:
logger.warning("Held-credit grant for %s failed: %s", email, e)
await sync_loomio_memberships()
await sync_matrix_space()
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}"
# /team/list has grown slow as usage accumulates; the httpx default 5s
# read timeout was too short (root cause of dashboard 500s on 2026-09-26).
async with httpx.AsyncClient(timeout=httpx.Timeout(connect=15.0, read=60.0, write=15.0, pool=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():
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 litellm_pin_reset_date(team_id: str, reset_at: datetime) -> None:
"""Pin a team's allowance reset to an exact moment (next OC payment date).
LiteLLM's /team/update silently DROPS budget_reset_at (verified v1.74.0),
so this writes the authoritative DB directly. Budgets reset to zero spend
at the pinned instant; the pinned date is honored across restarts.
"""
def _write() -> None:
conn = _pg_connect()
try:
cur = conn.cursor()
cur.execute(
'UPDATE "LiteLLM_TeamTable" SET budget_reset_at = %s WHERE team_id = %s',
(reset_at, team_id),
)
conn.commit()
finally:
conn.close()
await asyncio.get_event_loop().run_in_executor(None, _write)
logger.info("Pinned budget reset for %s to %s", team_id, reset_at.isoformat())
def next_reset_from_now() -> datetime:
"""The member's next allowance reset: 30 days from now, on the hour."""
return datetime.now(timezone.utc).replace(minute=0, second=0, microsecond=0) + timedelta(days=30)
async def sync_team_budget_from_balance(email: str) -> float:
"""Apply the member's effective budget to their LiteLLM team.
Effective budget = monthly allowance + remaining non-expiring credits
(credits drawn only after the allowance is exhausted this cycle). Reads
the member's live spend from LiteLLM to determine how much credit has
already been consumed in the current cycle. Returns the applied budget.
"""
balance = get_member_balance(email)
credits = get_credit_balance(email)
team_id = await litellm_get_or_create_team(email, balance)
# Live spend this cycle (DB path preferred; LiteLLM computes the reset).
team = get_team_budget_from_db(team_id) or {}
team_spend = team.get("spend")
eff = compute_max_budget(balance, credits, team_spend)
await litellm_update_team_budget(team_id, eff["max_budget"])
return eff["max_budget"]
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 ""
# Reveal the key ONCE (in the response). Persist only a fingerprint (last-4)
# and the LiteLLM hash — the plaintext is NOT stored. Revocation uses the
# hash, so the full key never needs to be retained.
fingerprint = sk_token[-4:] if sk_token else ""
conn = get_db()
conn.execute(
"INSERT INTO member_api_keys (email, name, sk_token, hash, fingerprint) VALUES (?, ?, ?, ?, ?) "
"ON CONFLICT(email, name) DO UPDATE SET sk_token='', hash=excluded.hash, fingerprint=excluded.fingerprint, "
"created_at=datetime('now')",
(email, name, "", h, fingerprint),
)
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, fingerprint, created_at FROM member_api_keys WHERE email = ? ORDER BY created_at DESC",
(email,),
).fetchall()
conn.close()
return [
{"name": r[0], "fingerprint": r[1] or "", "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}")
@limiter.limit("30/minute")
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). Default to
# False so a missing field never triggers a spurious re-provision; the
# nightly sweep backstops any missed first payment.
first_payment = data.get("firstPayment", False)
# Membership provision: only on payments to the "Membership" tier. A
# one-time donation (no tier, or any other tier) grants BACKER role on
# OC but is NOT a membership — don't provision it.
#
# Tier detection: the webhook payload carries only `data.order.TierId`
# (numeric) — NO tier name. The order description embeds it
# ("Financial contribution to Inference Cooperative (Credit pack)").
order = data.get("order") or {}
order_desc = order.get("description") or ""
m = re.search(r"\(([^)]+)\)\s*$", order_desc.strip())
tier_name = m.group(1).strip() if m else ""
amount_cents = data.get("amount") or 0
if not amount_cents:
amount_cents = order.get("totalAmount") or 0
# Credit-pack detection: one-time contributions to a "Credit pack"
# tier route to the member's NON-EXPIRING credit balance instead of
# provisioning/allowance. Detection: the OC tier name contains
# "credit" (case-insensitive). Recurring subscriptions never route
# here (a credit pack is a one-time purchase).
if "credit" in tier_name.lower() and slug and event_type == "order.processed":
members = await fetch_members()
email = next((m["email"] for m in members if m["slug"] == slug), None)
# The tier is VARIABLE-amount — no fixed face value. Credits are
# granted on the contribution amount EXCLUDING the OC platform
# tip (totalAmount - tip = amount), which we resolve via OC's
# public order API. Fallback: webhook totalAmount (slightly
# overgrants if a tip was included; logged).
tip_free_cents = await get_order_amount(order.get("idV2"))
if tip_free_cents is not None:
amount = round(tip_free_cents / 100.0, 2)
else:
amount = round(float(amount_cents) / 100.0, 2)
charged = round(float(amount_cents) / 100.0, 2)
if charged > amount:
logger.info(
"Credit pack: charged $%.2f including tip; granting tip-free $%.2f",
charged, amount,
)
oc_ref = f"oc:{slug}:{data.get('id', '')}"
if email:
# Member purchase → credits land immediately.
new_balance = add_credits(
email, amount,
reason="credit pack purchase (Open Collective)",
reference=oc_ref,
)
# Push the enlarged budget immediately.
await sync_team_budget_from_balance(email)
logger.info("Credit pack: +%s credits for %s (balance %s)",
amount, email, new_balance)
return JSONResponse({
"status": "credits_added", "email": email,
"credits_added": amount, "credit_balance": new_balance,
})
# Non-member purchase → HOLD. Record it, email the buyer with
# their two options (refund or join-and-apply). Credits are NOT
# granted: a membership is required to use them.
buyer_name = name or slug
buyer_email = ""
# The webhook payload strips emails; look up the payer via the
# admin token (bot account sees contributor emails).
try:
q = ('{ collective(slug: "%s") { members(limit: 100) { nodes '
'{ account { name slug ... on Individual { email } } } } } }' % OC_COLLECTIVE_SLUG)
async with httpx.AsyncClient(timeout=15.0) as client:
r = await client.post(OC_GRAPHQL_URL,
headers={"Personal-Token": OC_PERSONAL_TOKEN,
"User-Agent": "Mozilla/5.0 (X11; Linux x86_64)"},
json={"query": q})
if r.status_code == 200:
for n in r.json().get("data", {}).get("collective", {}).get("members", {}).get("nodes", []):
if (n.get("account") or {}).get("slug") == slug:
buyer_email = (n["account"].get("email") or "").strip()
buyer_name = n["account"].get("name") or buyer_name
break
except Exception as e:
logger.warning("Credit-pack hold: payer lookup failed for %s: %s", slug, e)
conn = get_db()
conn.execute(
"INSERT INTO held_credits (slug, email, name, amount, oc_reference) "
"VALUES (?, ?, ?, ?, ?)",
(slug, buyer_email or None, buyer_name, amount, oc_ref),
)
conn.commit()
conn.close()
logger.warning(
"Credit pack HELD from non-member %s (%.2f credits); buyer notified",
slug, amount,
)
if buyer_email:
send_email(
buyer_email,
"Inference Cooperative — about your credit pack",
f"Hi {buyer_name},\n\n"
f"Thank you for your ${amount:.2f} credit pack purchase — we've received it.\n\n"
"One thing to flag: credit packs are for members, and our records\n"
"don't show you as one yet. You have two options:\n\n"
"1. Join the co-op ($10-20/month, sliding scale):\n"
" https://opencollective.com/inference-cooperative/contribute\n"
f" Once you're a member, your ${amount:.2f} in credits will be\n"
" applied to your account automatically.\n\n"
"2. Request a refund:\n"
" Reply to this email or write to info@inference.coop and\n"
" we'll refund you in full.\n\n"
"Sorry for the friction — memberships keep the co-op\n"
"member-governed, which is rather the point of the place.\n\n"
"— The Inference Cooperative\n"
"https://inference.coop",
)
else:
logger.warning("Credit-pack hold: no buyer email found for %s; no notification sent", slug)
return JSONResponse({
"status": "held",
"note": "credit pack purchased without membership; hold recorded and buyer notified",
})
if slug and "membership" in tier_name.lower() and event_type == "order.processed":
# Every Membership-tier payment (first or recurring) re-pins the
# member's allowance reset to payment + 30 days — the reset
# follows the member's actual billing rhythm, not a global
# calendar (all members landing on Oct 1 otherwise).
members = await fetch_members()
member = next((m for m in members if m["slug"] == slug), None)
if member:
if first_payment:
try:
await provision_member(member["email"], member["name"], member["slug"])
except Exception as e:
logger.warning("Webhook: failed to provision %s: %s", member["email"], e)
try:
team_id = await litellm_get_or_create_team(
member["email"], get_member_balance(member["email"]))
await litellm_pin_reset_date(team_id, next_reset_from_now())
except Exception as e:
logger.warning("Webhook: failed to pin reset date for %s: %s",
member["email"], e)
elif slug and first_payment and event_type == "order.processed":
# One-time donation (no tier or non-membership tier): NOT a
# membership. Log it so we can see it, but don't provision.
logger.info("One-time donation from slug %s (tier=%r) — not a membership, not provisioning",
slug, tier_name)
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 listmonk_sync(member["email"], subscribe=False)
await sync_loomio_memberships()
await sync_matrix_space()
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")
@limiter.limit("60/minute")
async def list_models(request: Request):
"""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()
# Chat-list filter: exclude audio-only models (STT/TTS). They are
# registered and callable via /v1/audio/* for API users, but they are
# not chat models — listing them in Open WebUI would just confuse
# members. The gateway's model_info carries `mode` for audio models;
# fall back to a name check for safety.
audio_modes = {"audio_transcription", "audio_speech"}
audio_name_hints = ("whisper", "voxtral", "tts")
def is_audio_model(m) -> bool:
mode = ((m.get("info") or {}).get("meta") or {}).get("mode") or ""
if mode in audio_modes:
return True
mid = (m.get("id") or "").lower()
return any(h in mid for h in audio_name_hints)
payload["data"] = [m for m in payload.get("data", []) if not is_audio_model(m)]
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"])
@limiter.limit("120/minute")
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")
@limiter.limit("20/minute")
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()
await sync_matrix_space()
return {"status": "synced"}
@app.post("/admin/provision")
@limiter.limit("20/minute")
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")
@limiter.limit("20/minute")
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/set-credits")
@limiter.limit("20/minute")
async def admin_set_credits(request: Request):
"""Adjust a member's non-expiring credit balance (with ledger entry).
Body: {"email": "...", "add": 10.0, "reason": "goodwill"} → add/subtract
{"email": "...", "set": 0.0, "reason": "correction"} → set absolute
"""
_verify_admin_token(request)
body = await request.json()
email = (body.get("email") or "").strip().lower()
if not email:
raise HTTPException(400, "Missing email")
reason = (body.get("reason") or "admin adjustment").strip()
if "add" in body:
delta = float(body["add"])
new_balance = add_credits(email, delta, reason=reason)
elif "set" in body:
target = float(body["set"])
current = get_credit_balance(email)
new_balance = add_credits(email, target - current, reason=reason)
else:
raise HTTPException(400, "Provide 'add' or 'set'")
applied = await sync_team_budget_from_balance(email)
return {
"status": "updated", "email": email,
"credit_balance": new_balance, "team_budget": applied,
}
@app.get("/admin/credits/{email}")
@limiter.limit("30/minute")
async def admin_credit_ledger(email: str, request: Request):
"""Return a member's credit balance and full ledger (audit trail)."""
_verify_admin_token(request)
conn = get_db()
rows = conn.execute(
"SELECT delta, reason, reference, created_at FROM credit_ledger "
"WHERE email = ? ORDER BY created_at DESC",
(email,),
).fetchall()
conn.close()
return {
"email": email,
"credit_balance": get_credit_balance(email),
"ledger": [
{"delta": r[0], "reason": r[1], "reference": r[2], "created_at": r[3]}
for r in rows
],
}
@app.get("/admin/held-credits")
@limiter.limit("30/minute")
async def admin_held_credits(request: Request):
"""List held credit packs (non-member purchases awaiting refund or grant)."""
_verify_admin_token(request)
conn = get_db()
rows = conn.execute(
"SELECT id, slug, email, name, amount, oc_reference, status, created_at, resolved_at "
"FROM held_credits ORDER BY created_at DESC"
).fetchall()
conn.close()
return {
"held": [
{"id": r[0], "slug": r[1], "email": r[2], "name": r[3], "amount": r[4],
"reference": r[5], "status": r[6], "created_at": r[7], "resolved_at": r[8]}
for r in rows
]
}
@app.post("/admin/held-credits/{hold_id}/refund")
@limiter.limit("20/minute")
async def admin_held_refund(hold_id: int, request: Request):
"""Mark a held credit pack as refunded (the refund itself is done manually
in the Open Collective dashboard — this records it and updates the hold)."""
_verify_admin_token(request)
conn = get_db()
cur = conn.execute(
"UPDATE held_credits SET status = 'refunded', resolved_at = datetime('now') "
"WHERE id = ? AND status = 'held'",
(hold_id,),
)
conn.commit()
updated = cur.rowcount
conn.close()
if not updated:
raise HTTPException(404, "Held credit not found (or already resolved)")
return {"status": "refunded", "id": hold_id}
@app.get("/admin/overview")
@limiter.limit("30/minute")
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?per_page=100",
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)
# All-time spend per member team, from SpendLogs (the team `spend` counter
# is cycle-only and resets at each member's budget_reset_at — the panel
# needs both to avoid the "month change wiped the numbers" confusion).
all_time_spend: dict[str, float] = {}
try:
conn = _pg_connect()
try:
with conn.cursor() as cur:
cur.execute(
"""
SELECT t.team_alias, COALESCE(SUM(s.spend), 0)
FROM "LiteLLM_TeamTable" t
LEFT JOIN "LiteLLM_SpendLogs" s ON s."team_id" = t.team_id
WHERE t.team_alias LIKE 'member:%'
GROUP BY 1
"""
)
all_time_spend = {alias: float(v or 0) for alias, v in cur.fetchall()}
finally:
conn.close()
except Exception as e:
logger.warning("admin overview: all-time spend query failed: %s", e)
members = []
total_spend = 0.0
total_all_time = 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
member_all_time = all_time_spend.get(f"member:{email}", 0.0)
total_all_time += member_all_time
credits = get_credit_balance(email)
eff = compute_max_budget(
float(balance) if balance is not None else MEMBER_BUDGET,
credits,
spend,
)
members.append(
{
"email": email,
"all_time_spend": round(member_all_time, 2),
"username": usernames.get(user_id),
"activated": invite_accepted.get(user_id, False),
"active": bool(active),
"manual": email.lower() in MANUAL_MEMBERS,
"balance": float(balance) if balance is not None else MEMBER_BUDGET,
"credit_balance": credits,
"allowance_used": eff["allowance_used"],
"credits_used": eff["credits_used"],
"spend": spend,
"remaining": max(eff["max_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": round(total_spend, 2),
"total_all_time_spend": round(total_all_time, 2),
},
}
@app.get("/admin/model-usage")
@limiter.limit("30/minute")
async def admin_model_usage(request: Request, days: int = 0):
"""Co-op-wide per-model usage from LiteLLM spend logs.
Authenticated by the X-Admin-Token header. Aggregates spend-log rows
(across ALL members) by model: total spend, tokens, and request count,
sorted by spend descending. Optional ?days=N limits to the last N days
(0 or omitted = all-time). No per-member breakdown — this is co-op-level
so admins can see which models the membership actually uses.
Reads go straight to LiteLLM's Postgres when configured (fast, bounded);
otherwise falls back to the HTTP API path.
"""
_verify_admin_token(request)
# Fast path: direct Postgres aggregation.
if litellm_db_ready():
try:
return aggregate_model_usage(days=days)
except Exception as e:
logger.warning("admin model-usage: DB path failed (%s); falling back to HTTP", e)
# Fallback path: LiteLLM HTTP API.
cutoff = None
if days and days > 0:
from datetime import datetime, timedelta, timezone
cutoff = (datetime.now(timezone.utc) - timedelta(days=days)).strftime(
"%Y-%m-%dT%H:%M:%S"
)
models = {}
total_spend = 0.0
total_tokens = 0
total_requests = 0
try:
async with httpx.AsyncClient(timeout=30.0) as client:
r = await client.get(
f"{LITELLM_BASE}/spend/logs",
headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"},
)
r.raise_for_status()
rows = r.json()
if isinstance(rows, dict):
rows = rows.get("data", rows.get("logs", []))
for row in rows if isinstance(rows, list) else []:
# Prefer model_group: it carries our provider/model alias (e.g.
# greenpt/green-r, publicai/apertus-v1.5-8b) when present; the
# raw `model` field holds messy upstream/internal names. Older
# rows predating the alias convention have model_group empty,
# so fall back to `model` for those.
if cutoff:
started = str(row.get("startTime") or "")
if started and started < cutoff:
continue
model = row.get("model_group") or row.get("model") or "unknown"
m = models.setdefault(model, {"spend": 0.0, "tokens": 0, "requests": 0})
m["spend"] += float(row.get("spend") or 0.0)
m["tokens"] += int(row.get("total_tokens") or 0)
m["requests"] += 1
total_spend += float(row.get("spend") or 0.0)
total_tokens += int(row.get("total_tokens") or 0)
total_requests += 1
except Exception as e:
logger.warning("admin model-usage: spend/logs fetch failed: %s", e)
model_list = [
{"model": name, **stats}
for name, stats in sorted(models.items(), key=lambda kv: kv[1]["spend"], reverse=True)
]
return {
"models": model_list,
"totals": {
"total_spend": total_spend,
"total_tokens": total_tokens,
"total_requests": total_requests,
},
}
# --- 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, Cloudron username, or uid) to email.
Cloudron's OIDC `sub` claim is the user ID (uid-…), and email-invited users
have no username at all. So accept three identity shapes:
- email ("@") → return as-is
- uid-… → match on the user id
- username → match on the username field (legacy accounts)
"""
identity = identity.strip()
if "@" in identity:
return identity.lower()
# uid / username → email via Cloudron user list (one fetch covers both).
import urllib.request
req = urllib.request.Request(
f"{CLOUDRON_API}/api/v1/users?per_page=100",
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("id") == identity:
return (u.get("email") or "").lower()
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")
@limiter.limit("30/minute")
async def broker_list(request: Request):
email = _verify_broker_secret(request)
return {"keys": await broker_list_api_keys(email)}
@app.post("/broker/keys")
@limiter.limit("30/minute")
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}")
@limiter.limit("30/minute")
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, days: int = 0) -> dict:
"""Return the member's spend vs. balance, plus per-model and token detail.
This is the read-only primitive the member dashboard uses (no master key on
the client). Reads go straight to LiteLLM's Postgres when configured (fast,
bounded); otherwise falls back to the HTTP API path. `days` bounds the
per-model breakdown window (0 = all time); balance/spend always come from
the team's live totals.
"""
team_id = await litellm_get_or_create_team(email, get_member_balance(email))
allowance = get_member_balance(email)
credits = get_credit_balance(email)
spend = 0.0
reset_at = None
# Fast path: direct Postgres reads.
if litellm_db_ready():
try:
team = get_team_budget_from_db(team_id)
if team:
if team["max_budget"] is not None:
balance = float(team["max_budget"])
spend = float(team["spend"] or 0.0)
reset_at = team["budget_reset_at"]
usage = aggregate_model_usage(days=days, team_id=team_id)
# All-time spend for this member: total of their spend-log rows
# (the team `spend` counter resets with their allowance period).
all_time_spend = 0.0
try:
conn = _pg_connect()
try:
with conn.cursor() as cur:
cur.execute(
'SELECT COALESCE(SUM(spend), 0) FROM "LiteLLM_SpendLogs" WHERE team_id = %s',
(team_id,),
)
all_time_spend = float(cur.fetchone()[0] or 0.0)
finally:
conn.close()
except Exception as e:
logger.warning("broker_get_usage: all-time spend query failed: %s", e)
model_list = [
{"model": m["model"], "spend": m["spend"], "tokens": m["tokens"],
"calls": m["requests"]}
for m in usage["models"]
]
return {
"email": email,
"balance": balance,
"allowance": allowance,
"credits": credits,
"spend": spend,
"all_time_spend": round(all_time_spend, 2),
"remaining": max(balance - spend, 0.0),
"total_tokens": usage["totals"]["total_tokens"],
"reset_at": reset_at,
"models": model_list,
}
except Exception as e:
logger.warning("broker_get_usage: DB path failed (%s); falling back to HTTP", e)
# Fallback path: LiteLLM HTTP API.
# /team/list and /spend/logs have grown slow as usage accumulates; the
# httpx default 5s read timeout was too short and made /broker/usage 500.
client_timeout = httpx.Timeout(connect=15.0, read=60.0, write=15.0, pool=15.0)
async with httpx.AsyncClient(timeout=client_timeout) 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
# Prefer model_group: it carries our provider/model alias (e.g.
# greenpt/green-r) when present; the raw `model` field holds
# messy upstream/internal names. Older rows predating the alias
# convention have model_group empty, so fall back to `model`.
model = row.get("model_group") or 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)
# All-time spend from the same log rows (team counter is period-only).
all_time_spend = round(sum(m["spend"] for m in models.values()), 2)
# 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,
"allowance": allowance,
"credits": credits,
"spend": spend,
"all_time_spend": all_time_spend,
"remaining": max(balance - spend, 0.0),
"total_tokens": total_tokens,
"reset_at": reset_at,
"models": model_list,
}
@app.get("/broker/usage")
@limiter.limit("60/minute")
async def broker_usage(request: Request, days: int = 0):
email = _verify_broker_secret(request)
return await broker_get_usage(email, days=days)
@app.get("/broker/model-leaderboard")
@limiter.limit("30/minute")
async def broker_model_leaderboard(request: Request, days: int = 0):
"""Co-op-wide model popularity for the member dashboard.
Same aggregation as the admin model-usage view (spend/tokens/requests per
model, across all members) but served through the broker so members can
see which models the co-op uses most — social guidance for model choice.
No per-member data; aggregate totals only.
"""
_verify_broker_secret(request)
if litellm_db_ready():
try:
usage = aggregate_model_usage(days=days)
except Exception as e:
logger.warning("broker leaderboard: DB path failed (%s)", e)
usage = None
if usage:
return {
"models": [
{"model": m["model"], "spend": m["spend"],
"tokens": m["tokens"], "requests": m["requests"]}
for m in usage["models"][:10]
]
}
# Fallback: empty (rare; the DB path is the norm)
return {"models": []}
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}
active_by_email = {m["email"]: m 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)
await listmonk_sync(email, subscribe=False)
deactivated_count += 1
logger.info("Sweep: deactivated %s", email)
except Exception as e:
logger.warning("Sweep: failed to deactivate %s: %s", email, e)
# 2b. Re-invite never-activated members (provisioned but no account setup).
# A member who was provisioned (e.g. webhook fired while the OC token was
# missing the `email` scope, or the invite email went to spam) may sit
# active-but-unactivated indefinitely. Re-send their welcome email, but
# throttle by REINVITE_AFTER_DAYS so we don't spam them nightly.
reinvited_count = 0
conn = get_db()
reinvite_rows = conn.execute(
"SELECT email, cloudron_user_id, welcome_sent_at FROM members WHERE active = 1"
).fetchall()
conn.close()
now = datetime.now(timezone.utc)
for email, user_id, sent_at in reinvite_rows:
if not user_id:
continue
# Skip if already activated.
try:
if await cloudron_is_activated(user_id):
continue
except Exception as e:
logger.warning("Sweep: activation check failed for %s: %s", email, e)
continue
# Throttle: only re-invite if the last send was > REINVITE_AFTER_DAYS ago.
if sent_at:
try:
last = datetime.fromisoformat(sent_at)
# stored as UTC naive (datetime('now') is UTC); attach tz
if last.tzinfo is None:
last = last.replace(tzinfo=timezone.utc)
if (now - last).total_seconds() < REINVITE_AFTER_DAYS * 86400:
continue
except Exception:
pass
# Re-send the welcome email (reuses the stored key — no duplicate alias).
try:
name = active_by_email.get(email, {}).get("name") or email.split("@")[0]
await provision_member(email, name, "")
reinvited_count += 1
logger.info("Sweep: re-invited %s (never activated)", email)
except Exception as e:
logger.warning("Sweep: failed to re-invite %s: %s", email, e)
# 3. Sync Loomio.
await sync_loomio_memberships()
await sync_matrix_space()
return {
"status": "reconciled",
"provisioned": provisioned_count,
"deactivated": deactivated_count,
"reinvited": reinvited_count,
}
@app.post("/admin/sweep")
@limiter.limit("20/minute")
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"}