Coverage for netbox/search/backends.py: 48%

221 statements  

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

1import logging 

2from collections import defaultdict 

3 

4import netaddr 

5from django.conf import settings 

6from django.core.exceptions import ImproperlyConfigured 

7from django.db import DatabaseError, ProgrammingError, transaction 

8from django.db.models import F, Q, Window, prefetch_related_objects 

9from django.db.models.fields.related import ForeignKey 

10from django.db.models.functions import window 

11from django.utils.module_loading import import_string 

12from django.utils.translation import gettext_lazy as _ 

13from netaddr.core import AddrFormatError 

14 

15from core.models import ObjectType 

16from extras.models import CachedValue, CustomField 

17from netbox.registry import registry 

18from utilities.object_types import object_type_identifier 

19from utilities.querysets import RestrictedPrefetch 

20from utilities.string import title 

21 

22from . import FieldTypes, LookupTypes, get_indexer 

23 

24DEFAULT_LOOKUP_TYPE = LookupTypes.PARTIAL 

25MAX_RESULTS = 1000 

26 

27logger = logging.getLogger(__name__) 

28 

29 

30class SearchBackend: 

31 """ 

32 Base class for search backends. Subclasses must extend the `cache()`, `remove()`, and `clear()` 

33 methods below. 

34 """ 

35 _object_types = None 

36 

37 def get_object_types(self): 

38 """ 

39 Return a list of all registered object types, organized by category, suitable for populating a form's 

40 ChoiceField. 

41 """ 

42 if not self._object_types: 

43 

44 # Organize choices by category 

45 categories = defaultdict(dict) 

46 for label, idx in registry['search'].items(): 

47 categories[idx.get_category()][label] = _(title(idx.model._meta.verbose_name)) 

48 

49 # Compile a nested tuple of choices for form rendering 

50 results = ( 

51 ('', 'All Objects'), 

52 *[(category, list(choices.items())) for category, choices in categories.items()] 

53 ) 

54 

55 self._object_types = results 

56 

57 return self._object_types 

58 

59 def search(self, value, user=None, object_types=None, lookup=DEFAULT_LOOKUP_TYPE): 

60 """ 

61 Search cached object representations for the given value. 

62 """ 

63 raise NotImplementedError 

64 

65 # caching_handler() and removal_handler() are the default, synchronous signal receivers; they are 

66 # connected to post_save/post_delete from netbox.search.signals (wired from CoreConfig.ready()). 

67 # They are internal plumbing for signal dispatch, not a documented extension point: the public 

68 # backend contract is cache()/remove()/clear(). A backend that needs to do something other than 

69 # index inline (e.g. defer the work) overrides these in its subclass; see CachedValueSearchBackend. 

70 def caching_handler(self, sender, instance, created, **kwargs): 

71 """ 

72 Receiver for the post_save signal, responsible for caching object creation/changes. 

73 """ 

74 try: 

75 self.cache(instance, remove_existing=not created) 

76 except ProgrammingError as e: 

77 # The schema may be incomplete during migrations; skip caching. 

78 logger.warning(f"Skipping search cache update due to schema error: {e}") 

79 pass 

80 

81 def removal_handler(self, sender, instance, **kwargs): 

82 """ 

83 Receiver for the post_delete signal, responsible for caching object deletion. 

84 """ 

85 self.remove(instance) 

86 

87 def cache(self, instances, indexer=None, remove_existing=True): 

88 """ 

89 Create or update the cached representation of an instance. 

90 """ 

91 raise NotImplementedError 

92 

93 def remove(self, instance): 

94 """ 

95 Delete any cached representation of an instance. 

96 """ 

97 raise NotImplementedError 

98 

99 def clear(self, object_types=None): 

100 """ 

101 Delete *all* cached data (optionally filtered by object type). 

102 """ 

103 raise NotImplementedError 

104 

105 def count(self, object_types=None): 

106 """ 

107 Return a count of all cache entries (optionally filtered by object type). 

108 """ 

109 raise NotImplementedError 

110 

111 @property 

112 def size(self): 

113 """ 

114 Return a total number of cached entries. The meaning of this value will be 

115 backend-dependent. 

116 """ 

117 return None 

118 

119 

120class CachedValueSearchBackend(SearchBackend): 

121 

122 # These override the base's synchronous receivers to defer indexing past the response. They are 

123 # the seam where this backend captures the `using` alias Django passes to post_save/post_delete: 

124 # the deferred write runs after the transaction commits (and possibly in a worker), by which point 

125 # the originating routing context is gone, so the alias must be captured here and replayed on the 

126 # deferred write to keep cache entries in the originating schema (e.g. a branch schema under 

127 # netbox-branching). Deferral is internal to this backend; the public contract is unchanged. 

128 # 

129 # mark_for_deferred_indexing() etc. are imported inside each method rather than at module level: 

130 # this module's own top would import deferred.py *before* search_backend is defined further down 

131 # this same file, and deferred.py (plus jobs.py) need that singleton at their own module level. 

132 # A module-level import here would close that loop into a backends -> deferred -> backends 

133 # cycle. See #22485. 

134 def caching_handler(self, sender, instance, created, using=None, **kwargs): 

135 """ 

136 Receiver for the post_save signal, responsible for caching object creation/changes. 

137 """ 

138 from .deferred import OP_CACHE, mark_for_deferred_indexing 

139 

140 # Skip non-cacheable objects without scheduling any deferred work. 

141 try: 

142 indexer = get_indexer(instance) 

143 except KeyError: 

144 return 

145 

146 try: 

147 object_type = ObjectType.objects.get_for_model(indexer.model) 

148 except ProgrammingError as e: 

149 # The schema may be incomplete during migrations; skip caching. 

150 logger.warning(f"Skipping search cache update due to schema error: {e}") 

151 return 

152 

153 mark_for_deferred_indexing(object_type.pk, instance.pk, OP_CACHE, using=using) 

154 

155 def removal_handler(self, sender, instance, using=None, **kwargs): 

156 """ 

157 Receiver for the post_delete signal, responsible for caching object deletion. 

158 """ 

159 from .deferred import OP_REMOVE, mark_for_deferred_indexing 

160 

161 # Skip non-cacheable objects without scheduling any deferred work. 

162 try: 

163 indexer = get_indexer(instance) 

164 except KeyError: 

165 return 

166 

167 try: 

168 object_type = ObjectType.objects.get_for_model(indexer.model) 

169 except ProgrammingError as e: 

170 # The schema may be incomplete during migrations; skip caching. 

171 logger.warning(f"Skipping search cache update due to schema error: {e}") 

172 return 

173 

174 mark_for_deferred_indexing(object_type.pk, instance.pk, OP_REMOVE, using=using) 

175 

176 def search(self, value, user=None, object_types=None, lookup=DEFAULT_LOOKUP_TYPE): 

177 

178 # Build the filter used to find relevant CachedValue records 

179 query_filter = Q(**{f'value__{lookup}': value}) 

180 if object_types: 

181 # Limit results by object type 

182 query_filter &= Q(object_type__in=object_types) 

183 if lookup in (LookupTypes.STARTSWITH, LookupTypes.ENDSWITH): 

184 # "Starts/ends with" matches are valid only on string values 

185 query_filter &= Q(type=FieldTypes.STRING) 

186 elif lookup in (LookupTypes.PARTIAL, LookupTypes.EXACT): 

187 try: 

188 # If the value looks like an IP address, add extra filters for CIDR/INET values 

189 address = str(netaddr.IPNetwork(value.strip()).cidr) 

190 query_filter |= Q(type=FieldTypes.INET) & Q(value__net_host=address) 

191 if lookup == LookupTypes.PARTIAL: 

192 query_filter |= Q(type=FieldTypes.CIDR) & Q(value__net_contains_or_equals=address) 

193 except (AddrFormatError, ValueError): 

194 pass 

195 

196 # Construct the base queryset to retrieve matching results 

