Coverage for utilities/querysets.py: 54%

57 statements  

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

1from contextlib import nullcontext 

2 

3from django.conf import settings 

4from django.db import router, transaction 

5from django.db.models import Max, Prefetch, QuerySet 

6 

7from users.constants import CONSTRAINT_TOKEN_USER 

8from utilities.permissions import get_permission_for_model, permission_is_exempt, qs_filter_from_constraints 

9 

10__all__ = ( 

11 'RestrictedPrefetch', 

12 'RestrictedQuerySet', 

13 'chunked_update', 

14) 

15 

16 

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

25 

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. 

29 

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. 

40 

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) 

51 

52 model = queryset.model 

53 

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) 

58 

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

84 

85 return count 

86 

87 

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 

96 

97 super().__init__(lookup, queryset=queryset, to_attr=to_attr) 

98 

99 def get_current_querysets(self, level): 

100 params = { 

101 'user': self.restrict_user, 

102 'action': self.restrict_action, 

103 } 

104 

105 if querysets := super().get_current_querysets(level): 

106 return [qs.restrict(**params) for qs in querysets] 

107 

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 

112 

113 

114class RestrictedQuerySet(QuerySet): 

115 

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. 

120 

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) 

126 

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 

130 

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

134 

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) 

145 

146 return self