Coverage for daklib/policy_rpc.py: 87%

89 statements  

« prev     ^ index     » next       coverage.py v7.6.0, created at 2026-08-03 16:46 +0000

1# SPDX-License-Identifier: GPL-2.0-or-later 

2# © 2026, Ansgar 🙀 <ansgar@debian.org> 

3 

4import logging 

5from typing import override 

6 

7import grpc 

8from google.protobuf import empty_pb2 

9from sqlalchemy import select 

10from sqlalchemy.orm import Session, joinedload 

11 

12from dak.policyqueue.v1 import policyqueue_pb2, policyqueue_pb2_grpc 

13from daklib.dbconn import DBChange, DBConn, PolicyQueue, PolicyQueueUpload 

14from daklib.policy import PolicyQueueUploadHandler, section_matches_component 

15from daklib.rpc import to_timestamp 

16from daklib.rpc_auth import require_any_scope, require_scope 

17 

18logger = logging.getLogger(__name__) 

19 

20 

21class PolicyQueueServiceServicer(policyqueue_pb2_grpc.PolicyQueueServiceServicer): 

22 def __init__(self, conn: DBConn) -> None: 

23 self._conn = conn 

24 

25 _query = select(PolicyQueueUpload).options( 

26 joinedload(PolicyQueueUpload.changes), 

27 joinedload(PolicyQueueUpload.policy_queue), 

28 joinedload(PolicyQueueUpload.target_suite), 

29 ) 

30 

31 def _upload( 

32 self, session: Session, policy_queue: str, name: str, *, for_update=False 

33 ) -> PolicyQueueUpload | None: 

34 query = ( 

35 self._query.join(PolicyQueueUpload.changes) 

36 .join(PolicyQueueUpload.policy_queue) 

37 .where(PolicyQueue.queue_name == policy_queue, DBChange.changesname == name) 

38 ) 

39 if for_update: 

40 query = query.with_for_update(of=PolicyQueueUpload) 

41 return session.execute(query).scalar_one_or_none() 

42 

43 @override 

44 def ListUploads( 

45 self, request: policyqueue_pb2.ListUploadsRequest, context: grpc.ServicerContext 

46 ) -> policyqueue_pb2.ListUploadsResponse: 

47 require_scope(context, "policyqueue:read") 

48 logger.debug("list uploads: policy_queue=%s", request.policy_queue) 

49 with self._conn.session() as session: 

50 query = self._query.join(PolicyQueueUpload.policy_queue).where( 

51 PolicyQueue.queue_name == request.policy_queue 

52 ) 

53 uploads = session.execute(query).scalars() 

54 return policyqueue_pb2.ListUploadsResponse( 

55 uploads=[ 

56 policyqueue_pb2.Upload( 

57 name=u.changes.changesname, 

58 policy_queue=u.policy_queue.queue_name, 

59 target_suite=u.target_suite.suite_name, 

60 create_time=to_timestamp(u.changes.created), 

61 ) 

62 for u in uploads 

63 ], 

64 ) 

65 

66 @override 

67 def GetUpload( 

68 self, request: policyqueue_pb2.GetUploadRequest, context: grpc.ServicerContext 

69 ) -> policyqueue_pb2.Upload: 

70 require_scope(context, "policyqueue:read") 

71 logger.debug( 

72 "get upload: policy_queue=%s name=%s", 

73 request.policy_queue, 

74 request.name, 

75 ) 

76 with self._conn.session() as session: 

77 upload = self._upload(session, request.policy_queue, request.name) 

78 if upload is None: 78 ↛ 79line 78 didn't jump to line 79 because the condition on line 78 was never true

79 context.abort(grpc.StatusCode.NOT_FOUND, "not found") 

80 

81 handler = PolicyQueueUploadHandler(upload, session) 

82 missing_overrides = handler.missing_overrides() 

83 

84 return policyqueue_pb2.Upload( 

85 name=upload.changes.changesname, 

86 policy_queue=upload.policy_queue.queue_name, 

87 target_suite=upload.target_suite.suite_name, 

88 missing_overrides=[ 

89 policyqueue_pb2.Override( 

90 package=o["package"], 

91 component=o["component"], 

92 priority=o["priority"], 

93 section=o["section"], 

94 type=o["type"], 

95 ) 

96 for o in missing_overrides 

97 ], 

98 create_time=to_timestamp(upload.changes.created), 

99 ) 

100 

101 @override 

102 def AddOverrides( 

103 self, 

104 request: policyqueue_pb2.AddOverridesRequest, 

105 context: grpc.ServicerContext, 

106 ) -> empty_pb2.Empty: 

107 require_any_scope( 

108 context, ["policyqueue:write", f"policyqueue:write:{request.policy_queue}"] 

109 ) 

110 logger.info( 

111 "add overrides: policy_queue=%s name=%s count=%d", 

112 request.policy_queue, 

113 request.name, 

114 len(request.overrides), 

115 ) 

116 

117 mismatched = [ 

118 o 

119 for o in request.overrides 

120 if not section_matches_component(o.section, o.component) 

121 ] 

122 if mismatched: 

123 packages = ", ".join( 

124 f'{o.type}:{o.package} (section "{o.section}", component "{o.component}")' 

125 for o in mismatched 

126 ) 

127 context.abort( 

128 grpc.StatusCode.INVALID_ARGUMENT, 

129 f"section does not match component for {packages}", 

130 ) 

131 

132 with self._conn.session() as session: 

133 upload = self._upload( 

134 session, request.policy_queue, request.name, for_update=True 

135 ) 

136 if not upload: 136 ↛ 137line 136 didn't jump to line 137 because the condition on line 136 was never true

137 context.abort(grpc.StatusCode.NOT_FOUND, "not found") 

138 

139 handler = PolicyQueueUploadHandler(upload, session) 

140 missing_overrides = { 

141 (o["component"], o["type"], o["package"]) 

142 for o in handler.missing_overrides() 

143 } 

144 unneeded_overrides = [ 

145 o 

146 for o in request.overrides 

147 if (o.component, o.type, o.package) not in missing_overrides 

148 ] 

149 if unneeded_overrides: 149 ↛ 150line 149 didn't jump to line 150 because the condition on line 149 was never true

150 packages = ", ".join( 

151 f"{o.type}:{o.package}" for o in unneeded_overrides 

152 ) 

153 context.abort( 

154 grpc.StatusCode.FAILED_PRECONDITION, 

155 f"no new override needed for {packages}", 

156 ) 

157 

158 handler.add_overrides( 

159 [ 

160 { 

161 "component": o.component, 

162 "type": o.type, 

163 "package": o.package, 

164 "priority": o.priority, 

165 "section": o.section, 

166 } 

167 for o in request.overrides 

168 ], 

169 upload.target_suite, 

170 ) 

171 

172 return empty_pb2.Empty() 

173 

174 @override 

175 def AcceptUpload( 

176 self, 

177 request: policyqueue_pb2.AcceptUploadRequest, 

178 context: grpc.ServicerContext, 

179 ) -> empty_pb2.Empty: 

180 require_any_scope( 

181 context, ["policyqueue:write", f"policyqueue:write:{request.policy_queue}"] 

182 ) 

183 logger.info( 

184 "accept upload: policy_queue=%s name=%s", 

185 request.policy_queue, 

186 request.name, 

187 ) 

188 with self._conn.session() as session: 

189 upload = self._upload( 

190 session, request.policy_queue, request.name, for_update=True 

191 ) 

192 if not upload: 192 ↛ 193line 192 didn't jump to line 193 because the condition on line 192 was never true

193 context.abort(grpc.StatusCode.NOT_FOUND, "not found") 

194 

195 handler = PolicyQueueUploadHandler(upload, session) 

196 if missing_overrides := handler.missing_overrides(): 

197 packages = ", ".join( 

198 f"{o['type']}:{o['package']}" for o in missing_overrides 

199 ) 

200 context.abort( 

201 grpc.StatusCode.FAILED_PRECONDITION, 

202 f"missing overrides for {packages}", 

203 ) 

204 if (action := handler.get_action()) and action != "ACCEPT": 204 ↛ 205line 204 didn't jump to line 205 because the condition on line 204 was never true

205 context.abort( 

206 grpc.StatusCode.FAILED_PRECONDITION, 

207 f"upload is already processed with action {action}", 

208 ) 

209 

210 handler.accept() 

211 

212 return empty_pb2.Empty() 

213 

214 @override 

215 def RejectUpload( 

216 self, 

217 request: policyqueue_pb2.RejectUploadRequest, 

218 context: grpc.ServicerContext, 

219 ) -> empty_pb2.Empty: 

220 require_any_scope( 

221 context, ["policyqueue:write", f"policyqueue:write:{request.policy_queue}"] 

222 ) 

223 logger.info( 

224 "reject upload: policy_queue=%s name=%s rejected_by=%s", 

225 request.policy_queue, 

226 request.name, 

227 request.rejected_by, 

228 ) 

229 with self._conn.session() as session: 

230 upload = self._upload( 

231 session, request.policy_queue, request.name, for_update=True 

232 ) 

233 if not upload: 233 ↛ 234line 233 didn't jump to line 234 because the condition on line 233 was never true

234 context.abort(grpc.StatusCode.NOT_FOUND, "not found") 

235 

236 handler = PolicyQueueUploadHandler(upload, session) 

237 if (action := handler.get_action()) and action != "REJECT": 237 ↛ 238line 237 didn't jump to line 238 because the condition on line 237 was never true

238 context.abort( 

239 grpc.StatusCode.FAILED_PRECONDITION, 

240 f"upload is already processed with action {action}", 

241 ) 

242 

243 handler.reject(request.reason, rejected_by=request.rejected_by) 

244 

245 return empty_pb2.Empty()