fix: deduplicate scheduled tasks before enqueue
All checks were successful
Build TubeSync image / container (push) Successful in 3m56s

This commit is contained in:
wimby
2026-09-28 05:21:33 +02:00
parent 78b33544d3
commit 2f1d02c011
2 changed files with 92 additions and 1 deletions

View File

@@ -1,6 +1,6 @@
import uuid
from datetime import timedelta
from itertools import islice
from itertools import chain, islice
from django.db import connection, models, transaction
from django.utils import timezone
@@ -27,6 +27,18 @@ def th_schedule(cls, task_wrapper, /, *args, remove_duplicates=False, vn_args=()
if not ('delay' in kwargs or 'eta' in kwargs):
kwargs['delay'] = task_wrapper.settings.get('delay') or int()
task_obj = task_wrapper.s(*args, **kwargs)
if remove_duplicates:
waiting_tasks = chain(
task_wrapper.huey.pending(),
task_wrapper.huey.scheduled(),
)
duplicate_exists = any(
existing.name == task_obj.name and
existing.data == task_obj.data
for existing in waiting_tasks
)
if duplicate_exists:
return False
task_id = str(task_obj.id)
scheduled_at = task_wrapper.huey.scheduled_at_from_task(task_obj)
if scheduled_at:

View File

@@ -1,9 +1,12 @@
import os.path
from types import SimpleNamespace
from unittest.mock import Mock
from django.conf import settings
from django.test import TestCase, Client
from .testutils import prevent_request_warnings
from .utils import parse_database_connection_string, clean_filename
from .errors import DatabaseConnectionError
from .models.tasks import th_schedule
class ErrorPageTestCase(TestCase):
@@ -142,3 +145,79 @@ class UtilsTestCase(TestCase):
self.assertEqual(clean_filename('a a'), 'a a')
self.assertEqual(clean_filename('a\t\t\ta'), 'a a')
self.assertEqual(clean_filename('a\t\t\ta\t\t\t'), 'a a')
class TaskHistoryScheduleTestCase(TestCase):
def make_schedule_objects(self, pending=(), scheduled=()):
task = SimpleNamespace(
id='00000000-0000-0000-0000-000000000001',
__module__='sync.tasks',
name='download_media_metadata',
priority=60,
args=('media-id',),
kwargs={},
data=(('media-id',), {}),
)
huey = SimpleNamespace(
name='limited',
pending=Mock(return_value=list(pending)),
scheduled=Mock(return_value=list(scheduled)),
scheduled_at_from_task=Mock(return_value=None),
enqueue=Mock(),
)
wrapper = SimpleNamespace(
huey=huey,
settings={'delay': 60},
s=Mock(return_value=task),
)
objects = SimpleNamespace(update_or_create=Mock())
cls = SimpleNamespace(objects=objects)
return cls, wrapper, task
def test_remove_duplicates_skips_pending_duplicate(self):
existing = SimpleNamespace(
name='download_media_metadata',
data=(('media-id',), {}),
)
cls, wrapper, _ = self.make_schedule_objects(pending=(existing,))
result = th_schedule(
cls, wrapper, 'media-id',
remove_duplicates=True,
vn_fmt='metadata {}', vn_args=('media-id',),
)
self.assertFalse(result)
cls.objects.update_or_create.assert_not_called()
wrapper.huey.enqueue.assert_not_called()
def test_remove_duplicates_skips_scheduled_duplicate(self):
existing = SimpleNamespace(
name='download_media_metadata',
data=(('media-id',), {}),
)
cls, wrapper, _ = self.make_schedule_objects(scheduled=(existing,))
result = th_schedule(
cls, wrapper, 'media-id',
remove_duplicates=True,
vn_fmt='metadata {}', vn_args=('media-id',),
)
self.assertFalse(result)
cls.objects.update_or_create.assert_not_called()
wrapper.huey.enqueue.assert_not_called()
def test_remove_duplicates_enqueues_when_no_duplicate_exists(self):
cls, wrapper, task = self.make_schedule_objects()
result = th_schedule(
cls, wrapper, 'media-id',
remove_duplicates=True,
vn_fmt='metadata {}', vn_args=('media-id',),
)
self.assertTrue(result)
cls.objects.update_or_create.assert_called_once()
wrapper.huey.enqueue.assert_called_once_with(task)