Coverage for extras/jobs.py: 24%

155 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-10 18:35 +0000

1import logging 

2import traceback 

3from contextlib import ExitStack 

4 

5from django.apps import apps 

6from django.db import DEFAULT_DB_ALIAS, router, transaction 

7from django.utils.translation import gettext as _ 

8from django_pg_utils import advisory_lock 

9 

10from core.signals import clear_events 

11from dcim.models import Device 

12from extras.choices import CustomFieldStatusChoices 

13from extras.constants import CUSTOMFIELD_JOB_TIMEOUT 

14from extras.models import CustomField 

15from extras.models import Script as ScriptModel 

16from extras.scripts import _UNSET 

17from netbox.context_managers import event_tracking 

18from netbox.jobs import JobRunner 

19from netbox.registry import registry 

20from utilities.exceptions import AbortScript, AbortTransaction 

21 

22from .utils import is_report 

23 

24__all__ = ( 

25 'CustomFieldDataJob', 

26 'CustomFieldProvisioningJob', 

27 'CustomFieldPurgeJob', 

28 'RenderConfigContextJob', 

29 'ScriptJob', 

30 'provision_custom_field', 

31 'purge_custom_field', 

32) 

33 

34 

35# 

36# Config contexts 

37# 

38 

39RENDER_CONFIG_CONTEXT_CHUNK_SIZE = 500 

40 

41# Safety bound on the number of re-scan passes performed by RenderConfigContextJob.run() (see the 

42# loop there). Each pass re-queries for NULL caches, so any finite burst of concurrent 

43# invalidations is drained well within this limit; the cap only guards against an object whose 

44# cache is being invalidated faster than it can be rendered (pathological, unbounded churn). 

45RENDER_CONFIG_CONTEXT_MAX_PASSES = 100 

46 

47 

48class RenderConfigContextJob(JobRunner): 

49 """ 

50 Recompute the pre-rendered `_config_context_data` cache for a set of Devices or 

51 VirtualMachines. Enqueued (coalesced) by the invalidation helpers in extras/cache.py whenever 

52 an upstream change (ConfigContext, related object, or the object itself) NULLs a cache. 

53 

54 This is *not* a recurring system job: the initial post-upgrade population is handled by the 

55 `rebuild_config_context_cache` management command, and steady-state freshness is maintained by 

56 the invalidation signals. 

57 """ 

58 

59 class Meta: 

60 name = 'Render config context' 

61 

62 def run(self, model_label=None, pks=None, **kwargs): 

63 """ 

64 Args: 

65 model_label: 'dcim.device' or 'virtualization.virtualmachine'. If None, both are processed. 

66 pks: An iterable of object PKs to refresh. If None, refresh all objects whose cache is null. 

67 """ 

68 labels = (model_label,) if model_label is not None else ('dcim.device', 'virtualization.virtualmachine') 

69 pks = list(pks) if pks is not None else None 

70 

71 # Re-scan until a full pass renders nothing. An invalidation that commits while this job is 

72 # already RUNNING coalesces into this job — JobRunner.enqueue_once() treats RUNNING as an 

73 # enqueued state — so it will NOT schedule a follow-up job. If we rendered in a single pass, 

74 # any cache NULLed after the iterator moved past its row (or after the pass for its model 

75 # completed) would be left populated by no one, stranding it on the on-demand read path 

76 # indefinitely. Looping until a pass finds no renderable NULL caches guarantees those late 

77 # invalidations are picked up before this job finishes. 

78 total = 0 

79 for _pass in range(RENDER_CONFIG_CONTEXT_MAX_PASSES): 

80 rendered = sum(self._render_for_model(label, pks=pks) for label in labels) 

81 total += rendered 

82 # No progress this pass means either nothing is NULL or the only NULL rows are churning 

83 # under concurrent invalidation (each such invalidation enqueues its own follow-up), so 

84 # there is nothing more for us to safely do. 

85 if not rendered: 

86 break 

87 else: 

88 # The loop ran every pass without ever rendering nothing, meaning caches are being 

89 # invalidated about as fast as we can render them. This is pathological churn worth 

90 # surfacing: each lingering invalidation enqueues its own follow-up job, so the caches 

91 # are not stranded, but the sustained rate warrants investigation. 

