Merge pull request #1217 from tcely/patch-2

Remove the old `index_source_task` task
This commit is contained in:
meeb authored and GitHub committed 2025-07-17 17:44:15 +10:00
commit 9827de4869
9 files changed
+174 -139

No files matched your search

+3 -1
View File
@@ -115,7 +115,9 @@ jobs:
- name: Set up Django environment
run: |
mkdir -v -p ~/.config/TubeSync/config
mkdir -v -p ~/.config/TubeSync/downloads
sudo ln -v -s -f -T ~/.config/TubeSync/config /config
sudo ln -v -s -f -T ~/.config/TubeSync/downloads /downloads
cp -v -p tubesync/tubesync/local_settings.py.example tubesync/tubesync/local_settings.py
cp -v -a -t "${Python3_ROOT_DIR}"/lib/python3.*/site-packages/background_task/ patches/background_task/*
cp -v -a -t "${Python3_ROOT_DIR}"/lib/python3.*/site-packages/yt_dlp/ patches/yt_dlp/*
@@ -162,7 +164,7 @@ jobs:
--output-format github \
--ignore "${ignore_csv_list}"
- name: Run Django tests
run: cd tubesync && python3 -B -W default manage.py test --verbosity=2
run: cd tubesync && TUBESYNC_DEBUG=True python3 -B -W default manage.py test --no-input --buffer --verbosity=2
containerise:
if: ${{ !cancelled() && 'success' == needs.info.result }}
+27 -11
View File
@@ -12,20 +12,36 @@ def th_schedule(cls, task_wrapper, /, *args, remove_duplicates=False, vn_args=()
assert vn_fmt is not None, 'vn_fmt is required'
if vn_fmt is None:
return False
defaults = dict(
queue=task_wrapper.huey.name,
remove_duplicates=remove_duplicates,
verbose_name=str(vn_fmt).format(*vn_args),
)
# support using the delay setting from the decorator
if not ('delay' in kwargs or 'eta' in kwargs):
kwargs['delay'] = task_wrapper.settings.get('delay') or int()
result = task_wrapper.schedule(args=args, **kwargs)
try:
task_history = cls.objects.get(task_id=str(result.id))
except cls.DoesNotExist:
pass
else:
task_history.remove_duplicates = remove_duplicates
task_history.verbose_name = str(vn_fmt).format(*vn_args)
task_history.save()
return True
return False
task_obj = task_wrapper.s(*args, **kwargs)
task_id = str(task_obj.id)
scheduled_at = task_wrapper.huey.scheduled_at_from_task(task_obj)
if scheduled_at:
defaults['scheduled_at'] = scheduled_at
defaults['end_at'] = timezone.datetime.now(timezone.timezone.utc)
defaults['name'] = f'{task_obj.__module__}.{task_obj.name}'
defaults['priority'] = task_obj.priority
defaults['task_params'] = list((
list(task_obj.args),
repr(task_obj.kwargs),
))
cls.objects.update_or_create(
task_id=task_id,
defaults=defaults,
create_defaults={
'task_id': task_id,
**defaults,
},
)
task_wrapper.huey.enqueue(task_obj)
return True
class TaskHistoryQuerySet(models.QuerySet):
+20 -34
View File
@@ -17,16 +17,15 @@ from common.models import TaskHistory
from common.utils import glob_quote, mkdir_p
from .models import Source, Media, Metadata
from .tasks import (
delete_task_by_media, delete_task_by_source,
delete_task_by_media,
get_media_download_task, get_media_metadata_task, get_media_thumbnail_task,
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,
check_source_directory_exists, download_source_images, index_source,
download_media_file, download_media_metadata, download_media_image,
)
from .utils import delete_file
from .filtering import filter_media
from .choices import Val, YouTube_SourceType
@receiver(pre_save, sender=Source)
@@ -102,35 +101,34 @@ def source_pre_save(sender, instance, **kwargs):
)
if recreate_index_source_task:
# Indexing schedule has changed, recreate the indexing task
delete_task_by_source('sync.tasks.index_source_task', instance.pk)
verbose_name = _('Index media from source "{}"')
index_source_task(
TaskHistory.schedule(
index_source,
str(instance.pk),
repeat=0,
schedule=instance.task_run_at_dt,
verbose_name=verbose_name.format(instance.name),
eta=instance.task_run_at_dt,
remove_duplicates=True,
vn_fmt=_('Index media from source "{}"'),
vn_args=(instance.name,),
)
@receiver(post_save, sender=Source)
def source_post_save(sender, instance, created, **kwargs):
source = instance
# Check directory exists and create an indexing task for newly created sources
if created:
check_source_directory_exists(str(instance.pk))
if instance.source_type != Val(YouTube_SourceType.PLAYLIST) and instance.copy_channel_images:
download_source_images(str(instance.pk))
if instance.index_schedule > 0:
delete_task_by_source('sync.tasks.index_source_task', instance.pk)
log.info(f'Scheduling first media indexing for source: {instance.name}')
verbose_name = _('Index media from source "{}"')
index_source_task(
str(instance.pk),
repeat=0,
schedule=600,
verbose_name=verbose_name.format(instance.name),
check_source_directory_exists(str(source.pk))
if source.copy_channel_images and not source.is_playlist:
download_source_images(str(source.pk))
if source.is_active:
log.info(f'Scheduling first media indexing for source: {source.name}')
TaskHistory.schedule(
index_source,
str(source.pk),
delay=600,
vn_fmt=_('Index media from source "{}"'),
vn_args=(source.name,),
)
source = instance
TaskHistory.schedule(
save_all_media_for_source,
str(source.pk),
@@ -149,9 +147,6 @@ def source_pre_delete(sender, instance, **kwargs):
source = instance
log.info(f'Deactivating source: {instance.name}')
instance.deactivate()
log.info(f'Deleting tasks for source: {instance.name}')
delete_task_by_source('sync.tasks.index_source_task', instance.pk)
delete_task_by_source('sync.tasks.rename_all_media_for_source', instance.pk)
# Fetch the media source
sqs = Source.objects.filter(filter_text=str(source.pk))
@@ -171,15 +166,6 @@ def source_pre_delete(sender, instance, **kwargs):
))
@receiver(post_delete, sender=Source)
def source_post_delete(sender, instance, **kwargs):
# Triggered after a source is deleted
source = instance
log.info(f'Deleting tasks for removed source: {source.name}')
delete_task_by_source('sync.tasks.index_source_task', instance.pk)
delete_task_by_source('sync.tasks.rename_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):
+33 -53
View File
@@ -30,7 +30,7 @@ from common.huey import CancelExecution, dynamic_retry, register_huey_signals
from common.logger import log
from common.models import TaskHistory
from common.errors import (
BgTaskWorkerError, HueyConsumerError,
HueyConsumerError,
DownloadFailedException, FormatUnavailableError,
NoFormatException, NoMediaException, NoThumbnailException,
)
@@ -150,10 +150,19 @@ def get_source_completed_tasks(source_id, only_errors=False):
return CompletedTask.objects.filter(**q).order_by('-failed_at')
def get_model_task(model_pk, /, name=None, qs=None):
if qs is None:
qs = TaskHistory.objects.all()
if name is not None:
qs = qs.filter(name__endswith=name)
params_prefix = f'[["{model_pk}"'
qs = qs.filter(task_params__istartswith=params_prefix)
return qs[0] if qs.count() else False
def get_running_tasks(arg_dt=None, /):
return TaskHistory.objects.running(
now=arg_dt,
within=timezone.timedelta(hours=12),
within=timezone.timedelta(seconds=settings.MAX_RUN_TIME),
)
def get_running_task_by_name(arg_str, media_id, /):
@@ -164,10 +173,20 @@ def get_running_task_by_name(arg_str, 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)
#return get_running_task_by_name('download_media_file', media_id)
return get_model_task(
media_id,
name='download_media_file',
qs=get_running_tasks(),
)
def get_media_thumbnail_task(media_id):
return get_running_task_by_name('download_media_image', media_id)
#return get_running_task_by_name('download_media_image', media_id)
return get_model_task(
media_id,
name='download_media_image',
qs=get_running_tasks(),
)
def get_tasks(task_name, id=None, /, instance=None):
@@ -303,12 +322,14 @@ def schedule_indexing():
).clear()
# schedule a new indexing task
log.info(f'Scheduling an indexing task for source "{source.name}": {source.pk}')
vn_fmt = _('Index media from source "{}"')
index_source_task(
TaskHistory.schedule(
index_source,
str(source.pk),
repeat=0,
schedule=600,
verbose_name=vn_fmt.format(source.name),
delay=300,
expires=40*60,
remove_duplicates=True,
vn_fmt=_('Index media from source "{}"'),
vn_args=(source.name,),
)
@@ -359,9 +380,9 @@ def wait_for_errors(model, /, *, queue_name=None, task_name=None):
total_count += sum([ 1 if contains_http429(q, k) else 0 for k in q.all_results() ])
delay = 10 * total_count
time_str = seconds_to_timestr(delay)
db_down_path = Path('/run/service/huey-database/down')
fs_down_path = Path('/run/service/huey-filesystem/down')
log.info(f'waiting for errors: 429 ({time_str}): {model}')
db_down_path = Path('/run/service/tubesync-db-worker/down')
fs_down_path = Path('/run/service/tubesync-fs-worker/down')
while delay > 0:
# this happenes when the container is shutting down
# do not prevent that while we are delaying a task
@@ -372,7 +393,7 @@ def wait_for_errors(model, /, *, queue_name=None, task_name=None):
for task in tasks:
update_task_status(task, None)
if delay > 0:
raise BgTaskWorkerError(_('queue worker stopped'))
raise HueyConsumerError(_('queue consumer stopped'))
@db_task(priority=90, queue=Val(TaskQueue.FS))
@@ -735,8 +756,6 @@ def delete_media(media_id):
queue=Val(TaskQueue.DB),
):
media.delete()
return True
return False
@db_task(delay=60, priority=70, retries=5, retry_delay=60, queue=Val(TaskQueue.FS))
@@ -779,8 +798,6 @@ def save_media(media_id):
queue=Val(TaskQueue.DB),
):
save_model(media)
return True
return False
@db_task(delay=60, priority=60, retries=3, retry_delay=600, queue=Val(TaskQueue.LIMIT))
@@ -1261,43 +1278,6 @@ def wait_for_media_premiere(media_id):
if t[0]:
save_model(media)
@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 HueyConsumerError(_('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=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:
+81 -25
View File
@@ -13,8 +13,9 @@ from xml.etree import ElementTree
from django.conf import settings
from django.test import TestCase, Client, override_settings
from django.utils import timezone
from background_task.models import Task
from django_huey import DJANGO_HUEY, get_queue
from common.models import TaskHistory
from huey.consumer_options import ConsumerConfig
from .models import Source, Media
from .tasks import (
cleanup_old_media, check_source_directory_exists,
@@ -29,7 +30,29 @@ from .choices import (Val, Fallback, IndexSchedule, SourceResolution,
class FrontEndTestCase(TestCase):
@classmethod
def setUpClass(cls):
super().setUpClass()
cls._consumers = dict()
for qn, qc in DJANGO_HUEY.get('queues', dict()).items():
q = get_queue(qn)
consumer_opts = qc.get('consumer', {})
config = ConsumerConfig(**consumer_opts)
config.validate()
#consumer = q.create_consumer(**config.values)
#cls._consumers[qn] = consumer
#consumer.start()
q.immediate = True
q.immediate_use_memory = True
@classmethod
def tearDownClass(cls):
for qn, consumer in cls._consumers.items():
consumer.stop(graceful=True)
super().tearDownClass()
def setUp(self):
self.maxDiff = None
# Disable general logging for test case
logging.disable(logging.CRITICAL)
@@ -165,6 +188,14 @@ class FrontEndTestCase(TestCase):
self.assertTrue(checked_directory)
def test_source(self):
#logging.disable(logging.NOTSET)
def get_model_task(model_pk, /, name=None):
qs = TaskHistory.objects.all()
if name is not None:
qs = qs.filter(name__endswith=name)
params_prefix = f'[["{model_pk}"'
qs = qs.filter(task_params__istartswith=params_prefix)
return qs[0] if qs.count() else False
# Sources overview page
c = Client()
response = c.get('/sources')
@@ -185,6 +216,8 @@ class FrontEndTestCase(TestCase):
'filter_text': '.*',
'filter_seconds_min': int(True),
'index_schedule': 3600,
'download_media': False,
'index_videos': True,
'delete_old_media': False,
'days_to_keep': 14,
'source_resolution': '1080p',
@@ -210,11 +243,6 @@ class FrontEndTestCase(TestCase):
# Check that the SponsorBlock categories were saved
self.assertEqual(source.sponsorblock_categories.selected_choices,
expected_categories)
# Check a task was created to index the media for the new source
source_uuid = str(source.pk)
task = Task.objects.get_task('sync.tasks.index_source_task',
args=(source_uuid,))[0]
self.assertEqual(task.queue, Val(TaskQueue.NET))
# Run the check_source_directory_exists task
check_source_directory_exists.call_local(source_uuid)
# Check the source is now on the source overview page
@@ -224,6 +252,21 @@ class FrontEndTestCase(TestCase):
# Check the source detail page loads
response = c.get(f'/source/{source_uuid}')
self.assertEqual(response.status_code, 200)
# Check a task was created to index the media for the new source
index_task_qs = TaskHistory.objects.filter(
name='sync.tasks.index_source',
task_params__0__0=source_uuid,
).order_by('end_at')
self.assertNotEqual(list(), list(index_task_qs))
self.assertNotEqual(
list(),
[
th.__dict__ for th in TaskHistory.objects.all()
]
)
task = get_model_task(source_uuid, name='index_source')
self.assertNotEqual(False, task)
self.assertEqual(task.queue, get_queue(Val(TaskQueue.LIMIT)).name)
# save and refresh the Source
source.refresh_from_db()
source.sponsorblock_categories.selected_choices.append('sponsor')
@@ -305,8 +348,7 @@ class FrontEndTestCase(TestCase):
self.assertEqual(source.sponsorblock_categories.selected_choices,
expected_categories)
# Check a new task has been created by seeing if the pk has changed
new_task = Task.objects.get_task('sync.tasks.index_source_task',
args=(source_uuid,))[0]
new_task = index_task_qs.last()
self.assertNotEqual(task.pk, new_task.pk)
# Delete source confirmation page
response = c.get(f'/source-delete/{source_uuid}')
@@ -329,10 +371,6 @@ class FrontEndTestCase(TestCase):
# Check the source details page now 404s
response = c.get(f'/source/{source_uuid}')
self.assertEqual(response.status_code, 404)
# Check the indexing media task was removed
tasks = Task.objects.get_task('sync.tasks.index_source_task',
args=(source_uuid,))
self.assertFalse(tasks)
def test_media(self):
# Media overview page
@@ -420,24 +458,35 @@ class FrontEndTestCase(TestCase):
test_media3_pk = str(test_media3.pk)
# simulate the tasks consumer signals having already run
now_dt = timezone.now()
TaskHistory.objects.all().update(
TaskHistory.objects.filter(
name__startswith='sync.tasks.download_media_',
).update(
scheduled_at=before_dt,
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)
def get_model_task(model_pk, /, name=None):
qs = TaskHistory.objects.all()
if name is not None:
qs = qs.filter(name__endswith=name)
params_prefix = f'[["{model_pk}"'
qs = qs.filter(task_params__istartswith=params_prefix)
return 1 == qs.count()
name_suffix = 'download_media_file'
found_download_task1 = get_model_task(test_media1_pk, name_suffix)
found_download_task2 = get_model_task(test_media2_pk, name_suffix)
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)
name_suffix = 'download_media_image'
found_thumbnail_task1 = get_model_task(test_media1_pk, name_suffix)
found_thumbnail_task2 = get_model_task(test_media2_pk, name_suffix)
found_thumbnail_task3 = get_media_thumbnail_task(test_media3_pk)
self.assertTrue(found_download_task1)
self.assertTrue(found_download_task2)
self.assertTrue(found_download_task3)
self.assertTrue(not not found_download_task3)
self.assertTrue(found_thumbnail_task1)
self.assertTrue(found_thumbnail_task2)
self.assertTrue(found_thumbnail_task3)
self.assertTrue(not not found_thumbnail_task3)
# Check the media is listed on the media overview page
response = c.get('/media')
self.assertEqual(response.status_code, 200)
@@ -464,7 +513,9 @@ class FrontEndTestCase(TestCase):
response = c.get(f'/media/{test_media3_pk}')
self.assertEqual(response.status_code, 404)
# simulate the tasks consumer signals having already run
TaskHistory.objects.all().update(end_at=timezone.now())
TaskHistory.objects.filter(
name__startswith='sync.tasks.download_media_',
).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)
@@ -496,15 +547,20 @@ class FrontEndTestCase(TestCase):
metadata_filepath = settings.BASE_DIR / 'sync' / 'testdata' / 'metadata.json'
metadata = open(metadata_filepath, 'rt').read()
with open(metadata_filepath, 'rt') as file:
metadata = file.read()
metadata_hdr_filepath = settings.BASE_DIR / 'sync' / 'testdata' / 'metadata_hdr.json'
metadata_hdr = open(metadata_hdr_filepath, 'rt').read()
with open(metadata_hdr_filepath, 'rt') as file:
metadata_hdr = file.read()
metadata_60fps_filepath = settings.BASE_DIR / 'sync' / 'testdata' / 'metadata_60fps.json'
metadata_60fps = open(metadata_60fps_filepath, 'rt').read()
with open(metadata_60fps_filepath, 'rt') as file:
metadata_60fps = file.read()
metadata_60fps_hdr_filepath = settings.BASE_DIR / 'sync' / 'testdata' / 'metadata_60fps_hdr.json'
metadata_60fps_hdr = open(metadata_60fps_hdr_filepath, 'rt').read()
with open(metadata_60fps_hdr_filepath, 'rt') as file:
metadata_60fps_hdr = file.read()
metadata_20230629_filepath = settings.BASE_DIR / 'sync' / 'testdata' / 'metadata_2023-06-29.json'
metadata_20230629 = open(metadata_20230629_filepath, 'rt').read()
with open(metadata_20230629_filepath, 'rt') as file:
metadata_20230629 = file.read()
all_test_metadata = {
'boring': metadata,
'hdr': metadata_hdr,
+8 -10
View File
@@ -33,7 +33,7 @@ from .utils import delete_file, validate_url
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,
check_source_directory_exists, index_source, download_media_image,
)
from .choices import (Val, MediaServerType, SourceResolution, IndexSchedule,
YouTube_SourceType, youtube_long_source_types,
@@ -142,18 +142,16 @@ class SourcesView(ListView):
def get(self, *args, **kwargs):
if args[0].path.startswith("/source-sync-now/"):
sobj = Source.objects.get(pk=kwargs["pk"])
if sobj is None:
source = Source.objects.get(pk=kwargs["pk"])
if source is None:
return HttpResponseNotFound()
source = sobj
verbose_name = _('Index media from source "{}" once')
index_source_task(
TaskHistory.schedule(
index_source,
str(source.pk),
remove_existing_tasks=False,
repeat=0,
schedule=30,
verbose_name=verbose_name.format(source.name),
delay=30,
vn_fmt=_('Index media from source "{}" once'),
vn_args=(source.name,),
)
url = reverse_lazy('sync:sources')
url = append_uri_params(url, {'message': 'source-refreshed'})
@@ -21,8 +21,6 @@ SECRET_KEY = getenv('DJANGO_SECRET_KEY', 'tubesync-django-secret')
ALLOWED_HOSTS_STR = getenv('TUBESYNC_HOSTS', '*')
ALLOWED_HOSTS = ALLOWED_HOSTS_STR.split(',')
DEBUG_STR = getenv('TUBESYNC_DEBUG', False)
DEBUG = True if 'true' == DEBUG_STR.strip().lower() else False
FORCE_SCRIPT_NAME = getenv('DJANGO_FORCE_SCRIPT_NAME', DJANGO_URL_PREFIX)
@@ -7,7 +7,6 @@ DOWNLOADS_BASE_DIR = BASE_DIR
SECRET_KEY = 'example-secret-key'
DEBUG = False
DATABASES = {
+2 -2
View File
@@ -12,7 +12,7 @@ DOWNLOADS_BASE_DIR = BASE_DIR
VERSION = '0.15.7'
SECRET_KEY = ''
DEBUG = False
DEBUG = 'true' == getenv('TUBESYNC_DEBUG').strip().lower()
ALLOWED_HOSTS = []
@@ -53,7 +53,7 @@ FORCE_SCRIPT_NAME = None
DJANGO_HUEY = {
'default': TaskQueue.LIMIT.value,
'queues': dict(),
'verbose': None if 'true' == getenv('TUBESYNC_DEBUG', False).strip().lower() else False,
'verbose': None if DEBUG else False,
}
for queue_name in TaskQueue.values:
queues = DJANGO_HUEY['queues']