Merge pull request #1170 from tcely/patch-2

Use `TaskHistory` for reporting about tasks
This commit is contained in:
meeb authored and GitHub committed 2025-07-04 17:08:33 +10:00
commit f456c22b30
9 files changed
+208 -65

No files matched your search

+9
View File
@@ -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
@@ -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),
),
]
+5 -2
View File
@@ -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()
+6 -2
View File
@@ -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)
+52 -3
View File
@@ -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
+8 -3
View File
@@ -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:
@@ -16,16 +16,16 @@
<span class="collection-item">
{% if task.has_error %}
<i class="fas fa-exclamation-triangle"></i> <strong>{{ task.verbose_name }}</strong><br>
<i class="fas fa-exclamation-triangle"></i> <strong>{% if task.verbose_name %}{{ task.verbose_name }}{% else %}{{ task.name }}{% endif %}</strong><br>
Queue: &quot;{{ task.queue }}&quot;<br>
Error: &quot;{{ task.error_message }}&quot;<br>
{% else %}
<i class="fas fa-check"></i> <strong>{{ task.verbose_name }}</strong><br>
<i class="fas fa-check"></i> <strong>{% if task.verbose_name %}{{ task.verbose_name }}{% else %}{{ task.name }}{% endif %}</strong><br>
Queue: &quot;{{ task.queue }}&quot;<br>
{% endif %}
Task locked for: {{ task.run_at|sub:task.locked_at|timedelta }}<br>
<i class="far fa-clock"></i> Task locked at <strong>{{ task.locked_at|date:'Y-m-d H:i:s' }}</strong><br>
<i class="fas fa-hourglass-end"></i> Task ended at <strong>{{ task.run_at|date:'Y-m-d H:i:s' }}</strong>
Task locked for: {{ task.end_at|sub:task.start_at|timedelta }}<br>
<i class="far fa-clock"></i> Task locked at <strong>{{ task.start_at|date:'Y-m-d H:i:s' }}</strong><br>
<i class="fas fa-hourglass-end"></i> Task ended at <strong>{{ task.end_at|date:'Y-m-d H:i:s' }}</strong>
</span>
{% empty %}
<span class="collection-item no-items"><i class="fas fa-info-circle"></i> There have been no completed tasks{% if source %} that match the specified source filter{% endif %}.</span>
@@ -33,5 +33,5 @@
</div>
</div>
</div>
{% include 'pagination.html' with pagination=sources.paginator filter=source.pk %}
{% include 'pagination.html' with filter=source.pk %}
{% endblock %}
+11 -10
View File
@@ -23,9 +23,9 @@
</p>
<div class="collection">
{% for task in running %}
<a href="{% url task.url pk=task.instance.pk %}" class="collection-item">
<a href="{%if task.instance.pk %}{% url task.url pk=task.instance.pk %}{% else %}#{{ task.task_id }}{% endif %}" class="collection-item">
<i class="fas fa-running"></i> <strong>{{ task }}</strong><br>
<i class="far fa-clock"></i> Task started at <strong>{{ task.locked_at|date:'Y-m-d H:i:s' }}</strong>
<i class="far fa-clock"></i> Task started at <strong>{{ task.start_at|date:'Y-m-d H:i:s' }}</strong>
</a>
{% empty %}
<span class="collection-item no-items"><i class="fas fa-info-circle"></i> There are no running tasks.</span>
@@ -54,8 +54,8 @@
<i class="fas fa-exclamation-triangle"></i> <strong>{{ task }}</strong>, attempted {{ task.attempts }} time{{ task.attempts|pluralize }}<br>
Error: &quot;{{ task.error_message }}&quot;<br>
</a>
<i class="fas fa-history"></i> Task will be retried at <strong>{{ task.run_at|date:'Y-m-d H:i:s' }}</strong>
<a href="{% url 'sync:run-task' pk=task.pk %}" class="error-text">
<i class="fas fa-history"></i> Task will be retried at <strong>{{ task.scheduled_at|date:'Y-m-d H:i:s' }}</strong>
<a href="{% url 'sync:run-task' pk=task.task_id %}" class="error-text">
<i class="fas fa-undo"></i>
</a>
</div>
@@ -73,19 +73,20 @@
<p>
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 &quot;run at&quot; time.
run in the future at the specified &quot;scheduled at&quot; time.
</p>
<div class="collection">
{% for task in scheduled %}
<div class="collection-item">
<a href="{% url task.url pk=task.instance.pk %}">
<i class="fas fa-hourglass-start"></i> <strong>{{ task }}</strong><br>
{% if task.instance.index_schedule and task.repeat > 0 %}Scheduled to run {{ task.instance.get_index_schedule_display|lower }}.<br>{% endif %}
<i class="far fa-clock"></i> Task will run {% if task.run_now %}<strong>immediately</strong>{% else %}at <strong>{{ task.run_at|date:'Y-m-d H:i:s' }}</strong>
{% if task.instance.is_active and 'once' not in task.verbose_name %}Scheduled to run {{ task.instance.get_index_schedule_display|lower }}.<br>{% endif %}
<i class="far fa-clock"></i> Task will run {% if task.run_now %}<strong>immediately</strong>{% else %}at <strong>{{ task.scheduled_at|date:'Y-m-d H:i:s' }}</strong>
{% if '-' not in task.task_id %}
</a>
<a href="{% url 'sync:run-task' pk=task.pk %}">
<a href="{% url 'sync:run-task' pk=task.task_id %}">
<i class="far fa-play-circle"></i>
{% endif %}
{% endif %}{% endif %}
</a>
</div>
{% empty %}
@@ -94,7 +95,7 @@
</div>
</div>
</div>
{% include 'pagination.html' with pagination=sources.paginator filter=source.pk %}
{% include 'pagination.html' with filter=source.pk %}
<div class="row">
<div class="col s12">
<h2>Completed</h2>
+76 -39
View File
@@ -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():