Coverage for core/signals.py: 64%
190 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
2from threading import local
4from django.contrib.contenttypes.models import ContentType
5from django.core.exceptions import ObjectDoesNotExist, ValidationError
6from django.core.signals import request_finished
7from django.db import transaction
8from django.db.models import CASCADE, RESTRICT
9from django.db.models.fields.reverse_related import ManyToManyRel, ManyToOneRel
10from django.db.models.signals import m2m_changed, post_migrate, post_save, pre_delete
11from django.dispatch import Signal, receiver
12from django.utils.translation import gettext_lazy as _
13from django.utils.translation import ngettext
14from django_prometheus.models import model_deletes, model_inserts, model_updates
15from rq.timeouts import JobTimeoutException
17from core.choices import JobStatusChoices, ObjectChangeActionChoices
18from core.events import *
19from core.exceptions import SyncError
20from core.models import ObjectType
21from extras.events import enqueue_event
22from extras.models import Tag
23from extras.utils import run_validators
24from netbox.config import get_config
25from netbox.context import current_request, events_queue
26from netbox.models.features import ChangeLoggingMixin, get_model_features, model_is_public
27from utilities.data import get_config_value_ci
28from utilities.exceptions import AbortRequest
30from .models import ConfigRevision, DataSource, ObjectChange
32logger = logging.getLogger('netbox.core.signals')
34__all__ = (
35 'clear_events',
36 'job_end',
37 'job_start',
38 'post_sync',
39 'pre_sync',
40)
42# Job signals
43job_start = Signal()
44job_end = Signal()
46# DataSource signals
47pre_sync = Signal()
48post_sync = Signal()
50# Event signals
51clear_events = Signal()
54#
55# Object types
56#
58@receiver(post_migrate)
59def update_object_types(sender, **kwargs):
60 """
61 Create or update the corresponding ObjectType for each model within the migrated app.
62 """
63 for model in sender.get_models():
64 app_label, model_name = model._meta.label_lower.split('.')
66 # Determine whether model is public and its supported features
67 is_public = model_is_public(model)
68 features = get_model_features(model)
70 # Create/update the ObjectType for the model
71 try:
72 ot = ObjectType.objects.get_by_natural_key(app_label=app_label, model=model_name)
73 ot.public = is_public
74 ot.features = features
75 ot.save()
76 except ObjectDoesNotExist:
77 ObjectType.objects.create(
78 app_label=app_label,
79 model=model_name,
80 public=is_public,
81 features=features,
82 )
85#
86# Change logging & event handling
87#
89# Used to track received signals per object
90_signals_received = local()
93@receiver((post_save, m2m_changed))
94def handle_changed_object(sender, instance, **kwargs):
95 """
96 Fires when an object is created or updated.
97 """
98 m2m_changed = False
100 if not hasattr(instance, 'to_objectchange'):
101 return
103 # Get the current request, or bail if not set
104 request = current_request.get()
105 if request is None:
106 return
108 # Determine the type of change being made
109 if kwargs.get('created'):
110 event_type = OBJECT_CREATED
111 elif 'created' in kwargs:
112 event_type = OBJECT_UPDATED
113 elif kwargs.get('action') in ['post_add', 'post_remove'] and kwargs['pk_set']: 113 ↛ 115line 113 didn't jump to line 115 because the condition on line 113 was never true
114 # m2m_changed with objects added or removed
115 m2m_changed = True
116 event_type = OBJECT_UPDATED
117 elif kwargs.get('action') == 'post_clear':
118 # Handle clearing of an M2M field
119 if kwargs.get('model') == Tag and getattr(instance, '_prechange_snapshot', {}).get('tags'): 119 ↛ 122line 119 didn't jump to line 122 because the condition on line 119 was never true
120 # Handle generation of M2M changes for Tags which have a previous value (ignoring changes where the
121 # prechange snapshot is empty)
122 m2m_changed = True
123 event_type = OBJECT_UPDATED
124 else:
125 # Other endpoints are unimpacted as they send post_add and post_remove
126 # This will impact changes that utilize clear() however so we may want to give consideration for this branch
127 return
128 else:
129 return
131 # Create/update an ObjectChange record for this change
132 action = {
133 OBJECT_CREATED: ObjectChangeActionChoices.ACTION_CREATE,
134 OBJECT_UPDATED: ObjectChangeActionChoices.ACTION_UPDATE,
135 OBJECT_DELETED: ObjectChangeActionChoices.ACTION_DELETE,
136 }[event_type]
137 objectchange = instance.to_objectchange(action)
138 # If this is a many-to-many field change, check for a previous ObjectChange instance recorded
139 # for this object by this request and update it
140 if m2m_changed and ( 140 ↛ 147line 140 didn't jump to line 147 because the condition on line 140 was never true
141 prev_change := ObjectChange.objects.filter(
142 changed_object_type=ContentType.objects.get_for_model(instance),
143 changed_object_id=instance.pk,
144 request_id=request.id
145 ).first()
146 ):
147 prev_change.postchange_data = objectchange.postchange_data
148 prev_change.save()
149 elif objectchange and objectchange.has_changes:
150 objectchange.user = request.user
151 objectchange.request_id = request.id
152 objectchange.save()
154 # Ensure that we're working with fresh M2M assignments
155 if m2m_changed: 155 ↛ 156line 155 didn't jump to line 156 because the condition on line 155 was never true
156 instance.refresh_from_db()
158 # Enqueue the object for event processing
159 queue = events_queue.get()
160 enqueue_event(queue, instance, request, event_type)
161 events_queue.set(queue)
163 # Increment metric counters
164 if event_type == OBJECT_CREATED:
165 model_inserts.labels(instance._meta.model_name).inc()
166 elif event_type == OBJECT_UPDATED: 166 ↛ exitline 166 didn't return from function 'handle_changed_object' because the condition on line 166 was always true
167 model_updates.labels(instance._meta.model_name).inc()
170@receiver(pre_delete)
171def handle_deleted_object(sender, instance, **kwargs):
172 """
173 Fires when an object is deleted.
174 """
175 # Run any deletion protection rules for the object. Note that this must occur prior
176 # to queueing any events for the object being deleted, in case a validation error is
177 # raised, causing the deletion to fail.
178 model_name = f'{sender._meta.app_label}.{sender._meta.model_name}'
179 validators = get_config_value_ci(get_config().PROTECTION_RULES, model_name, default=[])
180 try:
181 run_validators(instance, validators)
182 except ValidationError as e:
183 raise AbortRequest(
184 _("Deletion is prevented by a protection rule: {message}").format(message=e)
185 )
187 # Get the current request, or bail if not set
188 request = current_request.get()
189 if request is None:
190 return
192 # Check whether we've already processed a pre_delete signal for this object. (This can
193 # happen e.g. when both a parent object and its child are deleted simultaneously, due
194 # to cascading deletion.)
195 if not hasattr(_signals_received, 'pre_delete'): 195 ↛ 196line 195 didn't jump to line 196 because the condition on line 195 was never true
196 _signals_received.pre_delete = set()
197 signature = (ContentType.objects.get_for_model(instance), instance.pk)
198 if signature in _signals_received.pre_delete: 198 ↛ 199line 198 didn't jump to line 199 because the condition on line 198 was never true
199 return
200 _signals_received.pre_delete.add(signature)
202 # Record an ObjectChange if applicable
203 if hasattr(instance, 'to_objectchange'):
204 if hasattr(instance, 'snapshot') and not getattr(instance, '_prechange_snapshot', None):
205 instance.snapshot()
206 objectchange = instance.to_objectchange(ObjectChangeActionChoices.ACTION_DELETE)
207 objectchange.user = request.user
208 objectchange.request_id = request.id
209 objectchange.save()
211 # Django does not automatically send an m2m_changed signal for the reverse direction of a
212 # many-to-many relationship (see https://code.djangoproject.com/ticket/17688), so we need to
213 # trigger one manually. We do this by checking for any reverse M2M relationships on the
214 # instance being deleted, and explicitly call .remove() on the remote M2M field to delete
215 # the association. This triggers an m2m_changed signal with the `post_remove` action type
216 # for the forward direction of the relationship, ensuring that the change is recorded.
217 # Similarly, for many-to-one relationships, we set the value on the related object to None
218 # and save it to trigger a change record on that object.
219 #
220 # Skip this for private models (e.g. CablePath) whose lifecycle is an internal
221 # implementation detail. Django's on_delete handlers (e.g. SET_NULL) already take
222 # care of the database integrity; recording changelog entries for the related
223 # objects would be spurious. (Ref: #21390)
224 if not getattr(instance, '_netbox_private', False):
225 for relation in instance._meta.related_objects:
226 if type(relation) not in [ManyToManyRel, ManyToOneRel]:
227 continue
228 related_model = relation.related_model
229 related_field_name = relation.remote_field.name
230 if not issubclass(related_model, ChangeLoggingMixin):
231 # We only care about triggering the m2m_changed signal for models which support
232 # change logging
233 continue
234 related_object_type = ContentType.objects.get_for_model(related_model)
235 for obj in related_model.objects.filter(**{related_field_name: instance.pk}):
236 # Skip any related object that is itself being deleted as part of this same
237 # operation (e.g. a sibling caught up in the same cascade). Its deletion has
238 # already been recorded, so nulling the FK and saving here would record an
239 # UPDATE ObjectChange *after* the object's DELETE, corrupting the changelog and
240 # breaking branch replay. (Ref: #22270)
241 #
242 # Note this is order-dependent: it only fires once the related object's own
243 # pre_delete has run (adding it to the set). If the cascade happens to delete
244 # this instance *before* the related object, the guard won't trigger and the
245 # related object still gets an UPDATE — but in the harmless UPDATE-then-DELETE
246 # order, not the corrupting DELETE-then-UPDATE order. Fully suppressing it in
247 # every ordering would require the complete deletion set, which isn't available
248 # from a pre_delete signal.
249 if (related_object_type, obj.pk) in _signals_received.pre_delete: 249 ↛ 251line 249 didn't jump to line 251 because the condition on line 249 was always true
250 continue
251 obj.snapshot() # Ensure the change record includes the "before" state
252 if type(relation) is ManyToManyRel:
253 getattr(obj, related_field_name).remove(instance)
254 elif type(relation) is ManyToOneRel and relation.null and relation.on_delete not in (CASCADE, RESTRICT):
255 setattr(obj, related_field_name, None)
256 obj.save()
258 # Enqueue the object for event processing
259 queue = events_queue.get()
260 enqueue_event(queue, instance, request, OBJECT_DELETED)
261 events_queue.set(queue)
263 # Increment metric counters
264 model_deletes.labels(instance._meta.model_name).inc()
267@receiver(request_finished)
268def clear_signal_history(sender, **kwargs):
269 """
270 Clear out the signals history once the request is finished.
271 """
272 _signals_received.pre_delete = set()
275@receiver(clear_events)
276def clear_events_queue(sender, **kwargs):
277 """
278 Delete any queued events (e.g. because of an aborted bulk transaction)
279 """
280 logger = logging.getLogger('events')
281 logger.info(f"Clearing {len(events_queue.get())} queued events ({sender})")
282 events_queue.set({})
285#
286# DataSource handlers
287#
289@receiver(post_save, sender=DataSource)
290def enqueue_sync_job(instance, created, **kwargs):
291 """
292 When a DataSource is saved, check its sync_interval and enqueue a sync job if appropriate.
293 """
294 from .jobs import SyncDataSourceJob
296 if instance.enabled and instance.sync_interval:
297 SyncDataSourceJob.enqueue_once(instance, interval=instance.sync_interval)
298 elif not created:
299 # Delete any previously scheduled recurring jobs for this DataSource
300 for job in SyncDataSourceJob.get_jobs(instance).defer('data').filter(
301 interval__isnull=False,
302 status=JobStatusChoices.STATUS_SCHEDULED
303 ):
304 # Call delete() per instance to ensure the associated background task is deleted as well
305 job.delete()
308# Keeps the aggregated error readable when a whole source fails at once
309_AUTO_SYNC_DETAIL_LIMIT = 10
312@receiver(post_sync)
313def auto_sync(instance, **kwargs):
314 """
315 Automatically synchronize any DataFiles with AutoSyncRecords after synchronizing a DataSource.
316 """
317 from .models import AutoSyncRecord
319 failure_count = 0
320 details = []
321 first_error = None
323 records = AutoSyncRecord.objects.filter(datafile__source=instance).order_by('pk')
324 for autosync in records.select_related('object_type').prefetch_related('object'):
325 # The object may be unresolvable or mid-failure, so identify by keys
326 target = f'{autosync.object_type.app_label}.{autosync.object_type.model} ID {autosync.object_id}'
327 if autosync.object_type.model_class() is None:
328 # Not an orphaned row, so leave it for remove_stale_contenttypes
329 logger.warning(f"Skipping AutoSyncRecord for uninstalled model {target}")
330 continue
331 try:
332 # The try must stay outside this, so the savepoint is rolled back before the handler runs
333 with transaction.atomic():
334 obj = autosync.object
335 if obj is None:
336 # The prefetch resolves through the default manager, so recheck with the base manager
337 try:
338 obj = autosync.object_type.get_object_for_this_type(pk=autosync.object_id)
339 except ObjectDoesNotExist:
340 logger.warning(f"Deleting stale AutoSyncRecord for {target}")
341 autosync.delete()
342 continue
343 obj.sync(save=True)
344 except JobTimeoutException:
345 # rq arms one alarm per job, so a timeout is not an ordinary per-object failure
346 raise
347 except Exception as e:
348 failure_count += 1
349 if first_error is None:
350 first_error = e
351 # Not capped, unlike the raised message below
352 logger.error(f"Error auto-syncing {target}: {e}", exc_info=True)
353 if len(details) < _AUTO_SYNC_DETAIL_LIMIT:
354 details.append(f'- {target}: {type(e).__name__}: {e}')
356 if first_error is not None:
357 summary = ngettext(
358 'Automatic synchronization failed for {count} object:',
359 'Automatic synchronization failed for {count} objects:',
360 failure_count,
361 ).format(count=failure_count)
362 if omitted := failure_count - len(details):
363 details.append(ngettext(
364 '{count} additional failure is not shown.',
365 '{count} additional failures are not shown.',
366 omitted,
367 ).format(count=omitted))
368 raise SyncError('\n'.join([summary, *details])) from first_error
371@receiver(post_save, sender=ConfigRevision)
372def update_config(sender, instance, **kwargs):
373 """
374 Update the cached NetBox configuration when a new ConfigRevision is created.
375 """
376 instance.activate()