Coverage for netbox/jobs.py: 36%
180 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 18:35 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 18:35 +0000
1import json
2import logging
3import os
4import traceback
5from abc import ABC, abstractmethod
6from datetime import timedelta
7from io import BytesIO
8from pathlib import Path
10from django.contrib.auth import get_user_model
11from django.core.exceptions import ImproperlyConfigured, PermissionDenied, ValidationError
12from django.core.handlers.wsgi import WSGIRequest
13from django.db.models import ProtectedError, RestrictedError
14from django.http import Http404
15from django.utils import timezone
16from django.utils.functional import classproperty
17from django.utils.module_loading import import_string
18from django.utils.translation import gettext_lazy as _
19from django_pg_utils import advisory_lock
20from rest_framework.exceptions import APIException
21from rq.timeouts import JobTimeoutException
23from core.choices import JobStatusChoices
24from core.exceptions import JobFailed
25from core.models import Job, ObjectType
26from netbox.constants import ADVISORY_LOCK_KEYS
27from netbox.registry import registry
28from utilities.exceptions import AbortRequest
29from utilities.request import apply_request_processors
31__all__ = (
32 'AsyncAPIJob',
33 'AsyncViewJob',
34 'JobRunner',
35 'system_job',
36)
38# The installation root, e.g. "/opt/netbox/". Used to strip absolute path
39# prefixes from traceback file paths before recording them in the job log.
40# jobs.py lives at <root>/netbox/netbox/jobs.py, so parents[2] is the root.
41_INSTALL_ROOT = str(Path(__file__).resolve().parents[2]) + os.sep
44def system_job(interval):
45 """
46 Decorator for registering a `JobRunner` class as system background job.
47 """
48 if type(interval) is not int: 48 ↛ 49line 48 didn't jump to line 49 because the condition on line 48 was never true
49 raise ImproperlyConfigured("System job interval must be an integer (minutes).")
51 def _wrapper(cls):
52 registry['system_jobs'][cls] = {
53 'interval': interval
54 }
55 return cls
57 return _wrapper
60class JobLogHandler(logging.Handler):
61 """
62 A logging handler which records entries on a Job.
63 """
64 def __init__(self, job, *args, **kwargs):
65 super().__init__(*args, **kwargs)
66 self.job = job
68 def emit(self, record):
69 # Enter the record in the log of the associated Job
70 self.job.log(record)
73class JobRunner(ABC):
74 """
75 Background Job helper class.
77 This class handles the execution of a background job. It is responsible for maintaining its state, reporting errors,
78 and scheduling recurring jobs.
79 """
81 class Meta:
82 pass
84 def __init__(self, job):
85 """
86 Args:
87 job: The specific `Job` this `JobRunner` is executing.
88 """
89 self.job = job
91 # Initiate the system logger
92 self.logger = logging.getLogger(f"netbox.jobs.{self.__class__.__name__}")
93 self.logger.setLevel(logging.DEBUG)
94 self.logger.addHandler(JobLogHandler(job))
96 @classproperty
97 def name(cls):
98 return getattr(cls.Meta, 'name', cls.__name__)
100 @abstractmethod
101 def run(self, *args, **kwargs):
102 """
103 Run the job.
105 A `JobRunner` class needs to implement this method to execute all commands of the job.
106 """
107 pass
109 @classmethod
110 def handle(cls, job, *args, **kwargs):
111 """
112 Handle the execution of a `Job`.
114 This method is called by the Job Scheduler to handle the execution of all job commands. It will maintain the
115 job's metadata and handle errors. For periodic jobs, a new job is automatically scheduled using its `interval`.
116 """
117 logger = logging.getLogger('netbox.jobs')
119 try:
120 job.start()
121 cls(job).run(*args, **kwargs)
122 job.terminate()
124 except JobFailed:
125 logger.warning(f"Job {job} failed")
126 job.terminate(status=JobStatusChoices.STATUS_FAILED)
128 except Exception as e:
129 tb_str = traceback.format_exc().replace(_INSTALL_ROOT, '')
130 tb_record = logging.makeLogRecord({
131 'levelno': logging.ERROR,
132 'levelname': 'ERROR',
133 'msg': tb_str,
134 })
135 job.log(tb_record)
136 job.terminate(status=JobStatusChoices.STATUS_ERRORED, error=repr(e))
137 if type(e) is JobTimeoutException:
138 logger.error(e)
140 # If the executed job is a periodic job, schedule its next execution at the specified interval.
141 finally:
142 if job.interval:
143 # Determine the new scheduled time. Cannot be earlier than one minute in the future.
144 new_scheduled_time = max(
145 (job.scheduled or job.started) + timedelta(minutes=job.interval),
146 timezone.now() + timedelta(minutes=1)
147 )
148 if job.object and getattr(job.object, "python_class", None):
149 kwargs["job_timeout"] = job.object.python_class.job_timeout
151 enqueue_kwargs = dict(
152 instance=job.object,
153 name=job.name,
154 user=job.user,
155 schedule_at=new_scheduled_time,
156 interval=job.interval,
157 notifications=job.notifications,
158 **kwargs,
159 )
161 # Reschedule the next occurrence. If the object's configuration has become invalid since this run was
162 # scheduled (e.g. a script's Meta.job_timeout was edited to an invalid value, see #22872), the enqueue
163 # will raise a ValidationError. Record it on this job and decline to reschedule rather than allowing an
164 # unhandled exception to escape the worker's finally block.
165 try:
166 if cls in registry['system_jobs']:
167 # System jobs are also scheduled by `enqueue_once()` at worker startup,
168 # which races with this finally block and can produce duplicate schedules
169 # (see #22232). Acquire the same advisory lock used by `enqueue_once()`
170 # and skip rescheduling if a successor is already enqueued.
171 #
172 # This branch is limited to system jobs because generic recurring jobs
173 # (e.g. scheduled scripts) may have multiple legitimate schedules sharing
174 # the same runner/object/interval but differing in their runtime kwargs.
175 with advisory_lock(ADVISORY_LOCK_KEYS['job-schedules']):
176 successor_exists = Job.objects.filter(
177 name=cls.name,
178 object_id__isnull=True,
179 status__in=JobStatusChoices.ENQUEUED_STATE_CHOICES,
180 interval=job.interval,
181 ).exclude(pk=job.pk).exists()
182 if not successor_exists:
183 cls.enqueue(**enqueue_kwargs)
184 else:
185 cls.enqueue(**enqueue_kwargs)
186 except ValidationError as e:
187 # The successor could not be scheduled because the object's configuration is now invalid. Record
188 # this against the (already-terminated) job without overwriting the outcome of the run that just
189 # completed — re-running terminate() here would clobber a successful run's status and fire a
190 # duplicate notification (see #22872).
191 error = _("Recurring job not rescheduled due to invalid configuration: {error}").format(
192 error='; '.join(e.messages)
193 )
194 logger.error(f"Job {job}: {error}")
195 job.log(logging.makeLogRecord({
196 'levelno': logging.ERROR,
197 'levelname': 'ERROR',
198 'msg': error,
199 }))
200 job.save()
202 @classmethod
203 def get_jobs(cls, instance=None):
204 """
205 Get all jobs of this `JobRunner` related to a specific instance.
206 """
207 jobs = Job.objects.filter(name=cls.name)
209 if instance: 209 ↛ 210line 209 didn't jump to line 210 because the condition on line 209 was never true
210 object_type = ObjectType.objects.get_for_model(instance, for_concrete_model=False)
211 jobs = jobs.filter(
212 object_type=object_type,
213 object_id=instance.pk,
214 )
216 return jobs
218 @classmethod
219 def enqueue(cls, *args, **kwargs):
220 """
221 Enqueue a new `Job`.
223 This method is a wrapper of `Job.enqueue()` using `handle()` as function callback. See its documentation for
224 parameters.
225 """
226 name = kwargs.pop('name', None) or cls.name
227 return Job.enqueue(cls.handle, name=name, *args, **kwargs)
229 @classmethod
230 @advisory_lock(ADVISORY_LOCK_KEYS['job-schedules'])
231 def enqueue_once(cls, instance=None, schedule_at=None, interval=None, *args, **kwargs):
232 """
233 Enqueue a new `Job` once, i.e. skip duplicate jobs.
235 Like `enqueue()`, this method adds a new `Job` to the job queue. However, if there's already a job of this
236 class scheduled for `instance`, the existing job will be updated if necessary. This ensures that a particular
237 schedule is only set up once at any given time, i.e. multiple calls to this method are idempotent.
239 Note that this does not forbid running additional jobs with the `enqueue()` method, e.g. to schedule an
240 immediate synchronization job in addition to a periodic synchronization schedule.
242 For additional parameters see `enqueue()`.
244 Args:
245 instance: The NetBox object to which this job pertains (optional)
246 schedule_at: Schedule the job to be executed at the passed date and time
247 interval: Recurrence interval (in minutes)
248 """
249 job = cls.get_jobs(instance).filter(status__in=JobStatusChoices.ENQUEUED_STATE_CHOICES).first()
250 if job: 250 ↛ 253line 250 didn't jump to line 253 because the condition on line 250 was never true
251 # If the job parameters haven't changed, don't schedule a new job and keep the current schedule. Otherwise,
252 # delete the existing job and schedule a new job instead.
253 if (not schedule_at or job.scheduled == schedule_at) and (job.interval == interval):
254 return job
255 job.delete()
257 return cls.enqueue(instance=instance, schedule_at=schedule_at, interval=interval, *args, **kwargs)
260class AsyncViewJob(JobRunner):
261 """
262 Execute a view as a background job.
263 """
264 class Meta:
265 name = 'Async View'
267 def run(self, view_cls, request, **kwargs):
268 view = view_cls.as_view()
269 request.job = self
271 # Apply all registered request processors (e.g. event_tracking)
272 with apply_request_processors(request):
273 view(request)
275 if self.job.error:
276 raise JobFailed()
279class AsyncAPIJob(JobRunner):
280 """
281 Execute a REST API bulk write (create/update/delete) as a background job.
283 The viewset's action method is re-invoked inside the worker against a reconstructed
284 request, so the synchronous and background code paths are identical (validation,
285 transaction semantics, object permissions, and change logging all behave the same).
286 The action's serialized Response is captured into the job's data.
287 """
288 class Meta:
289 name = 'Async API Request'
291 @staticmethod
292 def _build_request(request_copy, payload, scheme):
293 """
294 Reconstruct a real WSGIRequest from a copy_safe_request() snapshot, injecting the JSON
295 payload as the request body. DRF's Request wrapper requires a real HttpRequest to parse
296 the body, and the snapshot's host metadata (already correctly separated into
297 SERVER_NAME/SERVER_PORT/HTTP_HOST by the original WSGI layer) is carried verbatim so that
298 absolute URLs in the captured result (serializer hyperlink fields) point at the real
299 server. The scheme is applied separately, as copy_safe_request() does not capture it.
300 """
301 body = json.dumps(payload).encode('utf-8')
302 environ = {
303 'REQUEST_METHOD': request_copy.method,
304 'PATH_INFO': getattr(request_copy, 'path', '/') or '/',
305 'CONTENT_TYPE': 'application/json',
306 'CONTENT_LENGTH': str(len(body)),
307 'wsgi.input': BytesIO(body),
308 'wsgi.url_scheme': scheme,
309 'SERVER_PROTOCOL': 'HTTP/1.1',
310 # Sensible defaults (a real WSGI layer always sets these); overridden below by the
311 # snapshot's host metadata when present.
312 'SERVER_NAME': 'localhost',
313 'SERVER_PORT': '443' if scheme == 'https' else '80',
314 }
315 # Carry the host/forwarding metadata from the safe request copy (no host:port parsing
316 # needed: these were already split correctly when the original request was received).
317 for key in (
318 'HTTP_HOST', 'SERVER_NAME', 'SERVER_PORT',
319 'HTTP_X_FORWARDED_HOST', 'HTTP_X_FORWARDED_PORT', 'HTTP_X_FORWARDED_PROTO',
320 ):
321 if value := request_copy.META.get(key):
322 environ[key] = value
324 request = WSGIRequest(environ)
325 request.id = getattr(request_copy, 'id', None)
326 return request
328 def run(
329 self, viewset_class, action, payload, user_pk, request,
330 action_kwargs=None, scheme='http', **kwargs
331 ):
332 # Imported here to avoid a circular import (netbox.api.viewsets imports from this module).
333 from netbox.api.viewsets import HTTP_ACTIONS
335 action_kwargs = action_kwargs or {}
336 method = request.method
337 request_id = getattr(request, 'id', '') or ''
338 viewset_class = import_string(viewset_class)
340 # Re-fetch the requesting user. If the user no longer exists or is inactive, fail
341 # the job rather than running with stale identity (execution-time identity check).
342 User = get_user_model()
343 try:
344 user = User.objects.get(pk=user_pk)
345 except User.DoesNotExist:
346 self.job.error = "The requesting user no longer exists."
347 self.job.save()
348 raise JobFailed()
349 if not user.is_active:
350 self.job.error = "The requesting user is no longer active."
351 self.job.save()
352 raise JobFailed()
354 django_request = self._build_request(request, payload, scheme)
356 # Instantiate the viewset and apply the minimal scaffolding that DRF's as_view()
357 # normally sets, so initialize_request() can wire up parsers, etc.
358 viewset = viewset_class()
359 viewset.action_map = {method.lower(): action}
360 viewset.kwargs = {}
361 viewset.args = ()
362 viewset.action = action
363 viewset.format_kwarg = None
365 drf_request = viewset.initialize_request(django_request)
366 # Carry the authenticated user forward; we do not re-authenticate in the worker.
367 # Setting .user populates the request's user cache, so DRF never lazily invokes
368 # authentication (and nothing on the action path reads the authenticator/auth).
369 drf_request.user = user
370 drf_request.id = request_id
371 viewset.request = drf_request
373 # Re-apply object-level permission restriction exactly as BaseViewSet.initial() does.
374 if perm_action := HTTP_ACTIONS[method.upper()]:
375 viewset.queryset = viewset.queryset.restrict(user, perm_action)
377 # Execute the action method within the registered request processors so change
378 # logging and event rules fire (and are attributed to the original request_id).
379 #
380 # The synchronous path relies on NetBoxModelViewSet.dispatch() and DRF's
381 # handle_exception() to translate exceptions into HTTP responses. Because we invoke
382 # the action method directly (bypassing dispatch), we reproduce that translation here
383 # so the captured result matches what the synchronous API would have returned:
384 # - APIException (incl. ValidationError, PermissionDenied, Http404) -> handle_exception()
385 # - AbortRequest / ProtectedError / RestrictedError -> exception_to_response()
386 with apply_request_processors(drf_request):
387 try:
388 # A job queued before the endpoint opted out carries no query parameter to catch.
389 # Optional because run() accepts any importable viewset, not only mixin subclasses.
390 if check_background := getattr(viewset, 'check_background_enabled', None):
391 check_background()
392 response = getattr(viewset, action)(drf_request, **action_kwargs)
393 except (APIException, Http404, PermissionDenied) as e:
394 response = viewset.handle_exception(e)
395 except (AbortRequest, ProtectedError, RestrictedError) as e:
396 response = viewset.exception_to_response(e)
397 if response is None:
398 raise
400 # Capture the action's result for the polling client, in the same shape for both
401 # success and failure.
402 self.job.data = {
403 'status_code': response.status_code,
404 'data': response.data,
405 }
407 if response.status_code >= 400:
408 # A handled rejection (4xx), not a worker crash: record a concise summary and
409 # mark the job failed (JobRunner.handle reserves "errored" for unhandled crashes).
410 detail = response.data.get('detail') if isinstance(response.data, dict) else None
411 self.job.error = str(detail) if detail else f"Request failed with status {response.status_code}."
412 self.job.save()
413 raise JobFailed()
415 # On success, job.data is persisted by JobRunner.handle() -> job.terminate().