Coverage for netbox/jobs.py: 36%

180 statements  

« 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 

9 

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 

22 

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 

30 

31__all__ = ( 

32 'AsyncAPIJob', 

33 'AsyncViewJob', 

34 'JobRunner', 

35 'system_job', 

36) 

37 

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 

42 

43 

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).") 

50 

51 def _wrapper(cls): 

52 registry['system_jobs'][cls] = { 

53 'interval': interval 

54 } 

55 return cls 

56 

57 return _wrapper 

58 

59 

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 

67 

68 def emit(self, record): 

69 # Enter the record in the log of the associated Job 

70 self.job.log(record) 

71 

72 

73class JobRunner(ABC): 

74 """ 

75 Background Job helper class. 

76 

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 """ 

80 

81 class Meta: 

82 pass 

83 

84 def __init__(self, job): 

85 """ 

86 Args: 

87 job: The specific `Job` this `JobRunner` is executing. 

88 """ 

89 self.job = job 

90 

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)) 

95 

96 @classproperty 

97 def name(cls): 

98 return getattr(cls.Meta, 'name', cls.__name__) 

99 

100 @abstractmethod 

101 def run(self, *args, **kwargs): 

102 """ 

103 Run the job. 

104 

105 A `JobRunner` class needs to implement this method to execute all commands of the job. 

106 """ 

107 pass 

108 

109 @classmethod 

110 def handle(cls, job, *args, **kwargs): 

111 """ 

112 Handle the execution of a `Job`. 

113 

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') 

118 

119 try: 

120 job.start() 

121 cls(job).run(*args, **kwargs) 

122 job.terminate() 

123 

124 except JobFailed: 

125 logger.warning(f"Job {job} failed") 

126 job.terminate(status=JobStatusChoices.STATUS_FAILED) 

127 

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) 

139 

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 

150 

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 ) 

160 

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() 

201 

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) 

208 

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 ) 

215 

216 return jobs 

217 

218 @classmethod 

219 def enqueue(cls, *args, **kwargs): 

220 """ 

221 Enqueue a new `Job`. 

222 

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) 

228 

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. 

234 

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. 

238 

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. 

241 

242 For additional parameters see `enqueue()`. 

243 

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() 

256 

257 return cls.enqueue(instance=instance, schedule_at=schedule_at, interval=interval, *args, **kwargs) 

258 

259 

260class AsyncViewJob(JobRunner): 

261 """ 

262 Execute a view as a background job. 

263 """ 

264 class Meta: 

265 name = 'Async View' 

266 

267 def run(self, view_cls, request, **kwargs): 

268 view = view_cls.as_view() 

269 request.job = self 

270 

271 # Apply all registered request processors (e.g. event_tracking) 

272 with apply_request_processors(request): 

273 view(request) 

274 

275 if self.job.error: 

276 raise JobFailed() 

277 

278 

279class AsyncAPIJob(JobRunner): 

280 """ 

281 Execute a REST API bulk write (create/update/delete) as a background job. 

282 

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' 

290 

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 

323 

324 request = WSGIRequest(environ) 

325 request.id = getattr(request_copy, 'id', None) 

326 return request 

327 

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 

334 

335 action_kwargs = action_kwargs or {} 

336 method = request.method 

337 request_id = getattr(request, 'id', '') or '' 

338 viewset_class = import_string(viewset_class) 

339 

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() 

353 

354 django_request = self._build_request(request, payload, scheme) 

355 

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 

364 

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 

372 

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) 

376 

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 

399 

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 } 

406 

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() 

414 

415 # On success, job.data is persisted by JobRunner.handle() -> job.terminate().