92 self.logger.warning( 

93 f"Reached the maximum of {RENDER_CONFIG_CONTEXT_MAX_PASSES} render passes with caches " 

94 f"still being invalidated; config context caches may be churning under sustained " 

95 f"concurrent invalidation." 

96 ) 

97 

98 self.logger.info(f"Rendered config context for {total} object(s)") 

99 

100 def _render_for_model(self, model_label, pks): 

101 """ 

102 Render and cache config context for every object of the given model whose cache is 

103 currently NULL (optionally restricted to `pks`). Returns the number of objects written. 

104 """ 

105 Model = apps.get_model(model_label) 

106 qs = Model.objects.filter(_config_context_data__isnull=True) 

107 if pks is not None: 

108 qs = qs.filter(pk__in=list(pks)) 

109 

110 # Annotate so each instance's render() uses the same aggregated subquery the on-demand 

111 # path would use, avoiding N additional queries. 

112 qs = qs.annotate_config_context_data() 

113 

114 rendered = 0 

115 for obj in qs.iterator(chunk_size=RENDER_CONFIG_CONTEXT_CHUNK_SIZE): 

116 # Capture the generation we rendered against, then write the result back only if no 

117 # invalidation has bumped it in the meantime (compare-and-set). If a fresh invalidation 

118 # won the race, the row stays NULL with a higher generation and the follow-up job it 

119 # enqueued will re-render it — we never persist a stale value. 

120 generation = obj._config_context_generation 

121 data = obj.render_config_context() 

122 updated = Model.objects.filter( 

123 pk=obj.pk, 

124 _config_context_generation=generation, 

125 ).update(_config_context_data=data) 

126 rendered += updated 

127 

128 return rendered 

129 

130 

131# 

132# Custom fields 

133# 

134 

135 

136def provision_custom_field(pk, object_type_pks): 

137 """ 

138 Populate a new custom field's default value across the objects of the given types, then bring 

139 the field live. Returns True if the field was brought live. 

140 

141 The backfill is committed in batches, so an interruption leaves the field provisioning with some 

142 of its objects already updated. Running again completes it. 

143 

144 Args: 

145 pk: The primary key of the CustomField to provision 

146 object_type_pks: The primary keys of the object types to provision. Named explicitly, as 

147 only the caller which deferred the work knows which of the field's assignments are the 

148 new ones. 

149 """ 

150 # Taken on the connection the field is written on, as CustomField.delete() takes it, so that 

151 # the two are actually exclusive of one another. 

152 using = router.db_for_write(CustomField) 

153 with advisory_lock(CustomField.data_lock_key(pk), using=using): 

154 

155 # Rechecked now that the lock is held: where two jobs were enqueued for the same field, 

156 # whichever arrived first has left it in a state the other no longer matches. 

157 custom_field = CustomField.objects.filter(pk=pk, status=CustomFieldStatusChoices.STATUS_PROVISIONING).first() 

158 if custom_field is None: 

159 return False 

160 

161 # Restricted to the field's current assignments: a type unassigned since the job was 

162 # enqueued must not be provisioned, its data having been removed by remove_data(). That 

163 # method refuses an unassignment while the field is being provisioned, so this covers only 

164 # a change made without it -- through the m2m table directly, which emits no signal. 

165 object_types = custom_field.object_types.filter(pk__in=object_type_pks) 

166 custom_field.populate_initial_data(object_types, commit_per_batch=True) 

167 

168 # Applied via the queryset so that bringing the field live does not record a change of its 

169 # own, and cannot trip the guard in CustomField.clean(). 

170 activated = CustomField.objects.filter( 

171 pk=pk, status=CustomFieldStatusChoices.STATUS_PROVISIONING 

172 ).update(status=CustomFieldStatusChoices.STATUS_ACTIVE) 

173 CustomField.objects.clear_cache() 

174 

175 return bool(activated) 

176 

177 

178def purge_custom_field(pk): 

179 """ 

180 Remove a deleted custom field's data from all applicable objects, then remove the field itself. 

181 Returns True if the field was purged. 

182 

183 The row is dropped only once its data is gone: until then it reserves the field's name against a 

184 new field which would otherwise inherit the orphaned values. The removal is committed in batches, 

185 so an interruption leaves data behind for a later run to finish removing. 

186 

187 Args: 

188 pk: The primary key of the CustomField to purge 

189 """ 

190 # Taken on the connection the field is written on, as CustomField.delete() takes it, so that 

191 # the two are actually exclusive of one another. 

192 using = router.db_for_write(CustomField) 

193 with advisory_lock(CustomField.data_lock_key(pk), using=using): 

194 

195 # Rechecked now that the lock is held: where two jobs were enqueued for the same field, 

196 # whichever arrived first has left it in a state the other no longer matches. 

197 custom_field = CustomField.objects.filter(pk=pk, status=CustomFieldStatusChoices.STATUS_DELETING).first() 

198 if custom_field is None: 

199 return False 

200 

201 custom_field.remove_stale_data(custom_field.object_types.all(), commit_per_batch=True) 

202 custom_field._delete_row() 

203 

204 return True 

205 

206 

207class CustomFieldDataJob(JobRunner): 

208 """ 

209 Base class for the jobs which rewrite a custom field's stored data in bulk. 

210 

211 The field is passed by primary key rather than assigned to the job as its object. Job.clean() 

212 permits only models with the jobs feature there, and granting CustomField that feature would 

213 give it a cascading relation to its jobs -- so the purge job, whose last act is to remove the 

214 row, would delete the record of its own execution as it ran. 

215 """ 

216 @classmethod 

217 def enqueue_for(cls, custom_field, **kwargs): 

218 """ 

219 Enqueue this job for the given custom field, naming the field in the job's name and raising 

220 its timeout from the default (see CUSTOMFIELD_JOB_TIMEOUT). 

221 """ 

222 return cls.enqueue( 

223 name=f'{cls.name}: {custom_field}', 

224 custom_field_pk=custom_field.pk, 

225 job_timeout=CUSTOMFIELD_JOB_TIMEOUT, 

226 **kwargs, 

227 ) 

228 

229 

230class CustomFieldProvisioningJob(CustomFieldDataJob): 

231 """ 

232 Populate the default value of a newly created custom field. 

233 """ 

234 class Meta: 

235 name = 'Custom Field Provisioning' 

236 

237 def run(self, custom_field_pk, *args, object_type_pks, **kwargs): 

238 if provision_custom_field(custom_field_pk, object_type_pks): 

239 self.logger.info("Custom field provisioned") 

240 else: 

241 self.logger.info("Custom field is no longer awaiting provisioning; skipping") 

242 

243 

244class CustomFieldPurgeJob(CustomFieldDataJob): 

245 """ 

246 Purge the stored data of a deleted custom field, then delete the field. 

247 """ 

248 class Meta: 

249 name = 'Custom Field Purge' 

250 

251 def run(self, custom_field_pk, *args, **kwargs): 

252 if purge_custom_field(custom_field_pk): 

253 self.logger.info("Custom field data purged") 

254 else: 

255 self.logger.info("Custom field is no longer awaiting deletion; skipping") 

256 

257 

258# 

259# Scripts 

260# 

261 

262 

263class ScriptJob(JobRunner): 

264 """ 

265 Script execution job. 

266 

267 A wrapper for calling Script.run(). This performs error handling and provides a hook for committing changes. It 

268 exists outside the Script class to ensure it cannot be overridden by a script author. 

269 """ 

270 

271 class Meta: 

272 name = 'Run Script' 

273 

274 @classmethod 

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

276 """ 

277 Validate the script's execution parameters before enqueueing. This is the single choke point through which 

278 every script execution passes (interactive runs, the REST API, the runscript command, event-rule actions, and 

279 recurring reschedules), so validating here surfaces a misconfigured script as an actionable error rather than 

280 an unhandled exception at enqueue time (see #22872). 

281 

282 The values actually being enqueued are validated, not just the script's Meta defaults, so an explicit 

283 job_timeout or notifications supplied by the caller is checked too. 

284 """ 

285 # The instance may be passed positionally (JobRunner.enqueue() forwards it to Job.enqueue()'s first argument) 

286 # or by keyword. Resolve it for validation without consuming it, so the original arguments are forwarded to 

287 # super() unchanged and the inherited calling contract is preserved. 

288 instance = args[0] if args else kwargs.get('instance') 

289 script_class = getattr(instance, 'python_class', None) 

290 if script_class is not None: 

291 script_class.validate_meta( 

292 job_timeout=kwargs.get('job_timeout', _UNSET), 

293 notifications=kwargs.get('notifications', _UNSET), 

294 ) 

295 

296 return super().enqueue(*args, **kwargs) 

297 

298 def run_script(self, script, request, data, commit): 

299 """ 

300 Core script execution task. We capture this within a method to allow for conditionally wrapping it with the 

301 event_tracking context manager (which is bypassed if commit == False). 

302 

303 Args: 

304 request: The WSGI request associated with this execution (if any) 

305 data: A dictionary of data to be passed to the script upon execution 

306 commit: Passed through to Script.run() 

307 """ 

308 logger = logging.getLogger(f"netbox.scripts.{script.full_name}") 

309 logger.info(f"Running script (commit={commit})") 

310 

311 try: 

312 try: 

313 # A script can modify multiple models so need to do an atomic lock on 

314 # both the default database (for non ChangeLogged models) and potentially 

315 # any other database (for ChangeLogged models) 

316 changeloged_db = router.db_for_write(Device) 

317 with transaction.atomic(using=DEFAULT_DB_ALIAS): 

318 # If branch database is different from default, wrap in a second atomic transaction 

319 # Note: Don't add any extra code between the two atomic transactions, 

320 # otherwise the changes might get committed to the default database 

321 # if there are any raised exceptions. 

322 if changeloged_db != DEFAULT_DB_ALIAS: 

323 with transaction.atomic(using=changeloged_db): 

324 script.output = script.run(data, commit) 

325 if not commit: 

326 raise AbortTransaction() 

327 else: 

328 script.output = script.run(data, commit) 

329 if not commit: 

330 raise AbortTransaction() 

331 except AbortTransaction: 

332 script.log_info(message=_("Database changes have been reverted automatically.")) 

333 if script.failed: 

334 logger.warning("Script failed") 

335 

336 except Exception as e: 

337 if type(e) is AbortScript: 

338 msg = _("Script aborted with error: ") + str(e) 

339 if is_report(type(script)): 

340 script.log_failure(message=msg) 

341 else: 

342 script.log_failure(msg) 

343 logger.error(f"Script aborted with error: {e}") 

344 self.logger.error(f"Script aborted with error: {e}") 

345 

346 else: 

347 stacktrace = traceback.format_exc() 

348 script.log_failure( 

349 message=_("An exception occurred: ") + f"`{type(e).__name__}: {e}`\n```\n{stacktrace}\n```" 

350 ) 

351 logger.error(f"Exception raised during script execution: {e}") 

352 self.logger.error(f"Exception raised during script execution: {e}") 

353 

354 if type(e) is not AbortTransaction: 

355 script.log_info(message=_("Database changes have been reverted due to error.")) 

356 self.logger.info("Database changes have been reverted due to error.") 

357 

358 # Clear all pending events. Job termination (including setting the status) is handled by the job framework. 

359 if request: 

360 clear_events.send(request) 

361 raise 

362 

363 # Update the job data regardless of the execution status of the job. Successes should be reported as well as 

364 # failures. 

365 finally: 

366 self.job.data = script.get_job_data() 

367 

368 def run(self, data, request=None, commit=True, **kwargs): 

369 """ 

370 Run the script. 

371 

372 Args: 

373 job: The Job associated with this execution 

374 data: A dictionary of data to be passed to the script upon execution 

375 request: The WSGI request associated with this execution (if any) 

376 commit: Passed through to Script.run() 

377 """ 

378 script_model = ScriptModel.objects.get(pk=self.job.object_id) 

379 self.logger.debug(f"Found ScriptModel ID {script_model.pk}") 

380 script = script_model.python_class() 

381 self.logger.debug(f"Loaded script {script.full_name}") 

382 

383 # Add files to form data 

384 if request: 

385 files = request.FILES 

386 for field_name, fileobj in files.items(): 

387 data[field_name] = fileobj 

388 

389 # Add the current request as a property of the script 

390 script.request = request 

391 self.logger.debug(f"Request ID: {request.id if request else None}") 

392 

393 if commit: 

394 self.logger.info("Executing script (commit enabled)") 

395 else: 

396 self.logger.warning("Executing script (commit disabled)") 

397 

398 with ExitStack() as stack: 

399 for request_processor in registry['request_processors']: 

400 if not commit and request_processor is event_tracking: 

401 continue 

402 stack.enter_context(request_processor(request)) 

403 self.run_script(script, request, data, commit)