Merge pull request #1199 from tcely/tcely-download_media
Remove old `download_media` task
This commit is contained in:
5 files changed
+46
-57
No files matched your search
@@ -996,7 +996,12 @@ class Media(models.Model):
|
||||
if self.downloaded:
|
||||
return Val(MediaState.DOWNLOADED)
|
||||
if task:
|
||||
if task.locked_by_pid_running():
|
||||
def running(arg_task, /):
|
||||
if hasattr(arg_task, 'locked_by_pid_running'):
|
||||
return arg_task.locked_by_pid_running()
|
||||
from ..tasks import get_media_download_task
|
||||
return get_media_download_task(str(self.pk))
|
||||
if running(task):
|
||||
return Val(MediaState.DOWNLOADING)
|
||||
elif task.has_error():
|
||||
return Val(MediaState.ERROR)
|
||||
|
||||
@@ -22,7 +22,7 @@ from .tasks import (
|
||||
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_task,
|
||||
download_media, download_media_metadata, download_media_image,
|
||||
download_media_file, download_media_metadata, download_media_image,
|
||||
)
|
||||
from .utils import delete_file
|
||||
from .filtering import filter_media
|
||||
@@ -369,10 +369,12 @@ def media_post_save(sender, instance, created, **kwargs):
|
||||
downloaded = False
|
||||
if (instance.source.download_media and instance.can_download) and not (
|
||||
instance.skip or downloaded or existing_media_download_task):
|
||||
verbose_name = _('Downloading media for "{}"')
|
||||
download_media(
|
||||
str(instance.pk),
|
||||
verbose_name=verbose_name.format(instance.name),
|
||||
TaskHistory.schedule(
|
||||
download_media_file,
|
||||
str(media.pk),
|
||||
remove_duplicates=True,
|
||||
vn_fmt=_('Downloading media for "{}"'),
|
||||
vn_args=(media.name,),
|
||||
)
|
||||
# Save the instance if any changes were required
|
||||
if skip_changed or can_download_changed:
|
||||
@@ -385,7 +387,6 @@ def media_post_save(sender, instance, created, **kwargs):
|
||||
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', (str(instance.pk),))
|
||||
delete_task_by_media('sync.tasks.download_media_metadata', (str(instance.pk),))
|
||||
# Remove thumbnail file for deleted media
|
||||
if instance.thumb:
|
||||
|
||||
+14
-22
@@ -63,7 +63,7 @@ def map_task_to_instance(task, using_history=True):
|
||||
TASK_MAP = {
|
||||
'sync.tasks.index_source_task': Source,
|
||||
'sync.tasks.download_media_image': Media,
|
||||
'sync.tasks.download_media': Media,
|
||||
'sync.tasks.download_media_file': Media,
|
||||
'sync.tasks.download_media_metadata': Media,
|
||||
'sync.tasks.save_all_media_for_source': Source,
|
||||
'sync.tasks.rename_all_media_for_source': Source,
|
||||
@@ -160,13 +160,19 @@ def get_running_tasks(arg_dt=None, /):
|
||||
).order_by('end_at')
|
||||
return running_qs
|
||||
|
||||
def get_media_thumbnail_task(media_id):
|
||||
tqs = get_running_tasks().filter(
|
||||
name='sync.tasks.download_media_image',
|
||||
task_params__0__0=media_id,
|
||||
)
|
||||
def get_running_task_by_name(arg_str, media_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
|
||||
|
||||
def get_media_download_task(media_id):
|
||||
return get_running_task_by_name('download_media_file', media_id)
|
||||
|
||||
def get_media_thumbnail_task(media_id):
|
||||
return get_running_task_by_name('download_media_image', media_id)
|
||||
|
||||
|
||||
def get_tasks(task_name, id=None, /, instance=None):
|
||||
assert not (id is None and instance is None)
|
||||
@@ -177,9 +183,6 @@ def get_first_task(task_name, id=None, /, *, instance=None):
|
||||
tqs = get_tasks(task_name, id, instance).order_by('run_at')
|
||||
return tqs[0] if tqs.count() else False
|
||||
|
||||
def get_media_download_task(media_id):
|
||||
return get_first_task('sync.tasks.download_media', media_id)
|
||||
|
||||
def get_media_metadata_task(media_id):
|
||||
return get_first_task('sync.tasks.download_media_metadata', media_id)
|
||||
|
||||
@@ -979,12 +982,11 @@ def download_media_file(media_id, override=False):
|
||||
if not media.download_checklist(override):
|
||||
# any condition that needs to reschedule the task
|
||||
# should raise an exception to avoid this
|
||||
return False
|
||||
return
|
||||
|
||||
wait_for_errors(
|
||||
media,
|
||||
queue_name=Val(TaskQueue.LIMIT),
|
||||
task_name='sync.tasks.download_media',
|
||||
)
|
||||
with huey_lock_task(
|
||||
f'media:{media.uuid}',
|
||||
@@ -1038,7 +1040,6 @@ def download_media_file(media_id, override=False):
|
||||
media.write_nfo_file()
|
||||
# Schedule a task to update media servers
|
||||
schedule_media_servers_update()
|
||||
return True
|
||||
|
||||
|
||||
@db_task(delay=30, expires=210, priority=100, queue=Val(TaskQueue.NET))
|
||||
@@ -1113,7 +1114,7 @@ def rename_all_media_for_source(source_id):
|
||||
getattr(settings, 'RENAME_ALL_SOURCES', False)
|
||||
)
|
||||
if not create_rename_tasks:
|
||||
return None
|
||||
return
|
||||
mqs = Media.objects.filter(
|
||||
source=source,
|
||||
downloaded=True,
|
||||
@@ -1310,12 +1311,3 @@ def download_media_metadata(media_id):
|
||||
raise InvalidTaskError(str(e)) from e
|
||||
|
||||
|
||||
@background(schedule=dict(priority=30, run_at=60), queue=Val(TaskQueue.NET), remove_existing_tasks=True)
|
||||
def download_media(media_id, override=False):
|
||||
try:
|
||||
res = download_media_file(media_id, override)
|
||||
return res.get(blocking=True)
|
||||
except CancelExecution as e:
|
||||
raise InvalidTaskError(str(e)) from e
|
||||
|
||||
|
||||
+13
-18
@@ -18,7 +18,7 @@ from common.models import TaskHistory
|
||||
from .models import Source, Media
|
||||
from .tasks import (
|
||||
cleanup_old_media, check_source_directory_exists,
|
||||
get_media_thumbnail_task,
|
||||
get_media_download_task, get_media_thumbnail_task,
|
||||
)
|
||||
from .filtering import filter_media
|
||||
from .utils import filter_response
|
||||
@@ -421,26 +421,18 @@ class FrontEndTestCase(TestCase):
|
||||
now_dt = timezone.now()
|
||||
TaskHistory.objects.all().update(start_at=now_dt, end_at=now_dt)
|
||||
# Check the tasks to fetch the media thumbnails have been scheduled
|
||||
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)
|
||||
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)
|
||||
found_download_task1 = False
|
||||
found_download_task2 = False
|
||||
found_download_task3 = False
|
||||
q = {'task_name': 'sync.tasks.download_media'}
|
||||
for task in Task.objects.filter(**q):
|
||||
if test_media1_pk in task.task_params:
|
||||
found_download_task1 = True
|
||||
if test_media2_pk in task.task_params:
|
||||
found_download_task2 = True
|
||||
if test_media3_pk in task.task_params:
|
||||
found_download_task3 = True
|
||||
self.assertTrue(found_thumbnail_task1)
|
||||
self.assertTrue(found_thumbnail_task2)
|
||||
self.assertTrue(found_thumbnail_task3)
|
||||
self.assertTrue(found_download_task1)
|
||||
self.assertTrue(found_download_task2)
|
||||
self.assertTrue(found_download_task3)
|
||||
self.assertTrue(found_thumbnail_task1)
|
||||
self.assertTrue(found_thumbnail_task2)
|
||||
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)
|
||||
@@ -469,15 +461,18 @@ class FrontEndTestCase(TestCase):
|
||||
# simulate the tasks consumer signals having already run
|
||||
TaskHistory.objects.all().update(end_at=timezone.now())
|
||||
# Confirm any tasks have been deleted
|
||||
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)
|
||||
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.assertFalse(found_download_task1)
|
||||
self.assertFalse(found_download_task2)
|
||||
self.assertFalse(found_download_task3)
|
||||
self.assertFalse(found_thumbnail_task1)
|
||||
self.assertFalse(found_thumbnail_task2)
|
||||
self.assertFalse(found_thumbnail_task3)
|
||||
q = {'task_name': 'sync.tasks.download_media'}
|
||||
download_media_tasks = Task.objects.filter(**q)
|
||||
self.assertFalse(download_media_tasks)
|
||||
|
||||
def test_tasks(self):
|
||||
# Tasks overview page
|
||||
|
||||
+6
-10
@@ -30,11 +30,11 @@ from .forms import (ValidateSourceForm, ConfirmDeleteSourceForm, RedownloadMedia
|
||||
SkipMediaForm, EnableMediaForm, ResetTasksForm, ScheduleTaskForm,
|
||||
ConfirmDeleteMediaServerForm, SourceForm)
|
||||
from .utils import delete_file, validate_url
|
||||
from .tasks import (map_task_to_instance, get_error_message,
|
||||
get_source_completed_tasks, get_media_download_task,
|
||||
delete_task_by_media, index_source_task,
|
||||
download_media_image,
|
||||
check_source_directory_exists, migrate_queues)
|
||||
from .tasks import (
|
||||
map_task_to_instance, get_error_message, migrate_queues, delete_task_by_media,
|
||||
get_running_tasks, get_media_download_task, get_source_completed_tasks,
|
||||
check_source_directory_exists, index_source_task, download_media_image,
|
||||
)
|
||||
from .choices import (Val, MediaServerType, SourceResolution, IndexSchedule,
|
||||
YouTube_SourceType, youtube_long_source_types,
|
||||
youtube_help, youtube_validation_urls)
|
||||
@@ -875,11 +875,7 @@ class TasksView(ListView):
|
||||
scheduled_qs = get_waiting_tasks()
|
||||
# Huey removes running tasks,
|
||||
# so the waiting tasks will not include them.
|
||||
running_qs = TaskHistory.objects.filter(
|
||||
start_at=F('end_at'),
|
||||
scheduled_at__lte=F('end_at'),
|
||||
end_at__gte=now_dt-timezone.timedelta(hours=12),
|
||||
)
|
||||
running_qs = get_running_tasks(now_dt)
|
||||
errors_qs = scheduled_qs.filter(
|
||||
attempts__gt=0
|
||||
).exclude(last_error__exact='')
|
||||
|
||||
Reference in new issue
Block a user