Coverage for server / worker / student_identity_sync_task.py: 50%
34 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"""
2Periodically re-sync embedded student identity into class documents.
4WHY A PERIODIC TASK AND NOT A CALL ON SAVE. The real fix is for the writer to
5sync on write: update_portal_user_profile, which changes a student's
6first_name/last_name, lives in the ADMIN-STAFF API and would need one line to
7call sync_student_identity_in_classes. That repo is outside this service's
8scope, so this task exists to close the window rather than prevent it. When the
9direct call lands, this becomes what a reconciler should be — a safety net that
10normally finds nothing.
12BOUNDED BY updated_at, DELIBERATELY. Scanning every user against every class
13would be expensive and would grow with the estate. The Admin-Staff profile
14update stamps `updated_at`, so the set of accounts whose embedded copies COULD
15have drifted since the last run is exactly the set updated since then, plus a
16margin. A full sweep is available on demand by passing a wide window.
18IT IS SAFE TO RUN REPEATEDLY. The sync writes the account's current values over
19the embedded copy; running it twice changes nothing the first run did not
20already settle, and `modified_count` reports zero when every copy already
21matched.
23Developer: Allan Ninal
24Date: 2026-09-21
25"""
27import logging
28import os
29from datetime import datetime, timedelta, timezone
31from pymongo import MongoClient
33from server.services.common.student_identity_sync import sync_identity_in_classes
34from server.worker.celery_worker import TASK_QUEUE, celery
36logger = logging.getLogger(__name__)
38DB_NAME = os.getenv("DB_NAME", "teacher_student_db")
39MONGODB_URL = os.getenv("MONGODB_URL", "mongodb://localhost:27017")
40# TASK_QUEUE comes from celery_worker, which is also what the worker consumes
41# (task_default_queue). Re-deriving it from an env var here is how a task ends
42# up queued somewhere nothing is listening — it never runs, and nothing errors.
44# How often the sweep runs, and how far back it looks. The window is wider than
45# the interval on purpose: a run that fails or is delayed must not leave a gap
46# that the next run silently skips over.
47SYNC_INTERVAL_HOURS = int(os.getenv("STUDENT_IDENTITY_SYNC_INTERVAL_HOURS", "6"))
48SYNC_WINDOW_HOURS = int(os.getenv("STUDENT_IDENTITY_SYNC_WINDOW_HOURS", "24"))
50# Both roles are embedded in a class document — the teacher as a single object,
51# students as array elements — so both go stale and both need sweeping. `role`
52# is projected because the sync dispatches on it.
53_SWEPT_ROLES = ("student", "teacher")
55_IDENTITY_PROJECTION = {
56 "_id": 1,
57 "role": 1,
58 "first_name": 1,
59 "middle_name": 1,
60 "last_name": 1,
61 "email": 1,
62}
65@celery.task(name="sync_embedded_student_identity", bind=True, max_retries=2)
66def sync_embedded_student_identity(self, window_hours: int | None = None):
67 """Re-sync recently-updated teachers and students into their classes."""
68 window = window_hours if window_hours is not None else SYNC_WINDOW_HOURS
69 since = datetime.now(timezone.utc) - timedelta(hours=window)
71 client = MongoClient(MONGODB_URL)
72 try:
73 db = client[DB_NAME]
74 cursor = db["user_collection"].find(
75 {"role": {"$in": list(_SWEPT_ROLES)}, "updated_at": {"$gte": since}},
76 _IDENTITY_PROJECTION,
77 )
79 stats = {"examined": 0, "classes_updated": 0, "errors": 0}
80 for doc in cursor:
81 stats["examined"] += 1
82 try:
83 stats["classes_updated"] += sync_identity_in_classes(
84 db, doc["_id"], doc.get("role", ""), doc
85 )
86 except Exception: # noqa: BLE001 — one bad account must not stop the sweep
87 stats["errors"] += 1
88 logger.exception(
89 "student identity sync failed for %s", doc.get("_id")
90 )
92 logger.info(
93 "student identity sync: window=%sh examined=%s classes_updated=%s errors=%s",
94 window,
95 stats["examined"],
96 stats["classes_updated"],
97 stats["errors"],
98 )
99 return {"status": "completed", "window_hours": window, **stats}
100 finally:
101 client.close()
104celery.conf.beat_schedule = celery.conf.get("beat_schedule", {})
105celery.conf.beat_schedule["sync-embedded-student-identity"] = {
106 "task": "sync_embedded_student_identity",
107 "schedule": SYNC_INTERVAL_HOURS * 3600,
108 "options": {"queue": TASK_QUEUE}, # this service's queue, not the shared default
109}