Add email notifications (sendmail) + membership sweep (scheduler): welcome email on webhook, reconcile endpoint, daily sweep
This commit is contained in:
1 parent
12d36b8244
commit
c1866df2a4
4 files changed
+245
-1
No files matched your search
@@ -11,6 +11,13 @@
|
|||||||
"oidc": {
|
"oidc": {
|
||||||
"loginRedirectUri": "/",
|
"loginRedirectUri": "/",
|
||||||
"logoutRedirectUri": "/"
|
"logoutRedirectUri": "/"
|
||||||
|
},
|
||||||
|
"sendmail": {},
|
||||||
|
"scheduler": {
|
||||||
|
"membership_sweep": {
|
||||||
|
"schedule": "0 3 * * *",
|
||||||
|
"command": "/app/code/sweep.sh"
|
||||||
|
}
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
"manifestVersion": 2,
|
"manifestVersion": 2,
|
||||||
|
|||||||
@@ -8,6 +8,8 @@ RUN pip install --no-cache-dir -r requirements.txt
|
|||||||
|
|
||||||
# Copy app
|
# Copy app
|
||||||
COPY app/ ./app/
|
COPY app/ ./app/
|
||||||
|
COPY sweep.sh /app/code/sweep.sh
|
||||||
|
RUN chmod +x /app/code/sweep.sh
|
||||||
|
|
||||||
# Run as root (the base image default). Cloudron's sandboxing (read-only
|
# Run as root (the base image default). Cloudron's sandboxing (read-only
|
||||||
# rootfs, AppArmor, dropped capabilities) still applies to root inside the
|
# rootfs, AppArmor, dropped capabilities) still applies to root inside the
|
||||||
|
|||||||
+201
-1
@@ -18,6 +18,9 @@ import json
|
|||||||
import logging
|
import logging
|
||||||
import sqlite3
|
import sqlite3
|
||||||
import secrets
|
import secrets
|
||||||
|
import smtplib
|
||||||
|
from email.mime.text import MIMEText
|
||||||
|
from email.mime.multipart import MIMEMultipart
|
||||||
|
|
||||||
import httpx
|
import httpx
|
||||||
import jwt
|
import jwt
|
||||||
@@ -102,6 +105,17 @@ LOOMIO_ALWAYS_KEEP = [
|
|||||||
# Lives in /app/data (Cloudron localstorage addon persists this).
|
# Lives in /app/data (Cloudron localstorage addon persists this).
|
||||||
DB_PATH = os.environ.get("DB_PATH", "/app/data/members.db")
|
DB_PATH = os.environ.get("DB_PATH", "/app/data/members.db")
|
||||||
|
|
||||||
|
# --- Email (Cloudron sendmail addon) ---
|
||||||
|
# The sendmail addon exports SMTP credentials as CLOUDRON_MAIL_* env vars.
|
||||||
|
SMTP_SERVER = os.environ.get("CLOUDRON_MAIL_SMTP_SERVER", "")
|
||||||
|
SMTP_PORT = int(os.environ.get("CLOUDRON_MAIL_SMTP_PORT", "2525"))
|
||||||
|
SMTP_USERNAME = os.environ.get("CLOUDRON_MAIL_SMTP_USERNAME", "")
|
||||||
|
SMTP_PASSWORD = os.environ.get("CLOUDRON_MAIL_SMTP_PASSWORD", "")
|
||||||
|
MAIL_FROM = os.environ.get("CLOUDRON_MAIL_FROM", "info@inference.coop")
|
||||||
|
MAIL_FROM_NAME = os.environ.get("CLOUDRON_MAIL_FROM_DISPLAY_NAME", "Inference Cooperative")
|
||||||
|
# Public base URL for the portal (used in email links).
|
||||||
|
PORTAL_BASE = os.environ.get("PORTAL_BASE", "https://portal.inference.coop")
|
||||||
|
|
||||||
|
|
||||||
def _verify_portal_secret(request: Request) -> None:
|
def _verify_portal_secret(request: Request) -> None:
|
||||||
"""Reject requests that didn't come through LibreChat (shared secret)."""
|
"""Reject requests that didn't come through LibreChat (shared secret)."""
|
||||||
@@ -213,6 +227,106 @@ async def is_active_member(email: str) -> bool:
|
|||||||
return email.strip().lower() in members
|
return email.strip().lower() in members
|
||||||
|
|
||||||
|
|
||||||
|
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 { email name slug } } } } }'
|
||||||
|
% OC_COLLECTIVE_SLUG
|
||||||
|
)
|
||||||
|
async with httpx.AsyncClient() as client:
|
||||||
|
r = await client.post(
|
||||||
|
OC_GRAPHQL_URL,
|
||||||
|
headers={"Personal-Token": OC_PERSONAL_TOKEN},
|
||||||
|
json={"query": query},
|
||||||
|
)
|
||||||
|
if r.status_code != 200:
|
||||||
|
logger.warning("OC member fetch failed: %s", r.status_code)
|
||||||
|
return []
|
||||||
|
nodes = r.json().get("data", {}).get("collective", {}).get("members", {}).get("nodes", [])
|
||||||
|
result: list[dict] = []
|
||||||
|
for n in nodes:
|
||||||
|
role = n.get("role", "")
|
||||||
|
acct = n.get("account", {}) or {}
|
||||||
|
email = (acct.get("email") or "").strip().lower()
|
||||||
|
if email and role in ("BACKER", "ADMIN"):
|
||||||
|
result.append({
|
||||||
|
"email": email,
|
||||||
|
"name": acct.get("name") or email.split("@")[0],
|
||||||
|
"slug": acct.get("slug") or "",
|
||||||
|
"role": role,
|
||||||
|
})
|
||||||
|
return result
|
||||||
|
|
||||||
|
|
||||||
|
def send_email(to: str, subject: str, text_body: str, html_body: str | None = None) -> bool:
|
||||||
|
"""Send an email via Cloudron's sendmail addon (SMTP).
|
||||||
|
|
||||||
|
Returns True on success, False on failure (logged, not raised — email is
|
||||||
|
best-effort and must never break the provisioning flow).
|
||||||
|
"""
|
||||||
|
if not SMTP_SERVER:
|
||||||
|
logger.warning("SMTP not configured; skipping email to %s", to)
|
||||||
|
return False
|
||||||
|
try:
|
||||||
|
msg = MIMEMultipart("alternative")
|
||||||
|
msg["From"] = f"{MAIL_FROM_NAME} <{MAIL_FROM}>"
|
||||||
|
msg["To"] = to
|
||||||
|
msg["Subject"] = subject
|
||||||
|
msg.attach(MIMEText(text_body, "plain", "utf-8"))
|
||||||
|
if html_body:
|
||||||
|
msg.attach(MIMEText(html_body, "html", "utf-8"))
|
||||||
|
|
||||||
|
with smtplib.SMTP(SMTP_SERVER, SMTP_PORT, timeout=15) as s:
|
||||||
|
if SMTP_USERNAME:
|
||||||
|
s.login(SMTP_USERNAME, SMTP_PASSWORD)
|
||||||
|
s.sendmail(MAIL_FROM, [to], msg.as_string())
|
||||||
|
logger.info("Sent email to %s: %s", to, subject)
|
||||||
|
return True
|
||||||
|
except Exception as e:
|
||||||
|
logger.warning("Failed to send email to %s: %s", to, e)
|
||||||
|
return False
|
||||||
|
|
||||||
|
|
||||||
|
def welcome_email(name: str) -> tuple[str, str, str]:
|
||||||
|
"""Build the welcome email (subject, text, html) for a new member."""
|
||||||
|
join_url = f"{PORTAL_BASE}/join"
|
||||||
|
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, "
|
||||||
|
"visit the link below and connect your Open Collective account:\n\n"
|
||||||
|
f"{join_url}\n\n"
|
||||||
|
"This takes about a minute. You'll then get an account-setup email "
|
||||||
|
"from our platform.\n\n"
|
||||||
|
"Questions? Reply to this email or write to info@inference.coop.\n\n"
|
||||||
|
"— The Inference Cooperative\n"
|
||||||
|
)
|
||||||
|
html = (
|
||||||
|
f"<p>Hi {name},</p>"
|
||||||
|
"<p>Thanks for joining the <strong>Inference Cooperative</strong>!</p>"
|
||||||
|
"<p>To finish setting up your account and start using private AI chat, "
|
||||||
|
"click the button below and connect your Open Collective account:</p>"
|
||||||
|
f'<p><a href="{join_url}" style="display:inline-block;padding:12px 24px;'
|
||||||
|
'background:#5b8c5a;color:#fff;text-decoration:none;border-radius:8px;">'
|
||||||
|
"Finish setup</a></p>"
|
||||||
|
"<p>This takes about a minute. You'll then get an account-setup email "
|
||||||
|
"from our platform.</p>"
|
||||||
|
'<p>Questions? Reply to this email or write to '
|
||||||
|
'<a href="mailto:info@inference.coop">info@inference.coop</a>.</p>'
|
||||||
|
"<p>— The Inference Cooperative</p>"
|
||||||
|
)
|
||||||
|
return subject, text, html
|
||||||
|
|
||||||
|
|
||||||
async def get_active_member_emails() -> list[str]:
|
async def get_active_member_emails() -> list[str]:
|
||||||
"""Return the list of active member emails from the portal's own database.
|
"""Return the list of active member emails from the portal's own database.
|
||||||
|
|
||||||
@@ -488,9 +602,19 @@ async def opencollective_webhook(request: Request, token: str):
|
|||||||
# - payload has "firstPayment" boolean to distinguish new vs recurring
|
# - payload has "firstPayment" boolean to distinguish new vs recurring
|
||||||
# - "collective.transaction.created" is DEPRECATED (being removed)
|
# - "collective.transaction.created" is DEPRECATED (being removed)
|
||||||
if event_type in ("order.processed", "new.member", "collective.member.created"):
|
if event_type in ("order.processed", "new.member", "collective.member.created"):
|
||||||
# New or renewed member → store as pending; they complete via OAuth
|
# New or renewed member → store as pending; they complete via OAuth.
|
||||||
if slug:
|
if slug:
|
||||||
store_pending_member(slug, name)
|
store_pending_member(slug, name)
|
||||||
|
# Send the welcome email immediately so the member gets feedback that
|
||||||
|
# their contribution was received and knows the next step (/join).
|
||||||
|
# The webhook strips email, so look it up by slug via the admin token.
|
||||||
|
if slug:
|
||||||
|
members = await fetch_members()
|
||||||
|
for m in members:
|
||||||
|
if m["slug"] == slug:
|
||||||
|
subject, text, html = welcome_email(m["name"])
|
||||||
|
send_email(m["email"], subject, text, html)
|
||||||
|
break
|
||||||
return JSONResponse({"status": "pending", "name": name, "slug": slug})
|
return JSONResponse({"status": "pending", "name": name, "slug": slug})
|
||||||
|
|
||||||
if event_type in ("collective.member.deleted", "collective.transaction.deleted"):
|
if event_type in ("collective.member.deleted", "collective.transaction.deleted"):
|
||||||
@@ -815,6 +939,82 @@ async def admin_provision(token: str, request: Request):
|
|||||||
return {"status": "provisioned", "email": email, "user_id": user_id}
|
return {"status": "provisioned", "email": email, "user_id": user_id}
|
||||||
|
|
||||||
|
|
||||||
|
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}
|
||||||
|
|
||||||
|
# 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:
|
||||||
|
user_id = await cloudron_create_user(email, m["name"])
|
||||||
|
await cloudron_set_group(user_id)
|
||||||
|
await cloudron_set_active(user_id, True)
|
||||||
|
key_token = await litellm_create_key(email, MEMBER_BUDGET)
|
||||||
|
store_member(email, key_token, user_id, m["slug"])
|
||||||
|
await cloudron_send_invite(user_id, email)
|
||||||
|
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.
|
||||||
|
for email, user_id in provisioned.items():
|
||||||
|
if email not in active_emails:
|
||||||
|
try:
|
||||||
|
if INACTIVE_GROUP_ID and user_id:
|
||||||
|
await cloudron_set_group(user_id, INACTIVE_GROUP_ID)
|
||||||
|
deactivate_member(email)
|
||||||
|
deactivated_count += 1
|
||||||
|
logger.info("Sweep: deactivated %s", email)
|
||||||
|
except Exception as e:
|
||||||
|
logger.warning("Sweep: failed to deactivate %s: %s", email, e)
|
||||||
|
|
||||||
|
# 3. Sync Loomio.
|
||||||
|
await sync_loomio_memberships()
|
||||||
|
|
||||||
|
return {
|
||||||
|
"status": "reconciled",
|
||||||
|
"provisioned": provisioned_count,
|
||||||
|
"deactivated": deactivated_count,
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
@app.post("/admin/sweep/{token}")
|
||||||
|
async def admin_sweep(token: str):
|
||||||
|
"""Manually trigger a full membership reconciliation.
|
||||||
|
|
||||||
|
Protected by the same secret token as the webhook. Idempotent — safe to
|
||||||
|
call repeatedly. This is also what the Cloudron scheduler invokes on a
|
||||||
|
cron schedule (see the sweep script).
|
||||||
|
"""
|
||||||
|
if not WEBHOOK_TOKEN or not secrets.compare_digest(token, WEBHOOK_TOKEN):
|
||||||
|
raise HTTPException(401, "Invalid token")
|
||||||
|
|
||||||
|
result = await reconcile_memberships()
|
||||||
|
return result
|
||||||
|
|
||||||
|
|
||||||
@app.get("/")
|
@app.get("/")
|
||||||
async def index():
|
async def index():
|
||||||
return {"service": "Inference Cooperative Member Portal", "version": "0.1.0"}
|
return {"service": "Inference Cooperative Member Portal", "version": "0.1.0"}
|
||||||
@@ -0,0 +1,35 @@
|
|||||||
|
#!/bin/sh
|
||||||
|
# Membership reconciliation sweep — run by the Cloudron scheduler addon.
|
||||||
|
#
|
||||||
|
# Calls the portal's own /admin/sweep endpoint (localhost), which reconciles
|
||||||
|
# the portal against the live Open Collective member list: provisions missing
|
||||||
|
# members, deactivates lapsed members, and syncs Loomio.
|
||||||
|
#
|
||||||
|
# The WEBHOOK_TOKEN env var is shared with the main app (scheduler runs in the
|
||||||
|
# same launch environment), so we use it to authenticate the call.
|
||||||
|
#
|
||||||
|
# Uses Python (guaranteed present in the base image) rather than curl.
|
||||||
|
|
||||||
|
set -e
|
||||||
|
|
||||||
|
if [ -z "$WEBHOOK_TOKEN" ]; then
|
||||||
|
echo "sweep: WEBHOOK_TOKEN not set; aborting" >&2
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
|
||||||
|
python3 - "$WEBHOOK_TOKEN" <<'PY'
|
||||||
|
import sys, urllib.request, urllib.error
|
||||||
|
|
||||||
|
token = sys.argv[1]
|
||||||
|
url = f"http://127.0.0.1:8000/admin/sweep/{token}"
|
||||||
|
req = urllib.request.Request(url, method="POST", data=b"", headers={"Content-Type": "application/json"})
|
||||||
|
try:
|
||||||
|
with urllib.request.urlopen(req, timeout=60) as resp:
|
||||||
|
print("sweep:", resp.read().decode())
|
||||||
|
except urllib.error.HTTPError as e:
|
||||||
|
print(f"sweep: HTTP {e.code}: {e.read().decode()}", file=sys.stderr)
|
||||||
|
sys.exit(1)
|
||||||
|
except Exception as e:
|
||||||
|
print(f"sweep: request failed: {e}", file=sys.stderr)
|
||||||
|
sys.exit(1)
|
||||||
|
PY
|
||||||
Reference in new issue
Block a user