Coverage for utilities/querysets.py: 54%
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
1from contextlib import nullcontext
3from django.conf import settings
4from django.db import router, transaction
5from django.db.models import Max, Prefetch, QuerySet
7from users.constants import CONSTRAINT_TOKEN_USER
8from utilities.permissions import get_permission_for_model, permission_is_exempt, qs_filter_from_constraints
10__all__ = (
11 'RestrictedPrefetch',
12 'RestrictedQuerySet',
13 'chunked_update',
14)
17def chunked_update(queryset, chunk_size=None, commit_per_batch=False, **kwargs):
18 """
19 Perform a bulk UPDATE on the given queryset, optionally splitting it into batches of at most
20 `chunk_size` rows. Bounding the number of rows touched by each statement avoids exceeding the
21 database's statement timeout when updating very large tables. Batches are selected via keyset
22 pagination on the primary key and wrapped in a transaction so that the operation remains atomic,
23 as it would be when performed by a single UPDATE. Returns the total number of rows updated
24 (matching the return value of QuerySet.update()).
26 If `chunk_size` is None, it falls back to the BULK_UPDATE_CHUNK_SIZE configuration parameter
27 (5000 by default). If that is also None, a single unbounded UPDATE is issued, identical to
28 calling queryset.update(**kwargs) directly.
30 :param queryset: The QuerySet identifying the rows to update
31 :param chunk_size: The maximum number of rows to update per statement (defaults to
32 settings.BULK_UPDATE_CHUNK_SIZE)
33 :param commit_per_batch: Commit each batch independently rather than wrapping them all in a
34 single transaction, forfeiting the atomicity described above. Postgres holds a row lock on
35 every row updated until the transaction commits, which for a long-running job spanning a
36 large table means blocking concurrent edits for its whole run, so such callers commit as
37 they go. Only for updates which can safely be resumed, and only outside an enclosing atomic
38 block, which owns the commit regardless -- passed from within one, this silently degrades to
39 a single transaction rather than raising, as Django's TestCase wraps every test in one.
41 The callers which pass it (see extras.jobs) run in a worker, under a session-scoped advisory
42 lock which is taken through a bare cursor, so nothing in that path opens a transaction and
43 each batch does commit as intended.
44 """
45 if chunk_size is None: 45 ↛ 47line 45 didn't jump to line 47 because the condition on line 45 was always true
46 chunk_size = settings.BULK_UPDATE_CHUNK_SIZE
47 if chunk_size is not None and (type(chunk_size) is not int or chunk_size < 1): 47 ↛ 48line 47 didn't jump to line 48 because the condition on line 47 was never true
48 raise ValueError(f"chunk_size must be a positive integer or None (found {chunk_size!r})")
49 if chunk_size is None: 49 ↛ 50line 49 didn't jump to line 50 because the condition on line 49 was never true
50 return queryset.update(**kwargs)
52 model = queryset.model
54 # Pin the entire operation to a single write database so that the PK lookups, the UPDATE
55 # statements, and the enclosing transaction all use the same connection. This preserves an
56 # explicit .using() on the queryset and otherwise honors the router's write destination.
57 using = queryset._db or router.db_for_write(model)
59 count = 0
60 last_pk = 0
61 # Upper bound on the PKs to process. Established lazily (see below) only once a second batch is
62 # known to be needed, so the common single-batch case incurs no extra aggregate query.
63 max_pk = None
64 with nullcontext() if commit_per_batch else transaction.atomic(using=using):
65 while True:
66 batch = queryset.using(using).filter(pk__gt=last_pk).order_by('pk')
67 if max_pk is not None: 67 ↛ 68line 67 didn't jump to line 68 because the condition on line 67 was never true
68 batch = batch.filter(pk__lte=max_pk)
69 pks = list(batch.values_list('pk', flat=True)[:chunk_size])
70 if not pks: 70 ↛ 71line 70 didn't jump to line 71 because the condition on line 70 was never true
71 break
72 # Re-filter the original queryset by pk__in (rather than the model's default manager) so
73 # that its own filters are preserved and rows no longer matching them are left untouched.
74 count += queryset.using(using).filter(pk__in=pks).update(**kwargs)
75 last_pk = pks[-1]
76 # A batch shorter than chunk_size means the rows are exhausted; stop without issuing a
77 # trailing (empty) lookup. This keeps a single-batch update to one SELECT and one UPDATE.
78 if len(pks) < chunk_size: 78 ↛ 82line 78 didn't jump to line 82 because the condition on line 78 was always true
79 break
80 # A full batch means more rows may remain. Capture the current maximum PK as an upper
81 # bound (once) so that rows inserted while the operation runs cannot keep extending it.
82 if max_pk is None:
83 max_pk = queryset.using(using).aggregate(_max=Max('pk'))['_max']
85 return count
88class RestrictedPrefetch(Prefetch):
89 """
90 Extend Django's Prefetch to accept a user and action to be passed to the
91 `restrict()` method of the related object's queryset.
92 """
93 def __init__(self, lookup, user, action='view', queryset=None, to_attr=None):
94 self.restrict_user = user
95 self.restrict_action = action
97 super().__init__(lookup, queryset=queryset, to_attr=to_attr)
99 def get_current_querysets(self, level):
100 params = {
101 'user': self.restrict_user,
102 'action': self.restrict_action,
103 }
105 if querysets := super().get_current_querysets(level):
106 return [qs.restrict(**params) for qs in querysets]
108 # Bit of a hack. If no queryset is defined, pass through the dict of restrict()
109 # kwargs to be handled by the field. This is necessary e.g. for GenericForeignKey
110 # fields, which do not permit setting a queryset on a Prefetch object.
111 return params
114class RestrictedQuerySet(QuerySet):
116 def restrict(self, user, action='view'):
117 """
118 Filter the QuerySet to return only objects on which the specified user has been granted the specified
119 permission.
121 :param user: User instance
122 :param action: The action which must be permitted (e.g. "view" for "dcim.view_site"); default is 'view'
123 """
124 # Resolve the full name of the required permission
125 permission_required = get_permission_for_model(self.model, action)
127 # Bypass restriction for superusers and exempt views
128 if (user and user.is_active and user.is_superuser) or permission_is_exempt(permission_required): 128 ↛ 132line 128 didn't jump to line 132 because the condition on line 128 was always true
129 return self
131 # User is anonymous or has not been granted the requisite permission
132 if user is None or not user.is_authenticated or permission_required not in user.get_all_permissions():
133 return self.none()
135 # Filter the queryset to include only objects with allowed attributes
136 constraints = user._object_perm_cache[permission_required]
137 tokens = {
138 CONSTRAINT_TOKEN_USER: user,
139 }
140 if attrs := qs_filter_from_constraints(constraints, tokens):
141 # #8715: Avoid duplicates when JOIN on many-to-many fields without using DISTINCT.
142 # DISTINCT acts globally on the entire request, which may not be desirable.
143 allowed_objects = self.model.objects.filter(attrs)
144 return self.filter(pk__in=allowed_objects)
146 return self