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
« prev ^ index » next coverage.py v7.13.4, created at 2026-10-04 09:33 +0000
1import os
2import re
3import ssl
5from celery import Celery
6from dotenv import load_dotenv
7from server.utilities.redis_url import redis_url_or_default
9load_dotenv(".env")
11REDIS_URL = redis_url_or_default("redis://localhost:6379/0")
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 )
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)
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")
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)
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]