Convert save_all_media_for_source task

This commit is contained in:
tcely authored and GitHub committed 2025-07-04 14:52:32 -04:00
1 parent f456c22b30
commit d5edc33955
1 file changed
+65 -74
+65 -74
View File
@@ -1040,6 +1040,71 @@ def rename_all_media_for_source(source_id):
with atomic(durable=False):
media.rename_files()
@dynamic_retry(db_task, delay=600, priority=70, retries=15, queue=Val(TaskQueue.FS))
def save_all_media_for_source(source_id):
'''
Iterates all media items linked to a source and saves them to
trigger the post_save signal for every media item. Used when a
source has its parameters changed and all media needs to be
checked to see if its download status has changed.
'''
db.reset_queries()
try:
source = Source.objects.get(pk=source_id)
except Source.DoesNotExist as e:
# Task triggered but the source no longer exists, do nothing
log.error(f'Task save_all_media_for_source(pk={source_id}) called but no '
f'source exists with ID: {source_id}')
raise CancelExecution(_('no such source'), retry=False) from e
# Keep out of the way of the index task!
# SQLite will be locked for a while if we start
# a large source, which reschedules a more costly task.
if 'sqlite' == db_vendor:
index_task = get_source_index_task(source_id)
if index_task and index_task.locked_by_pid_running():
raise Exception(_('Indexing not completed'))
refresh_qs = Media.objects.all().only(
'pk',
'uuid',
'key',
'title', # for name property
).filter(
source=source,
can_download=False,
skip=False,
manual_skip=False,
downloaded=False,
metadata__isnull=False,
)
save_qs = Media.objects.all().only(
'pk',
'uuid',
).filter(
source=source,
)
saved_later = {
str(media.pk)
for media in qs_gen(refresh_qs)
if media.has_metadata
}
refresh_formats.map(saved_later)
# Trigger the post_save signal for each media item linked to this source as various
# flags may need to be recalculated
saved_now = {
str(media.pk)
for media in qs_gen(save_qs)
if str(media.pk) not in saved_later
}
# wait for tasks to complete
res = save_media.map(saved_now)
res.get(blocking=True)
# Old tasks system
from background_task import background # noqa: E402
from background_task.exceptions import InvalidTaskError # noqa: E402
@@ -1113,80 +1178,6 @@ def download_media(media_id, override=False):
raise InvalidTaskError(str(e)) from e
@background(schedule=dict(priority=30, run_at=600), queue=Val(TaskQueue.FS), remove_existing_tasks=True)
def save_all_media_for_source(source_id):
'''
Iterates all media items linked to a source and saves them to
trigger the post_save signal for every media item. Used when a
source has its parameters changed and all media needs to be
checked to see if its download status has changed.
'''
db.reset_queries()
try:
source = Source.objects.get(pk=source_id)
except Source.DoesNotExist as e:
# Task triggered but the source no longer exists, do nothing
log.error(f'Task save_all_media_for_source(pk={source_id}) called but no '
f'source exists with ID: {source_id}')
raise InvalidTaskError(_('no such source')) from e
refresh_qs = Media.objects.all().only(
'pk',
'uuid',
'key',
'title', # for name property
).filter(
source=source,
can_download=False,
skip=False,
manual_skip=False,
downloaded=False,
metadata__isnull=False,
)
save_qs = Media.objects.all().only(
'pk',
'uuid',
).filter(
source=source,
)
saved_later = set()
task = get_source_check_task(source_id)
if task:
task._verbose_name = remove_enclosed(
task.verbose_name, '[', ']', ' ',
valid='0123456789/,',
end=task.verbose_name.find('Check'),
)
tvn_format = '1/{:,}' + f'/{refresh_qs.count():,}'
for mn, media in enumerate(qs_gen(refresh_qs), start=1):
update_task_status(task, tvn_format.format(mn))
refresh_formats(str(media.pk))
saved_later.add(media.uuid)
# Keep out of the way of the index task!
# SQLite will be locked for a while if we start
# a large source, which reschedules a more costly task.
if 'sqlite' == db_vendor:
index_task = get_source_index_task(source_id)
if index_task and index_task.locked_by_pid_running():
raise Exception(_('Indexing not completed'))
# Trigger the post_save signal for each media item linked to this source as various
# flags may need to be recalculated
saved_now = set()
tvn_format = '2/{:,}' + f'/{save_qs.count():,}'
for mn, media in enumerate(qs_gen(save_qs), start=1):
if media.uuid not in saved_later:
update_task_status(task, tvn_format.format(mn))
saved_now.add(str(media.pk))
#save_model(media)
# Reset task.verbose_name to the saved value
update_task_status(task, None)
# wait for tasks to complete
res = save_media.map(saved_now)
res.get(blocking=True)
@background(schedule=dict(priority=0, run_at=60), queue=Val(TaskQueue.NET), remove_existing_tasks=True)
def wait_for_media_premiere(media_id):
try: