""" 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, timezone from email.mime.text import MIMEText from email.mime.multipart import MIMEMultipart import httpx import jwt from fastapi import FastAPI, Request, Response, HTTPException from fastapi.responses import JSONResponse, StreamingResponse 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") # 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", "") 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 :). 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 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") return conn def get_member_balance(email: str) -> float: """Return the member's current balance (USD), defaulting to MEMBER_BUDGET. The balance is the flexible allowance — it starts at MEMBER_BUDGET but can be adjusted (topped up by payments) without changing code. """ conn = get_db() row = conn.execute( "SELECT balance FROM members WHERE email = ?", (email,) ).fetchone() conn.close() if row and row[0] is not None: return float(row[0]) return MEMBER_BUDGET def set_member_balance(email: str, balance: float) -> None: """Set the member's balance (USD). Used by top-up / billing adjustments.""" conn = get_db() conn.execute( "UPDATE members SET balance = ? WHERE email = ?", (balance, email) ) conn.commit() conn.close() async def fetch_members() -> list[dict]: """Return the full active-member list: email, name, slug, role. Uses the bot's personal token (admin) so emails are exposed. This is the authoritative source of truth for reconciliation (provisioning missing members, deactivating lapsed ones). """ if not OC_PERSONAL_TOKEN: logger.warning("OC_PERSONAL_TOKEN not set; cannot fetch members") return [] query = ( '{ collective(slug: "%s") { members(limit: 100) { nodes { role account { name slug ... on Individual { email } } } } } }' % OC_COLLECTIVE_SLUG ) async with httpx.AsyncClient() as client: r = await client.post( OC_GRAPHQL_URL, headers={"Personal-Token": OC_PERSONAL_TOKEN}, json={"query": query}, ) if r.status_code != 200: logger.warning("OC member fetch failed: %s", r.status_code) return [] nodes = r.json().get("data", {}).get("collective", {}).get("members", {}).get("nodes", []) result: list[dict] = [] for n in nodes: role = n.get("role", "") acct = n.get("account", {}) or {} email = (acct.get("email") or "").strip().lower() if email and role in ("BACKER", "ADMIN"): result.append({ "email": email, "name": acct.get("name") or email.split("@")[0], "slug": acct.get("slug") or "", "role": role, }) return result def send_email(to: str, subject: str, text_body: str, html_body: str | None = None) -> bool: """Send an email via Cloudron's sendmail addon (SMTP). Returns True on success, False on failure (logged, not raised — email is best-effort and must never break the provisioning flow). """ if not SMTP_SERVER: logger.warning("SMTP not configured; skipping email to %s", to) return False try: msg = MIMEMultipart("alternative") msg["From"] = f"{MAIL_FROM_NAME} <{MAIL_FROM}>" msg["To"] = to msg["Subject"] = subject msg["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"""

Welcome, {name}

Thanks for joining the Inference Cooperative.

Finish setting up your account to start using private AI chat.

Finish setup

This takes about a minute.

Questions? Reply to this email or write to info@inference.coop.
The Inference Cooperative · inference.coop
By creating your account, you agree to our Terms of Service and Privacy Policy.
""" return subject, text, html async def get_active_member_emails() -> list[str]: """Return the list of active member emails from the portal's own database. The local `members` table is the authoritative source of member emails (populated during provisioning). Used for the Loomio sync. """ conn = get_db() rows = conn.execute( "SELECT email FROM members WHERE active = 1" ).fetchall() conn.close() return [r[0] for r in rows] async def sync_loomio_memberships() -> None: """Reconcile the Loomio group to the current active-member list. Uses POST /api/b2/memberships with remove_absent=1, which adds new members and removes anyone not in the list. Operator accounts (LOOMIO_ALWAYS_KEEP) are always included so they are never removed. """ if not LOOMIO_API_KEY or not LOOMIO_GROUP_ID: logger.warning("Loomio not configured; skipping membership sync") return emails = await get_active_member_emails() # Always keep operator accounts in the group. for keep in LOOMIO_ALWAYS_KEEP: if keep not in emails: emails.append(keep) async with httpx.AsyncClient() as client: r = await client.post( f"{LOOMIO_BASE}/api/b2/memberships", headers={"Authorization": f"Bearer {LOOMIO_API_KEY}"}, json={"group_id": int(LOOMIO_GROUP_ID), "emails": emails, "remove_absent": 1}, ) if r.status_code != 200: logger.warning("Loomio membership sync failed: %s %s", r.status_code, r.text[:200]) return result = r.json() logger.info( "Loomio sync: added=%s removed=%s", result.get("added_emails", []), result.get("removed_emails", []), ) def store_member(email: str, key_token: str, cloudron_user_id: str, slug: str = "") -> None: conn = get_db() # 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:` 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", 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", 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:` 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) await sync_loomio_memberships() logger.info("Provisioned member %s (user_id=%s)", email, user_id) return user_id async def litellm_get_or_create_team(email: str, budget: float) -> str: """Return the team_id for a member's budget team, creating it if needed. Each member gets a LiteLLM "team" (a budget container) with a $15/30d budget. The team holds the chat key + any member-managed API keys, so all spend draws from the one pool. Idempotent: looks up the existing team by alias before creating. Note that /team/new is NOT idempotent (it creates a new team each call), so we must check /team/list first. """ alias = f"member:{email}" async with httpx.AsyncClient() as client: r = await client.get( f"{LITELLM_BASE}/team/list", headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"}, ) r.raise_for_status() for t in r.json(): if t.get("team_alias") == alias: return t["team_id"] r = await client.post( f"{LITELLM_BASE}/team/new", headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"}, json={ "team_alias": alias, "max_budget": budget, "budget_duration": "30d", "models": MEMBER_MODELS, }, ) r.raise_for_status() return r.json()["team_id"] async def litellm_update_team_budget(team_id: str, budget: float) -> None: """Set a team's max_budget (the enforced spend cap). This is the primitive behind "payment → add to balance": a top-up adjusts the team's max_budget, which every key under the team (chat + API) draws from. Keys inherit the team budget (their own max_budget is ignored when a team_id is set), so the team is the single source of truth for allowance. """ async with httpx.AsyncClient() as client: r = await client.post( f"{LITELLM_BASE}/team/update", headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"}, json={"team_id": team_id, "max_budget": budget}, ) r.raise_for_status() async def sync_team_budget_from_balance(email: str) -> float: """Apply the member's stored balance to their team's max_budget. Reads the balance from the members table and pushes it to LiteLLM. Returns the applied budget. This keeps the portal DB and LiteLLM in sync whenever a balance changes (top-up, reset, pricing-model change). """ balance = get_member_balance(email) team_id = await litellm_get_or_create_team(email, balance) await litellm_update_team_budget(team_id, balance) return balance async def litellm_create_key(email: str, budget: float, team_id: str) -> str: """Create a fresh LiteLLM virtual key under the member's team. Returns the ACTUAL sk- key (only revealed at generation time). We always create a new key rather than reusing — the sk- value is unrecoverable from a stored hash, so "reuse" would store a useless hash (a bug that previously broke chat for three members). NOTE: the key does NOT set its own max_budget — it inherits the team's budget (LiteLLM ignores a key's max_budget when team_id is set anyway). The team is the single source of truth for allowance, so a balance change (top-up) updates the team once and every key follows. The `budget` arg is retained only for the team-creation fallback in the migration path. """ alias = f"member:{email}" async with httpx.AsyncClient() as client: r = await client.post( f"{LITELLM_BASE}/key/generate", headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"}, json={ "key_alias": alias, "team_id": team_id, "models": MEMBER_MODELS, }, ) r.raise_for_status() return r.json().get("key", "") async def litellm_disable_key(key_token: str) -> None: """Delete a member's LiteLLM key (on payment lapse). /key/delete accepts either the sk- key or its SHA256 hash, so this works for both current keys and legacy hash-keyed members. """ if not key_token: return async with httpx.AsyncClient() as client: await client.post( f"{LITELLM_BASE}/key/delete", headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"}, json={"keys": [key_token]}, ) # --- Broker: member-managed API keys (scoped to the member's team) --- def _api_key_alias(email: str, name: str) -> str: return f"api:{email}:{name}" async def broker_create_api_key(email: str, name: str) -> dict: """Create a member-managed API key under the member's team. Returns {name, sk_token, created_at}. The key draws from the member's team budget and is tracked in member_api_keys so it can be listed/revoked by name. """ if not name or not re.fullmatch(r"[A-Za-z0-9._-]{1,64}", name): raise HTTPException(400, "Invalid key name (1-64 chars, alphanumeric/._-)") balance = get_member_balance(email) team_id = await litellm_get_or_create_team(email, balance) alias = _api_key_alias(email, name) async with httpx.AsyncClient() as client: r = await client.post( f"{LITELLM_BASE}/key/generate", headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"}, json={ "key_alias": alias, "team_id": team_id, "models": MEMBER_MODELS, }, ) r.raise_for_status() data = r.json() sk_token = data.get("key", "") h = data.get("token") or data.get("token_id") or "" conn = get_db() conn.execute( "INSERT INTO member_api_keys (email, name, sk_token, hash) VALUES (?, ?, ?, ?) " "ON CONFLICT(email, name) DO UPDATE SET sk_token=excluded.sk_token, hash=excluded.hash, " "created_at=datetime('now')", (email, name, sk_token, h), ) conn.commit() row = conn.execute( "SELECT created_at FROM member_api_keys WHERE email = ? AND name = ?", (email, name), ).fetchone() conn.close() return {"name": name, "sk_token": sk_token, "created_at": row[0] if row else None} async def broker_list_api_keys(email: str) -> list[dict]: conn = get_db() rows = conn.execute( "SELECT name, sk_token, created_at FROM member_api_keys WHERE email = ? ORDER BY created_at DESC", (email,), ).fetchall() conn.close() return [ {"name": r[0], "sk_token": r[1], "created_at": r[2]} for r in rows ] async def broker_revoke_api_key(email: str, name: str) -> None: conn = get_db() row = conn.execute( "SELECT hash FROM member_api_keys WHERE email = ? AND name = ?", (email, name) ).fetchone() if not row: conn.close() raise HTTPException(404, "Key not found") conn.execute("DELETE FROM member_api_keys WHERE email = ? AND name = ?", (email, name)) conn.commit() conn.close() # Delete the underlying LiteLLM key by its hash. await litellm_disable_key(row[0]) # --- A. Open Collective webhook --- @app.post("/webhook/opencollective/{token}") @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) if slug and first_payment: members = await fetch_members() for m in members: if m["slug"] == slug: try: await provision_member(m["email"], m["name"], m["slug"]) except Exception as e: logger.warning("Webhook: failed to provision %s: %s", m["email"], e) break return JSONResponse({"status": "provisioned", "name": name, "slug": slug}) if event_type in ("collective.member.deleted", "collective.transaction.deleted"): # Lapsed member → move to the inactive group (lose access, keep data). # The webhook gives us the slug; we map it to the member's email/user. if slug: member = get_member_by_slug(slug) if member: # Move to inactive group (revokes chat access, keeps account + data) if INACTIVE_GROUP_ID and member.get("cloudron_user_id"): await cloudron_set_group(member["cloudron_user_id"], INACTIVE_GROUP_ID) # Mark inactive in our DB (key injector will refuse requests) deactivate_member(member["email"]) await listmonk_sync(member["email"], subscribe=False) await sync_loomio_memberships() logger.info("Deactivated member %s (slug=%s)", member["email"], slug) return JSONResponse({"status": "deactivated", "email": member["email"]}) return JSONResponse({"status": "deactivated", "note": "no matching member"}) return JSONResponse({"status": "ignored", "type": event_type}) # --- B. Key injector (proxy) --- @app.get("/v1/models") @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() for model in payload.get("data", []): model.setdefault("info", {}).setdefault("meta", {}) model["info"]["meta"].setdefault("capabilities", {})["web_search"] = True model["info"]["meta"].setdefault("defaultFeatureIds", []) if "web_search" not in model["info"]["meta"]["defaultFeatureIds"]: model["info"]["meta"]["defaultFeatureIds"].append("web_search") return JSONResponse(content=payload) @app.api_route("/v1/{path:path}", methods=["GET", "POST", "PUT", "DELETE", "PATCH"]) async def inject_key(request: Request, path: str): """Read the member's email from a header, inject their key, forward to LiteLLM.""" # Verify the request came through the chat frontend, so a direct caller # can't spoof the user-email header and use another member's key. # # Two supported auth paths: # 1. OpenWebUI signs the user identity as a JWT (X-OpenWebUI-User-Jwt) # with a shared secret — the email is tamper-proof, so no separate # secret header is needed. # 2. LibreChat sends a plain X-User-Email header plus an X-Portal-Secret # shared-secret header. email = "" jwt_header = request.headers.get("x-openwebui-user-jwt", "") if jwt_header and OWUI_JWT_SECRET: try: claims = jwt.decode(jwt_header, OWUI_JWT_SECRET, algorithms=["HS256"]) email = claims.get("email", "") except jwt.PyJWTError: raise HTTPException(401, "Invalid user identity token") else: _verify_portal_secret(request) email = ( request.headers.get("x-user-email", "") or request.headers.get("x-openwebui-user-email", "") ) if not email: raise HTTPException(401, "No member identity (user-email header)") # Look up the member's key from the persistent store member_key = get_member_key(email) if not member_key: raise HTTPException(403, "No active membership key") # Forward the request to LiteLLM with the member's key, STREAMING the # response back token-by-token (or byte-by-byte) rather than buffering it. # Buffering collapsed the model's stream and handed OpenWebUI the full # answer in one blob, so members saw nothing until generation finished — # which made chat feel slow even though time-to-first-byte was fine. body = await request.body() headers = dict(request.headers) headers["authorization"] = f"Bearer {member_key}" headers.pop("host", None) headers.pop("content-length", None) headers.pop("accept-encoding", None) # let httpx handle decompression timeout = httpx.Timeout(connect=15.0, read=600.0, write=30.0, pool=15.0) # Send the request upstream, but don't buffer the body — use a stream. client = httpx.AsyncClient(timeout=timeout) upstream_req = client.build_request( method=request.method, url=f"{LITELLM_BASE}/v1/{path}", headers=headers, content=body, ) upstream = await client.send(upstream_req, stream=True) async def gen(): try: async for chunk in upstream.aiter_raw(): yield chunk finally: await upstream.aclose() await client.aclose() # Pass through the content-type (critical for SSE streaming) and status. return StreamingResponse( gen(), status_code=upstream.status_code, headers={ "content-type": upstream.headers.get("content-type", "application/json"), "cache-control": "no-cache", }, ) # --- D. Admin / health --- @app.get("/health") async def health(): return {"status": "ok"} @app.post("/admin/sync-loomio") @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() 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.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", headers={"Authorization": f"Bearer {CLOUDRON_TOKEN}"}, ) r.raise_for_status() for u in r.json().get("users", []): uid = u.get("id") if u.get("username"): usernames[uid] = u["username"] invite_accepted[uid] = bool(u.get("inviteAccepted")) except Exception as e: logger.warning("admin overview: cloudron users fetch failed: %s", e) # Fetch team spend/budget from LiteLLM once. teams = {} try: async with httpx.AsyncClient(timeout=15.0) as client: r = await client.get( f"{LITELLM_BASE}/team/list", headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"}, ) r.raise_for_status() for t in r.json(): teams[t.get("team_alias")] = t except Exception as e: logger.warning("admin overview: team/list failed: %s", e) members = [] total_spend = 0.0 for email, active, balance, user_id in rows: team = teams.get(f"member:{email}", {}) spend = float(team.get("spend") or 0.0) total_spend += spend members.append( { "email": email, "username": usernames.get(user_id), "activated": invite_accepted.get(user_id, False), "active": bool(active), "balance": float(balance) if balance is not None else MEMBER_BUDGET, "spend": spend, "remaining": max((balance or MEMBER_BUDGET) - spend, 0.0), "reset_at": team.get("budget_reset_at"), "cloudron_user_id": user_id, } ) return { "members": members, "totals": { "members": len(members), "active": sum(1 for m in members if m["active"]), "total_spend": total_spend, }, } # --- Broker HTTP endpoints (called by the member dashboard) --- # The dashboard is Cloudron-SSO-gated; it authenticates the member and passes # their email with a shared secret. The portal trusts the email only because it # carries BROKER_SECRET. These endpoints expose ONLY the authenticated member's # own keys — never another member's, and never the chat key. def _verify_broker_secret(request: Request) -> str: """Return the authenticated member's email, or 401. The member dashboard sends X-Broker-Secret (shared secret) + X-Member-Email (or X-Member-User, a Cloudron username). Both must be present and the secret must match, or we refuse (fail closed). Cloudron's proxyAuth injects the USERNAME (e.g. "ntnsndr"), not the email. So when the dashboard sends a username (no "@"), we resolve it to an email via the Cloudron user list (we already hold a Cloudron admin token). """ if not BROKER_SECRET: raise HTTPException(503, "Broker secret not configured") provided = request.headers.get("x-broker-secret", "") if not secrets.compare_digest(provided, BROKER_SECRET): raise HTTPException(401, "Invalid broker secret") identity = ( request.headers.get("x-member-email") or request.headers.get("x-member-user") or "" ).strip() if not identity: raise HTTPException(401, "Missing member identity") email = resolve_member_email(identity) # The member must be active in our DB before they can manage keys. if not get_member_key(email): raise HTTPException(403, "No active membership") return email def resolve_member_email(identity: str) -> str: """Resolve a member identity (email OR Cloudron username) to their email. If it looks like an email (contains "@"), return it lowercased. Otherwise treat it as a Cloudron username and look up the corresponding email from the Cloudron user list (synchronous, using urllib since this runs outside the async request path of httpx). """ identity = identity.strip() if "@" in identity: return identity.lower() # Username → email via Cloudron user list. import urllib.request req = urllib.request.Request( f"{CLOUDRON_API}/api/v1/users", headers={"Authorization": f"Bearer {CLOUDRON_TOKEN}"}, ) with urllib.request.urlopen(req, timeout=15) as r: data = json.loads(r.read().decode()) users = data.get("users", []) for u in users: if u.get("username") == identity: return (u.get("email") or "").lower() raise HTTPException(403, f"No Cloudron user for identity {identity!r}") @app.get("/broker/keys") @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) -> dict: """Return the member's spend vs. balance, plus per-model and token detail. Reads the member's team (spend + max_budget) and aggregates their recent spend logs by model for token counts. This is the read-only primitive the member dashboard uses (no master key on the client). """ team_id = await litellm_get_or_create_team(email, get_member_balance(email)) balance = get_member_balance(email) spend = 0.0 reset_at = None async with httpx.AsyncClient() as client: # 1. Team spend + budget r = await client.get( f"{LITELLM_BASE}/team/list", headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"}, ) r.raise_for_status() for t in r.json(): if t.get("team_id") == team_id: balance = float(t.get("max_budget") or balance) spend = float(t.get("spend") or 0.0) reset_at = t.get("budget_reset_at") break # 2. Per-model token/spend breakdown from spend logs (filtered by team). # /spend/logs does not reliably filter by team_id server-side, so we # filter client-side on the returned rows. models = {} total_tokens = 0 try: r2 = await client.get( f"{LITELLM_BASE}/spend/logs", headers={"Authorization": f"Bearer {LITELLM_MASTER_KEY}"}, ) r2.raise_for_status() rows = r2.json() if isinstance(rows, dict): rows = rows.get("data", rows.get("logs", [])) for row in rows if isinstance(rows, list) else []: if row.get("team_id") != team_id: continue model = row.get("model") or "unknown" m = models.setdefault(model, {"spend": 0.0, "tokens": 0, "calls": 0}) m["spend"] += float(row.get("spend") or 0.0) m["tokens"] += int(row.get("total_tokens") or 0) m["calls"] += 1 total_tokens += int(row.get("total_tokens") or 0) except Exception as e: # spend logs are best-effort; don't fail the whole view logger.warning("spend/logs fetch failed: %s", e) # Sort models by spend descending (which models are draining the most). model_list = [ {"model": name, **stats} for name, stats in sorted(models.items(), key=lambda kv: kv[1]["spend"], reverse=True) ] return { "email": email, "balance": balance, "spend": spend, "remaining": max(balance - spend, 0.0), "total_tokens": total_tokens, "reset_at": reset_at, "models": model_list, } @app.get("/broker/usage") @limiter.limit("60/minute") async def broker_usage(request: Request): email = _verify_broker_secret(request) return await broker_get_usage(email) async def reconcile_memberships() -> dict: """Reconcile the portal against the live Open Collective member list. - Provisions members who are BACKER/ADMIN on OC but not yet in our DB (missed webhook, or contributed before the email step existed). - Deactivates members who are in our DB but no longer BACKER/ADMIN (lapsed or cancelled membership). - Syncs Loomio to the resulting active-member set. Idempotent and safe to run repeatedly. Returns a summary of actions taken. """ active = await fetch_members() active_emails = {m["email"] for m in active} 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() 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"}