diff --git a/netbox/netbox/jobs.py b/netbox/netbox/jobs.py index 3203adea1..e011d8bbf 100644 --- a/netbox/netbox/jobs.py +++ b/netbox/netbox/jobs.py @@ -5,6 +5,7 @@ from abc import ABC, abstractmethod from datetime import timedelta from pathlib import Path +from django.conf import settings from django.core.exceptions import ImproperlyConfigured from django.utils import timezone from django.utils.functional import classproperty @@ -22,6 +23,7 @@ __all__ = ( 'AsyncViewJob', 'JobRunner', 'reconcile_stale_system_jobs', + 'resolve_job_timeout', 'system_job', ) @@ -30,10 +32,11 @@ __all__ = ( # jobs.py lives at /netbox/netbox/jobs.py, so parents[2] is the root. _INSTALL_ROOT = str(Path(__file__).resolve().parents[2]) + os.sep -# Floor for the grace period before a "running" system job is treated as stranded. The window -# is max(interval, this floor), so a short-interval job still gets time to finish a legitimate -# run. See reconcile_stale_system_jobs() and issue #22714. -STALE_RUNNING_JOB_GRACE_MINUTES = 30 +# Margin added to a job's run timeout before a "running" system job is treated as stranded. A +# job still running this long past its RQ timeout can no longer be executing legitimately (its +# worker was killed), so the window is keyed on run duration rather than the recurrence interval. +# See reconcile_stale_system_jobs() and issue #22714. +STALE_RUNNING_JOB_GRACE_SECONDS = 600 def system_job(interval): @@ -52,6 +55,17 @@ def system_job(interval): return _wrapper +def resolve_job_timeout(job_class, job): + """ + Return the RQ run timeout (in seconds) a `JobRunner` runs under, falling back to + `RQ_DEFAULT_TIMEOUT` when it declares none. There is no single timeout contract across job + types: scripts expose `Meta.job_timeout` (class-level), some plugin runners expose an instance + `@property`, and the base class declares nothing. Instantiate the runner so an instance property + resolves correctly regardless of what it reads. + """ + return getattr(job_class(job), 'job_timeout', None) or settings.RQ_DEFAULT_TIMEOUT + + @advisory_lock(ADVISORY_LOCK_KEYS['job-schedules']) def reconcile_stale_system_jobs(job_class, interval): """ @@ -59,21 +73,31 @@ def reconcile_stale_system_jobs(job_class, interval): that was killed mid-run (issue #22714). Such a row is never reset, and because "running" is an enqueued state, `enqueue_once()` mistakes it for a live schedule and never re-arms it. - A running job is treated as stranded once its `started` timestamp is older than - `max(interval, STALE_RUNNING_JOB_GRACE_MINUTES)`. RQ is not consulted: a killed job's RQ - entry can outlive the worker, and the schedule may be recovered with RQ state wiped. + A running job is treated as stranded once its `started` timestamp is older than the job's own + RQ run timeout plus a margin: past that point the run can no longer be executing legitimately. + Keying on run duration rather than the recurrence `interval` lets recovery fire on the common + case (a worker killed and restarted moments later) while still tolerating a legitimately long + run that declares a large timeout. RQ is not consulted: a killed job's RQ entry can outlive the + worker, and the schedule may be recovered with RQ state wiped. """ - grace = max(interval, STALE_RUNNING_JOB_GRACE_MINUTES) - cutoff = timezone.now() - timedelta(minutes=grace) - - orphaned = Job.objects.filter( - name=job_class.name, - object_id__isnull=True, - interval=interval, - status=JobStatusChoices.STATUS_RUNNING, - started__lte=cutoff, + running = list( + Job.objects.filter( + name=job_class.name, + object_id__isnull=True, + interval=interval, + status=JobStatusChoices.STATUS_RUNNING, + ) ) - for job in orphaned: + if not running: + return + + # The timeout is a property of the runner, not the individual job, so resolve it once. + grace = resolve_job_timeout(job_class, running[0]) + STALE_RUNNING_JOB_GRACE_SECONDS + cutoff = timezone.now() - timedelta(seconds=grace) + + for job in running: + if not job.started or job.started > cutoff: + continue # STATUS_ERRORED (not FAILED) records an unexpected fault rather than a self-declared # failure, as handle() does for an unhandled exception. For an object-less, userless # system job, terminate() sends no notification and triggers no event rule. diff --git a/netbox/netbox/tests/test_jobs.py b/netbox/netbox/tests/test_jobs.py index 8577f5bfb..99a79c7f7 100644 --- a/netbox/netbox/tests/test_jobs.py +++ b/netbox/netbox/tests/test_jobs.py @@ -2,17 +2,18 @@ import uuid from datetime import timedelta from unittest.mock import patch +from django.conf import settings from django.test import TestCase from django.utils import timezone -from core.choices import JobStatusChoices +from core.choices import JobIntervalChoices, JobStatusChoices from core.exceptions import JobFailed from core.models import DataSource, Job from utilities.testing import disable_warnings from utilities.testing.mixins import RQQueueTestMixin from ..jobs import * -from ..jobs import _INSTALL_ROOT +from ..jobs import _INSTALL_ROOT, STALE_RUNNING_JOB_GRACE_SECONDS class TestJobRunner(JobRunner): @@ -33,8 +34,18 @@ class TestSystemJobRunner(JobRunner): pass -@system_job(interval=1) -class TestShortIntervalSystemJobRunner(JobRunner): +class TestClassTimeoutJobRunner(JobRunner): + job_timeout = 3600 + + def run(self, *args, **kwargs): + pass + + +class TestPropertyTimeoutJobRunner(JobRunner): + # Mirrors plugins (e.g. netbox-branching) that expose job_timeout as an instance property. + @property + def job_timeout(self): + return 7200 def run(self, *args, **kwargs): pass @@ -394,7 +405,7 @@ class ReconcileStaleJobsTestCase(BaseJobRunnerTestCase): """ def test_reconcile_terminates_orphaned_running_job(self): - """A running system job whose `started` is older than the grace window is an + """A running system job whose `started` predates the run timeout window is an orphan (its worker died) and must be moved to `errored`.""" orphan = Job.objects.create( name=TestSystemJobRunner.name, @@ -413,8 +424,8 @@ class ReconcileStaleJobsTestCase(BaseJobRunnerTestCase): self.assertEqual(orphan.error, "Worker terminated before job completed") def test_reconcile_preserves_recently_started_running_job(self): - """A running system job that started recently (within the grace window) is a - legitimately in-flight job and must NOT be reaped.""" + """A running system job that started recently (within the run timeout window) is a + legitimately in-flight job and must NOT be reaped, regardless of its interval.""" live = Job.objects.create( name=TestSystemJobRunner.name, status=JobStatusChoices.STATUS_RUNNING, @@ -429,59 +440,43 @@ class ReconcileStaleJobsTestCase(BaseJobRunnerTestCase): live.refresh_from_db() self.assertEqual(live.status, JobStatusChoices.STATUS_RUNNING) - def test_reconcile_grace_window_scales_with_interval(self): - """The grace window is max(interval, floor). A job whose interval exceeds the floor - gets the longer window: a 60-minute-interval job started 45 minutes ago is past the - 30-minute floor but still within its own interval, so it must be preserved.""" - within_interval = Job.objects.create( + def test_reconcile_window_is_run_timeout_not_interval(self): + """The window is keyed on the RQ run timeout, not the recurrence interval, so a + long-interval job whose worker was just killed is recovered promptly. A daily job + started well past its run timeout (but far short of a day) must be reaped.""" + grace = settings.RQ_DEFAULT_TIMEOUT + STALE_RUNNING_JOB_GRACE_SECONDS + orphan = Job.objects.create( name=TestSystemJobRunner.name, status=JobStatusChoices.STATUS_RUNNING, - interval=60, - scheduled=timezone.now() - timedelta(minutes=45), - started=timezone.now() - timedelta(minutes=45), + interval=JobIntervalChoices.INTERVAL_DAILY, + scheduled=timezone.now() - timedelta(seconds=grace + 60), + started=timezone.now() - timedelta(seconds=grace + 60), job_id=uuid.uuid4(), ) - reconcile_stale_system_jobs(TestSystemJobRunner, 60) - - within_interval.refresh_from_db() - self.assertEqual(within_interval.status, JobStatusChoices.STATUS_RUNNING) - - def test_reconcile_grace_floor_protects_short_interval_job(self): - """For a job whose interval is shorter than the floor, the floor governs the window. - A 1-minute-interval job started 20 minutes ago is well past its interval but within - the 30-minute floor, so it must NOT be reaped.""" - within_floor = Job.objects.create( - name=TestShortIntervalSystemJobRunner.name, - status=JobStatusChoices.STATUS_RUNNING, - interval=1, - scheduled=timezone.now() - timedelta(minutes=20), - started=timezone.now() - timedelta(minutes=20), - job_id=uuid.uuid4(), - ) - - reconcile_stale_system_jobs(TestShortIntervalSystemJobRunner, 1) - - within_floor.refresh_from_db() - self.assertEqual(within_floor.status, JobStatusChoices.STATUS_RUNNING) - - def test_reconcile_grace_floor_reaps_short_interval_orphan(self): - """A 1-minute-interval job started 40 minutes ago is past the 30-minute floor and is - an orphan, so it must be reaped despite its short interval.""" - orphan = Job.objects.create( - name=TestShortIntervalSystemJobRunner.name, - status=JobStatusChoices.STATUS_RUNNING, - interval=1, - scheduled=timezone.now() - timedelta(minutes=40), - started=timezone.now() - timedelta(minutes=40), - job_id=uuid.uuid4(), - ) - - reconcile_stale_system_jobs(TestShortIntervalSystemJobRunner, 1) + reconcile_stale_system_jobs(TestSystemJobRunner, JobIntervalChoices.INTERVAL_DAILY) orphan.refresh_from_db() self.assertEqual(orphan.status, JobStatusChoices.STATUS_ERRORED) + def test_reconcile_preserves_job_within_run_timeout(self): + """A job started just inside the run timeout window is still legitimately running + and must be preserved.""" + grace = settings.RQ_DEFAULT_TIMEOUT + STALE_RUNNING_JOB_GRACE_SECONDS + live = Job.objects.create( + name=TestSystemJobRunner.name, + status=JobStatusChoices.STATUS_RUNNING, + interval=JobIntervalChoices.INTERVAL_DAILY, + scheduled=timezone.now() - timedelta(seconds=grace - 60), + started=timezone.now() - timedelta(seconds=grace - 60), + job_id=uuid.uuid4(), + ) + + reconcile_stale_system_jobs(TestSystemJobRunner, JobIntervalChoices.INTERVAL_DAILY) + + live.refresh_from_db() + self.assertEqual(live.status, JobStatusChoices.STATUS_RUNNING) + def test_reconcile_ignores_instance_bound_job(self): """The sweep targets object-less system jobs only. An instance-bound job of the same runner class must be left alone even if it looks stale.""" @@ -529,3 +524,43 @@ class ReconcileStaleJobsTestCase(BaseJobRunnerTestCase): status__in=JobStatusChoices.ENQUEUED_STATE_CHOICES, ) self.assertEqual(enqueued.count(), 1) + + def test_reconcile_honors_longer_per_job_timeout(self): + """A runner that declares a large `job_timeout` gets a proportionally longer window. A + job started past the default window but within its own declared timeout must be preserved, + so a legitimately long run (e.g. a bulk branch archival) is not reaped.""" + interval = JobIntervalChoices.INTERVAL_DAILY + default_window = settings.RQ_DEFAULT_TIMEOUT + STALE_RUNNING_JOB_GRACE_SECONDS + # Older than the default window, but well within this runner's 3600s job_timeout. + started = timezone.now() - timedelta(seconds=default_window + 120) + live = Job.objects.create( + name=TestClassTimeoutJobRunner.name, + status=JobStatusChoices.STATUS_RUNNING, + interval=interval, + scheduled=started, + started=started, + job_id=uuid.uuid4(), + ) + + reconcile_stale_system_jobs(TestClassTimeoutJobRunner, interval) + + live.refresh_from_db() + self.assertEqual(live.status, JobStatusChoices.STATUS_RUNNING) + + +class ResolveJobTimeoutTestCase(BaseJobRunnerTestCase): + """ + Test resolution of a runner's RQ run timeout across the ad-hoc conventions (#22714). + """ + + def test_falls_back_to_rq_default(self): + job = Job(name=TestSystemJobRunner.name) + self.assertEqual(resolve_job_timeout(TestSystemJobRunner, job), settings.RQ_DEFAULT_TIMEOUT) + + def test_reads_class_attribute(self): + job = Job(name=TestClassTimeoutJobRunner.name) + self.assertEqual(resolve_job_timeout(TestClassTimeoutJobRunner, job), 3600) + + def test_reads_instance_property(self): + job = Job(name=TestPropertyTimeoutJobRunner.name) + self.assertEqual(resolve_job_timeout(TestPropertyTimeoutJobRunner, job), 7200)