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
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 18:35 +0000
1import logging
2from collections import defaultdict
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
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
22from . import FieldTypes, LookupTypes, get_indexer
24DEFAULT_LOOKUP_TYPE = LookupTypes.PARTIAL
25MAX_RESULTS = 1000
27logger = logging.getLogger(__name__)
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
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:
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))
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 )
55 self._object_types = results
57 return self._object_types
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
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
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)
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
93 def remove(self, instance):
94 """
95 Delete any cached representation of an instance.
96 """
97 raise NotImplementedError
99 def clear(self, object_types=None):
100 """
101 Delete *all* cached data (optionally filtered by object type).
102 """
103 raise NotImplementedError
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
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
120class CachedValueSearchBackend(SearchBackend):
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
140 # Skip non-cacheable objects without scheduling any deferred work.
141 try:
142 indexer = get_indexer(instance)
143 except KeyError:
144 return
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
153 mark_for_deferred_indexing(object_type.pk, instance.pk, OP_CACHE, using=using)
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
161 # Skip non-cacheable objects without scheduling any deferred work.
162 try:
163 indexer = get_indexer(instance)
164 except KeyError:
165 return
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
174 mark_for_deferred_indexing(object_type.pk, instance.pk, OP_REMOVE, using=using)
176 def search(self, value, user=None, object_types=None, lookup=DEFAULT_LOOKUP_TYPE):
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
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]
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)
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')
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 )
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
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}')
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)
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)
255 return ret
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
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]
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
277 buffer = []
278 counter = 0
279 for instance in instances:
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
284 # Determine the indexer
285 if indexer is None:
286 try:
287 indexer = get_indexer(instance)
288 except KeyError:
289 break
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 ]
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)
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 )
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 = []
320 # Final buffer flush
321 if buffer:
322 counter += len(manager.bulk_create(buffer))
324 return counter
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
334 qs = CachedValue.objects.filter(object_type_id=object_type_id, object_id__in=object_ids)
336 # Call _raw_delete() on the queryset to avoid first loading instances into memory
337 return qs._raw_delete(using=using or qs.db)
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
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)
350 return self._remove_by_id(object_type.pk, [instance.pk], using=using)
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 ))
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 ))
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
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
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.
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 )
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
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
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
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 )
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)
505 # Call _raw_delete() on the queryset to avoid first loading instances into memory
506 return qs._raw_delete(using=qs.db)
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()
514 @property
515 def size(self):
516 return CachedValue.objects.count()
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}")
528 # Initialize and return the backend instance
529 return backend_cls()
532search_backend = get_backend()