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
« 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>
4import logging
5from typing import override
7import grpc
8from google.protobuf import empty_pb2
9from sqlalchemy import select
10from sqlalchemy.orm import Session, joinedload
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
18logger = logging.getLogger(__name__)
21class PolicyQueueServiceServicer(policyqueue_pb2_grpc.PolicyQueueServiceServicer):
22 def __init__(self, conn: DBConn) -> None:
23 self._conn = conn
25 _query = select(PolicyQueueUpload).options(
26 joinedload(PolicyQueueUpload.changes),
27 joinedload(PolicyQueueUpload.policy_queue),
28 joinedload(PolicyQueueUpload.target_suite),
29 )
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()
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 )
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")
81 handler = PolicyQueueUploadHandler(upload, session)
82 missing_overrides = handler.missing_overrides()
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 )
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 )
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 )
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")
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 )
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 )
172 return empty_pb2.Empty()
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")
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 )
210 handler.accept()
212 return empty_pb2.Empty()
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")
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 )
243 handler.reject(request.reason, rejected_by=request.rejected_by)
245 return empty_pb2.Empty()