diff --git a/.github/workflows/ci.yaml b/.github/workflows/ci.yaml index 4c755057..02c10ddd 100644 --- a/.github/workflows/ci.yaml +++ b/.github/workflows/ci.yaml @@ -119,7 +119,6 @@ jobs: sudo ln -v -s -f -T ~/.config/TubeSync/config /config sudo ln -v -s -f -T ~/.config/TubeSync/downloads /downloads cp -v -p tubesync/tubesync/local_settings.py.example tubesync/tubesync/local_settings.py - cp -v -a -t "${Python3_ROOT_DIR}"/lib/python3.*/site-packages/background_task/ patches/background_task/* cp -v -a -t "${Python3_ROOT_DIR}"/lib/python3.*/site-packages/yt_dlp/ patches/yt_dlp/* cd tubesync && python3 -B manage.py collectstatic --no-input --link - name: Check with ruff diff --git a/Dockerfile b/Dockerfile index f509ab8b..32361d10 100644 --- a/Dockerfile +++ b/Dockerfile @@ -514,10 +514,6 @@ RUN --mount=type=tmpfs,target=/cache \ # Copy root COPY config/root / -# patch background_task -COPY patches/background_task/ \ - /usr/local/lib/python3/dist-packages/background_task/ - # patch yt_dlp COPY patches/yt_dlp/ \ /usr/local/lib/python3/dist-packages/yt_dlp/ diff --git a/Pipfile b/Pipfile index a23c3dbb..b2448a1c 100644 --- a/Pipfile +++ b/Pipfile @@ -14,7 +14,6 @@ pillow = "*" whitenoise = "*" gunicorn = "*" httptools = "*" -django-background-tasks = ">=1.2.8" django-basicauth = "*" psycopg = {extras = ["binary", "pool"], version = "*"} mysqlclient = "*" diff --git a/config/root/etc/s6-overlay/s6-rc.d/background-task-workers/contents.d/tubesync-network-worker b/config/root/etc/s6-overlay/s6-rc.d/background-task-workers/contents.d/tubesync-network-worker deleted file mode 100644 index 8b137891..00000000 --- a/config/root/etc/s6-overlay/s6-rc.d/background-task-workers/contents.d/tubesync-network-worker +++ /dev/null @@ -1 +0,0 @@ - diff --git a/config/root/etc/s6-overlay/s6-rc.d/background-task-workers/type b/config/root/etc/s6-overlay/s6-rc.d/background-task-workers/type deleted file mode 100644 index 757b4221..00000000 --- a/config/root/etc/s6-overlay/s6-rc.d/background-task-workers/type +++ /dev/null @@ -1 +0,0 @@ -bundle diff --git a/config/root/etc/s6-overlay/s6-rc.d/tubesync-network-worker/dependencies b/config/root/etc/s6-overlay/s6-rc.d/tubesync-network-worker/dependencies deleted file mode 100644 index 283e1305..00000000 --- a/config/root/etc/s6-overlay/s6-rc.d/tubesync-network-worker/dependencies +++ /dev/null @@ -1 +0,0 @@ -gunicorn \ No newline at end of file diff --git a/config/root/etc/s6-overlay/s6-rc.d/tubesync-network-worker/down-signal b/config/root/etc/s6-overlay/s6-rc.d/tubesync-network-worker/down-signal deleted file mode 100644 index d751378e..00000000 --- a/config/root/etc/s6-overlay/s6-rc.d/tubesync-network-worker/down-signal +++ /dev/null @@ -1 +0,0 @@ -SIGINT diff --git a/config/root/etc/s6-overlay/s6-rc.d/tubesync-network-worker/run b/config/root/etc/s6-overlay/s6-rc.d/tubesync-network-worker/run deleted file mode 100755 index 7f7bcd26..00000000 --- a/config/root/etc/s6-overlay/s6-rc.d/tubesync-network-worker/run +++ /dev/null @@ -1,6 +0,0 @@ -#!/command/with-contenv bash - -exec nice -n "${TUBESYNC_NICE:-1}" s6-setuidgid app \ - /usr/bin/python3 /app/manage.py process_tasks \ - --queue network --duration 43200 \ - --sleep "10.${RANDOM}" diff --git a/config/root/etc/s6-overlay/s6-rc.d/tubesync-network-worker/type b/config/root/etc/s6-overlay/s6-rc.d/tubesync-network-worker/type deleted file mode 100644 index 1780f9f4..00000000 --- a/config/root/etc/s6-overlay/s6-rc.d/tubesync-network-worker/type +++ /dev/null @@ -1 +0,0 @@ -longrun \ No newline at end of file diff --git a/config/root/etc/s6-overlay/s6-rc.d/user/contents.d/background-task-workers b/config/root/etc/s6-overlay/s6-rc.d/user/contents.d/background-task-workers deleted file mode 100644 index 8b137891..00000000 --- a/config/root/etc/s6-overlay/s6-rc.d/user/contents.d/background-task-workers +++ /dev/null @@ -1 +0,0 @@ - diff --git a/tubesync/common/migrations/0004_alter_taskhistory_task_id.py b/tubesync/common/migrations/0004_alter_taskhistory_task_id.py index 9d2dc45f..9b251e4e 100644 --- a/tubesync/common/migrations/0004_alter_taskhistory_task_id.py +++ b/tubesync/common/migrations/0004_alter_taskhistory_task_id.py @@ -1,6 +1,39 @@ # Generated by Django 5.2.4 on 2025-07-15 21:40 from django.db import migrations, models +from common.logger import log + + +def remove_duplicated_rows(apps, schema_editor): + def keep_which(task_id): + th_qs = TaskHistory.objects.all().order_by('-id') + tqs = th_qs.filter(task_id=task_id) + return tqs[0].id + + TaskHistory = apps.get_model("common", 'TaskHistory') + duplicates = set( + TaskHistory.objects.values( + 'task_id', + ).alias( + count=models.Count('id'), + ).filter( + count__gt=1, + ).values_list( + 'task_id', + flat=True, + ), + ) + + log.info(f'TaskHistory rows: {len(duplicates)=}') + for task_id, n in enumerate(duplicates, start=1): + keeping = keep_which(task_id) + log.debug(f'{n=}: {task_id=}: {keeping=}') + TaskHistory.objects.filter( + task_id=task_id, + ).exclude( + id=keeping, + ).delete() + log.info('TaskHistory rows: finished removing duplicates.') class Migration(migrations.Migration): @@ -10,6 +43,10 @@ class Migration(migrations.Migration): ] operations = [ + migrations.RunPython( + remove_duplicated_rows, + migrations.RunPython.noop, + ), migrations.AlterField( model_name='taskhistory', name='task_id', diff --git a/tubesync/common/templates/pagination.html b/tubesync/common/templates/pagination.html index e48b24d8..1bb0e2b5 100644 --- a/tubesync/common/templates/pagination.html +++ b/tubesync/common/templates/pagination.html @@ -3,7 +3,7 @@
diff --git a/tubesync/sync/models/source.py b/tubesync/sync/models/source.py index 42afd2ca..7930a214 100644 --- a/tubesync/sync/models/source.py +++ b/tubesync/sync/models/source.py @@ -5,7 +5,6 @@ from collections import deque as queue from pathlib import Path from django import db from django.conf import settings -from django.core.exceptions import SuspiciousOperation from django.core.validators import RegexValidator from django.utils import timezone from django.utils.text import slugify @@ -20,7 +19,7 @@ from ..choices import (Val, from ..fields import CommaSepChoiceField from ..youtube import ( get_media_info as get_youtube_media_info, - get_channel_image_info as get_youtube_channel_image_info, + get_image_info as get_youtube_image_info, ) from ._migrations import media_file_storage from ._private import _srctype_dict @@ -484,10 +483,7 @@ class Source(db.models.Model): @property def get_image_url(self): - if self.is_playlist: - raise SuspiciousOperation('This source is a playlist so it doesn\'t have thumbnail.') - - return get_youtube_channel_image_info(self.url) + return get_youtube_image_info(self.url) def directory_exists(self): diff --git a/tubesync/sync/signals.py b/tubesync/sync/signals.py index fdd65981..77da1650 100644 --- a/tubesync/sync/signals.py +++ b/tubesync/sync/signals.py @@ -6,20 +6,13 @@ from django.db import IntegrityError from django.db.models.signals import pre_save, post_save, pre_delete, post_delete from django.db.transaction import atomic, on_commit from django.dispatch import receiver -from django.utils import timezone from django.utils.translation import gettext_lazy as _ -from background_task.signals import ( - task_created, task_started, task_successful, task_rescheduled, task_failed, -) -from background_task.models import Task from common.logger import log from common.models import TaskHistory from common.utils import glob_quote, mkdir_p from .models import Source, Media, Metadata from .tasks import ( - delete_task_by_media, get_media_download_task, get_media_metadata_task, get_media_thumbnail_task, - map_task_to_instance, delete_all_media_for_source, rename_media, save_all_media_for_source, check_source_directory_exists, download_source_images, index_source, download_media_file, download_media_metadata, download_media_image, @@ -43,7 +36,7 @@ def source_pre_save(sender, instance, **kwargs): check_source_directory_exists.call_local(*args) existing_copy_channel_images = existing_source.copy_channel_images new_copy_channel_images = instance.copy_channel_images - if new_copy_channel_images and not (existing_copy_channel_images or instance.is_playlist): + if new_copy_channel_images and not existing_copy_channel_images: download_source_images(str(instance.pk)) existing_dirpath = existing_source.directory_path.resolve(strict=True) new_dirpath = instance.directory_path.resolve(strict=False) @@ -117,7 +110,7 @@ def source_post_save(sender, instance, created, **kwargs): # Check directory exists and create an indexing task for newly created sources if created: check_source_directory_exists(str(source.pk)) - if source.copy_channel_images and not source.is_playlist: + if source.copy_channel_images: download_source_images(str(source.pk)) if source.is_active: log.info(f'Scheduling first media indexing for source: {source.name}') @@ -166,114 +159,6 @@ def source_pre_delete(sender, instance, **kwargs): )) -@receiver(task_created, dispatch_uid='sync.signals.task_task_created') -@atomic(durable=False) -def task_task_created(sender, task=None, **kwargs): - if task is None: - return - task_obj = task - th, created = TaskHistory.objects.get_or_create( - task_id=str(task_obj.pk), - name=task_obj.task_name, - queue=task_obj.queue, - ) - th.scheduled_at = task_obj.run_at - th.priority = (100 - task_obj.priority) - th.repeat = task_obj.repeat - th.repeat_until = task_obj.repeat_until - th.task_params = list(task_obj.params()) - th.verbose_name = task_obj.verbose_name - th.save() - if created: - log.debug(f'Created a new task history record: {th.pk}: {th.verbose_name}') - - -@receiver(task_started, dispatch_uid='sync.signals.task_task_started') -@atomic(durable=False) -def task_task_started(sender, **kwargs): - locked_tasks = Task.objects.locked(timezone.now()) - for task_obj in locked_tasks: - th, created = TaskHistory.objects.get_or_create( - task_id=str(task_obj.pk), - name=task_obj.task_name, - queue=task_obj.queue, - ) - th.attempts += 1 - th.end_at = task_obj.locked_at - th.priority = (100 - task_obj.priority) - th.repeat = task_obj.repeat - th.repeat_until = task_obj.repeat_until - th.start_at = task_obj.locked_at - th.task_params = list(task_obj.params()) - th.verbose_name = task_obj.verbose_name - th.save() - if created: - log.debug(f'Started a new task history record: {th.pk}: {th.verbose_name}') - - -@receiver(task_rescheduled, dispatch_uid='sync.signals.task_task_rescheduled') -@atomic(durable=False) -def task_task_rescheduled(sender, task=None, **kwargs): - if task is None: - return - now_dt = timezone.now() - task_obj = task - th, created = TaskHistory.objects.get_or_create( - task_id=str(task_obj.pk), - name=task_obj.task_name, - queue=task_obj.queue, - ) - th.elapsed += ( - now_dt - task_obj.locked_at - ).total_seconds() - th.end_at = now_dt - th.scheduled_at = task_obj.run_at - th.start_at = task_obj.locked_at - th.save() - if created: - log.debug(f'Rescheduled a new task history record: {th.pk}: {th.verbose_name}') - -def merge_completed_task_into_history(task_id, task_obj): - th, created = TaskHistory.objects.get_or_create( - task_id=str(task_id), - name=task_obj.task_name, - queue=task_obj.queue, - ) - th.elapsed += ( - (task_obj.failed_at or task_obj.run_at) - task_obj.locked_at - ).total_seconds() - th.end_at = task_obj.run_at - th.failed_at = task_obj.failed_at - th.last_error = task_obj.last_error - th.repeat = task_obj.repeat - th.repeat_until = task_obj.repeat_until - th.start_at = task_obj.locked_at - th.verbose_name = task_obj.verbose_name - th.save() - - -@receiver(task_successful, dispatch_uid='sync.signals.task_task_successful') -@atomic(durable=False) -def task_task_successful(sender, task_id, completed_task, **kwargs): - merge_completed_task_into_history(task_id, completed_task) - - -@receiver(task_failed, dispatch_uid='sync.signals.task_task_failed') -@atomic(durable=False) -def task_task_failed(sender, task_id, completed_task, **kwargs): - merge_completed_task_into_history(task_id, completed_task) - # Triggered after a task fails by reaching its max retry attempts - obj, url = map_task_to_instance(completed_task, using_history=False) - if isinstance(obj, Source): - log.error(f'Permanent failure for source: {obj} task: {completed_task}') - obj.has_failed = True - obj.save() - - if isinstance(obj, Media) and completed_task.task_name == "sync.tasks.download_media_metadata": - log.error(f'Permanent failure for media: {obj} task: {completed_task}') - obj.skip = True - obj.save() - @receiver(post_save, sender=Media) def media_post_save(sender, instance, created, **kwargs): media = instance @@ -322,10 +207,12 @@ def media_post_save(sender, instance, created, **kwargs): # If the media is missing metadata schedule it to be downloaded if not (media.skip or media.has_metadata or existing_media_metadata_task): log.info(f'Scheduling task to download metadata for: {media.url}') - verbose_name = _('Downloading metadata for: {}: "{}"') - download_media_metadata( + TaskHistory.schedule( + download_media_metadata, str(media.pk), - verbose_name=verbose_name.format(media.key, media.name), + remove_duplicates=True, + vn_fmt=_('Downloading metadata for: {}: "{}"'), + vn_args=(media.key, media.name,), ) # If the media is missing a thumbnail schedule it to be downloaded (unless we are skipping this media) if not media.thumb_file_exists: @@ -379,9 +266,6 @@ def media_post_save(sender, instance, created, **kwargs): @receiver(pre_delete, sender=Media) def media_pre_delete(sender, instance, **kwargs): - # Triggered before media is deleted, delete any unlocked scheduled tasks - log.info(f'Deleting tasks for media: {instance.name}') - delete_task_by_media('sync.tasks.download_media_metadata', (str(instance.pk),)) # Remove thumbnail file for deleted media if instance.thumb: instance.thumb.delete(save=False) diff --git a/tubesync/sync/tasks.py b/tubesync/sync/tasks.py index 9cf3c954..ebf0d4aa 100644 --- a/tubesync/sync/tasks.py +++ b/tubesync/sync/tasks.py @@ -5,14 +5,12 @@ import os -import json import random import requests import time import uuid from collections import deque as queue from io import BytesIO -from hashlib import sha1 from pathlib import Path from datetime import timedelta from shutil import copyfile, rmtree @@ -46,15 +44,7 @@ db_vendor = db.connection.vendor register_huey_signals() -def get_hash(task_name, pk): - ''' - Create a background_task compatible hash for a Task or CompletedTask. - ''' - task_params = json.dumps(((str(pk),), {}), sort_keys=True) - return sha1(f'{task_name}{task_params}'.encode('utf-8')).hexdigest() - - -def map_task_to_instance(task, using_history=True): +def map_task_to_instance(task): ''' Reverse-maps a scheduled backgrond task to an instance. Requires the task name to be a known task function and the first argument to be a UUID. This is used @@ -67,25 +57,20 @@ def map_task_to_instance(task, using_history=True): 'sync.tasks.download_media_metadata': Media, 'sync.tasks.save_all_media_for_source': Source, 'sync.tasks.rename_all_media_for_source': Source, - 'sync.tasks.wait_for_media_premiere': Media, } MODEL_URL_MAP = { Source: 'sync:source', Media: 'sync:media-item', } # Unpack - task_func = task.name if using_history else task.task_name + task_func = task.name model = TASK_MAP.get(task_func, None) if not model: return None, None url = MODEL_URL_MAP.get(model, None) if not url: return None, None - task_args = task.task_params if using_history else None - try: - task_args = task_args or json.loads(task.task_params) - except (TypeError, ValueError, AttributeError): - return None, None + task_args = task.task_params if len(task_args) != 2: return None, None args, kwargs = task_args @@ -142,51 +127,43 @@ def update_task_status(task, status): def get_source_completed_tasks(source_id, only_errors=False): ''' - Returns a queryset of CompletedTask objects for a source by source ID. + Returns a queryset of TaskHistory objects for a source by source ID. ''' - q = {'task_params__istartswith': f'[["{source_id}"'} + qs = get_model_tasks(source_id) if only_errors: - q['failed_at__isnull'] = False - return CompletedTask.objects.filter(**q).order_by('-failed_at') + qs = qs.filter(failed_at__isnull=False) + return qs.order_by('-failed_at') -def get_model_task(model_pk, /, name=None, qs=None): +def get_model_tasks(model_pk, /, name=None, qs=None): if qs is None: qs = TaskHistory.objects.all() if name is not None: qs = qs.filter(name__endswith=name) - params_prefix = f'[["{model_pk}"' - qs = qs.filter(task_params__istartswith=params_prefix) - return qs[0] if qs.count() else False + #return qs.filter(task_params__0__0=model_pk) + return qs.filter(task_params__istartswith=f'[["{model_pk}"') def get_running_tasks(arg_dt=None, /): + max_run_time = getattr(settings, 'MAX_RUN_TIME', 3600) return TaskHistory.objects.running( now=arg_dt, - within=timezone.timedelta(seconds=settings.MAX_RUN_TIME), + within=timezone.timedelta(seconds=max_run_time), ) -def get_running_task_by_name(arg_str, media_id, /): +def get_running_tasks_by_name(arg_str, instance_id, /): name = arg_str if '.' not in name: name = f'sync.tasks.{name}' - tqs = get_running_tasks().filter(name=name, task_params__0__0=media_id) - return tqs[0] if tqs.count() else False + tqs = get_model_tasks(instance_id, qs=get_running_tasks()) + return tqs.filter(name=name) def get_media_download_task(media_id): - #return get_running_task_by_name('download_media_file', media_id) - return get_model_task( - media_id, - name='download_media_file', - qs=get_running_tasks(), - ) + tqs = get_running_tasks_by_name('download_media_file', media_id) + return tqs[0] if tqs.count() else False def get_media_thumbnail_task(media_id): - #return get_running_task_by_name('download_media_image', media_id) - return get_model_task( - media_id, - name='download_media_image', - qs=get_running_tasks(), - ) + tqs = get_running_tasks_by_name('download_media_image', media_id) + return tqs[0] if tqs.count() else False def get_source_index_task(source_id): #return get_running_task_by_name('index_source', source_id) @@ -200,33 +177,17 @@ def get_source_index_task(source_id): 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,),) + return get_running_tasks_by_name(str(task_name), arg) def get_first_task(task_name, id=None, /, *, instance=None): - tqs = get_tasks(task_name, id, instance).order_by('run_at') + tqs = get_tasks(task_name, id, instance).order_by('scheduled_at') return tqs[0] if tqs.count() else False def get_media_metadata_task(media_id): return get_first_task('sync.tasks.download_media_metadata', media_id) - -def delete_task_by_source(task_name, source_id): - now = timezone.now() - unlocked = Task.objects.unlocked(now) - qs = unlocked.filter( - task_name=task_name, - task_params__istartswith=f'[["{source_id}"', - ) - return qs.delete() - - -def delete_task_by_media(task_name, args): - max_run_time = getattr(settings, 'MAX_RUN_TIME', 3600) - now = timezone.now() - expires_at = now - timedelta(seconds=max_run_time) - task_qs = Task.objects.get_task(task_name, args=args) - unlocked = task_qs.filter(locked_by=None) | task_qs.filter(locked_at__lt=expires_at) - return unlocked.delete() +def get_source_index_task(source_id): + return get_first_task('sync.tasks.index_source', source_id) def cleanup_completed_tasks(): @@ -234,19 +195,9 @@ def cleanup_completed_tasks(): delta = timezone.now() - timedelta(days=days_to_keep) log.info(f'Deleting completed tasks older than {days_to_keep} days ' f'(run_at before {delta})') - CompletedTask.objects.filter(run_at__lt=delta).delete() TaskHistory.objects.filter(end_at__lt=delta).delete() -@atomic(durable=False) -def migrate_queues(): - return Task.objects.exclude( - queue=Val(TaskQueue.NET) - ).update( - queue=Val(TaskQueue.NET) - ) - - def save_model(instance): with atomic(durable=False): instance.save() @@ -282,23 +233,9 @@ def upcoming_media(): ), ) for media in qs_gen(qs): - media_id = str(media.pk) valid, hours = media.wait_for_premiere() if valid: save_model(media) - task = get_first_task('sync.tasks.wait_for_media_premiere', media_id) - if not task: - # create a task to update - when = media.published + timezone.timedelta(minutes=1) - vn_fmt = _('Waiting for the premiere of "{}" at: {}') - vn = vn_fmt.format( - media.key, - media.published.isoformat(' ', 'seconds'), - ) - wait_for_media_premiere(media_id, run_at=when, verbose_name=vn) - task = get_first_task('sync.tasks.wait_for_media_premiere', media_id) - if hours: - update_task_status(task, f'available in {hours} hours') log.debug(f'upcoming_media: wait_for_premiere: {media.key}: {valid=} {hours=}') @@ -364,7 +301,7 @@ def contains_http429(q, task_id, /): def wait_for_errors(model, /, *, queue_name=None, task_name=None): if task_name is None: task_name=tuple(( - 'sync.tasks.download_media', + 'sync.tasks.download_media_file', 'sync.tasks.download_media_metadata', )) elif isinstance(task_name, str): @@ -374,18 +311,10 @@ def wait_for_errors(model, /, *, queue_name=None, task_name=None): ft = get_first_task(tn, instance=model) if ft: tasks.append(ft) - window = timezone.timedelta(hours=3) + timezone.now() - tqs = Task.objects.filter( - task_name__in=task_name, - attempts__gt=0, - locked_at__isnull=True, - run_at__lte=window, - last_error__contains='HTTPError 429: Too Many Requests', - ) for task in tasks: update_task_status(task, 'paused (429)') - total_count = tqs.count() + total_count = int() if queue_name: from django_huey import get_queue q = get_queue(queue_name) @@ -661,11 +590,13 @@ def index_source(source_id): delay=65-(30*num), ) log.info(f'Scheduling task to download metadata for: {media.url}') - verbose_name = _('Downloading metadata for: "{}": {}') - download_media_metadata( + TaskHistory.schedule( + download_media_metadata, str(media.pk), - schedule=dict(priority=35), - verbose_name=verbose_name.format(media.key, media.name), + priority=65, + remove_duplicates=True, + vn_fmt=_('Downloading metadata for: "{}": {}'), + vn_args=(media.key, media.name,), ) # Reset task.verbose_name to the saved value update_task_status(task, None) @@ -725,10 +656,27 @@ def download_source_images(source_id): log.error(f'Task download_source_images(pk={source_id}) called but no ' f'source exists with ID: {source_id}') raise CancelExecution(_('no such source'), retry=False) from e - avatar, banner = source.get_image_url + avatar, banner, thumbnail = source.get_image_url log.info(f'Thumbnail URL for source with ID: {source_id} / {source} ' f'Avatar: {avatar} ' - f'Banner: {banner}') + f'Banner: {banner} ' + f'Thumbnail: {thumbnail}') + if thumbnail is not None: + url = thumbnail + i = get_remote_image(url) + image_file = BytesIO() + i.save(image_file, 'JPEG', quality=85, optimize=True, progressive=True) + + for file_name in ["thumbnail.jpg",]: + # Reset file pointer to the beginning for the next save + image_file.seek(0) + # Create a Django ContentFile from BytesIO stream + django_file = ContentFile(image_file.read()) + file_path = source.directory_path / file_name + with open(file_path, 'wb') as f: + f.write(django_file.read()) + i = image_file = None + if banner is not None: url = banner i = get_remote_image(url) @@ -828,7 +776,7 @@ def save_media(media_id): @db_task(delay=60, priority=60, retries=3, retry_delay=600, queue=Val(TaskQueue.LIMIT)) -def download_metadata(media_id): +def download_media_metadata(media_id): ''' Downloads the metadata for a media item. ''' @@ -841,7 +789,7 @@ def download_metadata(media_id): raise CancelExecution(_('no such media'), retry=False) from e if media.manual_skip: log.info(f'Task for ID: {media_id} / {media} skipped, due to task being manually skipped.') - return False + return source = media.source wait_for_errors( media, @@ -891,7 +839,6 @@ def download_metadata(media_id): if raise_exception: raise log.debug(str(e)) - return False else: keep_metadata_lock = True finally: @@ -934,7 +881,6 @@ def download_metadata(media_id): else: log.info(f'Saved {len(media.metadata_dumps())} bytes of metadata for: ' f'{source} / {media}: {media_id}') - return True finally: metadata_lock.acquired = False @@ -1288,29 +1234,3 @@ def delete_all_media_for_source(source_id, source_name, source_directory): rmtree(directory_path, True) -# Old tasks system -from background_task import background # noqa: E402 -from background_task.exceptions import InvalidTaskError # noqa: E402 -from background_task.models import Task, CompletedTask # noqa: E402 - - -@background(schedule=dict(priority=0, run_at=60), queue=Val(TaskQueue.NET), remove_existing_tasks=True) -def wait_for_media_premiere(media_id): - try: - media = Media.objects.get(pk=media_id) - except Media.DoesNotExist as e: - raise InvalidTaskError(_('no such media')) from e - else: - t = media.wait_for_premiere() - if t[0]: - save_model(media) - -@background(schedule=dict(priority=40, run_at=60), queue=Val(TaskQueue.NET), remove_existing_tasks=True) -def download_media_metadata(media_id): - try: - res = download_metadata(media_id) - return res.get(blocking=True) - except CancelExecution as e: - raise InvalidTaskError(str(e)) from e - - diff --git a/tubesync/sync/templates/sync/media.html b/tubesync/sync/templates/sync/media.html index d2d4e639..0ba43b04 100644 --- a/tubesync/sync/templates/sync/media.html +++ b/tubesync/sync/templates/sync/media.html @@ -9,19 +9,37 @@
{% if show_skipped %} - Hide skipped media + Hide skipped media {% else %} - Show skipped media + Show skipped media {% endif %}
{% if only_skipped %} - Only skipped media + Only skipped media {% else %} - Only skipped media + Only skipped media {% endif %}
+
+
+ +
+ +
+
+ + +
+
+
{% include 'infobox.html' with message=message %}
{% for m in media %} @@ -64,5 +82,5 @@
{% endfor %} -{% include 'pagination.html' with pagination=sources.paginator filter=source.pk show_skipped=show_skipped only_skipped=only_skipped%} +{% include 'pagination.html' %} {% endblock %} diff --git a/tubesync/sync/templates/sync/tasks.html b/tubesync/sync/templates/sync/tasks.html index 85f38c20..cdf84cc7 100644 --- a/tubesync/sync/templates/sync/tasks.html +++ b/tubesync/sync/templates/sync/tasks.html @@ -55,11 +55,9 @@ Error: "{{ task.error_message }}"
Task will be retried at {{ task.scheduled_at|date:'Y-m-d H:i:s' }} - {% if '-' not in task.task_id %} - + - {% endif %} {% empty %} There are no tasks with errors on this page. @@ -87,11 +85,10 @@ Priority: {{ task.priority }} Queue: {{ task.queue }}
Task will run {% if task.run_now %}immediately{% else %}at {{ task.scheduled_at|date:'Y-m-d H:i:s' }} - {% if '-' not in task.task_id %} - + - {% endif %}{% endif %} + {% endif %} {% empty %} diff --git a/tubesync/sync/tests.py b/tubesync/sync/tests.py index 90d20ca8..fc380a9c 100644 --- a/tubesync/sync/tests.py +++ b/tubesync/sync/tests.py @@ -15,7 +15,6 @@ from django.test import TestCase, Client, override_settings from django.utils import timezone from django_huey import DJANGO_HUEY, get_queue from common.models import TaskHistory -from huey.consumer_options import ConsumerConfig from .models import Source, Media from .tasks import ( cleanup_old_media, check_source_directory_exists, @@ -29,30 +28,19 @@ from .choices import (Val, Fallback, IndexSchedule, SourceResolution, class FrontEndTestCase(TestCase): + maxDiff = None @classmethod def setUpClass(cls): super().setUpClass() - cls._consumers = dict() - for qn, qc in DJANGO_HUEY.get('queues', dict()).items(): + # Use immediate mode to execute tasks in this process + for qn in DJANGO_HUEY.get('queues', dict()): q = get_queue(qn) - consumer_opts = qc.get('consumer', {}) - config = ConsumerConfig(**consumer_opts) - config.validate() - #consumer = q.create_consumer(**config.values) - #cls._consumers[qn] = consumer - #consumer.start() - q.immediate = True + # Set the storage variable before using the property. q.immediate_use_memory = True - - @classmethod - def tearDownClass(cls): - for qn, consumer in cls._consumers.items(): - consumer.stop(graceful=True) - super().tearDownClass() + q.immediate = False def setUp(self): - self.maxDiff = None # Disable general logging for test case logging.disable(logging.CRITICAL) @@ -189,13 +177,6 @@ class FrontEndTestCase(TestCase): def test_source(self): #logging.disable(logging.NOTSET) - def get_model_task(model_pk, /, name=None): - qs = TaskHistory.objects.all() - if name is not None: - qs = qs.filter(name__endswith=name) - params_prefix = f'[["{model_pk}"' - qs = qs.filter(task_params__istartswith=params_prefix) - return qs[0] if qs.count() else False # Sources overview page c = Client() response = c.get('/sources') @@ -257,15 +238,8 @@ class FrontEndTestCase(TestCase): name='sync.tasks.index_source', task_params__0__0=source_uuid, ).order_by('end_at') - self.assertNotEqual(list(), list(index_task_qs)) - self.assertNotEqual( - list(), - [ - th.__dict__ for th in TaskHistory.objects.all() - ] - ) - task = get_model_task(source_uuid, name='index_source') - self.assertNotEqual(False, task) + self.assertTrue(index_task_qs) + task = index_task_qs.last() self.assertEqual(task.queue, get_queue(Val(TaskQueue.LIMIT)).name) # save and refresh the Source source.refresh_from_db() @@ -466,27 +440,18 @@ class FrontEndTestCase(TestCase): end_at=now_dt, ) # Check the tasks to fetch the media thumbnails have been scheduled - def get_model_task(model_pk, /, name=None): - qs = TaskHistory.objects.all() - if name is not None: - qs = qs.filter(name__endswith=name) - params_prefix = f'[["{model_pk}"' - qs = qs.filter(task_params__istartswith=params_prefix) - return 1 == qs.count() - name_suffix = 'download_media_file' - found_download_task1 = get_model_task(test_media1_pk, name_suffix) - found_download_task2 = get_model_task(test_media2_pk, name_suffix) + found_download_task1 = get_media_download_task(test_media1_pk) + found_download_task2 = get_media_download_task(test_media2_pk) found_download_task3 = get_media_download_task(test_media3_pk) - name_suffix = 'download_media_image' - found_thumbnail_task1 = get_model_task(test_media1_pk, name_suffix) - found_thumbnail_task2 = get_model_task(test_media2_pk, name_suffix) + found_thumbnail_task1 = get_media_thumbnail_task(test_media1_pk) + found_thumbnail_task2 = get_media_thumbnail_task(test_media2_pk) found_thumbnail_task3 = get_media_thumbnail_task(test_media3_pk) self.assertTrue(found_download_task1) self.assertTrue(found_download_task2) - self.assertTrue(not not found_download_task3) + self.assertTrue(found_download_task3) self.assertTrue(found_thumbnail_task1) self.assertTrue(found_thumbnail_task2) - self.assertTrue(not not found_thumbnail_task3) + self.assertTrue(found_thumbnail_task3) # Check the media is listed on the media overview page response = c.get('/media') self.assertEqual(response.status_code, 200) diff --git a/tubesync/sync/views.py b/tubesync/sync/views.py index 26fd1644..9c49d4ab 100644 --- a/tubesync/sync/views.py +++ b/tubesync/sync/views.py @@ -22,16 +22,16 @@ from django.utils.translation import gettext_lazy as _ from common.models import TaskHistory from common.timestamp import timestamp_to_datetime from common.utils import append_uri_params, mkdir_p, multi_key_sort -from background_task.models import Task from django_huey import DJANGO_HUEY, get_queue from common.huey import h_q_reset_tasks +from common.logger import log from .models import Source, Media, MediaServer from .forms import (ValidateSourceForm, ConfirmDeleteSourceForm, RedownloadMediaForm, SkipMediaForm, EnableMediaForm, ResetTasksForm, ScheduleTaskForm, ConfirmDeleteMediaServerForm, SourceForm) from .utils import delete_file, validate_url from .tasks import ( - map_task_to_instance, get_error_message, migrate_queues, delete_task_by_media, + map_task_to_instance, get_error_message, get_running_tasks, get_media_download_task, get_source_completed_tasks, check_source_directory_exists, index_source, download_media_image, ) @@ -43,9 +43,6 @@ from . import youtube def get_waiting_tasks(): - background_task_ids = { - str(t.pk) for t in Task.objects.all() - } huey_queue_names = (DJANGO_HUEY or {}).get('queues', {}) huey_queues = list(map(get_queue, huey_queue_names)) huey_task_ids = { @@ -56,7 +53,7 @@ def get_waiting_tasks(): ) } return TaskHistory.objects.filter( - task_id__in=huey_task_ids.union(background_task_ids), + task_id__in=huey_task_ids, ) @@ -503,22 +500,32 @@ class MediaView(ListView): self.filter_source = None self.show_skipped = False self.only_skipped = False + self.query = None + self.search_description = False super().__init__(*args, **kwargs) def dispatch(self, request, *args, **kwargs): - filter_by = request.GET.get('filter', '') + def is_active(arg, /): + return str(arg).strip().lower() in ( + 'enable', 'enabled', 'on', 'true', 'yes', '1', + ) + def post_or_get(request, /, key, default=None): + return request.POST.get(key) or request.GET.get(key) or default + + filter_by = post_or_get(request, 'filter', '') if filter_by: try: self.filter_source = Source.objects.get(pk=filter_by) except Source.DoesNotExist: self.filter_source = None - show_skipped = request.GET.get('show_skipped', '').strip() - if show_skipped == 'yes': + show_skipped = post_or_get(request, 'show_skipped', '') + if is_active(show_skipped): self.show_skipped = True - if not self.show_skipped: - only_skipped = request.GET.get('only_skipped', '').strip() - if only_skipped == 'yes': - self.only_skipped = True + only_skipped = post_or_get(request, 'only_skipped', '') + if is_active(only_skipped): + self.only_skipped = True + self.query = post_or_get(request, 'query') + self.search_description = is_active(post_or_get(request, 'search_description')) return super().dispatch(request, *args, **kwargs) def get_queryset(self): @@ -526,6 +533,21 @@ class MediaView(ListView): if self.filter_source: q = q.filter(source=self.filter_source) + if self.query: + needle = self.query + if self.search_description: + q = q.filter( + Q(new_metadata__value__fulltitle__icontains=needle) | + Q(new_metadata__value__description__icontains=needle) | + Q(title__icontains=needle) | + Q(key__contains=needle) + ) + else: + q = q.filter( + Q(new_metadata__value__fulltitle__icontains=needle) | + Q(title__icontains=needle) | + Q(key__contains=needle) + ) if self.only_skipped: q = q.filter(Q(can_download=False) | Q(skip=True) | Q(manual_skip=True)) elif not self.show_skipped: @@ -543,6 +565,8 @@ class MediaView(ListView): data['source'] = self.filter_source data['show_skipped'] = self.show_skipped data['only_skipped'] = self.only_skipped + data['query'] = self.query or str() + data['search_description'] = self.search_description return data @@ -662,8 +686,6 @@ class MediaRedownloadView(FormView, SingleObjectMixin): return super().dispatch(request, *args, **kwargs) def form_valid(self, form): - # Delete any active download tasks for the media - delete_task_by_media('sync.tasks.download_media', (str(self.object.pk),)) # If the thumbnail file exists on disk, delete it if self.object.thumb_file_exists: delete_file(self.object.thumb.path) @@ -714,8 +736,6 @@ class MediaSkipView(FormView, SingleObjectMixin): return super().dispatch(request, *args, **kwargs) def form_valid(self, form): - # Delete any active download tasks for the media - delete_task_by_media('sync.tasks.download_media', (str(self.object.pk),)) # If the media file exists on disk, delete it if self.object.media_file_exists: # Delete all files which contains filename @@ -886,7 +906,6 @@ class TasksView(ListView): data['total_errors'] = errors_qs.count() data['scheduled'] = list() data['total_scheduled'] = scheduled_qs.count() - data['migrated'] = migrate_queues() data['wait_for_database_queue'] = False def add_to_task(task): @@ -901,37 +920,6 @@ class TasksView(ListView): return 'error' return True and obj - verbose_names = dict() - for task in Task.objects.filter(locked_by__isnull=False): - # There was broken logic in `Task.objects.locked()`, work around it. - # With that broken logic, the tasks never resume properly. - # This check unlocks the tasks without a running process. - # `task.locked_by_pid_running()` returns: - # - `True`: locked and PID exists - # - `False`: locked and PID does not exist - # - `None`: not `locked_by`, so there was no PID to check - locked_by_pid_running = task.locked_by_pid_running() - if locked_by_pid_running is False: - task.locked_by = None - # do not wait for the task to expire - task.locked_at = None - task.save() - task_id = str(task.pk) - verbose_names[task_id] = task.verbose_name - try: - task = TaskHistory.objects.get(task_id=task_id) - except TaskHistory.DoesNotExist: - # possibly create a new instance? - pass - else: - if locked_by_pid_running and add_to_task(task): - # Use the status if it is available - task.verbose_name = verbose_names.get(task_id) or task.verbose_name - data['running'].append(task) - elif locked_by_pid_running and 'wait_for_database_queue' in task.name: - data['wait_for_database_queue'] = True - verbose_names = None - for task in running_qs: if task in data['running']: continue @@ -1036,7 +1024,6 @@ class ResetTasks(FormView): def form_valid(self, form): # Delete all tasks - Task.objects.all().delete() huey_queue_names = (DJANGO_HUEY or {}).get('queues', {}) for queue_name in huey_queue_names: h_q_reset_tasks(queue_name) @@ -1059,7 +1046,7 @@ class TaskScheduleView(FormView, SingleObjectMixin): template_name = 'sync/task-schedule.html' form_class = ScheduleTaskForm - model = Task + model = TaskHistory errors = dict( invalid_when=_('The type ({}) was incorrect.'), when_before_now=_('The date and time must be in the future.'), @@ -1109,7 +1096,6 @@ class TaskScheduleView(FormView, SingleObjectMixin): ) def form_valid(self, form): - max_attempts = getattr(settings, 'MAX_ATTEMPTS', 15) when = form.cleaned_data.get('when') if not isinstance(when, self.now.__class__): @@ -1130,14 +1116,24 @@ class TaskScheduleView(FormView, SingleObjectMixin): if form.errors: return super().form_invalid(form) - self.object.attempts = max_attempts // 2 - self.object.run_at = max(self.now, when) - self.object.save() - TaskHistory.objects.filter( - task_id=str(self.object.pk), - ).update( - scheduled_at=self.object.run_at, - ) + huey_queue_names = (DJANGO_HUEY or {}).get('queues', {}) + huey_queues = list(map(get_queue, huey_queue_names)) + pk = self.object.pk + queue = self.object.queue + task_id = self.object.task_id + matching = { q for q in huey_queues if q.name == queue } + try: + q = matching.pop() + except KeyError as e: + msg = f'TaskScheduleView: queue not found: {pk=} {queue=}' + log.exception(msg, exc_info=e) + else: + self.object.scheduled_at = max(self.now, when) + if q and q.reschedule(task_id, self.object.scheduled_at): + self.object.save() + else: + msg = f'TaskScheduleView: task not found: {pk=} {task_id=}' + log.warning(msg) return super().form_valid(form) diff --git a/tubesync/sync/youtube.py b/tubesync/sync/youtube.py index c0e55f6d..7bfda388 100644 --- a/tubesync/sync/youtube.py +++ b/tubesync/sync/youtube.py @@ -81,7 +81,10 @@ def get_channel_id(url): else: return channel_id -def get_channel_image_info(url): +def get_image_info(url): + avatar_url = None + banner_url = None + thumbnail_url = None opts = get_yt_opts() opts.update({ 'skip_download': True, @@ -94,20 +97,53 @@ def get_channel_image_info(url): with yt_dlp.YoutubeDL(opts) as y: try: response = y.extract_info(url, download=False) - - avatar_url = None - banner_url = None + except yt_dlp.utils.DownloadError as e: + raise YouTubeError(f'Failed to extract info for "{url}": {e}') from e + else: + max_height = 0 for thumbnail in response['thumbnails']: + thumbnail_height = thumbnail.get('height') + try: + thumbnail_height = int(thumbnail_height) + except (TypeError, ValueError,): + thumbnail_height = int() if thumbnail['id'] == 'avatar_uncropped': avatar_url = thumbnail['url'] - if thumbnail['id'] == 'banner_uncropped': + elif thumbnail['id'] == 'banner_uncropped': banner_url = thumbnail['url'] - if banner_url is not None and avatar_url is not None: - break - - return avatar_url, banner_url - except yt_dlp.utils.DownloadError as e: - raise YouTubeError(f'Failed to extract channel info for "{url}": {e}') from e + elif thumbnail_height > max_height: + max_height = thumbnail_height + thumbnail_url = thumbnail['url'] + try: + entry_type = response['entries'][0].get('_type') + except IndexError: + # an empty entries list + pass + else: + if 'url' == entry_type: + del response['entries'] + elif 'playlist' == entry_type: + for playlist in response['entries']: + del playlist['entries'] + from .models import Metadata + t = Metadata.objects.defer('value').filter( + source__isnull=True, + media__isnull=True, + ).get_or_create( + key=response['id'], + site=response['extractor_key'], + ) + md = t[0] + field_defaults = { + f.attname: f.get_default() + for f in md._meta.fields + if f.has_default() + } + if 'retrieved' in field_defaults: + md.retrieved = field_defaults['retrieved'] + md.value = response + md.save() + return avatar_url, banner_url, thumbnail_url def _subscriber_only(msg='', response=None): diff --git a/tubesync/tubesync/local_settings.py.container b/tubesync/tubesync/local_settings.py.container index 566cfb24..f4205772 100644 --- a/tubesync/tubesync/local_settings.py.container +++ b/tubesync/tubesync/local_settings.py.container @@ -51,8 +51,13 @@ else: "OPTIONS": { "timeout": 10, "transaction_mode": "IMMEDIATE", + # PRAGMA locking_mode = NORMAL | EXCLUSIVE # PRAGMA journal_mode = DELETE | TRUNCATE | PERSIST | MEMORY | WAL | OFF + # DO NOT change locking_mode to EXCLUSIVE! + # This is a foot-gun, and invalidates a behavior the code relies on. + # journal_mode WAL offers increased concurrency, the default is DELETE. "init_command": """ + PRAGMA locking_mode = NORMAL; PRAGMA journal_mode = TRUNCATE; PRAGMA journal_size_limit = 67108864; PRAGMA legacy_alter_table = OFF; diff --git a/tubesync/tubesync/settings.py b/tubesync/tubesync/settings.py index 3360be3a..da2fab5d 100644 --- a/tubesync/tubesync/settings.py +++ b/tubesync/tubesync/settings.py @@ -25,7 +25,6 @@ INSTALLED_APPS = [ 'django.contrib.staticfiles', 'django.contrib.humanize', 'sass_processor', - 'background_task', 'django_huey', 'common', 'sync',