From 89f2f2b0d3f57b7a343840c162051358e9d9ab0d Mon Sep 17 00:00:00 2001 From: tcely Date: Tue, 11 Mar 2025 17:39:33 -0400 Subject: [PATCH 01/10] Show progress on tasks page --- tubesync/sync/tasks.py | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/tubesync/sync/tasks.py b/tubesync/sync/tasks.py index 4a1884d8..a75c269a 100644 --- a/tubesync/sync/tasks.py +++ b/tubesync/sync/tasks.py @@ -203,9 +203,12 @@ def index_source_task(source_id): # Got some media, update the last crawl timestamp source.last_crawl = timezone.now() source.save() - log.info(f'Found {len(videos)} media items for source: {source}') + num_videos = len(videos) + log.info(f'Found {num_videos} media items for source: {source}') fields = lambda f, m: m.get_metadata_field(f) - for video in videos: + tvn_format = '[{}' + f'/{num_videos}] {task.verbose_name}' + for vn, video in enumerate(videos, start=1): + task.verbose_name = tvn_format.format(vn) # Create or update each video as a Media object key = video.get(source.key_field, None) if not key: From 3ad0fad72e27ee38171f0cf809b40fb5fd4fcc03 Mon Sep 17 00:00:00 2001 From: tcely Date: Wed, 12 Mar 2025 15:38:52 -0400 Subject: [PATCH 02/10] Get task functions --- tubesync/sync/tasks.py | 30 +++++++++++++++--------------- 1 file changed, 15 insertions(+), 15 deletions(-) diff --git a/tubesync/sync/tasks.py b/tubesync/sync/tasks.py index a75c269a..1f7ec3ab 100644 --- a/tubesync/sync/tasks.py +++ b/tubesync/sync/tasks.py @@ -114,27 +114,26 @@ def get_source_completed_tasks(source_id, only_errors=False): q['failed_at__isnull'] = False return CompletedTask.objects.filter(**q).order_by('-failed_at') +def get_tasks(task_name, id=None, /, instance=None): + assert not (id is None and instance is None) + arg = str(id or instance.pk) + return Task.objects.get_task(str(task_name), args=(arg,),) + +def get_first_task(task_name, id=None, /, *, instance=None): + tqs = get_tasks(task_name, id, instance).order_by('run_at') + return tqs[0] if tqs.count() else False def get_media_download_task(media_id): - try: - return Task.objects.get_task('sync.tasks.download_media', - args=(str(media_id),))[0] - except IndexError: - return False + return get_first_task('sync.tasks.download_media', media_id) def get_media_metadata_task(media_id): - try: - return Task.objects.get_task('sync.tasks.download_media_metadata', - args=(str(media_id),))[0] - except IndexError: - return False + return get_first_task('sync.tasks.download_media_metadata', media_id) def get_media_premiere_task(media_id): - try: - return Task.objects.get_task('sync.tasks.wait_for_media_premiere', - args=(str(media_id),))[0] - except IndexError: - return False + return get_first_task('sync.tasks.wait_for_media_premiere', media_id) + +def get_source_index_task(source_id): + return get_first_task('sync.tasks.index_source_task', source_id) def delete_task_by_source(task_name, source_id): now = timezone.now() @@ -206,6 +205,7 @@ def index_source_task(source_id): num_videos = len(videos) log.info(f'Found {num_videos} media items for source: {source}') fields = lambda f, m: m.get_metadata_field(f) + task = get_source_index_task(source_id) tvn_format = '[{}' + f'/{num_videos}] {task.verbose_name}' for vn, video in enumerate(videos, start=1): task.verbose_name = tvn_format.format(vn) From c0355e8f696973b72f1710d353c52a395c0ece9d Mon Sep 17 00:00:00 2001 From: tcely Date: Wed, 12 Mar 2025 15:46:55 -0400 Subject: [PATCH 03/10] Reset task verbose name after the loop ends --- tubesync/sync/tasks.py | 9 +++++++-- 1 file changed, 7 insertions(+), 2 deletions(-) diff --git a/tubesync/sync/tasks.py b/tubesync/sync/tasks.py index 1f7ec3ab..c99d83aa 100644 --- a/tubesync/sync/tasks.py +++ b/tubesync/sync/tasks.py @@ -206,9 +206,12 @@ def index_source_task(source_id): log.info(f'Found {num_videos} media items for source: {source}') fields = lambda f, m: m.get_metadata_field(f) task = get_source_index_task(source_id) - tvn_format = '[{}' + f'/{num_videos}] {task.verbose_name}' + if task: + verbose_name = task.verbose_name + tvn_format = '[{}' + f'/{num_videos}] {verbose_name}' for vn, video in enumerate(videos, start=1): - task.verbose_name = tvn_format.format(vn) + if task: + task.verbose_name = tvn_format.format(vn) # Create or update each video as a Media object key = video.get(source.key_field, None) if not key: @@ -239,6 +242,8 @@ def index_source_task(source_id): log.info(f'Indexed new media: {source} / {media}') except IntegrityError as e: log.error(f'Index media failed: {source} / {media} with "{e}"') + if task: + task.verbose_name = verbose_name # Tack on a cleanup of old completed tasks cleanup_completed_tasks() # Tack on a cleanup of old media From 408a3e1c952560324d086d719606365f922b555a Mon Sep 17 00:00:00 2001 From: tcely Date: Wed, 12 Mar 2025 17:32:19 -0400 Subject: [PATCH 04/10] Save the updated `verbose_name` of the task --- tubesync/sync/tasks.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/tubesync/sync/tasks.py b/tubesync/sync/tasks.py index c99d83aa..3b3dce45 100644 --- a/tubesync/sync/tasks.py +++ b/tubesync/sync/tasks.py @@ -212,6 +212,7 @@ def index_source_task(source_id): for vn, video in enumerate(videos, start=1): if task: task.verbose_name = tvn_format.format(vn) + task.save(update_fields=('verbose_name')) # Create or update each video as a Media object key = video.get(source.key_field, None) if not key: @@ -243,7 +244,8 @@ def index_source_task(source_id): except IntegrityError as e: log.error(f'Index media failed: {source} / {media} with "{e}"') if task: - task.verbose_name = verbose_name + task.verbose_name = verbose_name + task.save(update_fields=('verbose_name')) # Tack on a cleanup of old completed tasks cleanup_completed_tasks() # Tack on a cleanup of old media From f99c8fc5963d3235f1346c6de164baf5ef806640 Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 13 Mar 2025 06:16:01 -0400 Subject: [PATCH 05/10] Use set since tuple is dangerous for strings Even in explicit form tuple used a collection of characters instead of a single string. I hate that Python has these little traps. --- tubesync/sync/tasks.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/tubesync/sync/tasks.py b/tubesync/sync/tasks.py index 3b3dce45..3bb6a329 100644 --- a/tubesync/sync/tasks.py +++ b/tubesync/sync/tasks.py @@ -212,7 +212,7 @@ def index_source_task(source_id): for vn, video in enumerate(videos, start=1): if task: task.verbose_name = tvn_format.format(vn) - task.save(update_fields=('verbose_name')) + task.save(update_fields={'verbose_name'}) # Create or update each video as a Media object key = video.get(source.key_field, None) if not key: @@ -245,7 +245,7 @@ def index_source_task(source_id): log.error(f'Index media failed: {source} / {media} with "{e}"') if task: task.verbose_name = verbose_name - task.save(update_fields=('verbose_name')) + task.save(update_fields={'verbose_name'}) # Tack on a cleanup of old completed tasks cleanup_completed_tasks() # Tack on a cleanup of old media From ec45f29e1d30bd147954e966a6fafa5b46ddc909 Mon Sep 17 00:00:00 2001 From: tcely Date: Sat, 15 Mar 2025 21:05:39 -0400 Subject: [PATCH 06/10] Use smaller transactions --- tubesync/sync/tasks.py | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/tubesync/sync/tasks.py b/tubesync/sync/tasks.py index 3bb6a329..702086fe 100644 --- a/tubesync/sync/tasks.py +++ b/tubesync/sync/tasks.py @@ -179,7 +179,6 @@ def cleanup_removed_media(source, videos): @background(schedule=300, remove_existing_tasks=True) -@atomic(durable=True) def index_source_task(source_id): ''' Indexes media available from a Source object. @@ -210,9 +209,6 @@ def index_source_task(source_id): verbose_name = task.verbose_name tvn_format = '[{}' + f'/{num_videos}] {verbose_name}' for vn, video in enumerate(videos, start=1): - if task: - task.verbose_name = tvn_format.format(vn) - task.save(update_fields={'verbose_name'}) # Create or update each video as a Media object key = video.get(source.key_field, None) if not key: @@ -229,8 +225,12 @@ def index_source_task(source_id): published_dt = media.metadata_published(timestamp) if published_dt is not None: media.published = published_dt + if task: + task.verbose_name = tvn_format.format(vn) try: with atomic(): + if task: + task.save(update_fields={'verbose_name'}) media.save() log.debug(f'Indexed media: {source} / {media}') # log the new media instances From 021f4b172ae0971712ca6661445813f32ab5254f Mon Sep 17 00:00:00 2001 From: tcely Date: Tue, 18 Mar 2025 20:05:14 -0400 Subject: [PATCH 07/10] Display progress for checking task --- tubesync/sync/tasks.py | 32 +++++++++++++++++++++++++++----- 1 file changed, 27 insertions(+), 5 deletions(-) diff --git a/tubesync/sync/tasks.py b/tubesync/sync/tasks.py index 702086fe..d89d3f66 100644 --- a/tubesync/sync/tasks.py +++ b/tubesync/sync/tasks.py @@ -132,6 +132,9 @@ def get_media_metadata_task(media_id): def get_media_premiere_task(media_id): return get_first_task('sync.tasks.wait_for_media_premiere', media_id) +def get_source_check_task(source_id): + return get_first_task('sync.tasks.save_all_media_for_source', source_id) + def get_source_index_task(source_id): return get_first_task('sync.tasks.index_source_task', source_id) @@ -605,6 +608,7 @@ def save_all_media_for_source(source_id): already_saved = set() mqs = Media.objects.filter(source=source) + task = get_source_check_task(source_id) refresh_qs = mqs.filter( can_download=False, skip=False, @@ -612,22 +616,40 @@ def save_all_media_for_source(source_id): downloaded=False, metadata__isnull=False, ) - for media in refresh_qs: + if task: + verbose_name = task.verbose_name + tvn_format = '[{}' + f'/{refresh_qs.count()}] {verbose_name}' + for mn, media in enumerate(refresh_qs, start=1): + if task: + task.verbose_name = tvn_format.format(mn) + with atomic(): + task.save(update_fields={'verbose_name'}) try: media.refresh_formats except YouTubeError as e: log.debug(f'Failed to refresh formats for: {source} / {media.key}: {e!s}') pass else: - media.save() + with atomic(): + media.save() already_saved.add(media.uuid) # Trigger the post_save signal for each media item linked to this source as various # flags may need to be recalculated - with atomic(): - for media in mqs: + if task: + tvn_format = '[{}' + f'/{mqs.count()}] {verbose_name}' + for mn, media in enumerate(mqs, start=1): + if task: + task.verbose_name = tvn_format.format(mn) + with atomic(): + task.save(update_fields={'verbose_name'}) if media.uuid not in already_saved: - media.save() + with atomic(): + media.save() + if task: + task.verbose_name = verbose_name + with atomic(): + task.save(update_fields={'verbose_name'}) @background(schedule=60, remove_existing_tasks=True) From d2458a297965428729cd876db3803953bec0dbce Mon Sep 17 00:00:00 2001 From: tcely Date: Tue, 18 Mar 2025 20:12:54 -0400 Subject: [PATCH 08/10] Keep transactions specific to task --- tubesync/sync/tasks.py | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/tubesync/sync/tasks.py b/tubesync/sync/tasks.py index d89d3f66..fd5d1800 100644 --- a/tubesync/sync/tasks.py +++ b/tubesync/sync/tasks.py @@ -230,11 +230,10 @@ def index_source_task(source_id): media.published = published_dt if task: task.verbose_name = tvn_format.format(vn) - try: with atomic(): - if task: task.save(update_fields={'verbose_name'}) - media.save() + try: + media.save() log.debug(f'Indexed media: {source} / {media}') # log the new media instances new_media_instance = ( @@ -248,7 +247,8 @@ def index_source_task(source_id): log.error(f'Index media failed: {source} / {media} with "{e}"') if task: task.verbose_name = verbose_name - task.save(update_fields={'verbose_name'}) + with atomic(): + task.save(update_fields={'verbose_name'}) # Tack on a cleanup of old completed tasks cleanup_completed_tasks() # Tack on a cleanup of old media From 1f72718f317b6f7c66301b316f2f8db611589709 Mon Sep 17 00:00:00 2001 From: tcely Date: Tue, 18 Mar 2025 20:15:17 -0400 Subject: [PATCH 09/10] fixup: indentation --- tubesync/sync/tasks.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tubesync/sync/tasks.py b/tubesync/sync/tasks.py index fd5d1800..cf0d99d4 100644 --- a/tubesync/sync/tasks.py +++ b/tubesync/sync/tasks.py @@ -231,7 +231,7 @@ def index_source_task(source_id): if task: task.verbose_name = tvn_format.format(vn) with atomic(): - task.save(update_fields={'verbose_name'}) + task.save(update_fields={'verbose_name'}) try: media.save() log.debug(f'Indexed media: {source} / {media}') From abae403a8fbed7fb5520a551f91174c2aba47916 Mon Sep 17 00:00:00 2001 From: tcely Date: Tue, 18 Mar 2025 21:12:18 -0400 Subject: [PATCH 10/10] Remove extra blank lines --- tubesync/sync/tasks.py | 2 -- 1 file changed, 2 deletions(-) diff --git a/tubesync/sync/tasks.py b/tubesync/sync/tasks.py index c510b8fd..4fcf8455 100644 --- a/tubesync/sync/tasks.py +++ b/tubesync/sync/tasks.py @@ -254,12 +254,10 @@ def index_source_task(source_id): priority=20, verbose_name=verbose_name.format(media.pk), ) - if task: task.verbose_name = verbose_name with atomic(): task.save(update_fields={'verbose_name'}) - # Tack on a cleanup of old completed tasks cleanup_completed_tasks() with atomic(durable=True):