This reverts commit 39a22734c1.
This commit is contained in:
committed by
GitHub
parent
e5e3fcdc1c
commit
c1fe3c3a93
@@ -5,11 +5,12 @@ and auto discover tasks in all installed django apps.
|
||||
Taken from: https://celery.readthedocs.org/en/latest/django/first-steps-with-django.html
|
||||
"""
|
||||
|
||||
|
||||
import os
|
||||
|
||||
from celery import Celery
|
||||
|
||||
from openedx.core.lib.celery.routers import route_task_queue
|
||||
from openedx.core.lib.celery.routers import AlternateEnvironmentRouter
|
||||
|
||||
# set the default Django settings module for the 'celery' program.
|
||||
os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'lms.envs.production')
|
||||
@@ -23,11 +24,14 @@ APP.config_from_object('django.conf:settings')
|
||||
APP.autodiscover_tasks()
|
||||
|
||||
|
||||
def route_task(name, args, kwargs, options, task=None, **kw): # pylint: disable=unused-argument
|
||||
class Router(AlternateEnvironmentRouter):
|
||||
"""
|
||||
Celery-defined method allowing for custom routing logic.
|
||||
|
||||
If None is returned from this method, default routing logic is used.
|
||||
An implementation of AlternateEnvironmentRouter, for routing tasks to non-cms queues.
|
||||
"""
|
||||
|
||||
return route_task_queue(name)
|
||||
@property
|
||||
def alternate_env_tasks(self):
|
||||
"""
|
||||
Defines alternate environment tasks, as a dict of form { task_name: alternate_queue }
|
||||
"""
|
||||
return {}
|
||||
|
||||
@@ -196,6 +196,8 @@ def perform_delegate_email_batches(entry_id, course_id, task_input, action_name)
|
||||
|
||||
total_recipients = combined_set.count()
|
||||
|
||||
routing_key = settings.BULK_EMAIL_ROUTING_KEY
|
||||
|
||||
# Weird things happen if we allow empty querysets as input to emailing subtasks
|
||||
# The task appears to hang at "0 out of 0 completed" and never finishes.
|
||||
if total_recipients == 0:
|
||||
@@ -215,6 +217,7 @@ def perform_delegate_email_batches(entry_id, course_id, task_input, action_name)
|
||||
initial_subtask_status.to_dict(),
|
||||
),
|
||||
task_id=subtask_id,
|
||||
routing_key=routing_key,
|
||||
)
|
||||
return new_subtask
|
||||
|
||||
|
||||
@@ -35,6 +35,7 @@ log = logging.getLogger(__name__)
|
||||
|
||||
|
||||
DEFAULT_LANGUAGE = 'en'
|
||||
ROUTING_KEY = getattr(settings, 'ACE_ROUTING_KEY', None)
|
||||
|
||||
|
||||
@task(base=LoggedTask)
|
||||
@@ -59,7 +60,7 @@ class ResponseNotification(BaseMessageType):
|
||||
pass
|
||||
|
||||
|
||||
@task(base=LoggedTask)
|
||||
@task(base=LoggedTask, routing_key=ROUTING_KEY)
|
||||
def send_ace_message(context):
|
||||
context['course_id'] = CourseKey.from_string(context['course_id'])
|
||||
|
||||
|
||||
@@ -18,9 +18,10 @@ from .models import EmailMarketingConfiguration
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
SAILTHRU_LIST_CACHE_KEY = "email.marketing.cache"
|
||||
ACE_ROUTING_KEY = getattr(settings, 'ACE_ROUTING_KEY', None)
|
||||
|
||||
|
||||
@task(bind=True)
|
||||
@task(bind=True, routing_key=ACE_ROUTING_KEY)
|
||||
def get_email_cookies_via_sailthru(self, user_email, post_parms):
|
||||
"""
|
||||
Adds/updates Sailthru cookie information for a new user.
|
||||
@@ -60,7 +61,7 @@ def get_email_cookies_via_sailthru(self, user_email, post_parms):
|
||||
return None
|
||||
|
||||
|
||||
@task(bind=True, default_retry_delay=3600, max_retries=24)
|
||||
@task(bind=True, default_retry_delay=3600, max_retries=24, routing_key=ACE_ROUTING_KEY)
|
||||
def update_user(self, sailthru_vars, email, site=None, new_user=False, activation=False):
|
||||
"""
|
||||
Adds/updates Sailthru profile information for a user.
|
||||
@@ -142,7 +143,7 @@ def is_default_site(site):
|
||||
return not site or site.get('id') == settings.SITE_ID
|
||||
|
||||
|
||||
@task(bind=True, default_retry_delay=3600, max_retries=24)
|
||||
@task(bind=True, default_retry_delay=3600, max_retries=24, routing_key=ACE_ROUTING_KEY)
|
||||
def update_user_email(self, new_email, old_email):
|
||||
"""
|
||||
Adds/updates Sailthru when a user email address is changed
|
||||
@@ -302,7 +303,7 @@ def _retryable_sailthru_error(error):
|
||||
return code == 9 or code == 43
|
||||
|
||||
|
||||
@task(bind=True)
|
||||
@task(bind=True, routing_key=ACE_ROUTING_KEY)
|
||||
def update_course_enrollment(self, email, course_key, mode, site=None):
|
||||
"""Adds/updates Sailthru when a user adds to cart/purchases/upgrades a course
|
||||
Args:
|
||||
|
||||
@@ -80,7 +80,7 @@ class Command(BaseCommand):
|
||||
"""
|
||||
Enqueue all tasks, in shuffled order.
|
||||
"""
|
||||
task_options = {'queue': options['routing_key']} if options.get('routing_key') else {}
|
||||
task_options = {'routing_key': options['routing_key']} if options.get('routing_key') else {}
|
||||
for seq_id, kwargs in enumerate(self._shuffled_task_kwargs(options)):
|
||||
kwargs['seq_id'] = seq_id
|
||||
result = tasks.compute_grades_for_course_v2.apply_async(kwargs=kwargs, **task_options)
|
||||
|
||||
@@ -92,19 +92,19 @@ class TestComputeGrades(SharedModuleStoreTestCase):
|
||||
# Order doesn't matter, but can't use a set because dicts aren't hashable
|
||||
expected = [
|
||||
({
|
||||
'queue': 'key',
|
||||
'routing_key': 'key',
|
||||
'kwargs': _kwargs(self.course_keys[0], 0)
|
||||
},),
|
||||
({
|
||||
'queue': 'key',
|
||||
'routing_key': 'key',
|
||||
'kwargs': _kwargs(self.course_keys[0], 2)
|
||||
},),
|
||||
({
|
||||
'queue': 'key',
|
||||
'routing_key': 'key',
|
||||
'kwargs': _kwargs(self.course_keys[3], 0)
|
||||
},),
|
||||
({
|
||||
'queue': 'key',
|
||||
'routing_key': 'key',
|
||||
'kwargs': _kwargs(self.course_keys[3], 2)
|
||||
},),
|
||||
]
|
||||
|
||||
@@ -48,7 +48,7 @@ RETRY_DELAY_SECONDS = 40
|
||||
SUBSECTION_GRADE_TIMEOUT_SECONDS = 300
|
||||
|
||||
|
||||
@task(base=LoggedPersistOnFailureTask)
|
||||
@task(base=LoggedPersistOnFailureTask, routing_key=settings.POLICY_CHANGE_GRADES_ROUTING_KEY)
|
||||
def compute_all_grades_for_course(**kwargs):
|
||||
"""
|
||||
Compute grades for all students in the specified course.
|
||||
@@ -69,7 +69,7 @@ def compute_all_grades_for_course(**kwargs):
|
||||
'batch_size': batch_size,
|
||||
})
|
||||
compute_grades_for_course_v2.apply_async(
|
||||
kwargs=kwargs, queue=settings.POLICY_CHANGE_GRADES_ROUTING_KEY
|
||||
kwargs=kwargs, routing_key=settings.POLICY_CHANGE_GRADES_ROUTING_KEY
|
||||
)
|
||||
|
||||
|
||||
@@ -131,6 +131,7 @@ def compute_grades_for_course(course_key, offset, batch_size, **kwargs): # pyli
|
||||
time_limit=SUBSECTION_GRADE_TIMEOUT_SECONDS,
|
||||
max_retries=2,
|
||||
default_retry_delay=RETRY_DELAY_SECONDS,
|
||||
routing_key=settings.POLICY_CHANGE_GRADES_ROUTING_KEY
|
||||
)
|
||||
def recalculate_course_and_subsection_grades_for_user(self, **kwargs): # pylint: disable=unused-argument
|
||||
"""
|
||||
@@ -170,7 +171,8 @@ def recalculate_course_and_subsection_grades_for_user(self, **kwargs): # pylint
|
||||
base=LoggedPersistOnFailureTask,
|
||||
time_limit=SUBSECTION_GRADE_TIMEOUT_SECONDS,
|
||||
max_retries=2,
|
||||
default_retry_delay=RETRY_DELAY_SECONDS
|
||||
default_retry_delay=RETRY_DELAY_SECONDS,
|
||||
routing_key=settings.RECALCULATE_GRADES_ROUTING_KEY
|
||||
)
|
||||
def recalculate_subsection_grade_v3(self, **kwargs):
|
||||
"""
|
||||
|
||||
@@ -162,6 +162,7 @@ def send_bulk_course_email(entry_id, _xmodule_instance_args):
|
||||
@task(
|
||||
name='lms.djangoapps.instructor_task.tasks.calculate_problem_responses_csv.v2',
|
||||
base=BaseInstructorTask,
|
||||
routing_key=settings.GRADES_DOWNLOAD_ROUTING_KEY,
|
||||
)
|
||||
def calculate_problem_responses_csv(entry_id, xmodule_instance_args):
|
||||
"""
|
||||
@@ -174,7 +175,7 @@ def calculate_problem_responses_csv(entry_id, xmodule_instance_args):
|
||||
return run_main_task(entry_id, task_fn, action_name)
|
||||
|
||||
|
||||
@task(base=BaseInstructorTask)
|
||||
@task(base=BaseInstructorTask, routing_key=settings.GRADES_DOWNLOAD_ROUTING_KEY)
|
||||
def calculate_grades_csv(entry_id, xmodule_instance_args):
|
||||
"""
|
||||
Grade a course and push the results to an S3 bucket for download.
|
||||
@@ -190,7 +191,7 @@ def calculate_grades_csv(entry_id, xmodule_instance_args):
|
||||
return run_main_task(entry_id, task_fn, action_name)
|
||||
|
||||
|
||||
@task(base=BaseInstructorTask)
|
||||
@task(base=BaseInstructorTask, routing_key=settings.GRADES_DOWNLOAD_ROUTING_KEY)
|
||||
def calculate_problem_grade_report(entry_id, xmodule_instance_args):
|
||||
"""
|
||||
Generate a CSV for a course containing all students' problem
|
||||
@@ -255,7 +256,7 @@ def calculate_may_enroll_csv(entry_id, xmodule_instance_args):
|
||||
return run_main_task(entry_id, task_fn, action_name)
|
||||
|
||||
|
||||
@task(base=BaseInstructorTask)
|
||||
@task(base=BaseInstructorTask, routing_key=settings.GRADES_DOWNLOAD_ROUTING_KEY)
|
||||
def generate_certificates(entry_id, xmodule_instance_args):
|
||||
"""
|
||||
Grade students and generate certificates.
|
||||
|
||||
@@ -15,6 +15,8 @@ from django.core.mail import EmailMessage
|
||||
from edxmako.shortcuts import render_to_string
|
||||
from openedx.core.djangoapps.site_configuration import helpers as configuration_helpers
|
||||
|
||||
ACE_ROUTING_KEY = getattr(settings, 'ACE_ROUTING_KEY', None)
|
||||
SOFTWARE_SECURE_VERIFICATION_ROUTING_KEY = getattr(settings, 'SOFTWARE_SECURE_VERIFICATION_ROUTING_KEY', None)
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
|
||||
@@ -71,7 +73,7 @@ class BaseSoftwareSecureTask(Task):
|
||||
)
|
||||
|
||||
|
||||
@task
|
||||
@task(routing_key=ACE_ROUTING_KEY)
|
||||
def send_verification_status_email(context):
|
||||
"""
|
||||
Spins a task to send verification status email to the learner
|
||||
@@ -97,6 +99,7 @@ def send_verification_status_email(context):
|
||||
bind=True,
|
||||
default_retry_delay=settings.SOFTWARE_SECURE_REQUEST_RETRY_DELAY,
|
||||
max_retries=settings.SOFTWARE_SECURE_RETRY_MAX_ATTEMPTS,
|
||||
routing_key=SOFTWARE_SECURE_VERIFICATION_ROUTING_KEY,
|
||||
)
|
||||
def send_request_to_ss_for_user(self, user_verification_id, copy_id_photo_from):
|
||||
"""
|
||||
|
||||
@@ -150,7 +150,7 @@ CELERY_QUEUES = {
|
||||
HIGH_MEM_QUEUE: {},
|
||||
}
|
||||
|
||||
CELERY_ROUTES = "{}celery.route_task".format(QUEUE_VARIANT)
|
||||
CELERY_ROUTES = "{}celery.Router".format(QUEUE_VARIANT)
|
||||
CELERYBEAT_SCHEDULE = {} # For scheduling tasks, entries can be added to this dict
|
||||
|
||||
# STATIC_ROOT specifies the directory where static files are
|
||||
@@ -974,62 +974,3 @@ COMPLETION_VIDEO_COMPLETE_PERCENTAGE = ENV_TOKENS.get('COMPLETION_VIDEO_COMPLETE
|
||||
COMPLETION_VIDEO_COMPLETE_PERCENTAGE)
|
||||
COMPLETION_VIDEO_COMPLETE_PERCENTAGE = ENV_TOKENS.get('COMPLETION_BY_VIEWING_DELAY_MS',
|
||||
COMPLETION_BY_VIEWING_DELAY_MS)
|
||||
|
||||
######################## CELERY ROUTING ########################
|
||||
|
||||
# Defines alternate environment tasks, as a dict of form { task_name: alternate_queue }
|
||||
ALTERNATE_ENV_TASKS = {}
|
||||
|
||||
# Defines the task -> alternate worker queue to be used when routing.
|
||||
EXPLICIT_QUEUES = {
|
||||
'openedx.core.djangoapps.content.course_overviews.tasks.async_course_overview_update': {
|
||||
'queue': GRADES_DOWNLOAD_ROUTING_KEY},
|
||||
'lms.djangoapps.bulk_email.tasks.send_course_email': {
|
||||
'queue': BULK_EMAIL_ROUTING_KEY},
|
||||
'openedx.core.djangoapps.heartbeat.tasks.sample_task': {
|
||||
'queue': HEARTBEAT_CELERY_ROUTING_KEY},
|
||||
'lms.djangoapps.instructor_task.tasks.calculate_grades_csv': {
|
||||
'queue': GRADES_DOWNLOAD_ROUTING_KEY},
|
||||
'lms.djangoapps.instructor_task.tasks.calculate_problem_grade_report': {
|
||||
'queue': GRADES_DOWNLOAD_ROUTING_KEY},
|
||||
'lms.djangoapps.instructor_task.tasks.generate_certificates': {
|
||||
'queue': GRADES_DOWNLOAD_ROUTING_KEY},
|
||||
'lms.djangoapps.email_marketing.tasks.get_email_cookies_via_sailthru': {
|
||||
'queue': ACE_ROUTING_KEY},
|
||||
'lms.djangoapps.email_marketing.tasks.update_user': {
|
||||
'queue': ACE_ROUTING_KEY},
|
||||
'lms.djangoapps.email_marketing.tasks.update_user_email': {
|
||||
'queue': ACE_ROUTING_KEY},
|
||||
'lms.djangoapps.email_marketing.tasks.update_course_enrollment': {
|
||||
'queue': ACE_ROUTING_KEY},
|
||||
'lms.djangoapps.verify_student.tasks.send_verification_status_email': {
|
||||
'queue': ACE_ROUTING_KEY},
|
||||
'lms.djangoapps.verify_student.tasks.send_ace_message': {
|
||||
'queue': ACE_ROUTING_KEY},
|
||||
'lms.djangoapps.verify_student.tasks.send_request_to_ss_for_user': {
|
||||
'queue': SOFTWARE_SECURE_VERIFICATION_ROUTING_KEY},
|
||||
'openedx.core.djangoapps.schedules.tasks._recurring_nudge_schedule_send': {
|
||||
'queue': ACE_ROUTING_KEY},
|
||||
'openedx.core.djangoapps.schedules.tasks._upgrade_reminder_schedule_send': {
|
||||
'queue': ACE_ROUTING_KEY},
|
||||
'openedx.core.djangoapps.schedules.tasks._course_update_schedule_send': {
|
||||
'queue': ACE_ROUTING_KEY},
|
||||
'openedx.core.djangoapps.schedules.tasks.v1.tasks.send_grade_to_credentials': {
|
||||
'queue': CREDENTIALS_GENERATION_ROUTING_KEY},
|
||||
'common.djangoapps.entitlements.tasks.expire_old_entitlements': {
|
||||
'queue': ENTITLEMENTS_EXPIRATION_ROUTING_KEY},
|
||||
'lms.djangoapps.grades.tasks.recalculate_course_and_subsection_grades_for_user': {
|
||||
'queue': POLICY_CHANGE_GRADES_ROUTING_KEY},
|
||||
'lms.djangoapps.grades.tasks.recalculate_subsection_grade_v3': {
|
||||
'queue': RECALCULATE_GRADES_ROUTING_KEY},
|
||||
'openedx.core.djangoapps.programs.tasks.v1.tasks.award_program_certificates': {
|
||||
'queue': PROGRAM_CERTIFICATES_ROUTING_KEY},
|
||||
'openedx.core.djangoapps.programs.tasks.v1.tasks.revoke_program_certificates': {
|
||||
'queue': PROGRAM_CERTIFICATES_ROUTING_KEY},
|
||||
'openedx.core.djangoapps.programs.tasks.v1.tasks.update_certificate_visible_date_on_course_update': {
|
||||
'queue': PROGRAM_CERTIFICATES_ROUTING_KEY},
|
||||
'openedx.core.djangoapps.programs.tasks.v1.tasks.award_course_certificate': {
|
||||
'queue': PROGRAM_CERTIFICATES_ROUTING_KEY},
|
||||
'openedx.core.djangoapps.coursegraph.dump_course_to_neo4j': {
|
||||
'queue': COURSEGRAPH_JOB_QUEUE},
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user