2599 lines
108 KiB
Python
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> · 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"}
|