197 queryset = CachedValue.objects.filter(query_filter).annotate( 

198 # Annotate the rank of each result for its object according to its weight 

199 row_number=Window( 

200 expression=window.RowNumber(), 

201 partition_by=[F('object_type'), F('object_id')], 

202 order_by=[F('weight').asc()], 

203 ) 

204 )[:MAX_RESULTS] 

205 

206 # Gather all ObjectTypes present in the search results (used for prefetching related 

207 # objects). This must be done before generating the final results list, which returns 

208 # a RawQuerySet. 

209 object_type_ids = set(queryset.values_list('object_type', flat=True)) 

210 object_types = ObjectType.objects.filter(pk__in=object_type_ids) 

211 

212 # Construct a Prefetch to pre-fetch only those related objects for which the 

213 # user has permission to view. 

214 if user: 

215 prefetch = (RestrictedPrefetch('object', user, 'view'), 'object_type') 

216 else: 

217 prefetch = ('object', 'object_type') 

218 

219 # Wrap the base query to return only the lowest-weight result for each object 

220 # Hat-tip to https://blog.oyam.dev/django-filter-by-window-function/ for the solution 

221 sql, params = queryset.query.sql_with_params() 

222 results = CachedValue.objects.prefetch_related(*prefetch).raw( 

223 f"SELECT * FROM ({sql}) t WHERE row_number = 1", 

224 params 

225 ) 

226 

227 # Iterate through each ObjectType represented in the search results and prefetch any 

228 # related objects necessary to render the prescribed display attributes (display_attrs). 

229 for object_type in object_types: 

230 model = object_type.model_class() 

231 indexer = registry['search'].get(object_type_identifier(object_type)) 

232 if not (display_attrs := getattr(indexer, 'display_attrs', None)): 

233 continue 

234 

235 # Add ForeignKey fields to prefetch list 

236 prefetch_fields = [] 

237 for attr in display_attrs: 

238 field = model._meta.get_field(attr) 

239 if type(field) is ForeignKey: 

240 prefetch_fields.append(f'object__{attr}') 

241 

242 # Compile a list of all CachedValues referencing this object type, and prefetch 

243 # any related objects 

244 if prefetch_fields: 

245 objects = [r for r in results if r.object_type == object_type] 

246 prefetch_related_objects(objects, *prefetch_fields) 

247 

248 # Omit any results pertaining to an object the user does not have permission to view 

249 ret = [] 

250 for r in results: 

251 if r.object is not None: 

252 r.name = str(r.object) 

253 ret.append(r) 

254 

255 return ret 

256 

257 # `using` here is a PostgreSQL/schema concern specific to this backend's deferred-write path (it 

258 # replays the originating alias so branch writes land in the branch schema). It is deliberately 

259 # NOT on the base cache()/remove() contract: a non-PostgreSQL backend (Redis, Solr, etc.) has no 

260 # such concept. Do not lift `using` onto the base for symmetry; doing so would leak this backend's 

261 # storage model into the generic contract. 

262 def cache(self, instances, indexer=None, remove_existing=True, using=None): 

263 custom_fields = None 

264 

265 # Convert a single instance to an iterable 

266 if not hasattr(instances, '__iter__'): 266 ↛ 267line 266 didn't jump to line 267 because the condition on line 266 was never true

267 instances = [instances] 

268 

269 # Determine the queryset manager used to write cache entries. When a 

270 # database alias is provided (e.g. by a deferred task replaying the alias 

271 # the originating write used), entries are written to that connection; 

272 # otherwise the configured router decides. `using` is expected to be a 

273 # concrete alias or falsy (None) per the caller's contract; a falsy value 

274 # defers to the router, which is the correct behavior either way. 

275 manager = CachedValue.objects.using(using) if using else CachedValue.objects 

276 

277 buffer = [] 

278 counter = 0 

279 for instance in instances: 

280 

281 # First item 

282 if not counter: 282 ↛ 298line 282 didn't jump to line 298 because the condition on line 282 was always true

283 

284 # Determine the indexer 

285 if indexer is None: 

286 try: 

287 indexer = get_indexer(instance) 

288 except KeyError: 

289 break 

290 

