Skip to content

Commit 3aa12fd

Browse files
committed
fix(scheduler): remove registry entries whose job model is gone
1 parent 67ad4ba commit 3aa12fd

4 files changed

Lines changed: 50 additions & 2 deletions

File tree

scheduler/helpers/queues/queue_logic.py

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -131,7 +131,10 @@ def clean_registries(self, timestamp: float | None = None) -> None:
131131

132132
for job_name, job_score in started_jobs:
133133
job = JobModel.get(job_name, connection=self.connection)
134-
if job is None or not job.has_failure_callback or job_score + job.timeout > before_score:
134+
if job is None:
135+
self.active_job_registry.delete(connection=self.connection, job_name=job_name)
136+
continue
137+
if not job.has_failure_callback or job_score + job.timeout > before_score:
135138
continue
136139

137140
logger.debug(f"Running failure callbacks for {job.name}")

scheduler/tests/test_internals.py

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,8 @@
66
from django.utils import timezone
77

88
from scheduler.helpers.callback import Callback, CallbackSetupError
9+
from scheduler.helpers.queues import get_queue
10+
from scheduler.helpers.utils import current_timestamp
911
from scheduler.models import TaskType, get_next_cron_time, get_scheduled_task
1012
from scheduler.tests.testtools import SchedulerBaseCase, task_factory
1113

@@ -47,6 +49,16 @@ def test_callback_bad_arguments(self):
4749
self.assertEqual(str(cm.exception), "Callback `func` must be a string or function, received 1")
4850

4951

52+
class TestCleanRegistries(SchedulerBaseCase):
53+
def test_active_registry_entry_without_job_is_removed(self):
54+
queue = get_queue("default")
55+
queue.active_job_registry.add(queue.connection, "orphan-job-name", current_timestamp() - 3600)
56+
57+
queue.clean_registries()
58+
59+
self.assertFalse(queue.active_job_registry.exists(queue.connection, "orphan-job-name"))
60+
61+
5062
class TestConfSettings(SchedulerBaseCase):
5163
@override_settings(SCHEDULER_CONFIG=[])
5264
def test_conf_settings__bad_scheduler_config(self):

scheduler/tests/test_worker/test_scheduler.py

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,9 @@
33
import time_machine
44
from django.utils import timezone
55

6+
from scheduler.helpers.utils import current_timestamp
67
from scheduler.models import TaskType
8+
from scheduler.redis_models import JobModel
79
from scheduler.settings import SCHEDULER_CONFIG
810
from scheduler.tests.testtools import SchedulerBaseCase, task_factory
911
from scheduler.worker import WorkerScheduler, create_worker
@@ -37,3 +39,32 @@ def test_scheduler_schedules_tasks(self):
3739
self.assertIsNotNone(task.job_name)
3840
self.assertTrue(task.rqueue.queued_job_registry.exists(task.rqueue.connection, task.job_name))
3941
self.assertFalse(task.rqueue.scheduled_job_registry.exists(task.rqueue.connection, task.job_name))
42+
43+
def test_scheduler_removes_scheduled_registry_entry_without_job(self):
44+
# arrange
45+
task = task_factory(TaskType.CRON)
46+
job_name = task.job_name
47+
self.assertIsNotNone(job_name)
48+
connection = task.rqueue.connection
49+
registry = task.rqueue.scheduled_job_registry
50+
connection.delete(JobModel.key_for(job_name))
51+
registry.add(connection, job_name, current_timestamp() - 10)
52+
53+
scheduler = WorkerScheduler([task.rqueue], worker_name="fake-worker")
54+
scheduler._acquire_locks()
55+
56+
# act
57+
scheduler.enqueue_scheduled_jobs()
58+
59+
# assert
60+
self.assertFalse(registry.exists(connection, job_name))
61+
self.assertFalse(task.rqueue.queued_job_registry.exists(connection, job_name))
62+
63+
# act: the next pass schedules the task again
64+
scheduler.enqueue_scheduled_jobs()
65+
66+
# assert
67+
task.refresh_from_db()
68+
self.assertIsNotNone(task.job_name)
69+
self.assertNotEqual(job_name, task.job_name)
70+
self.assertTrue(registry.exists(connection, task.job_name))

scheduler/worker/scheduler.py

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -148,9 +148,11 @@ def enqueue_scheduled_jobs(self) -> None:
148148
queue = get_queue(registry.name)
149149
jobs = JobModel.get_many(job_names, connection=self.connection)
150150
with self.connection.pipeline() as pipeline:
151-
for job in jobs:
151+
for job_name, job in zip(job_names, jobs):
152152
if job is not None:
153153
queue.enqueue_job(job, pipeline=pipeline, at_front=job.at_front)
154+
else:
155+
registry.delete(connection=pipeline, job_name=job_name)
154156
pipeline.execute()
155157
self.status = SchedulerStatus.STARTED
156158

0 commit comments

Comments
 (0)