Skip to content

g1t/services/work/src/reviews.rs

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