291 # Prefetch any associated custom fields (excluding those with a zero search weight) 

292 custom_fields = [ 

293 cf for cf in CustomField.objects.get_for_model(indexer.model) 

294 if cf.search_weight > 0 

295 ] 

296 

297 # Wipe out any previously cached values for the object 

298 if remove_existing: 298 ↛ 299line 298 didn't jump to line 299 because the condition on line 298 was never true

299 self.remove(instance, using=using) 

300 

301 # Generate cache data 

302 object_type = ObjectType.objects.get_for_model(indexer.model) 

303 for field in indexer.to_cache(instance, custom_fields=custom_fields): 

304 buffer.append( 

305 CachedValue( 

306 object_type=object_type, 

307 object_id=instance.pk, 

308 field=field.name, 

309 type=field.type, 

310 weight=field.weight, 

311 value=field.value 

312 ) 

313 ) 

314 

315 # Check whether the buffer needs to be flushed 

316 if len(buffer) >= 2000: 316 ↛ 317line 316 didn't jump to line 317 because the condition on line 316 was never true

317 counter += len(manager.bulk_create(buffer)) 

318 buffer = [] 

319 

320 # Final buffer flush 

321 if buffer: 

322 counter += len(manager.bulk_create(buffer)) 

323 

324 return counter 

325 

326 def _remove_by_id(self, object_type_id, object_ids, using=None): 

327 """ 

328 Delete cached values for the given content type and object IDs using a 

329 single raw DELETE. Shared by remove() and the deferred search task. 

330 """ 

331 if not object_ids: 331 ↛ 332line 331 didn't jump to line 332 because the condition on line 331 was never true

332 return None 

333 

334 qs = CachedValue.objects.filter(object_type_id=object_type_id, object_id__in=object_ids) 

335 

336 # Call _raw_delete() on the queryset to avoid first loading instances into memory 

337 return qs._raw_delete(using=using or qs.db) 

338 

339 def remove(self, instance, using=None): 

340 # Avoid attempting to query for non-cacheable objects 

341 try: 

342 indexer = get_indexer(instance) 

343 except KeyError: 

344 return None 

345 

346 # Use the indexer's (concrete) model to resolve the object type, matching 

347 # the content type that cache() writes entries under. 

348 object_type = ObjectType.objects.get_for_model(indexer.model) 

349 

350 return self._remove_by_id(object_type.pk, [instance.pk], using=using) 

351 

352 # Postgres SQLSTATEs indicating this backend's own CachedValue table (or the schema it lives in) 

353 # no longer exists. This happens when a branch is merged or deprovisioned (its schema, and every 

354 # table in it including its copy of CachedValue, dropped) between the time an update was enqueued 

355 # and when it is applied. There is nothing to write to and nothing worth retrying -- a later write 

356 # for a surviving branch's CachedValue table is unaffected -- so this is expected and safe to 

357 # skip. Any other DatabaseError (e.g. a deadlock, lost connection, or CachedValue itself being out 

358 # of sync with a not-yet-migrated deployment) is not expected and must propagate so the work fails 

359 # visibly, rather than silently dropping index updates. This set applies to every write this 

360 # backend makes to CachedValue (removals below, and the remove+insert inside the cache loop) -- 

361 # deliberately not to reads of the model being indexed; see _STALE_INDEX_TARGET_SQLSTATES for 

362 # that. 

363 _MISSING_CACHE_TABLE_SQLSTATES = frozenset(( 

364 '3F000', # invalid_schema_name 

365 '42P01', # undefined_table 

366 )) 

367 

368 # Additionally tolerated when reading the model being (re)indexed -- not when writing to 

369 # CachedValue (see _MISSING_CACHE_TABLE_SQLSTATES above, which this extends). A plugin whose 

370 # models are dynamically regenerated per branch (e.g. netbox-custom-objects) resolves 

371 # ObjectType.model_class() to a branch-unaware class pinned to main; if a branch has renamed or 

372 # removed a column since an update was enqueued, that stale class's column no longer matches the 

373 # branch's live table, and reading it raises undefined_column rather than undefined_table. 

