Coverage for netbox/search/deferred.py: 84%
57 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
3from django.db import DEFAULT_DB_ALIAS, connections, transaction
4from redis.exceptions import RedisError
6from netbox.constants import RQ_QUEUE_DEFAULT
7from netbox.search.backends import search_backend
8from netbox.search.jobs import SearchCacheJob
9from utilities.rqworker import any_workers_for_queue
11# This module is internal plumbing for the search signal handlers; nothing here
12# is part of the public/plugin API, so no symbols are exported via __all__.
14logger = logging.getLogger(__name__)
16# Operation markers stored in the per-transaction buffer
17OP_CACHE = 'cache'
18OP_REMOVE = 'remove'
20# Attributes used to tag a flush callback so we can recognize our own callbacks
21# among those registered on a connection and reach the batch they will flush.
22_FLUSH_ALIAS_ATTR = '_netbox_search_flush_alias'
23_FLUSH_BATCH_ATTR = '_netbox_search_flush_batch'
24# The savepoint stack active when a flush callback was registered; see
25# _pending_batch() for why buffering is scoped to it.
26_FLUSH_SCOPE_ATTR = '_netbox_search_flush_scope'
29def mark_for_deferred_indexing(object_type_id, pk, op, using=None):
30 """
31 Schedule a searchable object for deferred (re)indexing.
33 The work is coalesced per database connection and per transaction: repeated
34 operations on the same object collapse to a single entry (a deletion always
35 wins over a create/update), and a single flush is scheduled to run after the
36 transaction commits. When no transaction is open (autocommit), the indexing
37 runs synchronously.
39 Args:
40 object_type_id: PK of the object's ObjectType/ContentType.
41 pk: PK of the object.
42 op: OP_CACHE or OP_REMOVE.
43 using: The database alias the originating write used. Replayed verbatim
44 on the deferred write so the cache entries land in the same schema
45 (e.g. a branch schema under netbox-branching), regardless of any
46 routing context that may be unset by the time the flush runs.
47 """
48 # Fall back to the default alias when no originating alias was captured. This is correct for the
49 # common case (autocommit / non-branch writes route to the default connection), but if a write
50 # under a branch schema ever reached here without its alias, the deferred write would silently
51 # land in the main schema. Log at debug so that case is observable rather than invisible.
52 if not using: 52 ↛ 53line 52 didn't jump to line 53 because the condition on line 52 was never true
53 logger.debug("Search cache: no originating DB alias for object %s/%s; using default", object_type_id, pk)
54 alias = using or DEFAULT_DB_ALIAS
55 connection = connections[alias]
57 # No transaction in progress: index synchronously. Deferring would have
58 # nothing to defer past, and transaction.on_commit() in autocommit mode runs
59 # its callback immediately at registration (before we could populate the
60 # batch), so handle this case explicitly.
61 #
62 # On the transactional path below, transaction.on_commit(..., robust=True)
63 # ensures a flush failure can never propagate to the (already-committed)
64 # caller. The autocommit path has no such backstop, so guard it here: a
65 # broad catch is deliberate, because the originating write has committed and
66 # a search cache update must never turn a successful save into an error. The
67 # error is logged so a genuine indexing defect is still visible.
68 if not connection.in_atomic_block:
69 try:
70 _flush({(object_type_id, pk): op}, alias)
71 except Exception:
72 logger.exception("Search cache: error while indexing inline")
73 return
75 # Scope buffering to the current savepoint stack, not just the alias (see
76 # _pending_batch). May legitimately contain None entries for nested
77 # atomic(savepoint=False) blocks; matching is by equality, so that is fine.
78 scope = tuple(connection.savepoint_ids)
80 batch = _pending_batch(connection, alias, scope)
81 if batch is None:
82 batch = {}
84 def flush(batch=batch, alias=alias):
85 _flush(batch, alias)
87 setattr(flush, _FLUSH_ALIAS_ATTR, alias)
88 setattr(flush, _FLUSH_BATCH_ATTR, batch)
89 setattr(flush, _FLUSH_SCOPE_ATTR, scope)
90 # robust=True is required, not just belt-and-suspenders: Django runs
91 # on_commit callbacks synchronously as the atomic block exits (after the
92 # COMMIT), so an exception escaping the callback would propagate out of
93 # the view's transaction and become a 500 on an already-committed write.
94 # _flush handles the recoverable Redis fault itself; robust=True is the
95 # only thing that keeps any *other* failure here (logged by Django at
96 # ERROR) from surfacing as that post-commit 500.
97 transaction.on_commit(flush, using=alias, robust=True)
99 # Coalesce: a deletion supersedes any pending create/update for the object.
100 key = (object_type_id, pk)
101 if op == OP_REMOVE or batch.get(key) != OP_REMOVE: 101 ↛ exitline 101 didn't return from function 'mark_for_deferred_indexing' because the condition on line 101 was always true
102 batch[key] = op
105def _pending_batch(connection, alias, scope):
106 """
107 Return the batch dict of a flush callback already scheduled for the given
108 alias and savepoint scope on this connection's current transaction, or None
109 if there is none.
111 This scans `connection.run_on_commit` on each call rather than caching the
112 lookup elsewhere. That is intentional: the scan is bounded (run_on_commit
113 holds only the transaction's registered commit callbacks, not one per saved
114 object), and reading it fresh each time is what keeps the buffer correctly
115 scoped to the live transaction. Django clears run_on_commit on both commit
116 and rollback, so a rolled-back transaction's batch can never be found here.
118 Matching on `scope` (the savepoint stack active when the callback was
119 registered) as well as `alias` keeps each savepoint scope on its own
120 callback. Django prunes a callback when a savepoint in its registration
121 snapshot rolls back, so an op buffered inside a nested savepoint is dropped
122 with its callback if that savepoint rolls back -- it can never be found here
123 and reused by an outer scope.
124 """
125 for _sids, func, _robust in connection.run_on_commit:
126 if (
127 getattr(func, _FLUSH_ALIAS_ATTR, None) == alias
128 and getattr(func, _FLUSH_SCOPE_ATTR, None) == scope
129 ):
130 return getattr(func, _FLUSH_BATCH_ATTR)
131 return None
134def _flush(batch, using):
135 """
136 Dispatch a coalesced batch of dirty objects for (re)indexing.
138 `_flush` is the single guarded entry point for deferred indexing, reached
139 either directly (autocommit) or from a transaction.on_commit callback. By the
140 time it runs the originating write has already committed, so it must never
141 propagate an error back to the caller and turn a successful save into a 500.
143 The inline fallback is safe even during a broker outage: the search index
144 lives in PostgreSQL (the extras_cachedvalue table), so a Redis outage only
145 prevents backgrounding, not indexing itself.
147 If the broker fails mid-enqueue (after the probe succeeds), Job.enqueue() has
148 already saved a Job row before the Redis dispatch raised, so the fallback can
149 leave behind a PENDING Job that no worker will run. The index is still
150 correct (written inline); the stranded row is cosmetic and ages out via the
151 housekeeping job.
152 """
153 if not batch: 153 ↛ 154line 153 didn't jump to line 154 because the condition on line 153 was never true
154 return
156 cache_groups = {}
157 remove_groups = {}
158 for (object_type_id, pk), op in batch.items():
159 groups = remove_groups if op == OP_REMOVE else cache_groups
160 groups.setdefault(object_type_id, []).append(pk)
162 try:
163 # Both the worker-availability check and the job enqueue talk to Redis,
164 # and a worker can die between the two. Treat any Redis failure across the
165 # whole dispatch as "no worker available" and fall back to inline
166 # indexing (a PostgreSQL write that does not depend on Redis).
167 if any_workers_for_queue(RQ_QUEUE_DEFAULT): 167 ↛ 168line 167 didn't jump to line 168 because the condition on line 167 was never true
168 SearchCacheJob.enqueue(using=using, cache_groups=cache_groups, remove_groups=remove_groups)
169 return
170 except RedisError:
171 logger.warning("Search cache: broker unavailable; indexing inline", exc_info=True)
173 search_backend._apply_deferred_updates(using=using, cache_groups=cache_groups, remove_groups=remove_groups)