Coverage for server / worker / celery_worker.py: 100%

16 statements  

« prev     ^ index     » next       coverage.py v7.13.4, created at 2026-10-04 09:33 +0000

1import os 

2import re 

3import ssl 

4 

5from celery import Celery 

6from dotenv import load_dotenv 

7from server.utilities.redis_url import redis_url_or_default 

8 

9load_dotenv(".env") 

10 

11REDIS_URL = redis_url_or_default("redis://localhost:6379/0") 

12 

13 

14def _redis_db(url: str, db: int) -> str: 

15 """Point Celery at a SEPARATE Redis logical DB from the API response cache 

16 (which uses DB 0), so broker queues + task results never collide with cached 

17 payloads. Swaps a trailing ``/<n>`` db segment; appends one if absent.""" 

18 return ( 

19 re.sub(r"/\d+$", f"/{db}", url) if re.search(r"/\d+$", url) else f"{url}/{db}" 

20 ) 

21 

22 

23# Broker + result backend on DB 1 (the cache uses DB 0). Override via CELERY_BROKER_URL. 

24CELERY_BROKER_URL = os.getenv("CELERY_BROKER_URL", _redis_db(REDIS_URL, 1)) 

25# TLS params are required when the broker is reached over rediss:// (the live server). 

26_USE_SSL = ( 

27 {"ssl_cert_reqs": ssl.CERT_REQUIRED} 

28 if CELERY_BROKER_URL.startswith("rediss://") 

29 else None 

30) 

31 

32# Per-service queue. The teacher-student and staff-admin Celery apps SHARE the same Redis 

33# broker DB, so they must NOT share the default "celery" queue — otherwise each service's 

34# worker would consume (and reject as unregistered) the other's tasks. Each app dispatches 

35# to and each worker consumes only this queue. Keep it distinct from the staff-admin app's. 

36TASK_QUEUE = os.getenv("CELERY_TASK_QUEUE", "ts_celery") 

37 

38celery = Celery(__name__) 

39celery.conf.update( 

40 broker_url=CELERY_BROKER_URL, 

41 result_backend=CELERY_BROKER_URL, 

42 broker_use_ssl=_USE_SSL, 

43 redis_backend_use_ssl=_USE_SSL, 

44 task_default_queue=TASK_QUEUE, 

45 task_serializer="json", 

46 result_serializer="json", 

47 accept_content=["json"], 

48 timezone="UTC", 

49 enable_utc=True, 

50 task_acks_late=True, 

51 worker_prefetch_multiplier=1, 

52 broker_connection_retry_on_startup=True, 

53 # --- runtime bounds ------------------------------------------------------- 

54 # The only task here reconciles Mongo accounts against Auth0 and can soft-delete. 

55 # `task_acks_late=True` means the message is acknowledged only after the task 

56 # returns, so a worker that dies mid-run leaves the message on the queue to be 

57 # redelivered. That is the right trade for a task worth retrying, but it makes 

58 # two settings load-bearing rather than cosmetic. 

59 # 

60 # visibility_timeout is how long Redis waits before deciding a delivered message 

61 # was never handled and hands it to another worker. If it is SHORTER than the 

62 # task's worst case, a slow-but-healthy run is redelivered and the destructive 

63 # pass executes twice concurrently. The default is 3600s; the Auth0 export poll 

64 # alone has a ~3 minute ceiling plus per-account Mongo writes, so 3600s is 

65 # probably fine today — but "probably fine today" is what silently stops being 

66 # true. Pin it well above the hard time limit instead of inheriting a default. 

67 task_time_limit=30 * 60, # hard kill 

68 task_soft_time_limit=25 * 60, # raises SoftTimeLimitExceeded first, so it can log 

69 broker_transport_options={"visibility_timeout": 2 * 60 * 60}, 

70 # Results live in the SAME Redis as the response cache, on a box configured 

71 # `maxmemory 2gb` + `noeviction`. Under noeviction Redis does not discard old 

72 # keys when full — it starts REFUSING WRITES, which would take out the cache and 

73 # the broker together. Result rows must therefore expire on their own. 

74 result_expires=24 * 60 * 60, 

75 # Recycle workers periodically; a long-lived process doing HTTP + Mongo has no 

76 # reason to hold memory for weeks. 

77 worker_max_tasks_per_child=100, 

78) 

79 

80 

81# Auto-discover tasks (includes auth0_reconciliation beat schedule) 

82celery.conf.imports = [ 

83 "server.worker.auth0_reconciliation", 

84 "server.worker.student_identity_sync_task", 

85]