374 # Unlike a dropped schema, this is *not* self-healing on "the next reindex": model_class() 

375 # resolves the same stale, main-pinned class every time, so the object stays unindexed until the 

376 # branch is merged or reverted. Scoping this to the read only -- rather than folding it into 

377 # _MISSING_CACHE_TABLE_SQLSTATES broadly -- keeps a genuine defect in the write path (e.g. code 

378 # deployed against a database that has not yet been migrated, which would also raise 

379 # undefined_column) from being silently downgraded to a warning. 

380 # 

381 # This covers only the top-level read of the model's own columns, not a related object a search 

382 # index's to_cache() might lazily traverse into (e.g. a relational field indexed with a positive 

383 # search_weight): such a traversal issues its own query inside the write step below, which does 

384 # not tolerate undefined_column. A plugin whose indexed fields never reference another of its own 

385 # dynamically-regenerated models is unaffected; one that does would need a plugin-side fix (making 

386 # ObjectType.model_class() branch-aware) rather than a wider exemption here. 

387 _STALE_INDEX_TARGET_SQLSTATES = _MISSING_CACHE_TABLE_SQLSTATES | frozenset(( 

388 '42703', # undefined_column 

389 )) 

390 

391 def _is_missing_cache_table(self, exc): 

392 """ 

393 Return True if the given DatabaseError was caused by CachedValue's own schema/table no longer 

394 existing (vs. a transient error that should propagate). Covers writes to CachedValue; see 

395 _is_stale_index_target() for the read side of a deferred update. 

396 """ 

397 sqlstate = getattr(getattr(exc, '__cause__', None), 'sqlstate', None) 

398 return sqlstate in self._MISSING_CACHE_TABLE_SQLSTATES 

399 

400 def _is_stale_index_target(self, exc): 

401 """ 

402 Return True if the given DatabaseError was caused by the model being (re)indexed no longer 

403 matching what was expected when the update was enqueued -- its schema or table (as for 

404 _is_missing_cache_table()), or, for a model whose columns can themselves diverge per branch, 

405 an individual column. 

406 """ 

407 sqlstate = getattr(getattr(exc, '__cause__', None), 'sqlstate', None) 

408 return sqlstate in self._STALE_INDEX_TARGET_SQLSTATES 

409 

410 def _apply_deferred_updates(self, using=None, cache_groups=None, remove_groups=None, log=logger): 

411 """ 

412 Apply a coalesced batch of updates to the search cache. Private to this backend; called by the 

413 deferred-flush machinery (netbox.search.deferred) and the background job 

414 (netbox.search.jobs.SearchCacheJob), not part of the public backend contract. 

415 

416 The `using` alias captured when each object was saved/deleted is replayed here so entries are 

417 written to the originating database/schema (e.g. a branch schema under netbox-branching), 

418 regardless of any routing context that is no longer active by the time this runs. 

419 """ 

420 for object_type_id, pks in (remove_groups or {}).items(): 

421 # Resolved once and reused for both the atomic() block and the delete below: passing 

422 # `using` straight through when it's falsy would let transaction.atomic() default to 

423 # DEFAULT_DB_ALIAS while _remove_by_id() defers to the router (see cache()'s own comment 

424 # on this) -- the savepoint would then belong to a different connection than the one the 

425 # DELETE actually runs on, on any deployment where CachedValue is routed elsewhere. 

426 db = using or CachedValue.objects.db 

427 try: 

428 # A tolerated DatabaseError still needs a savepoint to roll back to: without one, 

429 # Postgres leaves the connection's enclosing transaction aborted (refusing every 

430 # further statement in it until a rollback) even though the Python exception was 

431 # caught, breaking every later iteration of this loop -- not just this one. 

432 with transaction.atomic(using=db): 

433 self._remove_by_id(object_type_id, pks, using=db) 

434 except DatabaseError as e: 

435 if not self._is_missing_cache_table(e): 

436 raise 

437 log.warning( 

438 f"Skipping search cache removal for object type {object_type_id}: " 

439 f"CachedValue's own table or schema no longer exists ({e})" 

440 ) 

