diff --git a/tubesync/common/huey.py b/tubesync/common/huey.py index 8a8a2ea6..6f34b5e9 100644 --- a/tubesync/common/huey.py +++ b/tubesync/common/huey.py @@ -5,6 +5,7 @@ from huey import ( CancelExecution, SqliteHuey as huey_SqliteHuey, signals, utils, ) +from .timestamp import datetime_to_timestamp, timestamp_to_datetime class SqliteHuey(huey_SqliteHuey): @@ -277,7 +278,15 @@ def historical_task(signal_name, task_obj, exception_obj=None, /, *, huey=None): th.verbose_name = f'{th.name} with: {model_instance.key}' if hasattr(model_instance, 'name'): th.verbose_name += f' / {model_instance.name}' + elif signal_name == signals.SIGNAL_SCHEDULED: + if huey.utc: + th.scheduled_at = task_obj.eta.replace(tzinfo=datetime.UTC) + else: # this path is unlikely + th.scheduled_at = timestamp_to_datetime( + datetime_to_timestamp(task_obj.eta, integer=False), + ).astimezone(tz=datetime.UTC) th.end_at = signal_dt + th.elapsed = history['elapsed'] th.save() # Registration of shared signal handlers diff --git a/tubesync/common/migrations/0002_taskhistory_elapsed_taskhistory_scheduled_at_and_more.py b/tubesync/common/migrations/0002_taskhistory_elapsed_taskhistory_scheduled_at_and_more.py new file mode 100644 index 00000000..68801a7e --- /dev/null +++ b/tubesync/common/migrations/0002_taskhistory_elapsed_taskhistory_scheduled_at_and_more.py @@ -0,0 +1,35 @@ +# Generated by Django 5.2.4 on 2025-07-03 21:51 + +import django.utils.timezone +from django.db import migrations, models + + +class Migration(migrations.Migration): + + dependencies = [ + ('common', '0001_initial'), + ] + + operations = [ + migrations.AddField( + model_name='taskhistory', + name='elapsed', + field=models.FloatField(default=float), + ), + migrations.AddField( + model_name='taskhistory', + name='scheduled_at', + field=models.DateTimeField(db_index=True, default=django.utils.timezone.now), + ), + migrations.AlterField( + model_name='taskhistory', + name='attempts', + field=models.IntegerField(db_index=True, default=int), + ), + migrations.AlterField( + model_name='taskhistory', + name='repeat', + field=models.BigIntegerField(default=int), + ), + ] + diff --git a/tubesync/common/models/tasks.py b/tubesync/common/models/tasks.py index 62f72c88..8c1e9ecb 100644 --- a/tubesync/common/models/tasks.py +++ b/tubesync/common/models/tasks.py @@ -50,6 +50,8 @@ class TaskHistory(models.Model): verbose_name = models.CharField(max_length=255, null=True, blank=True) start_at = models.DateTimeField(db_index=True, null=True, blank=True) + # when the task was scheduled to run + scheduled_at = models.DateTimeField(default=timezone.now, db_index=True) # when the task was completed end_at = models.DateTimeField(default=timezone.now, db_index=True) @@ -59,13 +61,14 @@ class TaskHistory(models.Model): queue = models.CharField(max_length=190, db_index=True, null=True, blank=True) # how many times the task has been tried - attempts = models.IntegerField(default=0, db_index=True) + attempts = models.IntegerField(default=int, db_index=True) # when the task last failed failed_at = models.DateTimeField(db_index=True, null=True, blank=True) # details of the error that occurred last_error = models.TextField(blank=True) - repeat = models.BigIntegerField(default=0) + elapsed = models.FloatField(default=float) + repeat = models.BigIntegerField(default=int) repeat_until = models.DateTimeField(null=True, blank=True) objects = TaskHistoryQuerySet.as_manager() diff --git a/tubesync/common/timestamp.py b/tubesync/common/timestamp.py index d8b69178..7a618960 100644 --- a/tubesync/common/timestamp.py +++ b/tubesync/common/timestamp.py @@ -1,8 +1,9 @@ import datetime +import math utc_tz = datetime.timezone.utc -posix_epoch = datetime.datetime.fromtimestamp(0, utc_tz) +posix_epoch = datetime.datetime.fromtimestamp(0, tz=utc_tz) def add_epoch(seconds): @@ -13,6 +14,9 @@ def add_epoch(seconds): def subtract_epoch(arg_dt, /): assert isinstance(arg_dt, datetime.datetime) + if arg_dt.utcoffset() is None: # naive + return arg_dt - datetime.datetime.fromtimestamp(0, tz=None) + utc_dt = arg_dt.astimezone(utc_tz) return utc_dt - posix_epoch @@ -22,7 +26,7 @@ def datetime_to_timestamp(arg_dt, /, *, integer=True): if not integer: return timestamp - return round(timestamp) + return math.ceil(timestamp) def timestamp_to_datetime(seconds, /): return add_epoch(seconds=seconds).astimezone(utc_tz) diff --git a/tubesync/sync/signals.py b/tubesync/sync/signals.py index f89776a8..f42c2c7c 100644 --- a/tubesync/sync/signals.py +++ b/tubesync/sync/signals.py @@ -8,7 +8,9 @@ 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_started, task_successful, task_failed +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 @@ -168,6 +170,28 @@ def source_post_delete(sender, instance, **kwargs): delete_task_by_source('sync.tasks.save_all_media_for_source', instance.pk) +@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): @@ -188,15 +212,40 @@ def task_task_started(sender, **kwargs): 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}') + 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 @@ -218,7 +267,7 @@ def task_task_successful(sender, task_id, completed_task, **kwargs): 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) + 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 diff --git a/tubesync/sync/tasks.py b/tubesync/sync/tasks.py index 684f9c15..d64a2471 100644 --- a/tubesync/sync/tasks.py +++ b/tubesync/sync/tasks.py @@ -51,7 +51,7 @@ def get_hash(task_name, pk): return sha1(f'{task_name}{task_params}'.encode('utf-8')).hexdigest() -def map_task_to_instance(task): +def map_task_to_instance(task, using_history=True): ''' 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 @@ -72,7 +72,12 @@ def map_task_to_instance(task): Media: 'sync:media-item', } # Unpack - task_func, task_args_str = task.task_name, task.task_params + task_args = None + if using_history: + task_func = task.name + task_args = task.task_params + else: + task_func, task_args_str = task.task_name, task.task_params model = TASK_MAP.get(task_func, None) if not model: return None, None @@ -80,7 +85,7 @@ def map_task_to_instance(task): if not url: return None, None try: - task_args = json.loads(task_args_str) + task_args = task_args or json.loads(task_args_str) except (TypeError, ValueError, AttributeError): return None, None if len(task_args) != 2: diff --git a/tubesync/sync/templates/sync/tasks-completed.html b/tubesync/sync/templates/sync/tasks-completed.html index ec2a0aa8..796d1ada 100644 --- a/tubesync/sync/templates/sync/tasks-completed.html +++ b/tubesync/sync/templates/sync/tasks-completed.html @@ -16,16 +16,16 @@ {% if task.has_error %} - {{ task.verbose_name }}
+ {% if task.verbose_name %}{{ task.verbose_name }}{% else %}{{ task.name }}{% endif %}
Queue: "{{ task.queue }}"
Error: "{{ task.error_message }}"
{% else %} - {{ task.verbose_name }}
+ {% if task.verbose_name %}{{ task.verbose_name }}{% else %}{{ task.name }}{% endif %}
Queue: "{{ task.queue }}"
{% endif %} - Task locked for: {{ task.run_at|sub:task.locked_at|timedelta }}
- Task locked at {{ task.locked_at|date:'Y-m-d H:i:s' }}
- Task ended at {{ task.run_at|date:'Y-m-d H:i:s' }} + Task locked for: {{ task.end_at|sub:task.start_at|timedelta }}
+ Task locked at {{ task.start_at|date:'Y-m-d H:i:s' }}
+ Task ended at {{ task.end_at|date:'Y-m-d H:i:s' }}
{% empty %} There have been no completed tasks{% if source %} that match the specified source filter{% endif %}. @@ -33,5 +33,5 @@ -{% include 'pagination.html' with pagination=sources.paginator filter=source.pk %} +{% include 'pagination.html' with filter=source.pk %} {% endblock %} diff --git a/tubesync/sync/templates/sync/tasks.html b/tubesync/sync/templates/sync/tasks.html index a074df85..74997164 100644 --- a/tubesync/sync/templates/sync/tasks.html +++ b/tubesync/sync/templates/sync/tasks.html @@ -23,9 +23,9 @@

{% for task in running %} - + {{ task }}
- Task started at {{ task.locked_at|date:'Y-m-d H:i:s' }} + Task started at {{ task.start_at|date:'Y-m-d H:i:s' }}
{% empty %} There are no running tasks. @@ -54,8 +54,8 @@ {{ task }}, attempted {{ task.attempts }} time{{ task.attempts|pluralize }}
Error: "{{ task.error_message }}"
- Task will be retried at {{ task.run_at|date:'Y-m-d H:i:s' }} - + Task will be retried at {{ task.scheduled_at|date:'Y-m-d H:i:s' }} +
@@ -73,19 +73,20 @@

Tasks which are scheduled to run in the future or are waiting in a queue to be processed. They can be waiting for an available worker to run immediately, or - run in the future at the specified "run at" time. + run in the future at the specified "scheduled at" time.

{% for task in scheduled %}
{{ task }}
- {% if task.instance.index_schedule and task.repeat > 0 %}Scheduled to run {{ task.instance.get_index_schedule_display|lower }}.
{% endif %} - Task will run {% if task.run_now %}immediately{% else %}at {{ task.run_at|date:'Y-m-d H:i:s' }} + {% if task.instance.is_active and 'once' not in task.verbose_name %}Scheduled to run {{ task.instance.get_index_schedule_display|lower }}.
{% endif %} + 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 %} @@ -94,7 +95,7 @@
-{% include 'pagination.html' with pagination=sources.paginator filter=source.pk %} +{% include 'pagination.html' with filter=source.pk %}

Completed

diff --git a/tubesync/sync/views.py b/tubesync/sync/views.py index 7afb2804..eb31a966 100644 --- a/tubesync/sync/views.py +++ b/tubesync/sync/views.py @@ -13,16 +13,17 @@ from django.core.exceptions import SuspiciousFileOperation from django.http import HttpResponse from django.urls import reverse_lazy from django.db import connection, IntegrityError -from django.db.models import Q, Count, Sum, When, Case +from django.db.models import F, Q, Count, Sum, When, Case from django.forms import Form, ValidationError from django.utils.text import slugify from django.utils._os import safe_join from django.utils import timezone 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, CompletedTask -from django_huey import DJANGO_HUEY +from background_task.models import Task +from django_huey import DJANGO_HUEY, get_queue from common.huey import h_q_reset_tasks from .models import Source, Media, MediaServer from .forms import (ValidateSourceForm, ConfirmDeleteSourceForm, RedownloadMediaForm, @@ -41,6 +42,24 @@ from . import signals # noqa 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 = { + str(t.id) for q in huey_queues for t in set( + q.pending() + ).union( + q.scheduled() + ) + } + return TaskHistory.objects.filter( + task_id__in=huey_task_ids.union(background_task_ids), + ) + + class DashboardView(TemplateView): ''' The dashboard shows non-interactive totals and summaries. @@ -62,8 +81,13 @@ class DashboardView(TemplateView): data['num_media'] = Media.objects.all().count() data['num_downloaded_media'] = Media.objects.filter(downloaded=True).count() # Tasks - data['num_tasks'] = Task.objects.all().count() - data['num_completed_tasks'] = CompletedTask.objects.all().count() + completed_qs = TaskHistory.objects.filter( + start_at__isnull=False, + end_at__gt=F('start_at'), + ) + waiting_qs = get_waiting_tasks() + data['num_tasks'] = waiting_qs.count() + data['num_completed_tasks'] = completed_qs.count() # Disk usage disk_usage = Media.objects.filter( downloaded=True, downloaded_filesize__isnull=False @@ -830,27 +854,23 @@ class TasksView(ListView): return super().dispatch(request, *args, **kwargs) def get_queryset(self): - qs = Task.objects.all() + qs = get_waiting_tasks() if self.filter_source: params_prefix=f'[["{self.filter_source.pk}"' qs = qs.filter(task_params__istartswith=params_prefix) - order = getattr(settings, - 'BACKGROUND_TASK_PRIORITY_ORDERING', - 'DESC' - ) - prefix = '-' if 'ASC' != order else '' - _priority = f'{prefix}priority' return qs.order_by( - _priority, - 'run_at', + 'priority', + 'end_at', ) def get_context_data(self, *args, **kwargs): data = super().get_context_data(*args, **kwargs) - now = timezone.now() - qs = Task.objects.all() - running_qs = qs.filter(locked_by__isnull=False) - scheduled_qs = qs.filter(locked_by__isnull=True) + now_dt = timezone.now() + scheduled_qs = get_waiting_tasks() + running_qs = scheduled_qs.filter( + start_at__isnull=False, + end_at__lte=F('start_at'), + ) errors_qs = scheduled_qs.filter( attempts__gt=0 ).exclude(last_error__exact='') @@ -867,19 +887,19 @@ class TasksView(ListView): data['wait_for_database_queue'] = False def add_to_task(task): + setattr(task, 'run_now', task.scheduled_at < now_dt) obj, url = map_task_to_instance(task) - if not obj: - return False - setattr(task, 'instance', obj) - setattr(task, 'url', url) - setattr(task, 'run_now', task.run_at < now) + if obj: + setattr(task, 'instance', obj) + setattr(task, 'url', url) if task.has_error(): error_message = get_error_message(task) setattr(task, 'error_message', error_message) return 'error' - return True - - for task in running_qs: + 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. @@ -893,11 +913,28 @@ class TasksView(ListView): # do not wait for the task to expire task.locked_at = None task.save() - if locked_by_pid_running and add_to_task(task): - data['running'].append(task) - elif locked_by_pid_running and 'wait_for_database_queue' in task.task_name: - data['wait_for_database_queue'] = True + 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 + add_to_task(task) + data['running'].append(task) + # show all the errors when they fit on one page if (data['total_errors'] + len(data['running'])) < self.paginate_by: for task in errors_qs: @@ -923,14 +960,10 @@ class TasksView(ListView): elif mapped: data['scheduled'].append(task) - order = getattr(settings, - 'BACKGROUND_TASK_PRIORITY_ORDERING', - 'DESC' - ) sort_keys = ( # key, reverse - ('run_at', False), - ('priority', 'ASC' != order), + ('scheduled_at', False), + ('priority', True), ('run_now', True), ) data['errors'] = multi_key_sort(data['errors'], sort_keys, attr=True) @@ -965,11 +998,14 @@ class CompletedTasksView(ListView): return super().dispatch(request, *args, **kwargs) def get_queryset(self): - qs = CompletedTask.objects.all() + qs = TaskHistory.objects.filter( + start_at__isnull=False, + end_at__gt=F('start_at'), + ) if self.filter_source: params_prefix=f'[["{self.filter_source.pk}"' qs = qs.filter(task_params__istartswith=params_prefix) - return qs.order_by('-run_at') + return qs.order_by('-end_at') def get_context_data(self, *args, **kwargs): data = super().get_context_data(*args, **kwargs) @@ -998,7 +1034,8 @@ class ResetTasks(FormView): def form_valid(self, form): # Delete all tasks Task.objects.all().delete() - for queue_name in (DJANGO_HUEY or {}).get('queues', {}): + huey_queue_names = (DJANGO_HUEY or {}).get('queues', {}) + for queue_name in huey_queue_names: h_q_reset_tasks(queue_name) # Iter all tasks for source in Source.objects.all():