Coverage for extras/events.py: 54%
109 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 collections import UserDict, defaultdict
4from django.conf import settings
5from django.utils.module_loading import import_string
6from django.utils.translation import gettext as _
8from core.events import *
9from core.models import ObjectType
10from netbox.models.features import has_feature
11from utilities.api import get_serializer_for_model
12from utilities.serialization import serialize_object
14from .conditions import AbsentData
15from .models import EventRule
17logger = logging.getLogger('netbox.events_processor')
20class EventContext(UserDict):
21 """
22 Dictionary-compatible wrapper for queued events that lazily serializes
23 ``event['data']`` on first access.
25 Backward-compatible with the plain-dict interface expected by existing
26 EVENTS_PIPELINE consumers. When the same object is enqueued more than once
27 in a single request, the serialization source is updated so consumers see
28 the latest state.
29 """
31 def __init__(self, *args, **kwargs):
32 super().__init__(*args, **kwargs)
34 # Track which model instance should be serialized if/when `data` is
35 # requested. This may be refreshed on duplicate enqueue, while leaving
36 # the public `object` entry untouched for compatibility.
37 self._serialization_source = None
38 if 'object' in self: 38 ↛ exitline 38 didn't return from function '__init__' because the condition on line 38 was always true
39 self._serialization_source = super().__getitem__('object')
41 def refresh_serialization_source(self, instance):
42 """
43 Point lazy serialization at a fresher instance, invalidating any
44 already-materialized ``data``.
45 """
46 self._serialization_source = instance
47 # UserDict.__contains__ checks the backing dict directly, so `in`
48 # does not trigger __getitem__'s lazy serialization.
49 if 'data' in self:
50 del self['data']
52 def freeze_data(self, instance):
53 """
54 Eagerly serialize and cache the payload for delete events, where the
55 object may become inaccessible after deletion.
56 """
57 super().__setitem__('data', serialize_for_event(instance))
58 self._serialization_source = None
60 def __getitem__(self, item):
61 if item == 'data' and 'data' not in self: 61 ↛ 67line 61 didn't jump to line 67 because the condition on line 61 was never true
62 # Materialize the payload only when an event consumer asks for it.
63 #
64 # On coalesced events, use the latest explicitly queued instance so
65 # webhooks/scripts/notifications observe the final queued state for
66 # that object within the request.
67 source = self._serialization_source or super().__getitem__('object')
68 super().__setitem__('data', serialize_for_event(source))
70 return super().__getitem__(item)
73def serialize_for_event(instance):
74 """
75 Return a serialized representation of the given instance suitable for use in a queued event.
76 """
77 serializer_class = get_serializer_for_model(instance.__class__)
78 serializer_context = {
79 'request': None,
80 }
81 serializer = serializer_class(instance, context=serializer_context)
83 return serializer.data
86def get_snapshots(instance, event_type):
87 """
88 Return a dictionary of pre- and post-change snapshots for the given instance.
89 """
90 if event_type == OBJECT_DELETED:
91 # Post-change snapshot must be empty for deleted objects
92 postchange_snapshot = None
93 elif hasattr(instance, '_postchange_snapshot'): 93 ↛ 96line 93 didn't jump to line 96 because the condition on line 93 was always true
94 # Use the cached post-change snapshot if one is available
95 postchange_snapshot = instance._postchange_snapshot
96 elif hasattr(instance, 'serialize_object'):
97 # Use model's serialize_object() method if defined
98 postchange_snapshot = instance.serialize_object()
99 else:
100 # Fall back to the serialize_object() utility function
101 postchange_snapshot = serialize_object(instance)
103 return {
104 'prechange': getattr(instance, '_prechange_snapshot', None),
105 'postchange': postchange_snapshot,
106 }
109def enqueue_event(queue, instance, request, event_type):
110 """
111 Enqueue (or coalesce) an event for a created/updated/deleted object.
113 Events are processed after the request completes.
114 """
115 # Bail if this type of object does not support event rules
116 if not has_feature(instance, 'event_rules'):
117 return
119 app_label = instance._meta.app_label
120 model_name = instance._meta.model_name
122 if instance.pk is None: 122 ↛ 123line 122 didn't jump to line 123 because the condition on line 122 was never true
123 raise ValueError(
124 _("Cannot enqueue an event for an unsaved {app_label}.{model} instance.").format(
125 app_label=app_label,
126 model=model_name,
127 )
128 )
129 key = f'{app_label}.{model_name}:{instance.pk}'
131 if key in queue: 131 ↛ 132line 131 didn't jump to line 132 because the condition on line 131 was never true
132 queue[key]['snapshots']['postchange'] = get_snapshots(instance, event_type)['postchange']
134 # If the object is being deleted, convert any prior update event into a
135 # delete event and freeze the payload before the object (or related
136 # rows) become inaccessible.
137 if event_type == OBJECT_DELETED:
138 queue[key]['event_type'] = event_type
139 else:
140 # Keep the public `object` entry stable for compatibility.
141 queue[key].refresh_serialization_source(instance)
142 else:
143 queue[key] = EventContext(
144 object_type=ObjectType.objects.get_for_model(instance),
145 object_id=instance.pk,
146 object=instance,
147 event_type=event_type,
148 snapshots=get_snapshots(instance, event_type),
149 request=request,
150 user=request.user,
151 )
153 # For delete events, eagerly serialize the payload before the row is gone.
154 # This covers both first-time enqueues and coalesced update→delete promotions.
155 if event_type == OBJECT_DELETED:
156 queue[key].freeze_data(instance)
159def process_event_rules(event_rules, object_type, event):
160 """
161 Process a list of EventRules against an event.
163 Notes on event sources:
164 - Object change events (created/updated/deleted) are enqueued via enqueue_event()
165 during an HTTP request. These events include a request object, and their payload is
166 always the serialized object.
167 - Job lifecycle events (JOB_STARTED/JOB_COMPLETED) are emitted by job_start/job_end
168 signal handlers and may not include a request context. Consumers must not assume
169 that a request is always present. Their payload is the job's `data` field, which is
170 nullable and (for a job which sets it directly) not guaranteed to be a dict.
171 """
172 if not event_rules: 172 ↛ 177line 172 didn't jump to line 177 because the condition on line 172 was always true
173 return
175 # Normalize object_type onto the event context so that an action's enqueue() can always read
176 # event_context['object_type']: job-lifecycle events pass it only as this parameter.
177 event['object_type'] = object_type
179 # Normalize the event payload to a dict or AbsentData once for all rules.
180 data = event['data']
181 if not isinstance(data, dict):
182 if data is not None:
183 logger.warning(
184 _('Ignoring invalid data payload on {event_type} event (got {data_type})').format(
185 event_type=event['event_type'],
186 data_type=type(data).__name__,
187 )
188 )
189 data = AbsentData()
191 for event_rule in event_rules:
193 # Merge snapshots and evaluate event rule conditions (if any).
194 condition_data = data.copy()
195 condition_data['snapshots'] = event.get('snapshots')
196 if not event_rule.eval_conditions(condition_data):
197 continue
199 # Guard against action_data that is valid JSON but not a dict
200 # (e.g. a bare string or number). Existing rows with bad data are
201 # tolerated at runtime; validation on EventRule.clean() prevents
202 # new ones.
203 if event_rule.action_data is None:
204 action_data = {}
205 elif isinstance(event_rule.action_data, dict):
206 action_data = event_rule.action_data
207 else:
208 logger.warning(
209 _('Ignoring invalid action_data on event rule "{rule}" (got {data_type})').format(
210 rule=event_rule,
211 data_type=type(event_rule.action_data).__name__,
212 )
213 )
214 action_data = {}
216 # Merge rule-specific action_data with the event payload.
217 # Copy to avoid mutating the rule's stored action_data dict.
218 event_data = {**action_data, **data}
220 action = event_rule.action_provider
221 if action is None:
222 # The plugin providing this action type may not be installed. Log and move on to the
223 # next rule rather than raising: one rule's unavailable action must not prevent any
224 # other rule in this batch from being processed.
225 logger.warning(
226 _('Skipping event rule "{rule}": action type "{action_type}" is not registered '
227 '(the providing plugin may not be installed).').format(
228 rule=event_rule, action_type=event_rule.action_type,
229 )
230 )
231 continue
233 try:
234 action.enqueue(
235 event_rule=event_rule,
236 event_context=event,
237 action_object=event_rule.action_object,
238 action_data=event_data,
239 )
240 except Exception:
241 # Isolate third-party bugs; a core action's own bugs should propagate instead.
242 if not action.is_plugin_provided:
243 raise
244 logger.exception(
245 _('Error processing event rule "{rule}" (action: {action_type})').format(
246 rule=event_rule, action_type=event_rule.action_type,
247 )
248 )
251def process_event_queue(events):
252 """
253 Flush a list of object representation to RQ for EventRule processing.
255 This is the default processor listed in EVENTS_PIPELINE.
256 """
257 events_cache = defaultdict(dict)
259 for event in events:
260 event_type = event['event_type']
261 object_type = event['object_type']
263 # Cache applicable Event Rules
264 if object_type not in events_cache[event_type]:
265 events_cache[event_type][object_type] = EventRule.objects.filter(
266 event_types__contains=[event['event_type']],
267 object_types=object_type,
268 enabled=True
269 )
270 event_rules = events_cache[event_type][object_type]
272 process_event_rules(
273 event_rules=event_rules,
274 object_type=object_type,
275 event=event,
276 )
279def flush_events(events):
280 """
281 Flush a list of object representations to RQ for event processing.
282 """
283 if events: 283 ↛ exitline 283 didn't return from function 'flush_events' because the condition on line 283 was always true
284 for name in settings.EVENTS_PIPELINE:
285 try:
286 func = import_string(name)
287 func(events)
288 except ImportError as e:
289 logger.error(_("Cannot import events pipeline {name} error: {error}").format(name=name, error=e))