From b162ba55a6990c83347e7074541e0abb0beecf01 Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 03:00:26 -0400 Subject: [PATCH 01/30] Update the dashboard to use `TaskHistory` --- tubesync/sync/views.py | 28 ++++++++++++++++++++++++---- 1 file changed, 24 insertions(+), 4 deletions(-) diff --git a/tubesync/sync/views.py b/tubesync/sync/views.py index 7afb2804..e1b14cca 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 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, @@ -62,8 +63,27 @@ 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'), + ) + background_task_ids = { + str(t.pk) for t in Task.objects.all() + } + huey_queues = list(map(get_queue, DJANGO_HUEY.get('queues', dict()))) + huey_task_ids = { + str(t.id) for q in huey_queues for t in set( + q.pending() + ).union( + q.scheduled() + ) + } + waiting_qs = TaskHistory.objects.filter( + start_at__isnull=True, + task_id__in=huey_task_ids.union(background_task_ids), + ) + 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 From c8e7aede03e3f3dcc5e9680ae679ed53670364e8 Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 03:17:40 -0400 Subject: [PATCH 02/30] Add and use the `get_waiting_tasks` function --- tubesync/sync/views.py | 44 ++++++++++++++++++++---------------------- 1 file changed, 21 insertions(+), 23 deletions(-) diff --git a/tubesync/sync/views.py b/tubesync/sync/views.py index e1b14cca..4f2f6e02 100644 --- a/tubesync/sync/views.py +++ b/tubesync/sync/views.py @@ -42,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_queues = list(map(get_queue, DJANGO_HUEY.get('queues', dict()))) + huey_task_ids = { + str(t.id) for q in huey_queues for t in set( + q.pending() + ).union( + q.scheduled() + ) + } + return TaskHistory.objects.filter( + start_at__isnull=True, + task_id__in=huey_task_ids.union(background_task_ids), + ) + + class DashboardView(TemplateView): ''' The dashboard shows non-interactive totals and summaries. @@ -67,21 +85,7 @@ class DashboardView(TemplateView): start_at__isnull=False, end_at__gt=F('start_at'), ) - background_task_ids = { - str(t.pk) for t in Task.objects.all() - } - huey_queues = list(map(get_queue, DJANGO_HUEY.get('queues', dict()))) - huey_task_ids = { - str(t.id) for q in huey_queues for t in set( - q.pending() - ).union( - q.scheduled() - ) - } - waiting_qs = TaskHistory.objects.filter( - start_at__isnull=True, - task_id__in=huey_task_ids.union(background_task_ids), - ) + waiting_qs = get_waiting_tasks() data['num_tasks'] = waiting_qs.count() data['num_completed_tasks'] = completed_qs.count() # Disk usage @@ -850,18 +854,12 @@ 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, + 'priority', 'run_at', ) From 8dc212496dc0ea8ac2eed116673c32bcd0c23b99 Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 03:27:16 -0400 Subject: [PATCH 03/30] fixup: use `end_at` with `TaskHistory` --- tubesync/sync/views.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tubesync/sync/views.py b/tubesync/sync/views.py index 4f2f6e02..9f214202 100644 --- a/tubesync/sync/views.py +++ b/tubesync/sync/views.py @@ -860,7 +860,7 @@ class TasksView(ListView): qs = qs.filter(task_params__istartswith=params_prefix) return qs.order_by( 'priority', - 'run_at', + 'end_at', ) def get_context_data(self, *args, **kwargs): From 91599cb53f3cb70d87bc9989991eb504c50da8f9 Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 03:51:15 -0400 Subject: [PATCH 04/30] Add a new `TaskHistory` row when a task is created --- tubesync/sync/signals.py | 29 +++++++++++++++++++++++++++-- 1 file changed, 27 insertions(+), 2 deletions(-) diff --git a/tubesync/sync/signals.py b/tubesync/sync/signals.py index f89776a8..c584c820 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_failed, +) from background_task.models import Task from common.logger import log from common.models import TaskHistory @@ -168,6 +170,29 @@ 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.end_at = task_obj.run_at + th.priority = (100 - task_obj.priority) + th.repeat = task_obj.repeat + th.repeat_until = task_obj.repeat_until + th.start_at = task_obj.run_at + 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,7 +213,7 @@ 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}') def merge_completed_task_into_history(task_id, task_obj): From ef0a3d9e822fff95f3ef1bb4774bf12deaa37ea2 Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 03:54:42 -0400 Subject: [PATCH 05/30] fixup: do not set `start_at` before it was started --- tubesync/sync/signals.py | 1 - 1 file changed, 1 deletion(-) diff --git a/tubesync/sync/signals.py b/tubesync/sync/signals.py index c584c820..0c420880 100644 --- a/tubesync/sync/signals.py +++ b/tubesync/sync/signals.py @@ -185,7 +185,6 @@ def task_task_created(sender, task=None, **kwargs): th.priority = (100 - task_obj.priority) th.repeat = task_obj.repeat th.repeat_until = task_obj.repeat_until - th.start_at = task_obj.run_at th.task_params = list(task_obj.params()) th.verbose_name = task_obj.verbose_name th.save() From eebdfca8244f623b92543c3c509d5343f61ba2b6 Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 04:41:24 -0400 Subject: [PATCH 06/30] Display scheduled tasks using `TaskHistory` --- tubesync/sync/views.py | 38 ++++++++++++++++++++++++-------------- 1 file changed, 24 insertions(+), 14 deletions(-) diff --git a/tubesync/sync/views.py b/tubesync/sync/views.py index 9f214202..dc4767bd 100644 --- a/tubesync/sync/views.py +++ b/tubesync/sync/views.py @@ -868,7 +868,7 @@ class TasksView(ListView): now = timezone.now() qs = Task.objects.all() running_qs = qs.filter(locked_by__isnull=False) - scheduled_qs = qs.filter(locked_by__isnull=True) + scheduled_qs = get_waiting_tasks() errors_qs = scheduled_qs.filter( attempts__gt=0 ).exclude(last_error__exact='') @@ -916,39 +916,49 @@ class TasksView(ListView): elif locked_by_pid_running and 'wait_for_database_queue' in task.task_name: data['wait_for_database_queue'] = True + running_task_ids = { str(t.pk) for t in data['running'] } # 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: - if task in data['running']: + if task.task_id in running_task_ids: continue - mapped = add_to_task(task) - if 'error' == mapped: + obj, url = map_task_to_instance(task) + if obj: + setattr(task, 'instance', obj) + setattr(task, 'url', url) + setattr(task, 'run_now', task.end_at < now) + if task.has_error(): + error_message = get_error_message(task) + setattr(task, 'error_message', error_message) data['errors'].append(task) - elif mapped: + elif obj: data['scheduled'].append(task) for task in data['tasks']: already_added = ( + task.task_id in running_task_ids or task in data['running'] or task in data['errors'] or task in data['scheduled'] ) if already_added: continue - mapped = add_to_task(task) - if 'error' == mapped: + obj, url = map_task_to_instance(task) + if obj: + setattr(task, 'instance', obj) + setattr(task, 'url', url) + setattr(task, 'run_now', task.end_at < now) + if task.has_error(): + error_message = get_error_message(task) + setattr(task, 'error_message', error_message) data['errors'].append(task) - elif mapped: + elif obj: data['scheduled'].append(task) - order = getattr(settings, - 'BACKGROUND_TASK_PRIORITY_ORDERING', - 'DESC' - ) sort_keys = ( # key, reverse - ('run_at', False), - ('priority', 'ASC' != order), + ('end_at', False), + ('priority', True), ('run_now', True), ) data['errors'] = multi_key_sort(data['errors'], sort_keys, attr=True) From 7b897fb7a8a007f087a73e873d7d361310a1bed9 Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 05:29:38 -0400 Subject: [PATCH 07/30] Adjust `map_task_to_instance` for `TaskHistory` --- tubesync/sync/tasks.py | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) 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: From f09a34f1e12f71c56a791533a41391bd840af2df Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 05:32:39 -0400 Subject: [PATCH 08/30] `add_to_task` still uses the old class --- tubesync/sync/views.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tubesync/sync/views.py b/tubesync/sync/views.py index dc4767bd..c78717ee 100644 --- a/tubesync/sync/views.py +++ b/tubesync/sync/views.py @@ -885,7 +885,7 @@ class TasksView(ListView): data['wait_for_database_queue'] = False def add_to_task(task): - obj, url = map_task_to_instance(task) + obj, url = map_task_to_instance(task, using_history=False) if not obj: return False setattr(task, 'instance', obj) From 8db622f5b8645c249b22fa7233db02162bf7edc1 Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 05:37:36 -0400 Subject: [PATCH 09/30] `task_task_failed` still uses old class --- tubesync/sync/signals.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tubesync/sync/signals.py b/tubesync/sync/signals.py index 0c420880..12648f1f 100644 --- a/tubesync/sync/signals.py +++ b/tubesync/sync/signals.py @@ -242,7 +242,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 From 407cace19ba2260e7d393046da69c1377e4882b6 Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 06:19:46 -0400 Subject: [PATCH 10/30] Include errors which have set `start_at` already --- tubesync/sync/views.py | 1 - 1 file changed, 1 deletion(-) diff --git a/tubesync/sync/views.py b/tubesync/sync/views.py index c78717ee..aac744a7 100644 --- a/tubesync/sync/views.py +++ b/tubesync/sync/views.py @@ -55,7 +55,6 @@ def get_waiting_tasks(): ) } return TaskHistory.objects.filter( - start_at__isnull=True, task_id__in=huey_task_ids.union(background_task_ids), ) From 7b196fa34c6e727fb1085431bd4d07c1c6110090 Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 07:56:47 -0400 Subject: [PATCH 11/30] Set `end_at` to the scheduled `eta` --- tubesync/common/huey.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/tubesync/common/huey.py b/tubesync/common/huey.py index 8a8a2ea6..9e46ebb4 100644 --- a/tubesync/common/huey.py +++ b/tubesync/common/huey.py @@ -278,6 +278,8 @@ def historical_task(signal_name, task_obj, exception_obj=None, /, *, huey=None): if hasattr(model_instance, 'name'): th.verbose_name += f' / {model_instance.name}' th.end_at = signal_dt + if signal_name == signals.SIGNAL_SCHEDULED: + th.end_at = task_obj.eta.replace(tzinfo=datetime.timezone.utc) th.save() # Registration of shared signal handlers From 064b0892667841f6b5322768cb5f49768a6fa3b5 Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 12:57:25 -0400 Subject: [PATCH 12/30] Update tasks.html --- tubesync/sync/templates/sync/tasks.html | 15 ++++++++------- 1 file changed, 8 insertions(+), 7 deletions(-) diff --git a/tubesync/sync/templates/sync/tasks.html b/tubesync/sync/templates/sync/tasks.html index a074df85..7569fb6a 100644 --- a/tubesync/sync/templates/sync/tasks.html +++ b/tubesync/sync/templates/sync/tasks.html @@ -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.end_at|date:'Y-m-d H:i:s' }} + @@ -80,12 +80,13 @@ {% empty %} @@ -94,7 +95,7 @@ -{% include 'pagination.html' with pagination=sources.paginator filter=source.pk %} +{% include 'pagination.html' with filter=source.pk %}

