From d5edc33955bd3f06d23ee9b6ae5e5138eebf1ecf Mon Sep 17 00:00:00 2001 From: tcely Date: Fri, 4 Jul 2025 14:52:32 -0400 Subject: [PATCH] Convert `save_all_media_for_source` task --- tubesync/sync/tasks.py | 139 +++++++++++++++++++---------------------- 1 file changed, 65 insertions(+), 74 deletions(-) diff --git a/tubesync/sync/tasks.py b/tubesync/sync/tasks.py index d64a2471..d14af4aa 100644 --- a/tubesync/sync/tasks.py +++ b/tubesync/sync/tasks.py @@ -1040,6 +1040,71 @@ def rename_all_media_for_source(source_id): with atomic(durable=False): media.rename_files() + +@dynamic_retry(db_task, delay=600, priority=70, retries=15, queue=Val(TaskQueue.FS)) +def save_all_media_for_source(source_id): + ''' + Iterates all media items linked to a source and saves them to + trigger the post_save signal for every media item. Used when a + source has its parameters changed and all media needs to be + checked to see if its download status has changed. + ''' + db.reset_queries() + try: + source = Source.objects.get(pk=source_id) + except Source.DoesNotExist as e: + # Task triggered but the source no longer exists, do nothing + log.error(f'Task save_all_media_for_source(pk={source_id}) called but no ' + f'source exists with ID: {source_id}') + raise CancelExecution(_('no such source'), retry=False) from e + + # Keep out of the way of the index task! + # SQLite will be locked for a while if we start + # a large source, which reschedules a more costly task. + if 'sqlite' == db_vendor: + index_task = get_source_index_task(source_id) + if index_task and index_task.locked_by_pid_running(): + raise Exception(_('Indexing not completed')) + + refresh_qs = Media.objects.all().only( + 'pk', + 'uuid', + 'key', + 'title', # for name property + ).filter( + source=source, + can_download=False, + skip=False, + manual_skip=False, + downloaded=False, + metadata__isnull=False, + ) + save_qs = Media.objects.all().only( + 'pk', + 'uuid', + ).filter( + source=source, + ) + saved_later = { + str(media.pk) + for media in qs_gen(refresh_qs) + if media.has_metadata + } + refresh_formats.map(saved_later) + + # Trigger the post_save signal for each media item linked to this source as various + # flags may need to be recalculated + saved_now = { + str(media.pk) + for media in qs_gen(save_qs) + if str(media.pk) not in saved_later + } + + # wait for tasks to complete + res = save_media.map(saved_now) + res.get(blocking=True) + + # Old tasks system from background_task import background # noqa: E402 from background_task.exceptions import InvalidTaskError # noqa: E402 @@ -1113,80 +1178,6 @@ def download_media(media_id, override=False): raise InvalidTaskError(str(e)) from e -@background(schedule=dict(priority=30, run_at=600), queue=Val(TaskQueue.FS), remove_existing_tasks=True) -def save_all_media_for_source(source_id): - ''' - Iterates all media items linked to a source and saves them to - trigger the post_save signal for every media item. Used when a - source has its parameters changed and all media needs to be - checked to see if its download status has changed. - ''' - db.reset_queries() - try: - source = Source.objects.get(pk=source_id) - except Source.DoesNotExist as e: - # Task triggered but the source no longer exists, do nothing - log.error(f'Task save_all_media_for_source(pk={source_id}) called but no ' - f'source exists with ID: {source_id}') - raise InvalidTaskError(_('no such source')) from e - - refresh_qs = Media.objects.all().only( - 'pk', - 'uuid', - 'key', - 'title', # for name property - ).filter( - source=source, - can_download=False, - skip=False, - manual_skip=False, - downloaded=False, - metadata__isnull=False, - ) - save_qs = Media.objects.all().only( - 'pk', - 'uuid', - ).filter( - source=source, - ) - saved_later = set() - task = get_source_check_task(source_id) - if task: - task._verbose_name = remove_enclosed( - task.verbose_name, '[', ']', ' ', - valid='0123456789/,', - end=task.verbose_name.find('Check'), - ) - tvn_format = '1/{:,}' + f'/{refresh_qs.count():,}' - for mn, media in enumerate(qs_gen(refresh_qs), start=1): - update_task_status(task, tvn_format.format(mn)) - refresh_formats(str(media.pk)) - saved_later.add(media.uuid) - - # Keep out of the way of the index task! - # SQLite will be locked for a while if we start - # a large source, which reschedules a more costly task. - if 'sqlite' == db_vendor: - index_task = get_source_index_task(source_id) - if index_task and index_task.locked_by_pid_running(): - raise Exception(_('Indexing not completed')) - - # Trigger the post_save signal for each media item linked to this source as various - # flags may need to be recalculated - saved_now = set() - tvn_format = '2/{:,}' + f'/{save_qs.count():,}' - for mn, media in enumerate(qs_gen(save_qs), start=1): - if media.uuid not in saved_later: - update_task_status(task, tvn_format.format(mn)) - saved_now.add(str(media.pk)) - #save_model(media) - # Reset task.verbose_name to the saved value - update_task_status(task, None) - # wait for tasks to complete - res = save_media.map(saved_now) - res.get(blocking=True) - - @background(schedule=dict(priority=0, run_at=60), queue=Val(TaskQueue.NET), remove_existing_tasks=True) def wait_for_media_premiere(media_id): try: