Coverage for server / worker / auth0_reconciliation.py: 91%
176 statements
« prev ^ index » next coverage.py v7.13.4, created at 2026-10-04 09:33 +0000
« prev ^ index » next coverage.py v7.13.4, created at 2026-10-04 09:33 +0000
1"""
2Auth0 ↔ MongoDB Reconciliation Task (Teacher-Student API)
4Periodic Celery task that fetches all users from Auth0 and reconciles them
5with MongoDB User documents (teacher/student roles only). Catches Auth0
6Dashboard edits/deletes that happen outside the webhook flow.
8Runs every 6 hours by default (configurable via AUTH0_RECONCILE_INTERVAL_HOURS).
9"""
11import gzip
12import json
13import logging
14import os
15import smtplib
16import time
17from datetime import datetime, timezone
18from email.mime.text import MIMEText
19from email.utils import parseaddr
21import httpx
22from dotenv import load_dotenv
23from pymongo import MongoClient
25from server.worker.celery_worker import celery, TASK_QUEUE
27load_dotenv()
29logger = logging.getLogger(__name__)
31# Auth0 Management API config
32AUTH0_DOMAIN = os.getenv("AUTH0_DOMAIN", "")
33AUTH0_MGMT_CLIENT_ID = os.getenv("AUTH0_MGMT_CLIENT_ID", "")
34AUTH0_MGMT_CLIENT_SECRET = os.getenv("AUTH0_MGMT_CLIENT_SECRET", "")
35AUTH0_MGMT_AUDIENCE = os.getenv(
36 "AUTH0_MGMT_AUDIENCE",
37 f"https://{AUTH0_DOMAIN}/api/v2/" if AUTH0_DOMAIN else "",
38)
40# MongoDB config
41DB_USER = os.getenv("DB_USER", "")
42DB_PASSWORD = os.getenv("DB_PASSWORD", "")
43DB_HOST = os.getenv("DB_HOST", "")
44DB_PORT = os.getenv("DB_PORT", "27017")
45DB_AUTH_SOURCE = os.getenv("DB_AUTH_SOURCE", "")
46DB_NAME = os.getenv("DB_NAME", "teacher_student_db")
47MONGODB_URL = os.getenv(
48 "MONGODB_URL",
49 f"mongodb://{DB_USER}:{DB_PASSWORD}@{DB_HOST}:{DB_PORT}/?authSource={DB_AUTH_SOURCE}&directConnection=true"
50 if DB_USER and DB_PASSWORD and DB_HOST
51 else "mongodb://localhost:27017/",
52)
54# Reconciliation interval (hours)
55RECONCILE_INTERVAL = int(os.getenv("AUTH0_RECONCILE_INTERVAL_HOURS", "6"))
58def _send_alert_email(subject: str, body: str) -> None:
59 """Best-effort email alert for notable/destructive reconcile outcomes.
61 FAIL-OPEN: never raise — a mail failure must not break reconciliation. Suppressed
62 under TESTING. Recipient via ALERT_EMAIL (default below).
64 From address: prefers MAIL_FROM, falling back to SMTP_USER.
66 That fallback used to be the only behaviour, because consumer Gmail rejects a
67 From that does not match the authenticated user. It does not survive the move to
68 Amazon SES, where SMTP_USER is an IAM **access key id** (``AKIA...``) rather than
69 an address — sending From that would be rejected outright, and because this
70 function is fail-open the alerts would simply stop arriving with nothing in the
71 logs but a warning nobody reads. MAIL_FROM must therefore be set wherever SES is
72 configured; it must also match the address pinned by the IAM policy's
73 ``ses:FromAddress`` condition.
74 """
75 if os.getenv("TESTING", "").strip().lower() in ("1", "true", "yes"):
76 return
77 to_addr = os.getenv("ALERT_EMAIL", "landix.ninal@gmail.com")
78 user = os.getenv("SMTP_USER", "")
79 pwd = os.getenv("SMTP_PASS") or os.getenv("SMTP_PASSWORD", "")
80 if not (user and pwd and to_addr):
81 logger.warning("Reconcile alert email skipped: SMTP not configured")
82 return
83 from_header = os.getenv("MAIL_FROM") or user
84 # Envelope sender must be the bare address, even when MAIL_FROM carries a
85 # display name such as ``EruditionTX <support@eruditionsys.com>``.
86 envelope_from = parseaddr(from_header)[1] or user
87 try:
88 msg = MIMEText(body)
89 msg["Subject"] = subject
90 msg["From"] = from_header
91 msg["To"] = to_addr
92 # Attribute the send to an SES configuration set when one is configured, so
93 # bounces reach the SNS destination. Harmless on non-SES providers, but only
94 # added when set so no stray header goes out.
95 cfg_set = os.getenv("SES_CONFIGURATION_SET", "").strip()
96 if cfg_set:
97 msg["X-SES-CONFIGURATION-SET"] = cfg_set
98 server = smtplib.SMTP(os.getenv("SMTP_HOST", "smtp.gmail.com"), int(os.getenv("SMTP_PORT", "587")), timeout=15)
99 server.starttls()
100 server.login(user, pwd)
101 server.sendmail(envelope_from, [to_addr], msg.as_string())
102 server.quit()
103 logger.info(f"Reconcile alert email sent to {to_addr}")
104 except Exception as e:
105 logger.warning(f"Reconcile alert email failed (non-fatal): {e}")
108def _get_management_token() -> str:
109 """Get an M2M token for Auth0 Management API (synchronous)."""
110 response = httpx.post(
111 f"https://{AUTH0_DOMAIN}/oauth/token",
112 json={
113 "client_id": AUTH0_MGMT_CLIENT_ID,
114 "client_secret": AUTH0_MGMT_CLIENT_SECRET,
115 "audience": AUTH0_MGMT_AUDIENCE,
116 "grant_type": "client_credentials",
117 },
118 timeout=10.0,
119 )
120 response.raise_for_status()
121 return response.json()["access_token"]
124def _fetch_all_auth0_users(token: str) -> list[dict]:
125 """Fetch ALL Auth0 users via the Management API user-export job.
127 The plain ``GET /api/v2/users`` endpoint hard-caps at 1000 results (offset
128 paging), so on a tenant with >1000 users every user past the first 1000 would
129 be (incorrectly) treated as deleted in Auth0 and soft-deleted. The bulk
130 user-export job has no such cap: it returns a gzipped JSON-lines file of every
131 user. Synchronous (start → poll → download), which is fine inside the worker.
132 """
133 headers = {"Authorization": f"Bearer {token}"}
135 # 1) Start the export job (only the fields the reconciliation actually reads).
136 start = httpx.post(
137 f"https://{AUTH0_DOMAIN}/api/v2/jobs/users-exports",
138 headers=headers,
139 json={
140 "format": "json",
141 "fields": [
142 {"name": "user_id"},
143 {"name": "email"},
144 {"name": "given_name"},
145 {"name": "family_name"},
146 {"name": "app_metadata"},
147 ],
148 },
149 timeout=15.0,
150 )
151 start.raise_for_status()
152 job_id = start.json()["id"]
154 # 2) Poll until the job finishes (export of a few thousand users takes seconds).
155 location = None
156 for _ in range(60): # ~3 min ceiling
157 time.sleep(3)
158 status_resp = httpx.get(
159 f"https://{AUTH0_DOMAIN}/api/v2/jobs/{job_id}",
160 headers=headers,
161 timeout=15.0,
162 )
163 status_resp.raise_for_status()
164 info = status_resp.json()
165 state = info.get("status")
166 if state == "completed":
167 location = info.get("location")
168 break
169 if state == "failed":
170 raise RuntimeError(f"Auth0 user-export job {job_id} failed: {info}")
171 if not location:
172 raise TimeoutError(f"Auth0 user-export job {job_id} did not complete in time")
174 # 3) Download + parse the gzipped JSON-lines export (one user object per line).
175 dl = httpx.get(location, timeout=60.0)
176 dl.raise_for_status()
177 try:
178 raw = gzip.decompress(dl.content).decode("utf-8")
179 except (OSError, EOFError):
180 raw = dl.text # already-decompressed by the transport
181 return [json.loads(line) for line in raw.splitlines() if line.strip()]
184@celery.task(name="auth0_reconcile_teacher_student_users", bind=True, max_retries=2)
185def reconcile_auth0_users(self):
186 """
187 Reconcile Auth0 users with MongoDB User collection (teacher/student only).
189 For each Auth0 user with teacher/student role:
190 - If they have a linked User (by auth0_user_id): update name/email if changed
191 - If no linked User exists but email matches: link the account
192 - If no linked User exists at all: create one
194 For each MongoDB User with auth0_user_id:
195 - If the auth0_user_id doesn't exist in Auth0: soft-delete the User
196 """
197 if not all([AUTH0_DOMAIN, AUTH0_MGMT_CLIENT_ID, AUTH0_MGMT_CLIENT_SECRET]):
198 logger.warning("Auth0 Management API not configured, skipping reconciliation")
199 return {"status": "skipped", "reason": "not_configured"}
201 try:
202 token = _get_management_token()
203 auth0_users = _fetch_all_auth0_users(token)
204 except Exception as e:
205 logger.error(f"Failed to fetch Auth0 users: {e}")
206 raise self.retry(exc=e, countdown=300)
208 # Connect to MongoDB (sync)
209 client = MongoClient(MONGODB_URL)
210 db = client[DB_NAME]
211 collection = db["user_collection"]
213 now = datetime.now(timezone.utc)
214 stats = {"created": 0, "updated": 0, "soft_deleted": 0, "unchanged": 0, "errors": 0, "soft_delete_aborted": 0}
216 # Dry-run gate: defaults TRUE so a freshly-deployed worker never mutates accounts
217 # until AUTH0_RECONCILE_DRY_RUN is explicitly set false. In dry-run we still compute
218 # and log every change ("WOULD have ...") and tally the stats, but issue no DB writes.
219 dry_run = os.getenv("AUTH0_RECONCILE_DRY_RUN", "true").strip().lower() in ("1", "true", "yes", "on")
220 log_prefix = "[DRY-RUN] WOULD have" if dry_run else "Reconciliation"
221 logger.info(f"Auth0 reconciliation starting (dry_run={dry_run})")
223 # Build a set of Auth0 user IDs for the deletion check
224 auth0_user_ids = set()
226 for auth0_user in auth0_users:
227 auth0_id = auth0_user.get("user_id", "")
228 email = auth0_user.get("email", "")
229 given_name = auth0_user.get("given_name", "")
230 family_name = auth0_user.get("family_name", "")
231 role = auth0_user.get("app_metadata", {}).get("role", "")
233 if not auth0_id or not email:
234 continue
236 # Only reconcile teacher/student roles for this API
237 if role and role not in ("teacher", "student"):
238 continue
240 auth0_user_ids.add(auth0_id)
242 # Find existing User by auth0_user_id
243 existing = collection.find_one({"auth0_user_id": auth0_id})
245 if existing:
246 # Check if any fields changed
247 updates = {}
248 if given_name and existing.get("first_name") != given_name:
249 updates["first_name"] = given_name
250 if family_name and existing.get("last_name") != family_name:
251 updates["last_name"] = family_name
252 if email and existing.get("email") != email:
253 updates["email"] = email
255 if updates:
256 updates["updated_at"] = now
257 if not dry_run:
258 collection.update_one(
259 {"_id": existing["_id"]},
260 {"$set": updates},
261 )
262 stats["updated"] += 1
263 logger.info(f"{log_prefix} updated User {existing['_id']} (auth0: {auth0_id})")
264 else:
265 stats["unchanged"] += 1
266 else:
267 # Check if a user exists by email (may have been created without auth0_user_id)
268 by_email = collection.find_one({"email": email})
269 if by_email:
270 # Link the existing user to Auth0
271 if not dry_run:
272 collection.update_one(
273 {"_id": by_email["_id"]},
274 {"$set": {
275 "auth0_user_id": auth0_id,
276 "auth_provider": "auth0",
277 "updated_at": now,
278 }},
279 )
280 stats["updated"] += 1
281 logger.info(f"{log_prefix} linked User {by_email['_id']} to Auth0 {auth0_id}")
282 else:
283 # Create new User from Auth0 user
284 if not role:
285 role = "student" # Default role for untagged Auth0 users
287 try:
288 if not dry_run:
289 collection.insert_one({
290 "first_name": given_name or "User",
291 "middle_name": None,
292 "last_name": family_name or "User",
293 "role": role,
294 "email": email,
295 "password": "",
296 "auth0_user_id": auth0_id,
297 "auth_provider": "auth0",
298 "status": "active",
299 "profile_picture": None,
300 "total_usage_time_in_minutes": 0,
301 "total_no_of_visits": 0,
302 "created_at": now,
303 "updated_at": now,
304 })
305 stats["created"] += 1
306 logger.info(f"{log_prefix} created User for Auth0 user {auth0_id} ({email})")
307 except Exception as e:
308 stats["errors"] += 1
309 logger.error(f"Failed to create User for {email}: {e}")
311 # Soft-delete pass — guarded by a SAFETY CAP. Gather the candidates first; an
312 # implausibly high orphan rate almost always means an incomplete Auth0 fetch (e.g.
313 # the list API's 1000-result cap), so we ABORT the whole pass rather than risk
314 # mass-deleting live accounts. linked_count is derived from the same scan (no extra
315 # query) so it stays correct even when the fetch is partial.
316 linked_count = 0
317 to_soft_delete = []
318 for user in collection.find(
319 {"auth0_user_id": {"$ne": None, "$exists": True}, "status": {"$ne": "deleted"}},
320 ):
321 if user.get("role") not in ("teacher", "student"):
322 continue
323 linked_count += 1
324 auth0_id = user.get("auth0_user_id")
325 if auth0_id and auth0_id not in auth0_user_ids:
326 to_soft_delete.append(user)
328 max_fraction = float(os.getenv("AUTH0_RECONCILE_MAX_SOFT_DELETE_FRACTION", "0.10"))
329 floor = int(os.getenv("AUTH0_RECONCILE_SOFT_DELETE_FLOOR", "25"))
330 if len(to_soft_delete) > floor and linked_count and (len(to_soft_delete) / linked_count) > max_fraction:
331 logger.error(
332 f"SAFETY CAP HIT: would soft-delete {len(to_soft_delete)}/{linked_count} linked users "
333 f"(> {max_fraction:.0%}) — aborting soft-delete pass (likely an incomplete Auth0 fetch). "
334 f"No users soft-deleted. dry_run={dry_run}"
335 )
336 stats["soft_delete_aborted"] = len(to_soft_delete)
337 else:
338 for user in to_soft_delete:
339 auth0_id = user.get("auth0_user_id")
340 if not dry_run:
341 collection.update_one(
342 {"_id": user["_id"]},
343 {"$set": {
344 "status": "deleted",
345 "updated_at": now,
346 }},
347 )
348 stats["soft_deleted"] += 1
349 logger.info(
350 f"{log_prefix} soft-deleted User {user['_id']} "
351 f"(Auth0 user {auth0_id} no longer exists)"
352 )
354 client.close()
356 logger.info(f"Auth0 reconciliation complete (dry_run={dry_run}): {stats}")
357 # Alertable summary: a soft-delete (or a cap-aborted batch) is the destructive/notable
358 # outcome — emit it at WARNING with a greppable marker so monitoring can alert on it,
359 # turning an otherwise-silent automated deletion into an auditable signal.
360 if stats["soft_deleted"] or stats["soft_delete_aborted"]:
361 alert = (
362 f"[RECONCILE-ALERT] dry_run={dry_run} soft_deleted={stats['soft_deleted']} "
363 f"soft_delete_aborted={stats['soft_delete_aborted']} created={stats['created']} "
364 f"updated={stats['updated']} — review (account deletions occurred or were blocked)."
365 )
366 logger.warning(alert)
367 _send_alert_email(
368 subject=f"[EruditionTX] Teacher-Student Auth0 reconcile alert (dry_run={dry_run})",
369 body=alert + f"\n\nFull stats: {stats}",
370 )
371 return {"status": "completed", "dry_run": dry_run, **stats}
374# Celery Beat schedule for periodic execution
375celery.conf.beat_schedule = celery.conf.get("beat_schedule", {})
376celery.conf.beat_schedule["auth0-reconcile-teacher-student-users"] = {
377 "task": "auth0_reconcile_teacher_student_users",
378 "schedule": RECONCILE_INTERVAL * 3600, # Convert hours to seconds
379 "options": {"queue": TASK_QUEUE}, # this service's queue (NOT the shared default)
380}