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

1""" 

2Periodically re-sync embedded student identity into class documents. 

3 

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. 

11 

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. 

17 

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. 

22 

23Developer: Allan Ninal 

24Date: 2026-09-21 

25""" 

26 

27import logging 

28import os 

29from datetime import datetime, timedelta, timezone 

30 

31from pymongo import MongoClient 

32 

33from server.services.common.student_identity_sync import sync_identity_in_classes 

34from server.worker.celery_worker import TASK_QUEUE, celery 

35 

36logger = logging.getLogger(__name__) 

37 

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. 

43 

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

49 

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

54 

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} 

63 

64 

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) 

70 

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 ) 

78 

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 ) 

91 

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

102 

103 

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}