Skip to content

g1t/services/work/src/reviews.rs

390 lines14,657 bytesCodeBlame
1//! What g1t adds around a pull request for many agents working at once:
2//! which files it changes, which other pull requests are changing the same
3//! files, whether it has fallen behind, and reviews written by a g1t agent.
4
5use std::collections::HashSet;
6
7use g1t_contracts::events::ReviewEvent;
8use g1t_contracts::repos::{BehindArgs, Comparison, GetByIdArgs, HeadArgs, Repo, RepoPath};
9use g1t_contracts::time::rfc3339;
10use g1t_contracts::work::*;
11use g1t_contracts::{FailureCode, Outcome, new_id};
12use g1t_kit::now_ms;
13use serde::Deserialize;
14use sha2::{Digest, Sha256};
15use worker::Result;
16use worker::wasm_bindgen::JsValue;
17
18use crate::rows::{PULL_COLUMNS, PullRow};
19use crate::{Work, optional_number};
20
21/// Beyond this a change is too large for per-file bookkeeping to be useful.
22const MAX_FILES: usize = 300;
23const MAX_REVIEW_COMMENTS: usize = 30;
24const MAX_REVIEW_CHARS: usize = 20_000;
25/// Who reviews by g1t's agent are attributed to: shown as `g1t`, which
26/// nobody can register or sign in as.
27pub(crate) use g1t_contracts::identity::{AGENT_ID, AGENT_NAME};
28
29#[derive(Deserialize)]
30struct OtherRow {
31 number: u32,
32 title: String,
33 issue_number: Option<u32>,
34 files: Option<String>,
35}
36
37#[derive(Deserialize)]
38struct ReviewRunRow {
39 id: String,
40 pull_id: String,
41 token_hash: String,
42 finished_at: Option<String>,
43}
44
45fn hash(token: &str) -> String {
46 hex::encode(Sha256::digest(token.as_bytes()))
47}
48
49impl Work {
50 pub(crate) async fn pull_by_id(&self, id: &str) -> Result<Option<Pull>> {
51 Ok(self
52 .db
53 .prepare(format!("SELECT {PULL_COLUMNS} FROM pulls WHERE id = ?"))
54 .bind(&[id.into()])?
55 .first::<PullRow>(None)
56 .await?
57 .map(Pull::from))
58 }
59
60 /// Works out which files a pull request changes and records them, so
61 /// that overlaps between pull requests can be found without comparing
62 /// every one of them each time.
63 pub(crate) async fn refresh_files(&self, pull: &Pull) -> Result<Vec<ChangedFile>> {
64 // Its author can read both the repository and the pull request's source.
65 let viewer = self.author_viewer(pull).await?;
66 let compared: Outcome<Comparison> =
67 g1t_kit::call(&self.repos, "compare", &pull.comparison(&viewer)).await?;
68 let files: Vec<ChangedFile> = match compared {
69 Outcome::Ok(comparison) => comparison
70 .files
71 .into_iter()
72 .take(MAX_FILES)
73 .map(|file| ChangedFile {
74 path: file.path,
75 additions: file.additions,
76 deletions: file.deletions,
77 })
78 .collect(),
79 // Nothing pushed yet, or nothing to compare against.
80 Outcome::Fail(_) => Vec::new(),
81 };
82 self.db
83 .prepare("UPDATE pulls SET files = ? WHERE id = ?")
84 .bind(&[
85 serde_json::to_string(&files)?.into(),
86 pull.id.as_str().into(),
87 ])?
88 .run()
89 .await?;
90 Ok(files)
91 }
92
93 /// Other pull requests in progress that change files this one changes.
94 pub(crate) async fn overlaps(&self, pull: &Pull) -> Result<Vec<Overlap>> {
95 if pull.files.is_empty() {
96 return Ok(Vec::new());
97 }
98 let mine: HashSet<&str> = pull.files.iter().map(|file| file.path.as_str()).collect();
99 let others = match self.prefetched_pull(&pull.id) {
100 Some(found) => found.rows::<OtherRow>(crate::prefetch::Slot::Others)?,
101 None => self.other_pulls(pull).await?,
102 };
103 Ok(others
104 .into_iter()
105 .filter_map(|other| {
106 let files: Vec<ChangedFile> =
107 serde_json::from_str(other.files.as_deref().unwrap_or("[]")).ok()?;
108 let paths: Vec<String> = files
109 .into_iter()
110 .map(|file| file.path)
111 .filter(|path| mine.contains(path.as_str()))
112 .collect();
113 (!paths.is_empty()).then_some(Overlap {
114 number: other.number,
115 title: other.title,
116 issue: other.issue_number,
117 paths,
118 })
119 })
120 .collect())
121 }
122
123 /// The other pull requests in progress in its repository.
124 async fn other_pulls(&self, pull: &Pull) -> Result<Vec<OtherRow>> {
125 self
126 .db
127 .prepare(
128 "SELECT number, title, issue_number, files FROM pulls
129 WHERE repo_id = ? AND id != ? AND status IN ('draft', 'open')
130 ORDER BY number LIMIT 200",
131 )
132 .bind(&[pull.repo_id.as_str().into(), pull.id.as_str().into()])?
133 .all()
134 .await?
135 .results::<OtherRow>()
136 }
137
138 /// Whether the branch a pull request would merge into has commits the
139 /// pull request does not, so that it must catch up before it can merge.
140 pub(crate) async fn is_behind(&self, repo_id: &str, pull: &Pull) -> Result<bool> {
141 if !pull.status.is_active() || pull.head_commit.is_none() {
142 return Ok(false);
143 }
144 let asked = BehindArgs {
145 source_id: pull
146 .fork_repo_id
147 .clone()
148 .unwrap_or_else(|| repo_id.to_owned()),
149 branch: pull.branch.clone(),
150 };
151 self.timing.rpc(g1t_kit::call(&self.repos, "behind", &asked)).await
152 }
153
154 pub(crate) async fn review_pending(&self, pull_id: &str) -> Result<bool> {
155 if let Some(found) = self.prefetched_pull(pull_id) {
156 return Ok(found.first::<crate::rows::ValueRow>(crate::prefetch::Slot::ReviewPending)?.is_some());
157 }
158 Ok(self
159 .db
160 .prepare(
161 "SELECT id AS value FROM review_runs
162 WHERE pull_id = ? AND finished_at IS NULL AND created_at > ? LIMIT 1",
163 )
164 // A sandbox that never reported is forgotten after a while.
165 .bind(&[
166 pull_id.into(),
167 rfc3339(now_ms().saturating_sub(crate::prefetch::REVIEW_PENDING_MS)).into(),
168 ])?
169 .first::<crate::rows::ValueRow>(None)
170 .await?
171 .is_some())
172 }
173
174 /// Begins a review of a pull request by a g1t agent and returns what a
175 /// sandbox needs to carry it out.
176 pub(crate) async fn start_review(&self, a: StartReviewArgs) -> Result<Outcome<ReviewJob>> {
177 let Some(pull) = self.pull_by_id(&a.pull_id).await? else {
178 return Ok(Outcome::fail(
179 FailureCode::NotFound,
180 "Pull request not found.",
181 ));
182 };
183 if pull.status != PullStatus::Open {
184 return Ok(Outcome::fail(
185 FailureCode::Conflict,
186 "A pull request can be reviewed once it is ready for review.",
187 ));
188 }
189 let viewer = self.author_viewer(&pull).await?;
190 let repo: Outcome<Repo> = g1t_kit::call(
191 &self.repos,
192 "get_by_id",
193 &GetByIdArgs {
194 id: pull.repo_id.clone(),
195 viewer,
196 },
197 )
198 .await?;
199 let Outcome::Ok(repo) = crate::retired::unless_archived(repo) else {
200 return Ok(Outcome::fail(
201 FailureCode::NotFound,
202 "Pull request not found.",
203 ));
204 };
205 let head: Option<String> = g1t_kit::call(
206 &self.repos,
207 "head",
208 &HeadArgs {
209 repo_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
210 branch: pull
211 .branch
212 .clone()
213 .unwrap_or_else(|| repo.default_branch.clone()),
214 },
215 )
216 .await?;
217 let Some(commit) = head else {
218 return Ok(Outcome::fail(
219 FailureCode::Conflict,
220 "This pull request has no commits to review.",
221 ));
222 };
223 let issue = match pull.issue {
224 Some(number) => self.issue(&pull.repo_id, number).await?,
225 None => None,
226 };
227
228 let now = now_ms();
229 let run_id = new_id("rvw", now);
230 let mut bytes = [0u8; 32];
231 getrandom::getrandom(&mut bytes).expect("no source of randomness");
232 let token = hex::encode(bytes);
233 self.db
234 .prepare(
235 "INSERT INTO review_runs (id, pull_id, token_hash, created_at) VALUES (?, ?, ?, ?)",
236 )
237 .bind(&[
238 run_id.as_str().into(),
239 pull.id.as_str().into(),
240 hash(&token).into(),
241 rfc3339(now).into(),
242 ])?
243 .run()
244 .await?;
245
246 let path = RepoPath {
247 namespace: repo.namespace,
248 name: repo.name,
249 };
250 let mut sensitive: Vec<String> = Vec::new();
251 for kind in pull.files.iter().filter_map(|file| crate::confidence::sensitive(&file.path)) {
252 if !sensitive.iter().any(|seen| seen == kind) {
253 sensitive.push(kind.to_owned());
254 }
255 }
256 Ok(Outcome::Ok(ReviewJob {
257 run_id,
258 token,
259 source: pull.fork.clone().unwrap_or_else(|| path.clone()),
260 commit,
261 repo: path,
262 default_branch: repo.default_branch,
263 number: pull.number,
264 title: pull.title,
265 description: pull.body.unwrap_or_default(),
266 issue,
267 author: pull.author,
268 files: pull.files,
269 sensitive,
270 }))
271 }
272
273 /// Records the review a sandbox's agent wrote: its verdict, its summary
274 /// and its comments on lines. The run's token is the only credential.
275 pub(crate) async fn report_review(&self, a: ReportReviewArgs) -> Result<Outcome<bool>> {
276 let run = self
277 .db
278 .prepare("SELECT id, pull_id, token_hash, finished_at FROM review_runs WHERE id = ?")
279 .bind(&[a.run_id.as_str().into()])?
280 .first::<ReviewRunRow>(None)
281 .await?;
282 let Some(run) = run.filter(|run| run.token_hash == hash(&a.token)) else {
283 return Ok(Outcome::fail(FailureCode::NotFound, "Review not found."));
284 };
285 if run.finished_at.is_some() {
286 return Ok(Outcome::fail(
287 FailureCode::Conflict,
288 "This review has already been reported.",
289 ));
290 }
291 let Some(pull) = self.pull_by_id(&run.pull_id).await? else {
292 return Ok(Outcome::fail(FailureCode::NotFound, "Review not found."));
293 };
294
295 let now = now_ms();
296 let timestamp = rfc3339(now);
297 let signature = a
298 .model
299 .as_deref()
300 .map(|model| format!("\n\n_Reviewed by g1t on {model}._"))
301 .unwrap_or_default();
302 let summary = match &a.error {
303 Some(error) => format!("The review could not be completed: {error}"),
304 None => {
305 let body: String = a.body.trim().chars().take(MAX_REVIEW_CHARS).collect();
306 format!("{body}{signature}")
307 }
308 };
309 // A failed review carries no verdict.
310 let verdict = a.verdict.filter(|_| a.error.is_none());
311 let known: HashSet<&str> = pull.files.iter().map(|file| file.path.as_str()).collect();
312
313 let insert = "INSERT INTO comments
314 (id, repo_id, number, author_id, author_name, body, path, line, verdict, created_at)
315 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)";
316 let mut statements = Vec::new();
317 // Line comments first, so the summary and verdict read as the conclusion.
318 for comment in a
319 .comments
320 .iter()
321 .filter(|comment| a.error.is_none() && !comment.body.trim().is_empty())
322 // A line in a file the pull request does not change is not a review of it.
323 .filter(|comment| known.is_empty() || known.contains(comment.path.as_str()))
324 .take(MAX_REVIEW_COMMENTS)
325 {
326 let body: String = comment.body.trim().chars().take(MAX_REVIEW_CHARS).collect();
327 statements.push(self.db.prepare(insert).bind(&[
328 new_id("cmt", now).into(),
329 pull.repo_id.as_str().into(),
330 pull.number.into(),
331 AGENT_ID.into(),
332 AGENT_NAME.into(),
333 body.into(),
334 comment.path.as_str().into(),
335 optional_number(Some(comment.line).filter(|line| *line > 0)),
336 JsValue::NULL,
337 timestamp.as_str().into(),
338 ])?);
339 }
340 statements.push(self.db.prepare(insert).bind(&[
341 new_id("cmt", now).into(),
342 pull.repo_id.as_str().into(),
343 pull.number.into(),
344 AGENT_ID.into(),
345 AGENT_NAME.into(),
346 summary.into(),
347 JsValue::NULL,
348 JsValue::NULL,
349 verdict.map_or(JsValue::NULL, |verdict| verdict.as_str().into()),
350 timestamp.as_str().into(),
351 ])?);
352 statements.push(
353 self.db
354 .prepare("UPDATE review_runs SET finished_at = ?, verdict = ? WHERE id = ?")
355 .bind(&[
356 timestamp.as_str().into(),
357 verdict.map_or(JsValue::NULL, |verdict| verdict.as_str().into()),
358 run.id.as_str().into(),
359 ])?,
360 );
361 statements.push(
362 self.db
363 .prepare("UPDATE pulls SET updated_at = ? WHERE id = ?")
364 .bind(&[timestamp.as_str().into(), pull.id.as_str().into()])?,
365 );
366 // The review g1t asked for itself is no longer under way.
367 statements.push(
368 self.db
369 .prepare(
370 "UPDATE pulls SET working_on = NULL, working_until = NULL
371 WHERE id = ? AND working_on = 'review'",
372 )
373 .bind(&[pull.id.as_str().into()])?,
374 );
375 self.db.batch(statements).await?;
376 self.publish_as(
377 "review.completed",
378 &pull.repo_id,
379 None,
380 ReviewEvent {
381 pull_id: pull.id.clone(),
382 repo_id: pull.repo_id.clone(),
383 number: pull.number,
384 verdict: verdict.map(Verdict::as_str),
385 },
386 )
387 .await?;
388 Ok(Outcome::Ok(true))
389 }
390}