Skip to content

Commit d980170

Browse files
committed
use more explicit name approach to mark 'target_queue' for specific celery queue tasks
1 parent bda7a43 commit d980170

3 files changed

Lines changed: 19 additions & 11 deletions

File tree

api/share/utils.py

Lines changed: 14 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -50,24 +50,31 @@ def is_qa_resource(resource):
5050
return has_qa_tags or has_qa_title
5151

5252

53-
def update_share(resource, urgent=True):
53+
def update_share(resource, target_queue=None):
54+
"""
55+
By default, tasks are routed to queue based on module routing in CeleryRouter,
56+
:param resource: osf resource that is needed to be reindexed
57+
:param target_queue: should be task queue attribute of CeleryConfig f.e 'task_low_queue' for bulk background
58+
passing 'target_queue' allows low-level queue task run (reindexing files after a user merge) even though
59+
related module path may be marked to work with task_high_queue.
60+
"""
5461
if not settings.SHARE_ENABLED:
5562
return
5663
if not hasattr(resource, 'guids'):
5764
logger.error(f'update_share called on non-guid resource: {resource}')
5865
return
59-
if urgent:
60-
_enqueue_update_share(resource)
66+
if target_queue is not None:
67+
_enqueue_update_share(resource, target_queue)
6168
else:
62-
_enqueue_update_share(resource, urgent=False)
69+
_enqueue_update_share(resource)
6370

6471

65-
def _enqueue_update_share(osfresource, urgent=True):
72+
def _enqueue_update_share(osfresource, target_queue=None):
6673
_osfguid_value = osfresource.guids.values_list('_id', flat=True).first()
6774
if not _osfguid_value:
6875
logger.warning(f'update_share skipping resource that has no guids: {osfresource}')
6976
return
70-
enqueue_task(task__update_share.s(_osfguid_value, urgent=urgent))
77+
enqueue_task(task__update_share.s(_osfguid_value, target_queue=target_queue))
7178

7279

7380
@celery_app.task(
@@ -76,7 +83,7 @@ def _enqueue_update_share(osfresource, urgent=True):
7683
max_retries=4,
7784
retry_backoff=True,
7885
)
79-
def task__update_share(self, guid: str, is_backfill=False, osfmap_partition_name='MAIN', urgent=True):
86+
def task__update_share(self, guid: str, is_backfill=False, osfmap_partition_name='MAIN', target_queue=None):
8087
"""
8188
Send SHARE/trove current metadata record(s) for the osf-guid-identified object
8289
"""

framework/celery_tasks/routers.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,8 +32,8 @@ def route_for_task(self, task, args=None, kwargs=None):
3232
:param str task: Of the form 'full.module.path.to.class.function'
3333
:returns dict: Tells celery into which queue to route this task.
3434
"""
35-
if kwargs and (kwargs.get('urgent') is False):
36-
return {'queue': CeleryConfig.task_low_queue}
35+
if kwargs and (target_queue := kwargs.get('target_queue')):
36+
return {'queue': target_queue}
3737
return {
3838
'queue': match_by_module(task)
3939
}

osf/models/user.py

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -68,6 +68,7 @@
6868
from website import filters
6969
from website.project import new_bookmark_collection
7070
from website.util.metrics import OsfSourceTags, unregistered_created_source_tag
71+
from website.settings import CeleryConfig
7172
from importlib import import_module
7273
from osf.models.notification_type import NotificationTypeEnum
7374
from osf.utils.requests import string_type_request_headers
@@ -894,12 +895,12 @@ def merge_user(self, user):
894895
)
895896
for file in nodes_files_to_reindex.iterator(chunk_size=100):
896897
try:
897-
update_share(file, urgent=False)
898+
update_share(file, target_queue=CeleryConfig.task_low_queue)
898899
except Exception as e:
899900
logger.exception(f'Failed to SHARE reindex file {file._id} during user merge: {e}')
900901
for file in preprints_files_to_reindex.iterator(chunk_size=100):
901902
try:
902-
update_share(file, urgent=False)
903+
update_share(file, target_queue=CeleryConfig.task_low_queue)
903904
except Exception as e:
904905
logger.exception(f'Failed to SHARE reindex preprints file {file._id} during user merge: {e}')
905906

0 commit comments

Comments
 (0)