Merge pull request #1153 from tcely/patch-4

Reorganize tasks
This commit is contained in:
meeb
2025-06-25 16:41:13 +10:00
committed by GitHub
9 changed files with 182 additions and 156 deletions

View File

@@ -1,6 +0,0 @@
#!/command/with-contenv bash
exec nice -n "${TUBESYNC_NICE:-1}" s6-setuidgid app \
/usr/bin/python3 /app/manage.py process_tasks \
--queue database --duration 86400 \
--sleep "30.${RANDOM}"

View File

@@ -20,13 +20,8 @@ from django import db
from django.conf import settings
from django.core.files.base import ContentFile
from django.core.files.uploadedfile import SimpleUploadedFile
from django.db import DatabaseError
from django.db.transaction import atomic
from django.utils import timezone
from django.utils.translation import gettext_lazy as _
from background_task import background
from background_task.exceptions import InvalidTaskError
from background_task.models import Task, CompletedTask
from django_huey import lock_task as huey_lock_task, task as huey_task # noqa
from django_huey import db_periodic_task, db_task, signal as huey_signal
from huey import crontab as huey_crontab, signals as huey_signals
@@ -42,6 +37,7 @@ from .models import Source, Media, MediaServer, Metadata
from .utils import get_remote_image, resize_image_to_height, filter_response
from .youtube import YouTubeError
atomic = db.transaction.atomic
db_vendor = db.connection.vendor
register_huey_signals()
@@ -132,7 +128,7 @@ def update_task_status(task, status):
task.verbose_name = f'[{status}] {task._verbose_name}'
try:
task.save(update_fields={'verbose_name'})
except DatabaseError as e:
except db.DatabaseError as e:
if 'Save with update_fields did not affect any rows.' == str(e):
pass
else:
@@ -208,7 +204,11 @@ def cleanup_completed_tasks():
@atomic(durable=False)
def migrate_queues():
tqs = Task.objects.all()
qs = tqs.exclude(queue__in=TaskQueue.values)
remaining_queues = list((
Val(TaskQueue.FS),
Val(TaskQueue.NET),
))
qs = tqs.exclude(queue__in=remaining_queues)
return qs.update(queue=Val(TaskQueue.NET))
@@ -224,6 +224,38 @@ def save_model(instance):
time.sleep(random.expovariate(arg))
@db_periodic_task(
huey_crontab(minute=40, strict=True,),
priority=100,
expires=15*60,
queue=Val(TaskQueue.DB),
)
def upcoming_media():
now = timezone.now()
next_hour = now + timezone.timedelta(hours=1, minutes=3)
previous_hour = now - timezone.timedelta(hours=1, minutes=1)
qs = Media.objects.filter(
manual_skip=True,
metadata__isnull=False,
published__isnull=False,
published__gte=previous_hour,
)
for media in qs_gen(qs):
valid, hours = media.wait_for_premiere()
if valid:
save_model(media)
vn_fmt = _('Waiting for the premiere of "{}" at: {}')
wait_for_media_premiere(
str(media.pk),
run_at=next_hour,
verbose_name=vn_fmt.format(
media.key,
media.published.isoformat(' ', 'seconds'),
),
)
log.debug(f'upcoming_media: wait_for_premiere: {media.key}: {valid=} {hours=}')
@db_periodic_task(
huey_crontab(minute=59, strict=True,),
priority=100,
@@ -432,25 +464,6 @@ def migrate_to_metadata(media_id):
media.save_to_metadata(field, value)
@background(schedule=dict(priority=0, run_at=0), queue=Val(TaskQueue.NET), remove_existing_tasks=False)
def wait_for_database_queue():
from common.huey import h_q_tuple
queue_name = Val(TaskQueue.DB)
consumer_down_path = Path(f'/run/service/huey-{queue_name}/down')
included_names = frozenset(('migrate_to_metadata',))
total_count = 1
while 0 < total_count:
if consumer_down_path.exists() and consumer_down_path.is_file():
raise BgTaskWorkerError(_('queue consumer stopped'))
time.sleep(5)
status_dict = h_q_tuple(queue_name)[2]
total_count = status_dict.get('pending', (0,))[0]
scheduled_tasks = status_dict.get('scheduled', (0,[]))[1]
total_count += sum(
[ 1 for t in scheduled_tasks if t.name.rsplit('.', 1)[-1] in included_names ],
)
@db_task(delay=30, priority=80, queue=Val(TaskQueue.LIMIT))
def index_source(source_id):
'''
@@ -484,15 +497,6 @@ def index_source(source_id):
# Got some media, update the last crawl timestamp
source.last_crawl = timezone.now()
save_model(source)
wait_for_database_queue(
priority=19, # the indexing task uses 20
verbose_name=_('Waiting for database tasks to complete'),
)
wait_for_database_queue(
priority=29, # the checking task uses 30
queue=Val(TaskQueue.FS),
verbose_name=_('Delaying checking all media for database tasks'),
)
delete_task_by_source('sync.tasks.save_all_media_for_source', source.pk)
num_videos = len(videos)
log.info(f'Found {num_videos} media items for source: {source}')
@@ -612,15 +616,6 @@ def index_source(source_id):
return True
@background(schedule=dict(priority=20, run_at=30), queue=Val(TaskQueue.NET), remove_existing_tasks=True)
def index_source_task(source_id):
try:
res = index_source(source_id)
return res.get(blocking=True)
except CancelExecution as e:
raise InvalidTaskError(str(e)) from e
@dynamic_retry(db_task, priority=100, retries=15, queue=Val(TaskQueue.FS))
def check_source_directory_exists(source_id):
'''
@@ -788,13 +783,6 @@ def download_metadata(media_id):
media.published = published_datetime
media.manual_skip = True
media.save()
verbose_name = _('Waiting for the premiere of "{}" at: {}')
wait_for_media_premiere(
str(media.pk),
repeat=Task.HOURLY,
repeat_until = published_datetime + timedelta(hours=1),
verbose_name=verbose_name.format(media.key, published_datetime.isoformat(' ', 'seconds')),
)
raise_exception = False
if raise_exception:
raise
@@ -837,15 +825,6 @@ def download_metadata(media_id):
return True
@background(schedule=dict(priority=40, run_at=60), queue=Val(TaskQueue.NET), remove_existing_tasks=True)
def download_media_metadata(media_id):
try:
res = download_metadata(media_id)
return res.get(blocking=True)
except CancelExecution as e:
raise InvalidTaskError(str(e)) from e
@dynamic_retry(db_task, delay=10, priority=90, retries=15, queue=Val(TaskQueue.NET))
def download_media_image(media_id, url):
'''
@@ -914,13 +893,6 @@ def on_complete_download_media_image(signal_name, task_obj, exception_obj=None,
if result is True:
huey.result(preserve=False, id=task_obj.id)
@background(schedule=dict(priority=10, run_at=10), queue=Val(TaskQueue.FS), remove_existing_tasks=True)
def download_media_thumbnail(media_id, url):
try:
return download_media_image.call_local(media_id, url)
except CancelExecution as e:
raise InvalidTaskError(str(e)) from e
@db_task(delay=60, priority=70, queue=Val(TaskQueue.LIMIT))
def download_media_file(media_id, override=False):
'''
@@ -977,15 +949,6 @@ def download_media_file(media_id, override=False):
return True
@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
@db_task(delay=30, expires=210, priority=100, queue=Val(TaskQueue.NET))
def rescan_media_server(mediaserver_id):
'''
@@ -1001,6 +964,148 @@ def rescan_media_server(mediaserver_id):
mediaserver.update()
@dynamic_retry(db_task, backoff_func=lambda n: (n*3600)+600, priority=50, retries=15, queue=Val(TaskQueue.LIMIT))
def refresh_formats(media_id):
try:
media = Media.objects.get(pk=media_id)
except Media.DoesNotExist as e:
raise CancelExecution(_('no such media'), retry=False) from e
else:
wait_for_errors(
media,
queue_name=Val(TaskQueue.LIMIT),
)
save, retry, msg = media.refresh_formats()
if save is not True:
log.warning(f'Refreshing formats for "{media.key}" failed: {msg}')
exc = CancelExecution(
_('failed to refresh formats for:'),
f'{media.key} / {media.uuid}:',
msg,
retry=retry,
)
# combine the strings
exc.args = (' '.join(map(str, exc.args)),)
# store instance details
exc.instance = dict(
key=media.key,
model='Media',
uuid=str(media.pk),
)
# store the function results
exc.reason = msg
exc.save = save
raise exc
log.info(f'Saving refreshed formats for "{media.key}": {msg}')
save_model(media)
@db_task(delay=300, priority=80, retries=5, retry_delay=600, queue=Val(TaskQueue.FS))
@atomic(durable=True)
def rename_all_media_for_source(source_id):
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 rename_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
# Check that the settings allow renaming
rename_sources_setting = getattr(settings, 'RENAME_SOURCES') or list()
create_rename_tasks = (
(
source.directory and
source.directory in rename_sources_setting
) or
getattr(settings, 'RENAME_ALL_SOURCES', False)
)
if not create_rename_tasks:
return None
mqs = Media.objects.filter(
source=source,
downloaded=True,
)
for media in qs_gen(mqs):
with huey_lock_task(
f'media:{media.uuid}',
queue=Val(TaskQueue.DB),
):
with atomic(durable=False):
media.rename_files()
# Old tasks system
from background_task import background # noqa: E402
from background_task.exceptions import InvalidTaskError # noqa: E402
from background_task.models import Task, CompletedTask # noqa: E402
@background(schedule=dict(priority=0, run_at=0), queue=Val(TaskQueue.FS), remove_existing_tasks=False)
def wait_for_database_queue():
from common.huey import h_q_tuple
queue_name = Val(TaskQueue.DB)
consumer_down_path = Path(f'/run/service/huey-{queue_name}/down')
included_names = frozenset(('migrate_to_metadata',))
total_count = 1
while 0 < total_count:
if consumer_down_path.exists() and consumer_down_path.is_file():
raise BgTaskWorkerError(_('queue consumer stopped'))
time.sleep(5)
status_dict = h_q_tuple(queue_name)[2]
total_count = status_dict.get('pending', (0,))[0]
scheduled_tasks = status_dict.get('scheduled', (0,[]))[1]
total_count += sum(
[ 1 for t in scheduled_tasks if t.name.rsplit('.', 1)[-1] in included_names ],
)
@background(schedule=dict(priority=20, run_at=30), queue=Val(TaskQueue.NET), remove_existing_tasks=True)
def index_source_task(source_id):
try:
res = index_source(source_id)
retval = res.get(blocking=True)
except CancelExecution as e:
raise InvalidTaskError(str(e)) from e
else:
if retval is not True:
return retval
wait_for_database_queue(
priority=29, # the checking task uses 30
queue=Val(TaskQueue.FS),
verbose_name=_('Delaying checking all media for database tasks'),
)
wait_for_database_queue(
priority=19, # the indexing task uses 20
queue=Val(TaskQueue.NET),
verbose_name=_('Waiting for database tasks to complete'),
)
return True
@background(schedule=dict(priority=40, run_at=60), queue=Val(TaskQueue.NET), remove_existing_tasks=True)
def download_media_metadata(media_id):
try:
res = download_metadata(media_id)
return res.get(blocking=True)
except CancelExecution as e:
raise InvalidTaskError(str(e)) from e
@background(schedule=dict(priority=10, run_at=10), queue=Val(TaskQueue.NET), remove_existing_tasks=True)
def download_media_thumbnail(media_id, url):
try:
return download_media_image.call_local(media_id, url)
except CancelExecution as e:
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
@background(schedule=dict(priority=30, run_at=600), queue=Val(TaskQueue.FS), remove_existing_tasks=True)
def save_all_media_for_source(source_id):
'''
@@ -1075,77 +1180,7 @@ def save_all_media_for_source(source_id):
res.get(blocking=True)
@dynamic_retry(db_task, backoff_func=lambda n: (n*3600)+600, priority=50, retries=15, queue=Val(TaskQueue.LIMIT))
def refresh_formats(media_id):
try:
media = Media.objects.get(pk=media_id)
except Media.DoesNotExist as e:
raise CancelExecution(_('no such media'), retry=False) from e
else:
wait_for_errors(
media,
queue_name=Val(TaskQueue.LIMIT),
)
save, retry, msg = media.refresh_formats()
if save is not True:
log.warning(f'Refreshing formats for "{media.key}" failed: {msg}')
exc = CancelExecution(
_('failed to refresh formats for:'),
f'{media.key} / {media.uuid}:',
msg,
retry=retry,
)
# combine the strings
exc.args = (' '.join(map(str, exc.args)),)
# store instance details
exc.instance = dict(
key=media.key,
model='Media',
uuid=str(media.pk),
)
# store the function results
exc.reason = msg
exc.save = save
raise exc
log.info(f'Saving refreshed formats for "{media.key}": {msg}')
save_model(media)
@db_task(delay=300, priority=80, retries=5, retry_delay=600, queue=Val(TaskQueue.FS))
@atomic(durable=True)
def rename_all_media_for_source(source_id):
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 rename_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
# Check that the settings allow renaming
rename_sources_setting = getattr(settings, 'RENAME_SOURCES') or list()
create_rename_tasks = (
(
source.directory and
source.directory in rename_sources_setting
) or
getattr(settings, 'RENAME_ALL_SOURCES', False)
)
if not create_rename_tasks:
return None
mqs = Media.objects.filter(
source=source,
downloaded=True,
)
for media in qs_gen(mqs):
with huey_lock_task(
f'media:{media.uuid}',
queue=Val(TaskQueue.DB),
):
with atomic(durable=False):
media.rename_files()
@background(schedule=dict(priority=0, run_at=60), queue=Val(TaskQueue.DB), remove_existing_tasks=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:
media = Media.objects.get(pk=media_id)