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

1""" 

2Auth0 ↔ MongoDB Reconciliation Task (Teacher-Student API) 

3 

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. 

7 

8Runs every 6 hours by default (configurable via AUTH0_RECONCILE_INTERVAL_HOURS). 

9""" 

10 

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 

20 

21import httpx 

22from dotenv import load_dotenv 

23from pymongo import MongoClient 

24 

25from server.worker.celery_worker import celery, TASK_QUEUE 

26 

27load_dotenv() 

28 

29logger = logging.getLogger(__name__) 

30 

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) 

39 

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) 

53 

54# Reconciliation interval (hours) 

55RECONCILE_INTERVAL = int(os.getenv("AUTH0_RECONCILE_INTERVAL_HOURS", "6")) 

56 

57 

58def _send_alert_email(subject: str, body: str) -> None: 

59 """Best-effort email alert for notable/destructive reconcile outcomes. 

60 

61 FAIL-OPEN: never raise — a mail failure must not break reconciliation. Suppressed 

62 under TESTING. Recipient via ALERT_EMAIL (default below). 

63 

64 From address: prefers MAIL_FROM, falling back to SMTP_USER. 

65 

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}") 

106 

107 

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"] 

122 

123 

124def _fetch_all_auth0_users(token: str) -> list[dict]: 

125 """Fetch ALL Auth0 users via the Management API user-export job. 

126 

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}"} 

134 

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"] 

153 

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") 

173 

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()] 

182 

183 

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). 

188 

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 

193 

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"} 

200 

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) 

207 

208 # Connect to MongoDB (sync) 

209 client = MongoClient(MONGODB_URL) 

210 db = client[DB_NAME] 

211 collection = db["user_collection"] 

212 

213 now = datetime.now(timezone.utc) 

214 stats = {"created": 0, "updated": 0, "soft_deleted": 0, "unchanged": 0, "errors": 0, "soft_delete_aborted": 0} 

215 

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})") 

222 

223 # Build a set of Auth0 user IDs for the deletion check 

224 auth0_user_ids = set() 

225 

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", "") 

232 

233 if not auth0_id or not email: 

234 continue 

235 

236 # Only reconcile teacher/student roles for this API 

237 if role and role not in ("teacher", "student"): 

238 continue 

239 

240 auth0_user_ids.add(auth0_id) 

241 

242 # Find existing User by auth0_user_id 

243 existing = collection.find_one({"auth0_user_id": auth0_id}) 

244 

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 

254 

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 

286 

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}") 

310 

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) 

327 

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 ) 

353 

354 client.close() 

355 

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} 

372 

373 

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}