Skip to content

Commit 39e1f5b

Browse files
committed
add files share reindexing tasks to low priority queue
1 parent 9d15292 commit 39e1f5b

3 files changed

Lines changed: 9 additions & 7 deletions

File tree

api/share/utils.py

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

5252

53-
def update_share(resource):
53+
def update_share(resource, urgent=True):
5454
if not settings.SHARE_ENABLED:
5555
return
5656
if not hasattr(resource, 'guids'):
5757
logger.error(f'update_share called on non-guid resource: {resource}')
5858
return
59-
_enqueue_update_share(resource)
59+
_enqueue_update_share(resource, urgent)
6060

6161

62-
def _enqueue_update_share(osfresource):
62+
def _enqueue_update_share(osfresource, urgent):
6363
_osfguid_value = osfresource.guids.values_list('_id', flat=True).first()
6464
if not _osfguid_value:
6565
logger.warning(f'update_share skipping resource that has no guids: {osfresource}')
6666
return
67-
enqueue_task(task__update_share.s(_osfguid_value))
67+
enqueue_task(task__update_share.s(_osfguid_value, urgent=urgent))
6868

6969

7070
@celery_app.task(
@@ -73,7 +73,7 @@ def _enqueue_update_share(osfresource):
7373
max_retries=4,
7474
retry_backoff=True,
7575
)
76-
def task__update_share(self, guid: str, is_backfill=False, osfmap_partition_name='MAIN'):
76+
def task__update_share(self, guid: str, is_backfill=False, osfmap_partition_name='MAIN', urgent=True):
7777
"""
7878
Send SHARE/trove current metadata record(s) for the osf-guid-identified object
7979
"""

framework/celery_tasks/routers.py

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +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}
3537
return {
3638
'queue': match_by_module(task)
3739
}

osf/models/user.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -894,12 +894,12 @@ def merge_user(self, user):
894894
)
895895
for file in nodes_files_to_reindex.iterator(chunk_size=100):
896896
try:
897-
update_share(file)
897+
update_share(file, urgent=False)
898898
except Exception as e:
899899
logger.exception(f'Failed to SHARE reindex file {file._id} during user merge: {e}')
900900
for file in preprints_files_to_reindex.iterator(chunk_size=100):
901901
try:
902-
update_share(file)
902+
update_share(file, urgent=False)
903903
except Exception as e:
904904
logger.exception(f'Failed to SHARE reindex preprints file {file._id} during user merge: {e}')
905905

0 commit comments

Comments
 (0)