Coverage for core/signals.py: 64%

190 statements  

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

1import logging 

2from threading import local 

3 

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 

16 

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 

29 

30from .models import ConfigRevision, DataSource, ObjectChange 

31 

32logger = logging.getLogger('netbox.core.signals') 

33 

34__all__ = ( 

35 'clear_events', 

36 'job_end', 

37 'job_start', 

38 'post_sync', 

39 'pre_sync', 

40) 

41 

42# Job signals 

43job_start = Signal() 

44job_end = Signal() 

45 

46# DataSource signals 

47pre_sync = Signal() 

48post_sync = Signal() 

49 

50# Event signals 

51clear_events = Signal() 

52 

53 

54# 

55# Object types 

56# 

57 

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

65 

66 # Determine whether model is public and its supported features 

67 is_public = model_is_public(model) 

68 features = get_model_features(model) 

69 

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 ) 

83 

84 

85# 

86# Change logging & event handling 

87# 

88 

89# Used to track received signals per object 

90_signals_received = local() 

91 

92 

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 

99 

100 if not hasattr(instance, 'to_objectchange'): 

101 return 

102 

103 # Get the current request, or bail if not set 

104 request = current_request.get() 

105 if request is None: 

106 return 

107 

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 

130 

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

153 

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

157 

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) 

162 

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

168 

169 

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 ) 

186 

187 # Get the current request, or bail if not set 

188 request = current_request.get() 

189 if request is None: 

190 return 

191 

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) 

201 

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

210 

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

257 

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) 

262 

263 # Increment metric counters 

264 model_deletes.labels(instance._meta.model_name).inc() 

265 

266 

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

273 

274 

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

283 

284 

285# 

286# DataSource handlers 

287# 

288 

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 

295 

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

306 

307 

308# Keeps the aggregated error readable when a whole source fails at once 

309_AUTO_SYNC_DETAIL_LIMIT = 10 

310 

311 

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 

318 

319 failure_count = 0 

320 details = [] 

321 first_error = None 

322 

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

355 

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 

369 

370 

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