Coverage for extras/events.py: 54%

109 statements  

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

1import logging 

2from collections import UserDict, defaultdict 

3 

4from django.conf import settings 

5from django.utils.module_loading import import_string 

6from django.utils.translation import gettext as _ 

7 

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 

13 

14from .conditions import AbsentData 

15from .models import EventRule 

16 

17logger = logging.getLogger('netbox.events_processor') 

18 

19 

20class EventContext(UserDict): 

21 """ 

22 Dictionary-compatible wrapper for queued events that lazily serializes 

23 ``event['data']`` on first access. 

24 

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 """ 

30 

31 def __init__(self, *args, **kwargs): 

32 super().__init__(*args, **kwargs) 

33 

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

40 

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'] 

51 

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 

59 

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

69 

70 return super().__getitem__(item) 

71 

72 

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) 

82 

83 return serializer.data 

84 

85 

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) 

102 

103 return { 

104 'prechange': getattr(instance, '_prechange_snapshot', None), 

105 'postchange': postchange_snapshot, 

106 } 

107 

108 

109def enqueue_event(queue, instance, request, event_type): 

110 """ 

111 Enqueue (or coalesce) an event for a created/updated/deleted object. 

112 

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 

118 

119 app_label = instance._meta.app_label 

120 model_name = instance._meta.model_name 

121 

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

130 

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'] 

133 

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 ) 

152 

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) 

157 

158 

159def process_event_rules(event_rules, object_type, event): 

160 """ 

161 Process a list of EventRules against an event. 

162 

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 

174 

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 

178 

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

190 

191 for event_rule in event_rules: 

192 

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 

198 

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 = {} 

215 

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} 

219 

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 

232 

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 ) 

249 

250 

251def process_event_queue(events): 

252 """ 

253 Flush a list of object representation to RQ for EventRule processing. 

254 

255 This is the default processor listed in EVENTS_PIPELINE. 

256 """ 

257 events_cache = defaultdict(dict) 

258 

259 for event in events: 

260 event_type = event['event_type'] 

261 object_type = event['object_type'] 

262 

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] 

271 

272 process_event_rules( 

273 event_rules=event_rules, 

274 object_type=object_type, 

275 event=event, 

276 ) 

277 

278 

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