441 

442 for object_type_id, pks in (cache_groups or {}).items(): 

443 try: 

444 object_type = ObjectType.objects.get(pk=object_type_id) 

445 except ObjectType.DoesNotExist: 

446 continue 

447 model = object_type.model_class() 

448 if model is None: 448 ↛ 449line 448 didn't jump to line 449 because the condition on line 448 was never true

449 continue 

450 

451 # Resolved once per group, for the same reason as the removal loop above -- and kept 

452 # separate for the read vs. the write, since a router could place the indexed model and 

453 # CachedValue on different aliases (in which case no single atomic() spans both anyway; 

454 # see the outer/inner split below). 

455 read_db = using or model._default_manager.db 

456 write_db = using or CachedValue.objects.db 

457 

458 try: 

459 # The outer atomic() makes the delete+insert pair below a single write. It does 

460 # *not* give the read a consistent snapshot with that write -- PostgreSQL's default 

461 # (and NetBox's) READ COMMITTED isolation takes a fresh snapshot per statement 

462 # regardless of transaction boundaries, so a concurrent update can still land 

463 # between the read and the write either way. That race is pre-existing and benign: 

464 # the concurrent save schedules its own deferred update, so the object is reindexed 

465 # again regardless of which value this pass happened to write. 

466 # 

467 # Nested inside it, the read gets its own savepoint so a tolerated failure there 

468 # (see _is_stale_index_target(), which tolerates a wider set of SQLSTATES than a 

469 # write to CachedValue does) rolls back only the read, without ever reaching -- or 

470 # needing to roll back -- the write. 

471 with transaction.atomic(using=write_db): 

472 try: 

473 with transaction.atomic(using=read_db): 

474 # Reading on `read_db` is required: a branch object's PK may be absent 

475 # (or refer to a different object) on the default connection. 

476 instances = list(model.objects.using(read_db).filter(pk__in=pks)) 

477 except DatabaseError as e: 

478 if not self._is_stale_index_target(e): 

479 raise 

480 log.warning( 

481 f"Skipping search cache update for object type {object_type_id}: the " 

482 f"indexed model no longer matches its table, e.g. a branch-diverged " 

483 f"column ({e})" 

484 ) 

485 continue 

486 

487 # Clear any stale entries for these objects, then re-insert. Wrapping both in 

488 # one transaction avoids leaving an object with no cache rows if execution fails 

489 # between the delete and the insert. 

490 self._remove_by_id(object_type_id, pks, using=write_db) 

491 self.cache(instances, remove_existing=False, using=write_db) 

492 except DatabaseError as e: 

493 if not self._is_missing_cache_table(e): 

494 raise 

495 log.warning( 

496 f"Skipping search cache update for object type {object_type_id}: " 

497 f"CachedValue's own table or schema no longer exists ({e})" 

498 ) 

499 

500 def clear(self, object_types=None): 

501 qs = CachedValue.objects.all() 

502 if object_types: 

503 qs = qs.filter(object_type__in=object_types) 

504 

505 # Call _raw_delete() on the queryset to avoid first loading instances into memory 

506 return qs._raw_delete(using=qs.db) 

507 

508 def count(self, object_types=None): 

509 qs = CachedValue.objects.all() 

510 if object_types: 510 ↛ 512line 510 didn't jump to line 512 because the condition on line 510 was always true

511 qs = qs.filter(object_type__in=object_types) 

512 return qs.count() 

513 

514 @property 

515 def size(self): 

516 return CachedValue.objects.count() 

517 

518 

519def get_backend(): 

520 """ 

521 Initializes and returns the configured search backend. 

522 """ 

523 try: 

524 backend_cls = import_string(settings.SEARCH_BACKEND) 

525 except AttributeError: 

526 raise ImproperlyConfigured(f"Failed to import configured SEARCH_BACKEND: {settings.SEARCH_BACKEND}") 

527 

528 # Initialize and return the backend instance 

529 return backend_cls() 

530 

531 

532search_backend = get_backend()