Completed

From ca09e7a6bed70f2ed1e43e56e9486eb75bdd90e6 Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 14:46:20 -0400 Subject: [PATCH 13/30] Handle naive `datetime.datetime` also --- tubesync/common/timestamp.py | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) 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) From 6d6a4d4ba5daa809038536a06f8b06f752848493 Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 14:56:18 -0400 Subject: [PATCH 14/30] Convert Huey `eta` to UTC always --- tubesync/common/huey.py | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/tubesync/common/huey.py b/tubesync/common/huey.py index 9e46ebb4..014104f2 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): @@ -279,7 +280,12 @@ def historical_task(signal_name, task_obj, exception_obj=None, /, *, huey=None): th.verbose_name += f' / {model_instance.name}' th.end_at = signal_dt if signal_name == signals.SIGNAL_SCHEDULED: - th.end_at = task_obj.eta.replace(tzinfo=datetime.timezone.utc) + if huey.utc: + th.end_at = task_obj.eta.replace(tzinfo=datetime.UTC) + else: # this path is unlikely + th.end_at = timestamp_to_datetime( + datetime_to_timestamp(task_obj.eta, integer=False), + ).astimezone(tz=datetime.UTC) th.save() # Registration of shared signal handlers From b5893b4604df94389b759d686fc0eb2e20b7c210 Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 17:44:29 -0400 Subject: [PATCH 15/30] Add `elapsed` and `scheduled_at` to `TaskHistory` --- tubesync/common/models/tasks.py | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) 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() From d195f43942e4e0f1abe11eb86fc20fb1db03ddd4 Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 18:00:15 -0400 Subject: [PATCH 16/30] Create 0002_taskhistory_elapsed_taskhistory_scheduled_at_and_more.py --- ...apsed_taskhistory_scheduled_at_and_more.py | 35 +++++++++++++++++++ 1 file changed, 35 insertions(+) create mode 100644 tubesync/common/migrations/0002_taskhistory_elapsed_taskhistory_scheduled_at_and_more.py 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), + ), + ] + From df3616dd03c181065145cde50cd22bc38e6c2b5e Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 18:13:34 -0400 Subject: [PATCH 17/30] Populate the new `TaskHistory` fields from huey --- tubesync/common/huey.py | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/tubesync/common/huey.py b/tubesync/common/huey.py index 014104f2..6f34b5e9 100644 --- a/tubesync/common/huey.py +++ b/tubesync/common/huey.py @@ -278,14 +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}' - th.end_at = signal_dt - if signal_name == signals.SIGNAL_SCHEDULED: + elif signal_name == signals.SIGNAL_SCHEDULED: if huey.utc: - th.end_at = task_obj.eta.replace(tzinfo=datetime.UTC) + th.scheduled_at = task_obj.eta.replace(tzinfo=datetime.UTC) else: # this path is unlikely - th.end_at = timestamp_to_datetime( + 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 From c8dfcb19834fada3f53f680c52d26ef3fc77b80d Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 19:10:32 -0400 Subject: [PATCH 18/30] Populate the new `TaskHistory` fields from background_task --- tubesync/sync/signals.py | 25 +++++++++++++++++++++++-- 1 file changed, 23 insertions(+), 2 deletions(-) diff --git a/tubesync/sync/signals.py b/tubesync/sync/signals.py index 12648f1f..5a32fe74 100644 --- a/tubesync/sync/signals.py +++ b/tubesync/sync/signals.py @@ -9,7 +9,7 @@ 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_failed, + task_created, task_started, task_successful, task_rescheduled, task_failed, ) from background_task.models import Task from common.logger import log @@ -181,7 +181,7 @@ def task_task_created(sender, task=None, **kwargs): name=task_obj.task_name, queue=task_obj.queue, ) - th.end_at = task_obj.run_at + 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 @@ -215,12 +215,33 @@ def task_task_started(sender, **kwargs): 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) + 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) th.end_at = task_obj.run_at th.failed_at = task_obj.failed_at th.last_error = task_obj.last_error From cffa34774f8407459c3b8190a40eff5033f6d9bf Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 20:02:04 -0400 Subject: [PATCH 19/30] Use the new `scheduled_at` on the tasks page --- tubesync/sync/templates/sync/tasks.html | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/tubesync/sync/templates/sync/tasks.html b/tubesync/sync/templates/sync/tasks.html index 7569fb6a..63fd3561 100644 --- a/tubesync/sync/templates/sync/tasks.html +++ b/tubesync/sync/templates/sync/tasks.html @@ -54,7 +54,7 @@ {{ task }}, attempted {{ task.attempts }} time{{ task.attempts|pluralize }}
Error: "{{ task.error_message }}"
- Task will be retried at {{ task.end_at|date:'Y-m-d H:i:s' }} + Task will be retried at {{ task.scheduled_at|date:'Y-m-d H:i:s' }} @@ -73,7 +73,7 @@

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 %} @@ -81,7 +81,7 @@ {{ task }}
{% 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.end_at|date:'Y-m-d H:i:s' }} + 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 %}
From 898f1d42cf04ea238d8986e10759826a649afba5 Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 20:51:20 -0400 Subject: [PATCH 20/30] Display running and completed tasks from `TaskHistory` --- tubesync/sync/views.py | 57 ++++++++++++++++++++++++++++-------------- 1 file changed, 38 insertions(+), 19 deletions(-) diff --git a/tubesync/sync/views.py b/tubesync/sync/views.py index aac744a7..110287ff 100644 --- a/tubesync/sync/views.py +++ b/tubesync/sync/views.py @@ -46,7 +46,8 @@ def get_waiting_tasks(): background_task_ids = { str(t.pk) for t in Task.objects.all() } - huey_queues = list(map(get_queue, DJANGO_HUEY.get('queues', dict()))) + 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() @@ -864,9 +865,11 @@ class TasksView(ListView): 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) + now_dt = timezone.now() + running_qs = TaskHistory.objects.filter( + start_at__isnull=False, + end_at__lte=F('start_at'), + ) scheduled_qs = get_waiting_tasks() errors_qs = scheduled_qs.filter( attempts__gt=0 @@ -889,14 +892,14 @@ class TasksView(ListView): return False setattr(task, 'instance', obj) setattr(task, 'url', url) - setattr(task, 'run_now', task.run_at < now) + setattr(task, 'run_now', task.run_at < now_dt) 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: + 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. @@ -910,22 +913,38 @@ 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) + 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): + data['running'].append(task) + elif locked_by_pid_running and 'wait_for_database_queue' in task.name: + data['wait_for_database_queue'] = True - running_task_ids = { str(t.pk) for t in data['running'] } + for task in running_qs: + if task in data['running']: + continue + setattr(task, 'run_now', task.end_at < now_dt) + obj, url = map_task_to_instance(task) + if obj: + setattr(task, 'instance', obj) + setattr(task, 'url', url) + 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: - if task.task_id in running_task_ids: + if task in data['running']: continue obj, url = map_task_to_instance(task) if obj: setattr(task, 'instance', obj) setattr(task, 'url', url) - setattr(task, 'run_now', task.end_at < now) + setattr(task, 'run_now', task.end_at < now_dt) if task.has_error(): error_message = get_error_message(task) setattr(task, 'error_message', error_message) @@ -935,7 +954,6 @@ class TasksView(ListView): for task in data['tasks']: already_added = ( - task.task_id in running_task_ids or task in data['running'] or task in data['errors'] or task in data['scheduled'] @@ -946,7 +964,7 @@ class TasksView(ListView): if obj: setattr(task, 'instance', obj) setattr(task, 'url', url) - setattr(task, 'run_now', task.end_at < now) + setattr(task, 'run_now', task.end_at < now_dt) if task.has_error(): error_message = get_error_message(task) setattr(task, 'error_message', error_message) @@ -956,7 +974,7 @@ class TasksView(ListView): sort_keys = ( # key, reverse - ('end_at', False), + ('scheduled_at', False), ('priority', True), ('run_now', True), ) @@ -992,11 +1010,11 @@ class CompletedTasksView(ListView): return super().dispatch(request, *args, **kwargs) def get_queryset(self): - qs = CompletedTask.objects.all() + qs = TaskHistory.objects.all() 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) @@ -1025,7 +1043,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(): From 815dcb99e8c098628276aa3b45e09451c20bc209 Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 20:58:16 -0400 Subject: [PATCH 21/30] fixup: remove unused import --- tubesync/sync/views.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tubesync/sync/views.py b/tubesync/sync/views.py index 110287ff..e62edeb6 100644 --- a/tubesync/sync/views.py +++ b/tubesync/sync/views.py @@ -22,7 +22,7 @@ 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 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 6c860387a21cfe5d1235e0cacb07139edebcd5b5 Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 21:03:39 -0400 Subject: [PATCH 22/30] Update tasks-completed.html --- tubesync/sync/templates/sync/tasks-completed.html | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/tubesync/sync/templates/sync/tasks-completed.html b/tubesync/sync/templates/sync/tasks-completed.html index ec2a0aa8..8dd5f2e7 100644 --- a/tubesync/sync/templates/sync/tasks-completed.html +++ b/tubesync/sync/templates/sync/tasks-completed.html @@ -23,9 +23,9 @@ {{ task.verbose_name }}
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 %} From ed7b53b6a6a6d3a8b32dfe2b5d1dd5a078e8e35b Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 21:59:09 -0400 Subject: [PATCH 23/30] Update tasks.html --- tubesync/sync/templates/sync/tasks.html | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/tubesync/sync/templates/sync/tasks.html b/tubesync/sync/templates/sync/tasks.html index 63fd3561..0d622925 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. From b814614eb3fb2a92bb222857724ab0a8a1f1c5b9 Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 22:18:36 -0400 Subject: [PATCH 24/30] Display the `task_id` for now --- tubesync/sync/templates/sync/tasks.html | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tubesync/sync/templates/sync/tasks.html b/tubesync/sync/templates/sync/tasks.html index 0d622925..74997164 100644 --- a/tubesync/sync/templates/sync/tasks.html +++ b/tubesync/sync/templates/sync/tasks.html @@ -23,7 +23,7 @@

{% for task in running %} - + {{ task }}
Task started at {{ task.start_at|date:'Y-m-d H:i:s' }}
From 1a61e3882ac0927406e50cd4190c2ef0a76d3ee8 Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 22:47:47 -0400 Subject: [PATCH 25/30] Fallback to `name` when `verbose_name` is unavailable --- tubesync/sync/templates/sync/tasks-completed.html | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/tubesync/sync/templates/sync/tasks-completed.html b/tubesync/sync/templates/sync/tasks-completed.html index 8dd5f2e7..796d1ada 100644 --- a/tubesync/sync/templates/sync/tasks-completed.html +++ b/tubesync/sync/templates/sync/tasks-completed.html @@ -16,11 +16,11 @@ {% 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.end_at|sub:task.start_at|timedelta }}
From 75ef10ed1f31a8eacba5a0ef7f24a5e35d4dbb5a Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 23:00:50 -0400 Subject: [PATCH 26/30] Use the same filter as the dashboard --- tubesync/sync/views.py | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/tubesync/sync/views.py b/tubesync/sync/views.py index e62edeb6..3037001d 100644 --- a/tubesync/sync/views.py +++ b/tubesync/sync/views.py @@ -1010,7 +1010,10 @@ class CompletedTasksView(ListView): return super().dispatch(request, *args, **kwargs) def get_queryset(self): - qs = TaskHistory.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) From 5106f276c5d73b1d669cea4848dd0a5362d66301 Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 23:36:04 -0400 Subject: [PATCH 27/30] Use `add_to_task` again now that everything uses `TaskHistory` --- tubesync/sync/views.py | 44 ++++++++++++++---------------------------- 1 file changed, 14 insertions(+), 30 deletions(-) diff --git a/tubesync/sync/views.py b/tubesync/sync/views.py index 3037001d..7cbbd3a4 100644 --- a/tubesync/sync/views.py +++ b/tubesync/sync/views.py @@ -868,6 +868,7 @@ class TasksView(ListView): now_dt = timezone.now() running_qs = TaskHistory.objects.filter( start_at__isnull=False, + end_at__gt=now_dt-timezone.timedelta(days=1), end_at__lte=F('start_at'), ) scheduled_qs = get_waiting_tasks() @@ -887,17 +888,16 @@ class TasksView(ListView): data['wait_for_database_queue'] = False def add_to_task(task): - obj, url = map_task_to_instance(task, using_history=False) - if not obj: - return False - setattr(task, 'instance', obj) - setattr(task, 'url', url) - setattr(task, 'run_now', task.run_at < now_dt) + setattr(task, 'run_now', task.scheduled_at < now_dt) + obj, url = map_task_to_instance(task) + 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 + return True and obj for task in Task.objects.filter(locked_by__isnull=False): # There was broken logic in `Task.objects.locked()`, work around it. @@ -928,11 +928,7 @@ class TasksView(ListView): for task in running_qs: if task in data['running']: continue - setattr(task, 'run_now', task.end_at < now_dt) - obj, url = map_task_to_instance(task) - if obj: - setattr(task, 'instance', obj) - setattr(task, 'url', url) + add_to_task(task) data['running'].append(task) # show all the errors when they fit on one page @@ -940,16 +936,10 @@ class TasksView(ListView): for task in errors_qs: if task in data['running']: continue - obj, url = map_task_to_instance(task) - if obj: - setattr(task, 'instance', obj) - setattr(task, 'url', url) - setattr(task, 'run_now', task.end_at < now_dt) - if task.has_error(): - error_message = get_error_message(task) - setattr(task, 'error_message', error_message) + mapped = add_to_task(task) + if 'error' == mapped: data['errors'].append(task) - elif obj: + elif mapped: data['scheduled'].append(task) for task in data['tasks']: @@ -960,16 +950,10 @@ class TasksView(ListView): ) if already_added: continue - obj, url = map_task_to_instance(task) - if obj: - setattr(task, 'instance', obj) - setattr(task, 'url', url) - setattr(task, 'run_now', task.end_at < now_dt) - if task.has_error(): - error_message = get_error_message(task) - setattr(task, 'error_message', error_message) + mapped = add_to_task(task) + if 'error' == mapped: data['errors'].append(task) - elif obj: + elif mapped: data['scheduled'].append(task) sort_keys = ( From 0a205231c9d2dbc8e71680e624fadbd64a758690 Mon Sep 17 00:00:00 2001 From: tcely Date: Thu, 3 Jul 2025 23:53:54 -0400 Subject: [PATCH 28/30] Add the seconds to the `float`, not a `datetime.timedelta` --- tubesync/sync/signals.py | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/tubesync/sync/signals.py b/tubesync/sync/signals.py index 5a32fe74..f42c2c7c 100644 --- a/tubesync/sync/signals.py +++ b/tubesync/sync/signals.py @@ -227,7 +227,9 @@ def task_task_rescheduled(sender, task=None, **kwargs): name=task_obj.task_name, queue=task_obj.queue, ) - th.elapsed += (now_dt - task_obj.locked_at) + 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 @@ -241,7 +243,9 @@ def merge_completed_task_into_history(task_id, task_obj): name=task_obj.task_name, queue=task_obj.queue, ) - th.elapsed += ((task_obj.failed_at or task_obj.run_at) - task_obj.locked_at) + 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 From 16008cf7b1e8dc66b904c9a533385d51378074aa Mon Sep 17 00:00:00 2001 From: tcely Date: Fri, 4 Jul 2025 00:24:15 -0400 Subject: [PATCH 29/30] `scheduled_qs` now also includes running and errors --- tubesync/sync/views.py | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/tubesync/sync/views.py b/tubesync/sync/views.py index 7cbbd3a4..c842facd 100644 --- a/tubesync/sync/views.py +++ b/tubesync/sync/views.py @@ -866,12 +866,11 @@ class TasksView(ListView): def get_context_data(self, *args, **kwargs): data = super().get_context_data(*args, **kwargs) now_dt = timezone.now() - running_qs = TaskHistory.objects.filter( + scheduled_qs = get_waiting_tasks() + running_qs = scheduled_qs.filter( start_at__isnull=False, - end_at__gt=now_dt-timezone.timedelta(days=1), end_at__lte=F('start_at'), ) - scheduled_qs = get_waiting_tasks() errors_qs = scheduled_qs.filter( attempts__gt=0 ).exclude(last_error__exact='') From 0a36c3a791500a27102c8207091f8f5ccb21c7ba Mon Sep 17 00:00:00 2001 From: tcely Date: Fri, 4 Jul 2025 01:21:09 -0400 Subject: [PATCH 30/30] Check for updates to `verbose_name` for locked tasks --- tubesync/sync/views.py | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/tubesync/sync/views.py b/tubesync/sync/views.py index c842facd..eb31a966 100644 --- a/tubesync/sync/views.py +++ b/tubesync/sync/views.py @@ -897,7 +897,8 @@ class TasksView(ListView): setattr(task, 'error_message', error_message) 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. @@ -913,6 +914,7 @@ class TasksView(ListView): 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: @@ -920,9 +922,12 @@ class TasksView(ListView): 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']: