Coverage for extras/jobs.py: 24%
155 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 logging
2import traceback
3from contextlib import ExitStack
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
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
22from .utils import is_report
24__all__ = (
25 'CustomFieldDataJob',
26 'CustomFieldProvisioningJob',
27 'CustomFieldPurgeJob',
28 'RenderConfigContextJob',
29 'ScriptJob',
30 'provision_custom_field',
31 'purge_custom_field',
32)
35#
36# Config contexts
37#
39RENDER_CONFIG_CONTEXT_CHUNK_SIZE = 500
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
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.
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 """
59 class Meta:
60 name = 'Render config context'
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
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 )
98 self.logger.info(f"Rendered config context for {total} object(s)")
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))
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()
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
128 return rendered
131#
132# Custom fields
133#
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.
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.
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):
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
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)
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()
175 return bool(activated)
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.
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.
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):
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
201 custom_field.remove_stale_data(custom_field.object_types.all(), commit_per_batch=True)
202 custom_field._delete_row()
204 return True
207class CustomFieldDataJob(JobRunner):
208 """
209 Base class for the jobs which rewrite a custom field's stored data in bulk.
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 )
230class CustomFieldProvisioningJob(CustomFieldDataJob):
231 """
232 Populate the default value of a newly created custom field.
233 """
234 class Meta:
235 name = 'Custom Field Provisioning'
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")
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'
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")
258#
259# Scripts
260#
263class ScriptJob(JobRunner):
264 """
265 Script execution job.
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 """
271 class Meta:
272 name = 'Run Script'
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).
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 )
296 return super().enqueue(*args, **kwargs)
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).
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})")
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")
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}")
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}")
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.")
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
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()
368 def run(self, data, request=None, commit=True, **kwargs):
369 """
370 Run the script.
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}")
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
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}")
393 if commit:
394 self.logger.info("Executing script (commit enabled)")
395 else:
396 self.logger.warning("Executing script (commit disabled)")
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)